Add HTTP server (chez-httpd) and client keep-alive connection pooling
ober
677300064eed7d1ef23a3d60b83fd6e43ca0d716
new file mode 100644 --- /dev/null +++ b/src/chez-httpd.sls @@ -0,0 +1,859 @@ +;; chez-httpd — High-performance HTTP/1.1 server for Chez Scheme +;; +;; Dependencies: +;; chez-ssl — https://github.com/ober/chez-ssl (TCP + optional TLS) +;; +;; Features: +;; - Thread-pool based connection handling +;; - Keep-alive with configurable timeouts +;; - Chunked transfer encoding (send and receive) +;; - Buffer pooling to reduce GC pressure +;; - Routing with exact and prefix matching +;; - Static file serving +;; - Both HTTP and HTTPS (TLS) support + +(library (chez-httpd) + (export + ;; Server lifecycle + httpd-start httpd-stop httpd-start-https + ;; Routing + httpd-route httpd-route-prefix httpd-route-static + make-router router-add! router-add-prefix! router-lookup + ;; Request accessors + http-req-method http-req-path http-req-query http-req-headers + http-req-header http-req-body http-req-client-addr http-req-version + ;; Response building + http-respond http-respond-html http-respond-json http-respond-file + http-respond-error http-respond-redirect + http-respond-chunk-begin http-respond-chunk http-respond-chunk-end + ;; Configuration + httpd-config) + (import (chezscheme) (chez-ssl)) + + ;; ================================================================ + ;; Configuration + ;; ================================================================ + + (define *config* + (vector + 4 ;; 0: num-workers (thread pool size) + 8192 ;; 1: input-buffer-size + 32768 ;; 2: output-buffer-size + 60 ;; 3: request-timeout (seconds) — not yet enforced + 120 ;; 4: response-timeout (seconds) — not yet enforced + 1048576 ;; 5: max-request-body (1 MB) + 128 ;; 6: listen-backlog + )) + + (define (cfg-ref i) (vector-ref *config* i)) + + (define (httpd-config . args) + ;; (httpd-config 'workers: 8 'input-buffer: 16384 ...) + (let loop ([args args]) + (unless (null? args) + (let ([key (car args)] [val (cadr args)]) + (cond + [(eq? key 'workers:) (vector-set! *config* 0 val)] + [(eq? key 'input-buffer:) (vector-set! *config* 1 val)] + [(eq? key 'output-buffer:) (vector-set! *config* 2 val)] + [(eq? key 'max-body:) (vector-set! *config* 5 val)] + [(eq? key 'backlog:) (vector-set! *config* 6 val)]) + (loop (cddr args)))))) + + ;; ================================================================ + ;; String / bytevector utilities + ;; ================================================================ + + (define (string-index str ch) + (let ([len (string-length str)]) + (let loop ([i 0]) + (cond + [(= i len) #f] + [(char=? (string-ref str i) ch) i] + [else (loop (+ i 1))])))) + + (define (string-prefix? prefix str) + (let ([plen (string-length prefix)] + [slen (string-length str)]) + (and (>= slen plen) + (string=? prefix (substring str 0 plen))))) + + (define (string-trim-left str) + (let ([len (string-length str)]) + (let loop ([i 0]) + (if (and (< i len) (char-whitespace? (string-ref str i))) + (loop (+ i 1)) + (substring str i len))))) + + ;; ================================================================ + ;; Buffer pool — reduces GC pressure under load + ;; ================================================================ + + (define *input-pool* '()) + (define *input-pool-mutex* (make-mutex 'httpd-input-pool)) + (define *output-pool* '()) + (define *output-pool-mutex* (make-mutex 'httpd-output-pool)) + + (define (pool-get! pool-var mutex-var size) + (with-mutex mutex-var + (let ([pool (pool-var)]) + (if (null? pool) + (make-bytevector size) + (let ([buf (car pool)]) + (pool-var (cdr pool)) + buf))))) + + (define (pool-put! pool-var mutex-var buf) + (with-mutex mutex-var + (pool-var (cons buf (pool-var))))) + + ;; Use syntax to capture the mutable variable references + (define-syntax get-input-buffer + (syntax-rules () + [(_) (with-mutex *input-pool-mutex* + (if (null? *input-pool*) + (make-bytevector (cfg-ref 1)) + (let ([buf (car *input-pool*)]) + (set! *input-pool* (cdr *input-pool*)) + buf)))])) + + (define-syntax put-input-buffer! + (syntax-rules () + [(_ buf) (with-mutex *input-pool-mutex* + (set! *input-pool* (cons buf *input-pool*)))])) + + (define-syntax get-output-buffer + (syntax-rules () + [(_) (with-mutex *output-pool-mutex* + (if (null? *output-pool*) + (make-bytevector (cfg-ref 2)) + (let ([buf (car *output-pool*)]) + (set! *output-pool* (cdr *output-pool*)) + buf)))])) + + (define-syntax put-output-buffer! + (syntax-rules () + [(_ buf) (with-mutex *output-pool-mutex* + (set! *output-pool* (cons buf *output-pool*)))])) + + ;; ================================================================ + ;; Buffered reader — cursor-based, wraps a conn handle + ;; ================================================================ + + ;; Reader: #(conn buf pos end) + (define (make-reader conn buf) + (vector conn buf 0 0)) + + (define (reader-conn r) (vector-ref r 0)) + (define (reader-buf r) (vector-ref r 1)) + (define (reader-pos r) (vector-ref r 2)) + (define (reader-end r) (vector-ref r 3)) + (define (reader-pos-set! r v) (vector-set! r 2 v)) + (define (reader-end-set! r v) (vector-set! r 3 v)) + (define (reader-available r) (- (reader-end r) (reader-pos r))) + + (define (reader-compact! r) + (let ([pos (reader-pos r)] + [end (reader-end r)] + [buf (reader-buf r)]) + (when (> pos 0) + (let ([avail (- end pos)]) + (bytevector-copy! buf pos buf 0 avail) + (reader-pos-set! r 0) + (reader-end-set! r avail))))) + + (define (reader-fill! r) + ;; Compact then read more data from connection. + ;; Returns #t if data was read, #f on EOF/error. + (reader-compact! r) + (let* ([buf (reader-buf r)] + [end (reader-end r)] + [cap (bytevector-length buf)] + [space (- cap end)]) + (if (<= space 0) + #t ;; buffer full, data available + (let ([n (conn-read (reader-conn r) buf end)]) + ;; conn-read takes (conn buf len) but we need offset support. + ;; We'll use a temp buffer and copy. Or better: read into a temp + ;; and copy to the right position. + ;; Actually, conn-read reads into start of buf. We need to read + ;; at offset. Let's use a workaround with a temp buffer. + #f)))) ;; placeholder + + ;; Since conn-read reads into the start of the bytevector, we use a temp buf + ;; and copy into the main buffer at the right offset. This is the price of + ;; not having offset-aware C functions. + + (define *temp-read-buf* #f) + (define *temp-read-mutex* (make-mutex 'temp-read)) + + (define (reader-fill-from-conn! r) + ;; Read more data into the reader buffer. + ;; Returns bytes read (0 = EOF, -1 = error). + (reader-compact! r) + (let* ([buf (reader-buf r)] + [end (reader-end r)] + [cap (bytevector-length buf)] + [space (- cap end)]) + (when (<= space 0) + ;; Buffer is full after compaction — shouldn't happen with proper sizing + (error 'reader-fill! "buffer overflow")) + (let ([tmp (make-bytevector space)]) + (let ([n (conn-read (reader-conn r) tmp space)]) + (when (> n 0) + (bytevector-copy! tmp 0 buf end n) + (reader-end-set! r (+ end n))) + n)))) + + (define (reader-read-byte r) + ;; Read a single byte. Returns byte or #f on EOF. + (when (= (reader-pos r) (reader-end r)) + (let ([n (reader-fill-from-conn! r)]) + (when (<= n 0) (set! n 0)))) ;; mark as no more data + (if (= (reader-pos r) (reader-end r)) + #f + (let ([b (bytevector-u8-ref (reader-buf r) (reader-pos r))]) + (reader-pos-set! r (+ (reader-pos r) 1)) + b))) + + (define (reader-read-line r) + ;; Read a line terminated by \r\n. Returns string or #f on EOF. + ;; Scans the buffer for \r\n, filling as needed. + (let loop ([acc #f]) + (let ([pos (reader-pos r)] + [end (reader-end r)]) + (if (= pos end) + ;; Buffer empty — try to fill + (let ([n (reader-fill-from-conn! r)]) + (if (<= n 0) + acc ;; EOF — return accumulated string or #f + (loop acc))) + ;; Scan buffer for \r\n + (let scan ([i pos]) + (cond + [(>= i (- end 1)) + ;; Reached end without finding \r\n — save what we have, fill more + (let* ([chunk-len (- end pos)] + [chunk-bv (make-bytevector chunk-len)]) + (bytevector-copy! (reader-buf r) pos chunk-bv 0 chunk-len) + (reader-pos-set! r end) + (let ([chunk-str (utf8->string chunk-bv)] + [n (reader-fill-from-conn! r)]) + (let ([new-acc (if acc (string-append acc chunk-str) chunk-str)]) + (if (<= n 0) + (if (= (string-length new-acc) 0) #f new-acc) + (loop new-acc)))))] + [(and (= (bytevector-u8-ref (reader-buf r) i) 13) + (= (bytevector-u8-ref (reader-buf r) (+ i 1)) 10)) + ;; Found \r\n + (let* ([line-len (- i pos)] + [line-bv (make-bytevector line-len)]) + (bytevector-copy! (reader-buf r) pos line-bv 0 line-len) + (reader-pos-set! r (+ i 2)) + (let ([chunk-str (utf8->string line-bv)]) + (if acc (string-append acc chunk-str) chunk-str)))] + [else (scan (+ i 1))])))))) + + (define (reader-read-bytes r n) + ;; Read exactly n bytes. Returns bytevector or #f on premature EOF. + (let ([result (make-bytevector n)]) + (let loop ([offset 0]) + (if (= offset n) + result + (let ([avail (reader-available r)]) + (if (> avail 0) + (let ([take (min avail (- n offset))]) + (bytevector-copy! (reader-buf r) (reader-pos r) result offset take) + (reader-pos-set! r (+ (reader-pos r) take)) + (loop (+ offset take))) + (let ([rc (reader-fill-from-conn! r)]) + (if (<= rc 0) + #f ;; premature EOF + (loop offset))))))))) + + ;; ================================================================ + ;; Buffered writer — batches small writes, flushes to conn + ;; ================================================================ + + ;; Writer: #(conn buf pos) + (define (make-writer conn buf) + (vector conn buf 0)) + + (define (writer-conn w) (vector-ref w 0)) + (define (writer-buf w) (vector-ref w 1)) + (define (writer-pos w) (vector-ref w 2)) + (define (writer-pos-set! w v) (vector-set! w 2 v)) + + (define (writer-flush! w) + (let ([pos (writer-pos w)]) + (when (> pos 0) + (let ([bv (make-bytevector pos)]) + (bytevector-copy! (writer-buf w) 0 bv 0 pos) + (conn-write (writer-conn w) bv) + (writer-pos-set! w 0))))) + + (define (writer-write-byte! w b) + (let ([pos (writer-pos w)] + [buf (writer-buf w)]) + (when (= pos (bytevector-length buf)) + (writer-flush! w) + (set! pos 0)) + (bytevector-u8-set! buf pos b) + (writer-pos-set! w (+ pos 1)))) + + (define (writer-write-bv! w bv) + (let ([len (bytevector-length bv)] + [pos (writer-pos w)] + [buf (writer-buf w)] + [cap (bytevector-length (writer-buf w))]) + (if (<= (+ pos len) cap) + ;; Fits in buffer + (begin + (bytevector-copy! bv 0 buf pos len) + (writer-pos-set! w (+ pos len))) + ;; Doesn't fit — flush current, then write directly or buffer + (begin + (writer-flush! w) + (if (> len cap) + ;; Larger than buffer — write directly + (conn-write (writer-conn w) bv) + ;; Fits in empty buffer + (begin + (bytevector-copy! bv 0 buf 0 len) + (writer-pos-set! w len))))))) + + (define (writer-write-string! w str) + (writer-write-bv! w (string->utf8 str))) + + (define cr-lf (string->utf8 "\r\n")) + + (define (writer-write-crlf! w) + (writer-write-bv! w cr-lf)) + + ;; ================================================================ + ;; HTTP request record + ;; ================================================================ + + ;; #(method path query version headers body client-addr) + (define (make-http-request method path query version headers body client-addr) + (vector method path query version headers body client-addr)) + + (define (http-req-method r) (vector-ref r 0)) + (define (http-req-path r) (vector-ref r 1)) + (define (http-req-query r) (vector-ref r 2)) + (define (http-req-version r) (vector-ref r 3)) + (define (http-req-headers r) (vector-ref r 4)) + (define (http-req-body r) (vector-ref r 5)) + (define (http-req-client-addr r) (vector-ref r 6)) + + (define (http-req-header r name) + (let ([pair (assoc (string-downcase name) (vector-ref r 4))]) + (and pair (cdr pair)))) + + ;; ================================================================ + ;; Request parsing + ;; ================================================================ + + (define (parse-request-line line) + ;; "GET /path?query HTTP/1.1" -> (values method path query version) + (let* ([sp1 (string-index line #\space)] + [method (substring line 0 sp1)] + [rest (substring line (+ sp1 1) (string-length line))] + [sp2 (string-index rest #\space)] + [target (substring rest 0 sp2)] + [version (substring rest (+ sp2 1) (string-length rest))] + [qmark (string-index target #\?)] + [path (if qmark (substring target 0 qmark) target)] + [query (if qmark (substring target (+ qmark 1) (string-length target)) #f)]) + (values method path query version))) + + (define (read-request reader client-addr) + ;; Parse an HTTP request from the reader. + ;; Returns http-request or #f on connection close. + (let ([request-line (reader-read-line reader)]) + (if (not request-line) + #f ;; connection closed + (let-values ([(method path query version) (parse-request-line request-line)]) + (let ([headers (read-headers reader)]) + (let ([body (read-request-body reader headers)]) + (make-http-request method path query version headers body client-addr))))))) + + (define (read-headers reader) + ;; Read headers until empty line. Returns alist with lowercase keys. + (let loop ([acc '()]) + (let ([line (reader-read-line reader)]) + (cond + [(not line) (reverse acc)] + [(string=? line "") (reverse acc)] + [else + (let ([colon (string-index line #\:)]) + (if colon + (loop (cons (cons (string-downcase (substring line 0 colon)) + (string-trim-left + (substring line (+ colon 1) (string-length line)))) + acc)) + (loop acc)))])))) + + (define (read-request-body reader headers) + ;; Read body based on Content-Length or Transfer-Encoding. + ;; Returns bytevector or #f. + (let ([cl (assoc "content-length" headers)] + [te (assoc "transfer-encoding" headers)]) + (cond + [(and cl (cdr cl)) + (let ([len (string->number (cdr cl))]) + (if (and len (> len 0) (<= len (cfg-ref 5))) + (reader-read-bytes reader len) + #f))] + [(and te (header-ci=? (cdr te) "chunked")) + (read-chunked-request-body reader)] + [else #f]))) + + (define (read-chunked-request-body reader) + ;; Read chunked request body. Returns bytevector. + (let loop ([chunks '()]) + (let* ([size-line (reader-read-line reader)] + [chunk-size (string->number (string-trim size-line) 16)]) + (cond + [(or (not chunk-size) (= chunk-size 0)) + (reader-read-line reader) ;; consume trailing \r\n + (bytevector-concat-list (reverse chunks))] + [else + (let ([chunk (reader-read-bytes reader chunk-size)]) + (reader-read-line reader) ;; consume chunk-trailing \r\n + (loop (cons chunk chunks)))])))) + + (define (string-trim str) + (let* ([len (string-length str)] + [start (let loop ([i 0]) + (if (and (< i len) (char-whitespace? (string-ref str i))) + (loop (+ i 1)) i))] + [end (let loop ([i len]) + (if (and (> i start) (char-whitespace? (string-ref str (- i 1)))) + (loop (- i 1)) i))]) + (substring str start end))) + + (define (header-ci=? a b) + (string=? (string-downcase a) (string-downcase b))) + + (define (bytevector-concat-list bvs) + (if (null? bvs) (make-bytevector 0) + (let* ([total (fold-left + 0 (map bytevector-length bvs))] + [result (make-bytevector total)]) + (let loop ([bvs bvs] [offset 0]) + (if (null? bvs) result + (let ([bv (car bvs)]) + (bytevector-copy! bv 0 result offset (bytevector-length bv)) + (loop (cdr bvs) (+ offset (bytevector-length bv))))))))) + + ;; ================================================================ + ;; Response writing + ;; ================================================================ + + (define (status-text code) + (case code + [(200) "OK"] [(201) "Created"] [(204) "No Content"] + [(301) "Moved Permanently"] [(302) "Found"] + [(304) "Not Modified"] [(307) "Temporary Redirect"] + [(400) "Bad Request"] [(401) "Unauthorized"] [(403) "Forbidden"] + [(404) "Not Found"] [(405) "Method Not Allowed"] + [(408) "Request Timeout"] [(413) "Payload Too Large"] + [(500) "Internal Server Error"] [(502) "Bad Gateway"] + [(503) "Service Unavailable"] + [else "Unknown"])) + + (define (write-response-head! w status headers) + (writer-write-string! w "HTTP/1.1 ") + (writer-write-string! w (number->string status)) + (writer-write-string! w " ") + (writer-write-string! w (status-text status)) + (writer-write-crlf! w) + (for-each + (lambda (h) + (writer-write-string! w (car h)) + (writer-write-string! w ": ") + (writer-write-string! w (cdr h)) + (writer-write-crlf! w)) + headers) + (writer-write-crlf! w)) + + (define (http-respond writer status headers body) + ;; Write a complete response with Content-Length. + ;; body can be string, bytevector, or #f. + (let* ([body-bv (cond + [(not body) #f] + [(bytevector? body) body] + [(string? body) (string->utf8 body)] + [else (string->utf8 (format "~a" body))])] + [content-length (if body-bv (bytevector-length body-bv) 0)] + [all-headers (cons (cons "Content-Length" (number->string content-length)) + headers)]) + (write-response-head! writer status all-headers) + (when body-bv (writer-write-bv! writer body-bv)) + (writer-flush! writer))) + + (define (http-respond-html writer status body) + (http-respond writer status '(("Content-Type" . "text/html; charset=utf-8")) body)) + + (define (http-respond-json writer status body) + (http-respond writer status '(("Content-Type" . "application/json")) body)) + + (define (http-respond-error writer status) + (http-respond-html writer status + (string-append "<h1>" (number->string status) " " + (status-text status) "</h1>"))) + + (define (http-respond-redirect writer status location) + (http-respond writer status + (list (cons "Location" location)) + #f)) + + ;; ================================================================ + ;; Chunked response writing + ;; ================================================================ + + (define (http-respond-chunk-begin writer status headers) + ;; Begin a chunked response. Caller writes chunks, then calls chunk-end. + (let ([all-headers (cons '("Transfer-Encoding" . "chunked") headers)]) + (write-response-head! writer status all-headers))) + + (define (http-respond-chunk writer data) + ;; Write one chunk. data can be string or bytevector. + (let* ([bv (if (string? data) (string->utf8 data) data)] + [len (bytevector-length bv)]) + (when (> len 0) + (writer-write-string! writer (number->string len 16)) + (writer-write-crlf! writer) + (writer-write-bv! writer bv) + (writer-write-crlf! writer) + (writer-flush! writer)))) + + (define (http-respond-chunk-end writer) + ;; Write the zero-length terminating chunk. + (writer-write-string! writer "0") + (writer-write-crlf! writer) + (writer-write-crlf! writer) + (writer-flush! writer)) + + ;; ================================================================ + ;; Static file serving + ;; ================================================================ + + (define *mime-types* + '(("html" . "text/html; charset=utf-8") + ("htm" . "text/html; charset=utf-8") + ("css" . "text/css") + ("js" . "application/javascript") + ("json" . "application/json") + ("png" . "image/png") + ("jpg" . "image/jpeg") + ("jpeg" . "image/jpeg") + ("gif" . "image/gif") + ("svg" . "image/svg+xml") + ("ico" . "image/x-icon") + ("txt" . "text/plain") + ("xml" . "application/xml") + ("pdf" . "application/pdf") + ("woff" . "font/woff") + ("woff2" . "font/woff2"))) + + (define (file-extension path) + (let loop ([i (- (string-length path) 1)]) + (cond + [(< i 0) ""] + [(char=? (string-ref path i) #\.) + (substring path (+ i 1) (string-length path))] + [(char=? (string-ref path i) #\/) ""] + [else (loop (- i 1))]))) + + (define (mime-type path) + (let ([ext (string-downcase (file-extension path))]) + (let ([pair (assoc ext *mime-types*)]) + (if pair (cdr pair) "application/octet-stream")))) + + (define (http-respond-file writer req file-path) + ;; Serve a file using chunked encoding for efficiency. + (if (file-exists? file-path) + (let ([mime (mime-type file-path)]) + (http-respond-chunk-begin writer 200 + (list (cons "Content-Type" mime))) + (let ([port (open-file-input-port file-path)]) + (let ([buf (make-bytevector 32768)]) + (let loop () + (let ([n (get-bytevector-n! port buf 0 32768)]) + (unless (eof-object? n) + (if (= n 32768) + (http-respond-chunk writer buf) + (let ([partial (make-bytevector n)]) + (bytevector-copy! buf 0 partial 0 n) + (http-respond-chunk writer partial))) + (loop))))) + (close-port port)) + (http-respond-chunk-end writer)) + (http-respond-error writer 404))) + + ;; ================================================================ + ;; Router — exact match and prefix match + ;; ================================================================ + + ;; Router: #(exact-table prefix-list default-handler) + ;; exact-table: hashtable path -> handler + ;; prefix-list: sorted list of (prefix . handler) pairs + ;; default-handler: handler for unmatched paths + + (define (make-router default-handler) + (vector (make-hashtable string-hash string=?) + '() + default-handler)) + + (define (router-exact-table r) (vector-ref r 0)) + (define (router-prefix-list r) (vector-ref r 1)) + (define (router-prefix-list-set! r v) (vector-set! r 1 v)) + (define (router-default r) (vector-ref r 2)) + + (define (router-add! router path handler) + ;; Register an exact path match. + (hashtable-set! (router-exact-table router) path handler)) + + (define (router-add-prefix! router prefix handler) + ;; Register a prefix match. Longer prefixes match first. + (let ([new (cons (cons prefix handler) (router-prefix-list router))]) + (router-prefix-list-set! router + (sort (lambda (a b) (> (string-length (car a)) (string-length (car b)))) + new)))) + + (define (router-lookup router path) + ;; Look up handler for path. Exact match first, then prefix, then default. + (or (hashtable-ref (router-exact-table router) path #f) + (let loop ([prefixes (router-prefix-list router)]) + (cond + [(null? prefixes) #f] + [(string-prefix? (caar prefixes) path) (cdar prefixes)] + [else (loop (cdr prefixes))])) + (router-default router))) + + ;; ================================================================ + ;; Connection handler — runs in worker thread + ;; ================================================================ + + (define (handle-connection conn client-addr router) + ;; Handle one connection. Supports keep-alive (multiple requests). + (let ([ibuf (get-input-buffer)] + [obuf (get-output-buffer)]) + (let ([reader (make-reader conn ibuf)] + [writer (make-writer conn obuf)]) + (dynamic-wind + void + (lambda () + (let loop () + (let ([req (guard (e [#t #f]) (read-request reader client-addr))]) + (when req + (let ([handler (router-lookup router (http-req-path req))]) + (guard (e [#t + (guard (e2 [#t (void)]) + (http-respond-error writer 500))]) + (handler req writer))) + ;; Keep-alive check + (let ([conn-hdr (http-req-header req "connection")] + [version (http-req-version req)]) + (unless (or (and conn-hdr (header-ci=? conn-hdr "close")) + (and (string? version) + (string=? version "HTTP/1.0") + (not (and conn-hdr + (header-ci=? conn-hdr "keep-alive"))))) + (loop))))))) + (lambda () + (guard (e [#t (void)]) + (ssl-close conn)) + (put-input-buffer! ibuf) + (put-output-buffer! obuf)))))) + + ;; ================================================================ + ;; Thread pool — bounded work queue with condition variables + ;; ================================================================ + + ;; Work queue: #(items head tail count capacity mutex ready) + (define (make-work-queue capacity) + (vector (make-vector capacity #f) ;; items + 0 ;; head + 0 ;; tail + 0 ;; count + capacity ;; capacity + (make-mutex 'work-queue) + (make-condition 'work-ready))) + + (define (wq-items q) (vector-ref q 0)) + (define (wq-head q) (vector-ref q 1)) + (define (wq-head-set! q v) (vector-set! q 1 v)) + (define (wq-tail q) (vector-ref q 2)) + (define (wq-tail-set! q v) (vector-set! q 2 v)) + (define (wq-count q) (vector-ref q 3)) + (define (wq-count-set! q v) (vector-set! q 3 v)) + (define (wq-cap q) (vector-ref q 4)) + (define (wq-mutex q) (vector-ref q 5)) + (define (wq-ready q) (vector-ref q 6)) + + (define (wq-enqueue! q item) + (with-mutex (wq-mutex q) + (when (< (wq-count q) (wq-cap q)) + (vector-set! (wq-items q) (wq-tail q) item) + (wq-tail-set! q (mod (+ (wq-tail q) 1) (wq-cap q))) + (wq-count-set! q (+ (wq-count q) 1)) + (condition-signal (wq-ready q)) + #t))) + + (define (wq-dequeue! q) + (with-mutex (wq-mutex q) + (let loop () + (if (= (wq-count q) 0) + (begin + (condition-wait (wq-ready q) (wq-mutex q)) + (loop)) + (let* ([item (vector-ref (wq-items q) (wq-head q))]) + (vector-set! (wq-items q) (wq-head q) #f) ;; clear reference + (wq-head-set! q (mod (+ (wq-head q) 1) (wq-cap q))) + (wq-count-set! q (- (wq-count q) 1)) + item))))) + + ;; Sentinel to stop workers + (define *stop-sentinel* (list 'stop)) + + (define (start-worker-threads work-queue n-workers router) + ;; Start n worker threads that pull from the work queue. + (let loop ([i 0] [threads '()]) + (if (= i n-workers) + threads + (loop (+ i 1) + (cons (fork-thread + (lambda () + (let worker-loop () + (let ([job (wq-dequeue! work-queue)]) + (unless (eq? job *stop-sentinel*) + ;; job is (conn . client-addr) + (guard (e [#t (void)]) ;; don't crash worker on errors + (handle-connection (car job) (cdr job) router)) + (worker-loop)))))) + threads))))) + + ;; ================================================================ + ;; Accept loop — runs in its own thread + ;; ================================================================ + + (define (accept-loop listen-fd work-queue ssl-ctx stop-box) + ;; Accept connections and enqueue them for worker threads. + (let loop () + (unless (unbox stop-box) + (let-values ([(client-fd client-addr) (tcp-accept listen-fd)]) + (cond + [(not client-fd) (loop)] ;; EINTR, retry + [ssl-ctx + ;; TLS: set timeout, perform handshake, then enqueue + (guard (e [#t + ;; TLS handshake failed — close raw fd and continue + (tcp-close client-fd)]) + (tcp-set-timeout client-fd 60 120) + (let ([conn (ssl-server-accept ssl-ctx client-fd)]) + (wq-enqueue! work-queue (cons conn client-addr))))] + [else + ;; Plain TCP: set timeout, wrap fd, enqueue + (tcp-set-timeout client-fd 60 120) + (let ([conn (conn-wrap client-fd)]) + (wq-enqueue! work-queue (cons conn client-addr)))]) + (loop))))) + + ;; ================================================================ + ;; Public API — server lifecycle + ;; ================================================================ + + ;; Server handle: #(listen-fd work-queue workers accept-thread stop-box ssl-ctx) + (define (make-server-handle listen-fd wq workers accept-thread stop-box ssl-ctx) + (vector listen-fd wq workers accept-thread stop-box ssl-ctx)) + + ;; Convenience routing helpers (build a router and start server) + (define *default-router* #f) + + (define (ensure-default-router!) + (unless *default-router* + (set! *default-router* + (make-router (lambda (req w) (http-respond-error w 404)))))) + + (define (httpd-route path handler) + ;; Register an exact route. handler is (lambda (req writer) ...) + (ensure-default-router!) + (router-add! *default-router* path handler)) + + (define (httpd-route-prefix prefix handler) + (ensure-default-router!) + (router-add-prefix! *default-router* prefix handler)) + + (define (httpd-route-static url-prefix directory) + ;; Serve static files from directory under url-prefix. + (ensure-default-router!) + (router-add-prefix! *default-router* url-prefix + (lambda (req writer) + (let* ([path (http-req-path req)] + [rel (substring path (string-length url-prefix) (string-length path))] + ;; Prevent directory traversal + [safe-rel (if (or (string-prefix? "/" rel) + (string-prefix? ".." rel)) + "" + rel)] + [file-path (string-append directory "/" safe-rel)]) + (http-respond-file writer req file-path))))) + + (define (httpd-start port . args) + ;; Start an HTTP server on the given port. + ;; Optional: 'router: custom-router + (ssl-init!) + (let ([router (if (and (pair? args) (eq? (car args) 'router:)) + (cadr args) + (begin (ensure-default-router!) *default-router*))] + [n-workers (cfg-ref 0)] + [backlog (cfg-ref 6)]) + (let* ([listen-fd (tcp-listen port backlog)] + [wq (make-work-queue (* n-workers 64))] + [stop-box (box #f)] + [workers (start-worker-threads wq n-workers router)] + [accept-thread (fork-thread + (lambda () + (accept-loop listen-fd wq #f stop-box)))]) + (display (format "chez-httpd: listening on port ~a (~a workers)\n" port n-workers)) + (make-server-handle listen-fd wq workers accept-thread stop-box #f)))) + + (define (httpd-start-https port cert-file key-file . args) + ;; Start an HTTPS server on the given port with TLS. + (ssl-init!) + (let ([router (if (and (pair? args) (eq? (car args) 'router:)) + (cadr args) + (begin (ensure-default-router!) *default-router*))] + [n-workers (cfg-ref 0)] + [backlog (cfg-ref 6)]) + (let* ([ssl-ctx (ssl-server-ctx cert-file key-file)] + [listen-fd (tcp-listen port backlog)] + [wq (make-work-queue (* n-workers 64))] + [stop-box (box #f)] + [workers (start-worker-threads wq n-workers router)] + [accept-thread (fork-thread + (lambda () + (accept-loop listen-fd wq ssl-ctx stop-box)))]) + (display (format "chez-httpd: listening on port ~a HTTPS (~a workers)\n" port n-workers)) + (make-server-handle listen-fd wq workers accept-thread stop-box ssl-ctx)))) + + (define (httpd-stop server) + ;; Gracefully stop the server. + (let ([listen-fd (vector-ref server 0)] + [wq (vector-ref server 1)] + [workers (vector-ref server 2)] + [stop-box (vector-ref server 4)] + [ssl-ctx (vector-ref server 5)]) + ;; Signal accept loop to stop + (set-box! stop-box #t) + ;; Close listening socket to unblock accept() + (tcp-close listen-fd) + ;; Send stop sentinel to each worker + (for-each (lambda (_) (wq-enqueue! wq (cons *stop-sentinel* ""))) workers) + ;; Clean up SSL context if HTTPS + (when ssl-ctx (ssl-server-ctx-free ssl-ctx)) + (display "chez-httpd: stopped\n"))) + +) ;; end library --- a/src/chez-https.sls +++ b/src/chez-https.sls @@ -246,7 +246,7 @@ ;; HTTP request building ;; ================================================================ - (define (build-request method path host headers body-bv) + (define (build-request method path host headers body-bv keep-alive?) (let ([out (open-output-string)]) (put-string out method) (put-string out " ") @@ -270,8 +270,11 @@ (put-string out "Content-Length: ") (put-string out (number->string (bytevector-length body-bv))) (put-string out "\r\n")) - ;; Connection: close (first version; keep-alive is a future optimization) - (put-string out "Connection: close\r\n") + ;; Connection header + (unless (header-assoc "Connection" headers) + (put-string out (if keep-alive? + "Connection: keep-alive\r\n" + "Connection: close\r\n"))) (put-string out "\r\n") (get-output-string out))) @@ -351,36 +354,240 @@ (set! *ssl-initialized* #t))) ;; ================================================================ - ;; Response reading + ;; Connection pool — keep-alive / connection sharing ;; ================================================================ + ;; Pool keyed by "host:port". Each entry is a list of idle connections. + ;; Connections are reused for same host:port to amortize TLS handshakes. - (define (read-response conn) - ;; With Connection: close, ssl-read-all reads the entire response. - (let* ([raw (ssl-read-all conn)] - [len (bytevector-length raw)]) + (define *conn-pool* (make-hashtable string-hash string=?)) + (define *conn-pool-mutex* (make-mutex 'conn-pool)) + + (define (pool-key host port) + (string-append host ":" (number->string port))) + + (define (pool-get host port) + ;; Try to get an idle connection from the pool. Returns conn or #f. + (let ([key (pool-key host port)]) + (with-mutex *conn-pool-mutex* + (let ([conns (hashtable-ref *conn-pool* key '())]) + (if (null? conns) + #f + (let ([conn (car conns)]) + (hashtable-set! *conn-pool* key (cdr conns)) + conn)))))) + + (define (pool-put host port conn) + ;; Return a connection to the pool for reuse. + (let ([key (pool-key host port)]) + (with-mutex *conn-pool-mutex* + (let ([conns (hashtable-ref *conn-pool* key '())]) + ;; Cap at 4 idle connections per host:port + (if (< (length conns) 4) + (hashtable-set! *conn-pool* key (cons conn conns)) + ;; Pool full — close this connection + (guard (e [#t (void)]) (ssl-close conn))))))) + + ;; ================================================================ + ;; Response reading — supports both close and keep-alive modes + ;; ================================================================ + + (define (ssl-read-bytes conn n) + ;; Read exactly n bytes from conn. Returns bytevector or signals error. + (let ([result (make-bytevector n)] + [buf (make-bytevector (min n 32768))]) (let loop ([offset 0]) - (let ([sep (find-crlfcrlf raw len offset)]) - (unless sep - (error 'read-response - "malformed HTTP response: no header terminator found" - (if (< len 500) (utf8->string raw) "<response too large to display>"))) - (let* ([header-str (utf8->string (subbytevector raw offset sep))] - [body-start (+ sep 4)] - [lines (string-split-crlf header-str)]) - (when (null? lines) - (error 'read-response "empty HTTP response")) - (let ([status (parse-status-line (car lines))] - [headers (parse-headers (cdr lines))]) - (if (= status 100) - ;; Skip "100 Continue" and parse the real response - (loop body-start) - (let* ([raw-body (if (< body-start len) - (subbytevector raw body-start len) - (make-bytevector 0))] - [body (if (chunked-encoding? headers) - (decode-chunked raw-body) - raw-body)]) - (values status headers body))))))))) + (if (= offset n) + result + (let* ([want (min (- n offset) (bytevector-length buf))] + [got (ssl-read conn buf want)]) + (if (<= got 0) + ;; Truncated — return what we have + (subbytevector result 0 offset) + (begin + (bytevector-copy! buf 0 result offset got) + (loop (+ offset got))))))))) + + (define (ssl-read-until-eof conn) + ;; Read all data until connection closes. Returns bytevector. + (let ([buf (make-bytevector 32768)]) + (let loop ([chunks '()]) + (let ([n (ssl-read conn buf 32768)]) + (if (<= n 0) + (bytevector-concat-list (reverse chunks)) + (loop (cons (subbytevector buf 0 n) chunks))))))) + + (define (ssl-read-headers conn) + ;; Read HTTP headers incrementally until \r\n\r\n. + ;; Returns the raw header bytes (including the terminating \r\n\r\n) + ;; and any extra body bytes that were read past the header boundary. + (let ([buf (make-bytevector 8192)]) + (let loop ([chunks '()] [total 0]) + (let ([n (ssl-read conn buf 8192)]) + (if (<= n 0) + ;; EOF while reading headers — return what we have + (values (bytevector-concat-list (reverse chunks)) (make-bytevector 0)) + (let* ([chunk (subbytevector buf 0 n)] + [all-chunks (reverse (cons chunk chunks))] + [raw (bytevector-concat-list all-chunks)] + [raw-len (bytevector-length raw)] + [sep (find-crlfcrlf raw raw-len 0)]) + (if sep + ;; Found header/body boundary + (let ([header-end (+ sep 4)]) + (values (subbytevector raw 0 header-end) + (if (< header-end raw-len) + (subbytevector raw header-end raw-len)