(in-package #:nekod.socket) ;;; Backpressure controls — prevent unbounded memory growth from infinite streams (defstruct backpressure-config (maximum-read-buffer (* 1024 1024) :type integer) (maximum-event-size (* 256 1024) :type integer) (maximum-header-size (* 64 1024) :type integer) (consumer-queue-depth 1024 :type integer) (high-water-mark 768 :type integer) (low-water-mark 256 :type integer)) (defvar *backpressure* (make-backpressure-config)) (defstruct backpressure-state (queue-size 0 :type integer) (paused nil :type boolean) (events-dropped 0 :type integer) (total-bytes-buffered 0 :type integer)) (defvar *bp-state* (make-backpressure-state)) (defun check-event-size (size) "Signal event-too-large if event exceeds maximum." (when (> size (backpressure-config-maximum-event-size *backpressure*)) (error 'nekod:event-too-large :message "Event exceeds maximum size" :size size :limit (backpressure-config-maximum-event-size *backpressure*)))) (defun check-header-size (size) "Signal http-header-too-large if header exceeds maximum." (when (> size (backpressure-config-maximum-header-size *backpressure*)) (error 'nekod:http-header-too-large :message "Header exceeds maximum size" :size size :limit (backpressure-config-maximum-header-size *backpressure*)))) (defun should-pause-p () "Return T if consumer queue is at high water mark." (>= (backpressure-state-queue-size *bp-state*) (backpressure-config-high-water-mark *backpressure*))) (defun should-resume-p () "Return T if consumer queue has drained to low water mark." (<= (backpressure-state-queue-size *bp-state*) (backpressure-config-low-water-mark *backpressure*))) (defun enqueue-event (event-string) "Attempt to enqueue event. Returns :accepted, :paused, or :dropped." (let ((size (length event-string))) (check-event-size size) (cond ((>= (backpressure-state-queue-size *bp-state*) (backpressure-config-consumer-queue-depth *backpressure*)) (incf (backpressure-state-events-dropped *bp-state*)) :dropped) ((should-pause-p) (setf (backpressure-state-paused *bp-state*) t) (incf (backpressure-state-queue-size *bp-state*)) (incf (backpressure-state-total-bytes-buffered *bp-state*) size) :paused) (t (incf (backpressure-state-queue-size *bp-state*)) (incf (backpressure-state-total-bytes-buffered *bp-state*) size) (when (and (backpressure-state-paused *bp-state*) (should-resume-p)) (setf (backpressure-state-paused *bp-state*) nil)) :accepted)))) (defun dequeue-event () "Mark one event as consumed." (when (> (backpressure-state-queue-size *bp-state*) 0) (decf (backpressure-state-queue-size *bp-state*)) (when (and (backpressure-state-paused *bp-state*) (should-resume-p)) (setf (backpressure-state-paused *bp-state*) nil)) t)) (defun reset-backpressure () "Reset backpressure state." (setf *bp-state* (make-backpressure-state)))