production hardening: semaphore, metrics, admission control (Phase 4)
ober
52955f0e5a8ba464778c49cd43224b6eab46cfb4
--- a/lib/std/fiber.sls +++ b/lib/std/fiber.sls @@ -84,6 +84,13 @@ with-fibers + ;; Semaphore + make-fiber-semaphore + fiber-semaphore? + fiber-semaphore-acquire! + fiber-semaphore-release! + fiber-semaphore-try-acquire! + ;; Low-level primitives for I/O integration (used by std net io) wake-fiber! fiber-gate-set! @@ -1240,4 +1247,77 @@ (fiber-group-fibers group)) (raise exn)))))])) + ;; ========================================================================= + ;; Fiber-aware semaphore + ;; ========================================================================= + ;; + ;; Counting semaphore that parks fibers instead of blocking OS threads. + ;; Used for admission control (max concurrent connections, etc.). + + (define-record-type fiber-semaphore + (fields + (mutable count) + (immutable mutex) + (mutable waiters)) ;; list of (fiber . gate) pairs + (protocol + (lambda (new) + (lambda (max-count) + (new max-count (make-mutex) '()))))) + + ;; Acquire a permit. If none available, park the fiber until one is released. + (define (fiber-semaphore-acquire! sem) + (let ([mx (fiber-semaphore-mutex sem)]) + (mutex-acquire mx) + (let ([c (fiber-semaphore-count sem)]) + (cond + [(> c 0) + ;; Permit available — take it + (fiber-semaphore-count-set! sem (- c 1)) + (mutex-release mx)] + [else + ;; No permits — park the fiber + (let* ([f (fiber-self)] + [gate (box 'channel)]) + (fiber-semaphore-waiters-set! sem + (append (fiber-semaphore-waiters sem) (list (cons f gate)))) + (mutex-release mx) + ;; Park + (fiber-gate-set! f gate) + (set-timer 1) + (spin-until-gate gate) + (fiber-gate-set! f #f))])))) + + ;; Release a permit. If fibers are waiting, wake the first one. + (define (fiber-semaphore-release! sem) + (let ([mx (fiber-semaphore-mutex sem)]) + (mutex-acquire mx) + (let ([waiters (fiber-semaphore-waiters sem)]) + (cond + [(null? waiters) + ;; No waiters — increment count + (fiber-semaphore-count-set! sem (+ (fiber-semaphore-count sem) 1)) + (mutex-release mx)] + [else + ;; Wake the first waiter + (let ([entry (car waiters)]) + (fiber-semaphore-waiters-set! sem (cdr waiters)) + (mutex-release mx) + ;; Open their gate and re-enqueue + (set-box! (cdr entry) 'done) + (wake-fiber! (car entry)))])))) + + ;; Try to acquire without blocking. Returns #t on success, #f if no permits. + (define (fiber-semaphore-try-acquire! sem) + (let ([mx (fiber-semaphore-mutex sem)]) + (mutex-acquire mx) + (let ([c (fiber-semaphore-count sem)]) + (cond + [(> c 0) + (fiber-semaphore-count-set! sem (- c 1)) + (mutex-release mx) + #t] + [else + (mutex-release mx) + #f])))) + ) ;; end library --- a/lib/std/net/fiber-httpd.sls +++ b/lib/std/net/fiber-httpd.sls @@ -29,6 +29,7 @@ (export ;; Server lifecycle fiber-httpd-start + fiber-httpd-start* fiber-httpd-stop! fiber-httpd? fiber-httpd-listen-port @@ -44,7 +45,20 @@ ;; Router make-router router-add! router-dispatch - route-get route-post route-put route-delete) + route-get route-post route-put route-delete + + ;; Metrics / production + fiber-httpd-metrics + httpd-metrics? + httpd-metrics-connections-active + httpd-metrics-connections-total + httpd-metrics-requests-total + httpd-metrics-errors-total + httpd-metrics-start-time + + ;; Middleware + wrap-health-check + wrap-metrics-endpoint) (import (chezscheme) (std fiber) @@ -332,6 +346,43 @@ (define (route-put r path handler) (router-add! r "PUT" path handler)) (define (route-delete r path handler) (router-add! r "DELETE" path handler)) + ;; ========== Metrics ========== + + (define-record-type httpd-metrics + (fields + (mutable connections-active) ;; current open connections + (mutable connections-total) ;; total connections accepted + (mutable requests-total) ;; total requests served + (mutable errors-total) ;; total 5xx responses + (immutable start-time) ;; (current-time) at server start + (immutable metrics-mutex)) + (protocol + (lambda (new) + (lambda () + (new 0 0 0 0 (current-time) (make-mutex)))))) + + (define (metrics-inc-active! m) + (with-mutex (httpd-metrics-metrics-mutex m) + (httpd-metrics-connections-active-set! m + (+ (httpd-metrics-connections-active m) 1)) + (httpd-metrics-connections-total-set! m + (+ (httpd-metrics-connections-total m) 1)))) + + (define (metrics-dec-active! m) + (with-mutex (httpd-metrics-metrics-mutex m) + (httpd-metrics-connections-active-set! m + (max 0 (- (httpd-metrics-connections-active m) 1))))) + + (define (metrics-inc-requests! m) + (with-mutex (httpd-metrics-metrics-mutex m) + (httpd-metrics-requests-total-set! m + (+ (httpd-metrics-requests-total m) 1)))) + + (define (metrics-inc-errors! m) + (with-mutex (httpd-metrics-metrics-mutex m) + (httpd-metrics-errors-total-set! m + (+ (httpd-metrics-errors-total m) 1)))) + ;; ========== Server ========== (define-record-type fiber-httpd @@ -340,21 +391,29 @@ (immutable runtime) (immutable poller) (mutable running?) - (mutable accept-fiber)) + (mutable accept-fiber) + (immutable metrics) + (immutable max-connections) ;; #f = unlimited + (immutable conn-semaphore)) ;; fiber-semaphore or #f (sealed #t)) ;; Connection handler: one fiber per connection, keep-alive loop - (define (handle-connection fd poller handler) + (define (handle-connection fd poller handler metrics) (let loop () (let ([req (read-request fd poller)]) (when req + (metrics-inc-requests! metrics) (let ([resp (guard (exn [#t + (metrics-inc-errors! metrics) (respond-text 500 (if (message-condition? exn) (condition-message exn) "Internal Server Error"))]) (handler req))]) (when (response? resp) + ;; Track 5xx errors + (when (>= (response-status resp) 500) + (metrics-inc-errors! metrics)) (write-response fd poller resp) ;; Keep-alive: check Connection header (let ([conn (request-header req "connection")]) @@ -362,24 +421,42 @@ (loop)))))))) (fiber-tcp-close fd)) - ;; Accept loop + ;; Accept loop with optional admission control (define (accept-loop listen-fd poller handler server) - (guard (exn [#t (void)]) ;; catch cancellation and all errors - (let loop () - (when (fiber-httpd-running? server) - (let ([client-fd (fiber-tcp-accept listen-fd poller)]) - (fiber-spawn* - (lambda () (handle-connection client-fd poller handler)) - "http-conn") - (loop)))))) - - ;; Start the server + (let ([metrics (fiber-httpd-metrics server)] + [sem (fiber-httpd-conn-semaphore server)]) + (guard (exn [#t (void)]) ;; catch cancellation and all errors + (let loop () + (when (fiber-httpd-running? server) + ;; Backpressure: if max-connections set, acquire permit + (when sem (fiber-semaphore-acquire! sem)) + (let ([client-fd (fiber-tcp-accept listen-fd poller)]) + (metrics-inc-active! metrics) + (fiber-spawn* + (lambda () + (guard (exn [#t (void)]) ;; ensure cleanup + (handle-connection client-fd poller handler metrics)) + (metrics-dec-active! metrics) + (when sem (fiber-semaphore-release! sem))) + "http-conn") + (loop))))))) + + ;; Start server with defaults (define (fiber-httpd-start port handler) + (fiber-httpd-start* port handler #f)) + + ;; Start server with max-connections limit. + ;; max-conn: #f = unlimited, integer = max concurrent connections + (define (fiber-httpd-start* port handler max-conn) (let* ([rt (make-fiber-runtime)] - [poller (make-io-poller rt)]) + [poller (make-io-poller rt)] + [metrics (make-httpd-metrics)] + [sem (and max-conn (> max-conn 0) + (make-fiber-semaphore max-conn))]) (io-poller-start! poller) (let-values ([(listen-fd listen-port) (fiber-tcp-listen "0.0.0.0" port)]) - (let ([srv (make-fiber-httpd listen-fd listen-port rt poller #t #f)]) + (let ([srv (make-fiber-httpd listen-fd listen-port rt poller #t #f + metrics max-conn sem)]) ;; Spawn accept loop (let ([af (fiber-spawn rt (lambda () (accept-loop listen-fd poller handler srv)) @@ -403,4 +480,42 @@ ;; Give background thread time to exit (sleep (make-time 'time-duration 150000000 0))) + ;; ========== Middleware helpers ========== + + ;; Wrap a handler to serve /health automatically + (define (wrap-health-check handler server) + (lambda (req) + (if (and (string=? (request-method req) "GET") + (string=? (request-path-only req) "/health")) + (let ([m (fiber-httpd-metrics server)]) + (respond-json 200 + (format "{\"status\":\"ok\",\"connections\":~a,\"requests\":~a}" + (httpd-metrics-connections-active m) + (httpd-metrics-requests-total m)))) + (handler req)))) + + ;; Wrap a handler to serve /metrics in a simple text format + (define (wrap-metrics-endpoint handler server) + (lambda (req) + (if (and (string=? (request-method req) "GET") + (string=? (request-path-only req) "/metrics")) + (let* ([m (fiber-httpd-metrics server)] + [uptime (let ([now (time-second (current-time))] + [start (time-second (httpd-metrics-start-time m))]) + (- now start))]) + (respond-text 200 + (format (string-append + "# Jerboa fiber-httpd metrics\n" + "connections_active ~a\n" + "connections_total ~a\n" + "requests_total ~a\n" + "errors_total ~a\n" + "uptime_seconds ~a\n") + (httpd-metrics-connections-active m) + (httpd-metrics-connections-total m) + (httpd-metrics-requests-total m) + (httpd-metrics-errors-total m) + uptime))) + (handler req)))) + ) ;; end library --- a/tests/test-fiber-httpd.ss +++ b/tests/test-fiber-httpd.ss @@ -35,7 +35,7 @@ (unless val (error 'assert msg))])) ;; Helper: send raw HTTP request over a fiber-aware TCP connection -;; and read the full response. Reads once (sufficient for small responses). +;; and read the full response. Accumulates reads until headers + full body. (define (http-request-raw fd poller method path body) (let* ([body-bv (if body (string->bytevector body (make-transcoder (utf-8-codec))) @@ -56,14 +56,67 @@ ;; Send body if present (when body-bv (fiber-tcp-write fd body-bv (bytevector-length body-bv) poller)) - ;; Read response — single read for small responses, then done - (let ([buf (make-bytevector 16384)]) - (let ([n (fiber-tcp-read fd buf 16384 poller)]) - (if (<= n 0) "" - (bytevector->string - (let ([b (make-bytevector n)]) - (bytevector-copy! buf 0 b 0 n) b) - (make-transcoder (utf-8-codec)))))))) + ;; Read response — accumulate into a growing buffer + (let ([buf (make-bytevector 16384)] + [tmp (make-bytevector 4096)]) + (let loop ([total 0]) + (let ([n (fiber-tcp-read fd tmp (min 4096 (- 16384 total)) poller)]) + (cond + [(<= n 0) + ;; EOF — return accumulated data + (bytevector->string + (let ([b (make-bytevector total)]) + (bytevector-copy! buf 0 b 0 total) b) + (make-transcoder (utf-8-codec)))] + [else + (bytevector-copy! tmp 0 buf total n) + (let ([new-total (+ total n)]) + ;; Check if we have full headers + (let ([hdr-end (find-crlf-crlf buf new-total)]) + (if hdr-end + ;; Parse Content-Length and check if body is complete + (let ([cl (extract-content-length buf hdr-end)]) + (if (>= new-total (+ hdr-end (or cl 0))) + ;; Full response received + (bytevector->string + (let ([b (make-bytevector new-total)]) + (bytevector-copy! buf 0 b 0 new-total) b) + (make-transcoder (utf-8-codec))) + (loop new-total))) + (loop new-total))))])))))) ;; keep reading + +;; Find \r\n\r\n in bytevector, return offset after it (start of body) +(define (find-crlf-crlf buf len) + (let loop ([i 0]) + (cond + [(> (+ i 3) len) #f] + [(and (= (bytevector-u8-ref buf i) 13) + (= (bytevector-u8-ref buf (+ i 1)) 10) + (= (bytevector-u8-ref buf (+ i 2)) 13) + (= (bytevector-u8-ref buf (+ i 3)) 10)) + (+ i 4)] + [else (loop (+ i 1))]))) + +;; Extract Content-Length value from headers in bytevector +(define (extract-content-length buf header-end) + (let* ([hdr-str (bytevector->string + (let ([b (make-bytevector header-end)]) + (bytevector-copy! buf 0 b 0 header-end) b) + (make-transcoder (utf-8-codec)))] + [lc (string-downcase hdr-str)]) + (let ([idx (string-search-helper lc "content-length:" 0)]) + (if idx + (let* ([start (+ idx 15)] ;; skip "content-length:" + [end (string-search-helper hdr-str "\r\n" start)]) + (and end (string->number (string-trim-ws (substring hdr-str start end))))) + 0)))) + +(define (string-trim-ws s) + (let ([len (string-length s)]) + (let ([start (let loop ([i 0]) + (if (and (< i len) (char=? (string-ref s i) #\space)) + (loop (+ i 1)) i))]) + (substring s start len)))) ;; Helper: parse response status code from raw response (define (response-status-code resp) new file mode 100644 --- /dev/null +++ b/tests/test-production.ss @@ -0,0 +1,264 @@ +;;; Tests for Phase 4: Production hardening +;;; Tests semaphore, metrics, health check, and max-connections. + +(import (chezscheme)) +(import (std fiber)) +(import (std net io)) +(import (std net fiber-httpd)) + +(define test-count 0) +(define pass-count 0) + +(define-syntax test + (syntax-rules () + [(_ name body ...) + (begin + (set! test-count (+ test-count 1)) + (guard (exn [#t + (display "FAIL: ") (display name) (newline) + (display " Error: ") + (display (if (message-condition? exn) (condition-message exn) exn)) + (newline)]) + body ... + (set! pass-count (+ pass-count 1)) + (display "PASS: ") (display name) (newline)))])) + +(define-syntax assert-equal + (syntax-rules () + [(_ got expected msg) + (unless (equal? got expected) + (error 'assert msg (list 'got: got 'expected: expected)))])) + +(define-syntax assert-true + (syntax-rules () + [(_ val msg) + (unless val (error 'assert msg))])) + +;; Helper: simple single-read HTTP request +(define (http-get fd poller path) + (let* ([req-str (string-append + "GET " path " HTTP/1.1\r\n" + "Host: localhost\r\nConnection: close\r\n\r\n")] + [req-bv (string->bytevector req-str (make-transcoder (utf-8-codec)))]) + (fiber-tcp-write fd req-bv (bytevector-length req-bv) poller) + (let ([buf (make-bytevector 8192)] + [tmp (make-bytevector 4096)]) + (let loop ([total 0]) + (let ([n (fiber-tcp-read fd tmp (min 4096 (- 8192 total)) poller)]) + (cond + [(<= n 0) + (bytevector->string + (let ([b (make-bytevector total)]) + (bytevector-copy! buf 0 b 0 total) b) + (make-transcoder (utf-8-codec)))] + [else + (bytevector-copy! tmp 0 buf total n) + (loop (+ total n))])))))) + +(define (response-body-text resp) + (let ([idx (string-search-raw resp "\r\n\r\n" 0)]) + (if idx + (substring resp (+ idx 4) (string-length resp)) + ""))) + +(define (response-status-code resp) + (let ([sp1 (string-index-raw resp #\space 0)]) + (and sp1 + (let ([sp2 (string-index-raw resp #\space (+ sp1 1))]) + (and sp2 (string->number (substring resp (+ sp1 1) sp2))))))) + +(define (string-search-raw s needle start) + (let ([slen (string-length s)] + [nlen (string-length needle)]) + (let loop ([i start]) + (cond + [(> (+ i nlen) slen) #f] + [(string=? (substring s i (+ i nlen)) needle) i] + [else (loop (+ i 1))])))) + +(define (string-index-raw s ch start) + (let loop ([i start]) + (cond + [(= i (string-length s)) #f] + [(char=? (string-ref s i) ch) i] + [else (loop (+ i 1))]))) + +;; ========================================================================= +;; Test 1: Fiber semaphore — basic acquire/release +;; ========================================================================= + +(test "semaphore: basic acquire/release" + (let ([rt (make-fiber-runtime 2)] + [sem (make-fiber-semaphore 2)] + [results (make-vector 3 #f)]) + (fiber-spawn rt + (lambda () + (fiber-semaphore-acquire! sem) + (vector-set! results 0 #t) + (fiber-semaphore-release! sem)) + "f1") + (fiber-spawn rt + (lambda () + (fiber-semaphore-acquire! sem) + (vector-set! results 1 #t) + (fiber-semaphore-release! sem)) + "f2") + (fiber-spawn rt + (lambda () + (fiber-semaphore-acquire! sem) + (vector-set! results 2 #t) + (fiber-semaphore-release! sem)) + "f3") + (fiber-runtime-run! rt) + (assert-true (vector-ref results 0) "fiber 1 ran") + (assert-true (vector-ref results 1) "fiber 2 ran") + (assert-true (vector-ref results 2) "fiber 3 ran"))) + +;; ========================================================================= +;; Test 2: Fiber semaphore — try-acquire +;; ========================================================================= + +(test "semaphore: try-acquire" + (let ([rt (make-fiber-runtime 2)] + [sem (make-fiber-semaphore 1)] + [r1 (box #f)] + [r2 (box #f)]) + (fiber-spawn rt + (lambda () + ;; First acquire should succeed + (set-box! r1 (fiber-semaphore-try-acquire! sem)) + ;; Second should fail (only 1 permit) + (set-box! r2 (fiber-semaphore-try-acquire! sem)) + (fiber-semaphore-release! sem)) + "try-fiber") + (fiber-runtime-run! rt) + (assert-true (unbox r1) "first try-acquire succeeds") + (assert-true (not (unbox r2)) "second try-acquire fails"))) + +;; ========================================================================= +;; Test 3: Metrics tracking +;; ========================================================================= + +(test "metrics: connections and requests tracked" + (let* ([handler (lambda (req) (respond-text 200 "ok"))] + [srv (fiber-httpd-start 0 handler)] + [port (fiber-httpd-listen-port srv)] + [m (fiber-httpd-metrics srv)]) + + (sleep (make-time 'time-duration 100000000 0)) + + ;; Send 3 requests + (let ([rt (make-fiber-runtime 2)]) + (with-io-poller rt poller + (do ([i 0 (+ i 1)]) + ((= i 3)) + (fiber-spawn rt + (lambda () + (let ([fd (fiber-tcp-connect "127.0.0.1" port poller)]) + (http-get fd poller "/test") + (fiber-tcp-close fd))) + (string-append "req-" (number->string i)))) + (fiber-runtime-run! rt))) + + ;; Give server time to process + (sleep (make-time 'time-duration 100000000 0)) + + (let ([total (httpd-metrics-connections-total m)] + [reqs (httpd-metrics-requests-total m)]) + (fiber-httpd-stop! srv) + (assert-equal total 3 "3 connections total") + (assert-equal reqs 3 "3 requests total")))) + +;; ========================================================================= +;; Test 4: Health check endpoint +;; ========================================================================= + +(test "health check endpoint" + (let* ([base-handler (lambda (req) (respond-text 200 "app"))] + [srv (fiber-httpd-start 0 base-handler)] + [port (fiber-httpd-listen-port srv)] + [wrapped (wrap-health-check base-handler srv)]) + ;; Note: wrap-health-check returns a new handler, but we need to start + ;; the server with it. For this test, start a new server with wrapped handler. + (fiber-httpd-stop! srv) + + (let* ([srv2 (fiber-httpd-start 0 (wrap-health-check base-handler srv))]) + ;; Actually the srv reference is for the old server — metrics won't match. + ;; Let me just verify the wrap works by calling it directly. + (fiber-httpd-stop! srv2)) + + ;; Test the handler function directly + (let ([health-req (make-request "GET" "/health" "HTTP/1.1" '() #f)] + [app-req (make-request "GET" "/app" "HTTP/1.1" '() #f)]) + ;; Health check response + (let ([resp (wrapped health-req)]) + (assert-equal (response-status resp) 200 "health status 200") + (assert-true (string? (response-body resp)) "health body is string")) + ;; Regular request + (let ([resp (wrapped app-req)]) + (assert-equal (response-body resp) "app" "regular handler works"))))) + +;; ========================================================================= +;; Test 5: Metrics endpoint +;; ========================================================================= + +(test "metrics endpoint" + (let* ([handler (lambda (req) (respond-text 200 "ok"))] + [srv (fiber-httpd-start 0 handler)] + [wrapped (wrap-metrics-endpoint handler srv)]) + (fiber-httpd-stop! srv) + + (let ([metrics-req (make-request "GET" "/metrics" "HTTP/1.1" '() #f)]) + (let ([resp (wrapped metrics-req)]) + (assert-equal (response-status resp) 200 "metrics status 200") + (let ([body (response-body resp)]) + (assert-true (string? body) "metrics body is string") + (assert-true (> (string-length body) 0) "metrics body non-empty")))))) + +;; ========================================================================= +;; Test 6: Max connections (admission control) +;; ========================================================================= + +(test "max-connections admission control" + (let* ([handler (lambda (req) + ;; Slow handler — holds the connection + (fiber-sleep 50) + (respond-text 200 "ok"))] + [srv (fiber-httpd-start* 0 handler 3)] ;; max 3 concurrent + [port (fiber-httpd-listen-port srv)]) + + (sleep (make-time 'time-duration 100000000 0)) + + ;; Send 5 requests — only 3 should be active simultaneously + (let ([rt (make-fiber-runtime 4)] + [results (make-vector 5 #f)]) + (with-io-poller rt poller + (do ([i 0 (+ i 1)]) + ((= i 5)) + (let ([idx i]) + (fiber-spawn rt + (lambda () + (let ([fd (fiber-tcp-connect "127.0.0.1" port poller)]) + (let ([resp (http-get fd poller "/slow")]) + (vector-set! results idx + (= (response-status-code resp) 200))) + (fiber-tcp-close fd))) + (string-append "slow-" (number->string idx))))) + (fiber-runtime-run! rt)) + + ;; All 5 should eventually complete (just not all at once) + (fiber-httpd-stop! srv) + (let ([ok (do ([i 0 (+ i 1)] [c 0 (+ c (if (vector-ref results i) 1 0))]) + ((= i 5) c))]) + (assert-equal ok 5 "all 5 eventually served"))))) + +;; ========================================================================= +;; Summary +;; ========================================================================= +(newline) +(display "=========================================") (newline) +(display "Results: ") (display pass-count) (display "/") +(display test-count) (display " passed") (newline) +(display "=========================================") (newline) +(when (< pass-count test-count) + (exit 1))