Download socket/backpressure.lisp from Snapkitty/nekomata: direct link, hf CLI and curl.
- Browser
- Download file 3.17 kB
-
https://huggingface.co/Snapkitty/nekomata/resolve/main/socket/backpressure.lisp
- Command line
-
hf download hf://Snapkitty/nekomata/socket/backpressure.lisp
-
curl -L -o backpressure.lisp https://huggingface.co/Snapkitty/nekomata/resolve/main/socket/backpressure.lisp
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))) | |