nekomata / socket /backpressure.lisp
SNAPKITTYWEST's picture
push from SNAPKITTYWEST/nekomata
3b70664 verified
Raw History Blame Contribute Delete
3.17 kB
(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)))