sendfile + connpool: zero-copy I/O and fiber-aware connection pooling (Phase 6)
ober
d5b831d4772db44b4c930d1200a42746d9a5ec2d
new file mode 100644 --- /dev/null +++ b/lib/std/net/connpool.sls @@ -0,0 +1,120 @@ +#!chezscheme +;;; (std net connpool) — Fiber-aware connection pool +;;; +;;; Maintains a pool of reusable TCP connections. Fibers acquire a +;;; connection, use it, and return it. The pool manages the lifecycle. +;;; +;;; API: +;;; (make-conn-pool host port poller max-size) — create pool +;;; (conn-pool-acquire! pool) — get a connection (parks if full) +;;; (conn-pool-release! pool fd) — return connection to pool +;;; (conn-pool-discard! pool fd) — close and remove connection +;;; (conn-pool-close! pool) — close all connections +;;; (conn-pool-size pool) — current pool size +;;; (with-pooled-connection pool body ...) — scoped acquire/release + +(library (std net connpool) + (export + make-conn-pool + conn-pool? + conn-pool-acquire! + conn-pool-release! + conn-pool-discard! + conn-pool-close! + conn-pool-size + with-pooled-connection) + + (import (chezscheme) + (std fiber) + (std net io)) + + ;; ========== Connection pool ========== + + (define-record-type conn-pool + (fields + (immutable host) + (immutable port) + (immutable poller) + (immutable max-size) + (mutable idle) ;; list of idle fd's + (mutable active) ;; count of checked-out connections + (immutable mutex) + (immutable semaphore)) ;; fiber-semaphore for max-size + (protocol + (lambda (new) + (lambda (host port poller max-size) + (new host port poller max-size + '() 0 (make-mutex) + (make-fiber-semaphore max-size)))))) + + (define (conn-pool-size pool) + (with-mutex (conn-pool-mutex pool) + (+ (length (conn-pool-idle pool)) + (conn-pool-active pool)))) + + ;; Acquire a connection from the pool. + ;; Returns an idle connection if available, or creates a new one. + ;; Parks the fiber if the pool is at max capacity. + (define (conn-pool-acquire! pool) + ;; Acquire semaphore permit (parks if at max) + (fiber-semaphore-acquire! (conn-pool-semaphore pool)) + (let ([mx (conn-pool-mutex pool)]) + (mutex-acquire mx) + (let ([idle (conn-pool-idle pool)]) + (cond + [(not (null? idle)) + ;; Reuse an idle connection + (let ([fd (car idle)]) + (conn-pool-idle-set! pool (cdr idle)) + (conn-pool-active-set! pool (+ (conn-pool-active pool) 1)) + (mutex-release mx) + fd)] + [else + ;; Create a new connection + (conn-pool-active-set! pool (+ (conn-pool-active pool) 1)) + (mutex-release mx) + (fiber-tcp-connect (conn-pool-host pool) + (conn-pool-port pool) + (conn-pool-poller pool))])))) + + ;; Return a connection to the pool for reuse. + (define (conn-pool-release! pool fd) + (let ([mx (conn-pool-mutex pool)]) + (mutex-acquire mx) + (conn-pool-idle-set! pool (cons fd (conn-pool-idle pool))) + (conn-pool-active-set! pool (max 0 (- (conn-pool-active pool) 1))) + (mutex-release mx)) + ;; Release semaphore permit + (fiber-semaphore-release! (conn-pool-semaphore pool))) + + ;; Close a connection and remove it from the pool (e.g., on error). + (define (conn-pool-discard! pool fd) + (let ([mx (conn-pool-mutex pool)]) + (mutex-acquire mx) + (conn-pool-active-set! pool (max 0 (- (conn-pool-active pool) 1))) + (mutex-release mx)) + (fiber-tcp-close fd) + ;; Release semaphore permit + (fiber-semaphore-release! (conn-pool-semaphore pool))) + + ;; Close all connections in the pool. + (define (conn-pool-close! pool) + (let ([mx (conn-pool-mutex pool)]) + (mutex-acquire mx) + (for-each (lambda (fd) (fiber-tcp-close fd)) (conn-pool-idle pool)) + (conn-pool-idle-set! pool '()) + (mutex-release mx))) + + ;; Convenience macro: acquire, use, release (or discard on error). + (define-syntax with-pooled-connection + (syntax-rules () + [(_ pool fd body ...) + (let ([fd (conn-pool-acquire! pool)]) + (guard (exn [#t + (conn-pool-discard! pool fd) + (raise exn)]) + (let ([result (begin body ...)]) + (conn-pool-release! pool fd) + result)))])) + +) ;; end library new file mode 100644 --- /dev/null +++ b/lib/std/net/sendfile.sls @@ -0,0 +1,107 @@ +#!chezscheme +;;; (std net sendfile) — Zero-copy file serving via sendfile(2) +;;; +;;; Uses the Linux sendfile(2) syscall to transfer file data directly +;;; from the kernel page cache to a socket, bypassing userspace copies. +;;; Fiber-aware: parks the fiber on EAGAIN for non-blocking sockets. +;;; +;;; API: +;;; (fiber-sendfile sock-fd path poller) — send entire file +;;; (fiber-sendfile* sock-fd path offset count poller) — send range + +(library (std net sendfile) + (export + fiber-sendfile + fiber-sendfile*) + + (import (chezscheme) + (std fiber) + (std net io)) + + ;; FFI + (define _libc-loaded + (let ((v (getenv "JEMACS_STATIC"))) + (if (and v (not (string=? v "")) (not (string=? v "0"))) + #f + (load-shared-object #f)))) + + (define c-open (foreign-procedure "open" (string int) int)) + (define c-close (foreign-procedure "close" (int) int)) + (define c-fstat (foreign-procedure "__fxstat" (int int void*) int)) + (define c-sendfile (foreign-procedure "sendfile" (int int void* size_t) ssize_t)) + + ;; errno + (define c-errno-location + (cond + ((foreign-entry? "__errno_location") + (foreign-procedure "__errno_location" () void*)) + (else (foreign-procedure "__errno_location" () void*)))) + (define (get-errno) (foreign-ref 'int (c-errno-location) 0)) + (define EAGAIN 11) + (define EINTR 4) + (define O_RDONLY 0) + + ;; Get file size using fstat + (define (file-size-fd fd) + ;; struct stat is 144 bytes on x86_64 Linux + ;; st_size is at offset 48 + (let ([buf (foreign-alloc 144)]) + ;; __fxstat version 1 = STAT_VER_LINUX on x86_64 + (let ([rc (c-fstat 1 fd buf)]) + (if (< rc 0) + (begin (foreign-free buf) -1) + (let ([size (foreign-ref 'long buf 48)]) ;; st_size at offset 48 + (foreign-free buf) + size))))) + + ;; Send entire file to socket using sendfile(2). + ;; Returns total bytes sent. + (define (fiber-sendfile sock-fd path poller) + (let ([file-fd (c-open path O_RDONLY)]) + (when (< file-fd 0) + (error 'fiber-sendfile "open() failed" path)) + (let ([size (file-size-fd file-fd)]) + (when (< size 0) + (c-close file-fd) + (error 'fiber-sendfile "fstat() failed" path)) + (let ([result (fiber-sendfile-loop sock-fd file-fd 0 size poller)]) + (c-close file-fd) + result)))) + + ;; Send a range of a file. + ;; Returns total bytes sent. + (define (fiber-sendfile* sock-fd path offset count poller) + (let ([file-fd (c-open path O_RDONLY)]) + (when (< file-fd 0) + (error 'fiber-sendfile* "open() failed" path)) + (let ([result (fiber-sendfile-loop sock-fd file-fd offset count poller)]) + (c-close file-fd) + result))) + + ;; Internal: sendfile loop with fiber parking on EAGAIN + (define (fiber-sendfile-loop sock-fd file-fd offset count poller) + (let ([off-buf (foreign-alloc 8)]) + (foreign-set! 'long off-buf 0 offset) + (let loop ([remaining count] [total 0]) + (if (<= remaining 0) + (begin (foreign-free off-buf) total) + (let ([rc (c-sendfile sock-fd file-fd off-buf remaining)]) + (cond + [(> rc 0) + (loop (- remaining rc) (+ total rc))] + [(= rc 0) + ;; EOF + (foreign-free off-buf) + total] + [else + (let ([e (get-errno)]) + (cond + [(or (= e EAGAIN) (= e EINTR)) + ;; Socket buffer full — park fiber until writable + (fiber-wait-writable sock-fd poller) + (loop remaining total)] + [else + (foreign-free off-buf) + total]))])))))) + +) ;; end library new file mode 100644 --- /dev/null +++ b/tests/test-advanced.ss @@ -0,0 +1,211 @@ +;;; Tests for Phase 6: Advanced optimizations +;;; Tests sendfile and connection pooling. + +(import (chezscheme)) +(import (std fiber)) +(import (std net io)) +(import (std net sendfile)) +(import (std net connpool)) +(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))])) + +;; ========================================================================= +;; Test 1: sendfile — serve a file over TCP +;; ========================================================================= + +(test "sendfile: serve file over TCP" + ;; Create a test file + (let ([path "/tmp/jerboa-test-sendfile.txt"] + [content "Hello from sendfile! This is zero-copy I/O.\n"]) + (let ([p (open-file-output-port path + (file-options no-fail) + (buffer-mode block) + (make-transcoder (utf-8-codec)))]) + (put-string p content) + (close-output-port p)) + + (let ([rt (make-fiber-runtime 4)] + [result-box (box #f)]) + (with-io-poller rt poller + (let-values ([(listen-fd listen-port) (fiber-tcp-listen "127.0.0.1" 0)]) + ;; Server: accept, sendfile, close + (fiber-spawn rt + (lambda () + (let ([client-fd (fiber-tcp-accept listen-fd poller)]) + (fiber-sendfile client-fd path poller) + (fiber-tcp-close client-fd))) + "sendfile-server") + + ;; Client: connect, read all, verify + (fiber-spawn rt + (lambda () + (fiber-sleep 20) + (let ([fd (fiber-tcp-connect "127.0.0.1" listen-port poller)]) + (let ([buf (make-bytevector 4096)] + [tmp (make-bytevector 4096)]) + (let loop ([total 0]) + (let ([n (fiber-tcp-read fd tmp (min 4096 (- 4096 total)) poller)]) + (cond + [(<= n 0) + (set-box! result-box + (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))])))) + (fiber-tcp-close fd))) + "sendfile-client") + + (fiber-runtime-run! rt) + (fiber-tcp-close listen-fd))) + + (assert-equal (unbox result-box) content "sendfile content matches") + (delete-file path)))) + +;; ========================================================================= +;; Test 2: Connection pool — acquire and release +;; ========================================================================= + +(test "connpool: acquire, use, release" + ;; Start a simple echo server + (let* ([handler (lambda (req) (respond-text 200 "pooled-ok"))] + [srv (fiber-httpd-start 0 handler)] + [port (fiber-httpd-listen-port srv)] + [result-box (box #f)]) + + (sleep (make-time 'time-duration 100000000 0)) + + (let ([rt (make-fiber-runtime 2)]) + (with-io-poller rt poller + (let ([pool (make-conn-pool "127.0.0.1" port poller 5)]) + (fiber-spawn rt + (lambda () + ;; Acquire a connection + (let ([fd (conn-pool-acquire! pool)]) + ;; Send an HTTP request + (let* ([req-str "GET / HTTP/1.1\r\nHost: 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 4096)]) + (let ([n (fiber-tcp-read fd buf 4096 poller)]) + (when (> n 0) + (set-box! result-box #t))))) + ;; Discard since server sent Connection: close + (conn-pool-discard! pool fd))) + "pool-client") + (fiber-runtime-run! rt) + (conn-pool-close! pool)))) + + (fiber-httpd-stop! srv) + (assert-true (unbox result-box) "got response via pool"))) + +;; ========================================================================= +;; Test 3: Connection pool — with-pooled-connection macro +;; ========================================================================= + +(test "connpool: with-pooled-connection" + (let* ([handler (lambda (req) (respond-text 200 "macro-ok"))] + [srv (fiber-httpd-start 0 handler)] + [port (fiber-httpd-listen-port srv)] + [result-box (box #f)]) + + (sleep (make-time 'time-duration 100000000 0)) + + (let ([rt (make-fiber-runtime 2)]) + (with-io-poller rt poller + (let ([pool (make-conn-pool "127.0.0.1" port poller 5)]) + (fiber-spawn rt + (lambda () + (with-pooled-connection pool fd + (let* ([req-str "GET / HTTP/1.1\r\nHost: 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 4096)]) + (let ([n (fiber-tcp-read fd buf 4096 poller)]) + (set-box! result-box (> n 0))))))) + "macro-client") + (fiber-runtime-run! rt) + (conn-pool-close! pool)))) + + (fiber-httpd-stop! srv) + (assert-true (unbox result-box) "macro worked"))) + +;; ========================================================================= +;; Test 4: Connection pool — max-size enforcement +;; ========================================================================= + +(test "connpool: pool size tracking" + (let* ([handler (lambda (req) (respond-text 200 "ok"))] + [srv (fiber-httpd-start 0 handler)] + [port (fiber-httpd-listen-port srv)]) + + (sleep (make-time 'time-duration 100000000 0)) + + (let ([rt (make-fiber-runtime 4)]) + (with-io-poller rt poller + (let ([pool (make-conn-pool "127.0.0.1" port poller 3)] + [results (make-vector 3 #f)]) + ;; Spawn 3 fibers — exactly at pool max + ;; Use manual acquire/discard since server sends Connection: close + (do ([i 0 (+ i 1)]) + ((= i 3)) + (let ([idx i]) + (fiber-spawn rt + (lambda () + (let ([fd (conn-pool-acquire! pool)]) + (let* ([req-str "GET / HTTP/1.1\r\nHost: 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 4096)]) + (let ([n (fiber-tcp-read fd buf 4096 poller)]) + (vector-set! results idx (> n 0))))) + (conn-pool-discard! pool fd))) + (string-append "pool-" (number->string idx))))) + (fiber-runtime-run! rt) + (conn-pool-close! pool) + + ;; All 3 should complete + (let ([ok (do ([i 0 (+ i 1)] [c 0 (+ c (if (vector-ref results i) 1 0))]) + ((= i 3) c))]) + (fiber-httpd-stop! srv) + (assert-equal ok 3 "all 3 served through pool"))))))) + +;; ========================================================================= +;; 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))