fiber-io: epoll-integrated I/O core for green threads — (std net io)
ober
6a7d5228eb7cf4a2841699aca17c04de8ff32dfa
--- a/jerboa-native-rs/src/epoll.rs +++ b/jerboa-native-rs/src/epoll.rs @@ -93,3 +93,43 @@ pub extern "C" fn jerboa_epoll_close(epfd: i32) -> i32 { rc }) } + +// ---------- eventfd for poller wakeup ---------- + +#[no_mangle] +pub extern "C" fn jerboa_eventfd_create() -> i32 { + ffi_wrap(|| { + let fd = unsafe { libc::eventfd(0, libc::EFD_NONBLOCK | libc::EFD_CLOEXEC) }; + if fd < 0 { + crate::panic::set_last_error(format!( + "eventfd: {}", + std::io::Error::last_os_error() + )); + } + fd + }) +} + +/// Write 1 to an eventfd to wake a blocked epoll_wait. +#[no_mangle] +pub extern "C" fn jerboa_eventfd_signal(fd: i32) -> i32 { + ffi_wrap(|| { + let val: u64 = 1; + let rc = unsafe { + libc::write(fd, &val as *const u64 as *const libc::c_void, 8) + }; + if rc < 0 { -1 } else { 0 } + }) +} + +/// Read and clear an eventfd (drain the counter). +#[no_mangle] +pub extern "C" fn jerboa_eventfd_drain(fd: i32) -> i32 { + ffi_wrap(|| { + let mut val: u64 = 0; + let rc = unsafe { + libc::read(fd, &mut val as *mut u64 as *mut libc::c_void, 8) + }; + if rc < 0 { -1 } else { 0 } + }) +} --- a/lib/std/fiber.sls +++ b/lib/std/fiber.sls @@ -82,7 +82,12 @@ current-fiber-runtime current-fiber - with-fibers) + with-fibers + + ;; Low-level primitives for I/O integration (used by std net io) + wake-fiber! + fiber-gate-set! + spin-until-gate) (import (chezscheme) (std misc cpu)) new file mode 100644 --- /dev/null +++ b/lib/std/net/io.sls @@ -0,0 +1,509 @@ +#!chezscheme +;;; (std net io) — Fiber-Aware I/O Core +;;; +;;; Integrates Linux epoll with the fiber scheduler so that socket I/O +;;; parks fibers instead of blocking OS threads. This is the foundation +;;; for high-scalability networking. +;;; +;;; Architecture: +;;; - A dedicated poller thread runs epoll_wait in a loop +;;; - Each fd maps to a poll-desc holding parked reader/writer fibers +;;; - fiber-wait-readable/writable park the current fiber and register +;;; with epoll; the poller wakes them via wake-fiber! +;;; - Uses edge-triggered epoll (EPOLLET) — one notification per state +;;; change, matching Go's proven netpoller model +;;; - An eventfd wakes the poller when new fds are registered +;;; +;;; API: +;;; (make-io-poller rt) — create poller for a fiber runtime +;;; (io-poller-start! poller) — start poller thread +;;; (io-poller-stop! poller) — stop poller thread +;;; (fiber-wait-readable fd poller) — park fiber until fd is readable +;;; (fiber-wait-writable fd poller) — park fiber until fd is writable +;;; (poller-register-fd! poller fd events) — register fd with epoll +;;; (poller-unregister-fd! poller fd) — remove fd from epoll +;;; +;;; (fiber-tcp-accept srv poller) — accept, parking fiber on EAGAIN +;;; (fiber-tcp-read fd buf n poller) — read, parking fiber on EAGAIN +;;; (fiber-tcp-write fd buf n poller)— write, parking fiber on EAGAIN +;;; (fiber-tcp-connect addr port poller) — non-blocking connect +;;; +;;; (with-io-poller rt body ...) — convenience macro + +(library (std net io) + (export + ;; Poller lifecycle + make-io-poller + io-poller? + io-poller-start! + io-poller-stop! + + ;; Core primitives + fiber-wait-readable + fiber-wait-writable + poller-register-fd! + poller-unregister-fd! + + ;; Fiber-aware TCP operations + fiber-tcp-accept + fiber-tcp-read + fiber-tcp-write + fiber-tcp-connect + + ;; Convenience + with-io-poller + fiber-tcp-listen + fiber-tcp-close) + + (import (chezscheme) + (std fiber) + (std os epoll-native)) + + ;; ========== FFI for raw socket ops ========== + + (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-socket (foreign-procedure "socket" (int int int) int)) + (define c-bind (foreign-procedure "bind" (int void* int) int)) + (define c-listen (foreign-procedure "listen" (int int) int)) + (define c-accept (foreign-procedure "accept" (int void* void*) int)) + (define c-connect (foreign-procedure "connect" (int void* int) int)) + (define c-close (foreign-procedure "close" (int) int)) + (define c-setsockopt (foreign-procedure "setsockopt" (int int int void* int) int)) + (define c-read (foreign-procedure "read" (int u8* size_t) ssize_t)) + (define c-write (foreign-procedure "write" (int u8* size_t) ssize_t)) + (define c-htons (foreign-procedure "htons" (unsigned-short) unsigned-short)) + (define c-inet-pton (foreign-procedure "inet_pton" (int string void*) int)) + (define c-getsockname (foreign-procedure "getsockname" (int void* void*) int)) + (define c-fcntl (foreign-procedure "fcntl" (int int int) int)) + + ;; errno + (define c-errno-location + (cond + ((foreign-entry? "__errno_location") + (foreign-procedure "__errno_location" () void*)) + ((foreign-entry? "__error") + (foreign-procedure "__error" () void*)) + ((foreign-entry? "__errno") + (foreign-procedure "__errno" () void*)) + (else (foreign-procedure "__errno_location" () void*)))) + (define (get-errno) (foreign-ref 'int (c-errno-location) 0)) + + (define EINTR 4) + (define EAGAIN 11) + (define EINPROGRESS 115) + + ;; Constants + (define AF_INET 2) + (define SOCK_STREAM 1) + (define SOCK_NONBLOCK #x800) + (define SOL_SOCKET 1) + (define SO_REUSEADDR 2) + (define SO_REUSEPORT 15) + (define SOCKADDR_IN_SIZE 16) + (define F_GETFL 3) + (define F_SETFL 4) + (define O_NONBLOCK #x800) + + (define (set-nonblocking! fd) + (let ([flags (c-fcntl fd F_GETFL 0)]) + (c-fcntl fd F_SETFL (bitwise-ior flags O_NONBLOCK)))) + + ;; ========== sockaddr_in helpers ========== + + (define (make-sockaddr-in address port) + (let ([buf (foreign-alloc SOCKADDR_IN_SIZE)]) + (do ([i 0 (fx+ i 1)]) ((fx= i SOCKADDR_IN_SIZE)) + (foreign-set! 'unsigned-8 buf i 0)) + (foreign-set! 'unsigned-short buf 0 AF_INET) + (foreign-set! 'unsigned-short buf 2 (c-htons port)) + (when (= (c-inet-pton AF_INET address (+ buf 4)) 0) + (foreign-free buf) + (error 'make-sockaddr-in "invalid address" address)) + buf)) + + (define (sockaddr-in-port buf) + (let ([hi (foreign-ref 'unsigned-8 buf 2)] + [lo (foreign-ref 'unsigned-8 buf 3)]) + (+ (* hi 256) lo))) + + ;; ========== Poll descriptor ========== + ;; + ;; Maps an fd to the fibers waiting on it. + + (define-record-type poll-desc + (fields + (immutable fd) + (mutable events) ;; currently registered epoll events + (mutable reader-fiber) ;; fiber parked for read, or #f + (mutable writer-fiber) ;; fiber parked for write, or #f + (mutable error?) ;; set #t on EPOLLERR/EPOLLHUP + (immutable pd-mutex)) ;; protects reader/writer fields + (protocol + (lambda (new) + (lambda (fd) + (new fd 0 #f #f #f (make-mutex)))))) + + ;; ========== FD table ========== + ;; + ;; Growable vector mapping fd -> poll-desc or #f. + ;; fd values on Linux are small non-negative integers. + + (define-record-type fd-table + (fields + (mutable vec) + (mutable size) + (immutable ft-mutex)) + (protocol + (lambda (new) + (lambda () + (new (make-vector 4096 #f) 4096 (make-mutex)))))) + + (define (ft-ensure-capacity! ft fd) + (when (fx>= fd (fd-table-size ft)) + (let* ([old-sz (fd-table-size ft)] + [new-sz (fx* 2 (fxmax old-sz (fx+ fd 1)))] + [old-vec (fd-table-vec ft)] + [new-vec (make-vector new-sz #f)]) + (do ([i 0 (fx+ i 1)]) ((fx= i old-sz)) + (vector-set! new-vec i (vector-ref old-vec i))) + (fd-table-vec-set! ft new-vec) + (fd-table-size-set! ft new-sz)))) + + (define (ft-ref ft fd) + (if (fx< fd (fd-table-size ft)) + (vector-ref (fd-table-vec ft) fd) + #f)) + + (define (ft-set! ft fd pd) + (ft-ensure-capacity! ft fd) + (vector-set! (fd-table-vec ft) fd pd)) + + (define (ft-remove! ft fd) + (when (fx< fd (fd-table-size ft)) + (vector-set! (fd-table-vec ft) fd #f))) + + ;; ========== IO Poller ========== + + (define-record-type io-poller + (fields + (immutable epfd) ;; epoll file descriptor + (immutable wakefd) ;; eventfd for waking the poller + (immutable fdt) ;; fd-table + (immutable runtime) ;; fiber-runtime (for wake-fiber!) + (mutable poller-thread) ;; OS thread handle + (mutable running?)) ;; shutdown flag + (protocol + (lambda (new) + (lambda (rt) + (let ([epfd (epoll-create)] + [wakefd (eventfd-create)] + [fdt (make-fd-table)]) + ;; Register the wakefd with epoll so epoll_wait unblocks + ;; when we signal it + (epoll-add! epfd wakefd (bitwise-ior EPOLLIN EPOLLET)) + (new epfd wakefd fdt rt #f #f)))))) + + ;; ========== Poller thread ========== + + (define *max-events* 256) + + (define (poller-loop poller) + (let ([epfd (io-poller-epfd poller)] + [wakefd (io-poller-wakefd poller)] + [fdt (io-poller-fdt poller)]) + (let loop () + (when (io-poller-running? poller) + ;; Block until events (or wakefd signal). + ;; Use 100ms timeout as safety net so we re-check running?. + (let ([events (epoll-wait epfd *max-events* 100)]) + (for-each + (lambda (ev) + (let ([fd (car ev)] + [mask (cdr ev)]) + (cond + ;; wakefd — just drain it, new registrations handled next iteration + [(fx= fd wakefd) (eventfd-drain wakefd)] + ;; Real fd — look up poll-desc and wake parked fibers + [else + (let ([pd (with-mutex (fd-table-ft-mutex fdt) (ft-ref fdt fd))]) + (when pd + (let ([pdmx (poll-desc-pd-mutex pd)]) + (mutex-acquire pdmx) + ;; Error/hangup — wake both reader and writer + (when (or (not (zero? (bitwise-and mask EPOLLERR))) + (not (zero? (bitwise-and mask EPOLLHUP)))) + (poll-desc-error?-set! pd #t)) + ;; Readable — wake reader + (when (or (not (zero? (bitwise-and mask EPOLLIN))) + (poll-desc-error? pd)) + (let ([rf (poll-desc-reader-fiber pd)]) + (when rf + (poll-desc-reader-fiber-set! pd #f) + (mutex-release pdmx) + (wake-fiber! rf) + (mutex-acquire pdmx)))) + ;; Writable — wake writer + (when (or (not (zero? (bitwise-and mask EPOLLOUT))) + (poll-desc-error? pd)) + (let ([wf (poll-desc-writer-fiber pd)]) + (when wf + (poll-desc-writer-fiber-set! pd #f) + (mutex-release pdmx) + (wake-fiber! wf) + (mutex-acquire pdmx)))) + (mutex-release pdmx))))]))) + events)) + (loop))))) + + (define (io-poller-start! poller) + (io-poller-running?-set! poller #t) + (io-poller-poller-thread-set! poller + (fork-thread (lambda () (poller-loop poller))))) + + (define (io-poller-stop! poller) + (io-poller-running?-set! poller #f) + (eventfd-signal (io-poller-wakefd poller)) + ;; Give poller thread time to exit + (sleep (make-time 'time-duration 100000000 0)) + (epoll-close (io-poller-epfd poller)) + (c-close (io-poller-wakefd poller))) + + ;; ========== Register/unregister fds ========== + + (define (poller-register-fd! poller fd events) + (let ([fdt (io-poller-fdt poller)] + [epfd (io-poller-epfd poller)]) + (let ([pd (make-poll-desc fd)]) + (poll-desc-events-set! pd events) + (with-mutex (fd-table-ft-mutex fdt) + (ft-set! fdt fd pd)) + (epoll-add! epfd fd (bitwise-ior events EPOLLET)) + (eventfd-signal (io-poller-wakefd poller)) + pd))) + + (define (poller-unregister-fd! poller fd) + (let ([fdt (io-poller-fdt poller)] + [epfd (io-poller-epfd poller)]) + (guard (e [#t (void)]) ;; ignore if fd already removed + (epoll-remove! epfd fd)) + (with-mutex (fd-table-ft-mutex fdt) + (ft-remove! fdt fd)))) + + ;; ========== Internal: ensure fd is registered, get its poll-desc ========== + + ;; ensure-poll-desc! — register fd with epoll if not already tracked. + ;; Uses level-triggered mode (no EPOLLET, no ONESHOT). + ;; The poller fires repeatedly while an fd is ready — if no fiber + ;; is registered, the event is simply ignored. When a fiber IS + ;; registered, it gets woken on the next poller iteration. + (define (ensure-poll-desc! poller fd) + (let ([fdt (io-poller-fdt poller)]) + (with-mutex (fd-table-ft-mutex fdt) + (or (ft-ref fdt fd) + (let ([pd (make-poll-desc fd)]) + (ft-set! fdt fd pd) + (epoll-add! (io-poller-epfd poller) fd + (bitwise-ior EPOLLIN EPOLLOUT)) + pd))))) + + ;; ========== fiber-wait-readable / fiber-wait-writable ========== + ;; + ;; Strategy: Level-triggered epoll. The poller fires continuously + ;; for ready fds. When a fiber parks for read/write, the next + ;; poller iteration will see the fd is ready (if it is), find the + ;; registered fiber, and wake it. No lost events possible. + ;; + ;; To avoid busy-wake on fds that are always writable, the poller + ;; only wakes fibers that are actually registered (non-#f). + + (define (fiber-wait-readable fd poller) + (let ([f (fiber-self)] + [pd (ensure-poll-desc! poller fd)]) + (fiber-check-cancelled!) + (let ([pdmx (poll-desc-pd-mutex pd)] + [gate (box 'channel)]) + ;; Register fiber as reader + (mutex-acquire pdmx) + (poll-desc-reader-fiber-set! pd f) + (mutex-release pdmx) + ;; Signal poller in case it's sleeping in epoll_wait + (eventfd-signal (io-poller-wakefd poller)) + ;; Park the fiber + (fiber-gate-set! f gate) + (set-timer 1) + (spin-until-gate gate) + (fiber-gate-set! f #f) + (fiber-check-cancelled!) + (void)))) + + (define (fiber-wait-writable fd poller) + (let ([f (fiber-self)] + [pd (ensure-poll-desc! poller fd)]) + (fiber-check-cancelled!) + (let ([pdmx (poll-desc-pd-mutex pd)] + [gate (box 'channel)]) + ;; Register fiber as writer + (mutex-acquire pdmx) + (poll-desc-writer-fiber-set! pd f) + (mutex-release pdmx) + ;; Signal poller + (eventfd-signal (io-poller-wakefd poller)) + ;; Park the fiber + (fiber-gate-set! f gate) + (set-timer 1) + (spin-until-gate gate) + (fiber-gate-set! f #f) + (fiber-check-cancelled!) + (void)))) + + ;; ========== Fiber-aware TCP operations ========== + + ;; ---------- fiber-tcp-listen ---------- + + (define fiber-tcp-listen + (case-lambda + [(address port) (fiber-tcp-listen address port 4096)] + [(address port backlog) + (let ([fd (c-socket AF_INET SOCK_STREAM 0)]) + (when (< fd 0) (error 'fiber-tcp-listen "socket() failed")) + ;; SO_REUSEADDR + SO_REUSEPORT + (let ([one (foreign-alloc 4)]) + (foreign-set! 'int one 0 1) + (c-setsockopt fd SOL_SOCKET SO_REUSEADDR one 4) + (c-setsockopt fd SOL_SOCKET SO_REUSEPORT one 4) + (foreign-free one)) + ;; Bind + (let ([addr (make-sockaddr-in address port)]) + (let ([rc (c-bind fd addr SOCKADDR_IN_SIZE)]) + (foreign-free addr) + (when (< rc 0) (c-close fd) + (error 'fiber-tcp-listen "bind() failed" address port)))) + ;; Listen + (when (< (c-listen fd backlog) 0) + (c-close fd) + (error 'fiber-tcp-listen "listen() failed")) + ;; Non-blocking + (set-nonblocking! fd) + ;; Get actual port + (let ([buf (foreign-alloc SOCKADDR_IN_SIZE)] + [len (foreign-alloc 4)]) + (foreign-set! 'int len 0 SOCKADDR_IN_SIZE) + (c-getsockname fd buf len) + (let ([p (sockaddr-in-port buf)]) + (foreign-free buf) + (foreign-free len) + (values fd p))))])) + + (define (fiber-tcp-close fd) + (c-close fd) + (void)) + + ;; ---------- fiber-tcp-accept ---------- + ;; + ;; Non-blocking accept. On EAGAIN, parks the fiber until the listen + ;; fd becomes readable. Returns client fd (already non-blocking). + + (define (fiber-tcp-accept listen-fd poller) + (let loop () + (let ([client-fd (c-accept listen-fd 0 0)]) + (cond + [(fx>= client-fd 0) + (set-nonblocking! client-fd) + client-fd] + [(let ([e (get-errno)]) (or (= e EAGAIN) (= e EINTR))) + (fiber-wait-readable listen-fd poller) + (loop)] + [else (error 'fiber-tcp-accept "accept() failed" (get-errno))])))) + + ;; ---------- fiber-tcp-read ---------- + ;; + ;; Read up to n bytes into bytevector buf starting at offset 0. + ;; Returns number of bytes read (0 = EOF, negative = error). + ;; Parks fiber on EAGAIN. + + (define (fiber-tcp-read fd buf n poller) + (let loop () + (let ([rc (c-read fd buf n)]) + (cond + [(fx> rc 0) rc] + [(fx= rc 0) 0] ;; EOF + [(let ([e (get-errno)]) (or (= e EAGAIN) (= e EINTR))) + (fiber-wait-readable fd poller) + (loop)] + [else rc])))) ;; error + + ;; ---------- fiber-tcp-write ---------- + ;; + ;; Write n bytes from bytevector buf. Returns total bytes written. + ;; Parks fiber on EAGAIN. Loops until all bytes sent or error. + + (define (fiber-tcp-write fd buf n poller) + (let loop ([written 0]) + (if (fx= written n) written + (let* ([remaining (fx- n written)] + [tmp (if (fx= written 0) buf + (let ([b (make-bytevector remaining)]) + (bytevector-copy! buf written b 0 remaining) b))] + [rc (c-write fd tmp remaining)]) + (cond + [(fx> rc 0) (loop (fx+ written rc))] + [(let ([e (get-errno)]) (or (= e EAGAIN) (= e EINTR))) + (fiber-wait-writable fd poller) + (loop written)] + [else written]))))) ;; partial write on error + + ;; ---------- fiber-tcp-connect ---------- + ;; + ;; Non-blocking connect. Parks fiber while connect is in progress. + ;; Returns the connected fd. + + (define (fiber-tcp-connect address port poller) + (let ([fd (c-socket AF_INET SOCK_STREAM 0)]) + (when (< fd 0) (error 'fiber-tcp-connect "socket() failed")) + (set-nonblocking! fd) + (let ([addr (make-sockaddr-in address port)]) + (let ([rc (c-connect fd addr SOCKADDR_IN_SIZE)]) + (foreign-free addr) + (cond + [(fx>= rc 0) fd] ;; connected immediately + [(= (get-errno) EINPROGRESS) + ;; Connection in progress — wait for writable + (fiber-wait-writable fd poller) + ;; Check SO_ERROR to see if connect succeeded + (let ([err-buf (foreign-alloc 4)] + [len-buf (foreign-alloc 4)]) + (foreign-set! 'int len-buf 0 4) + (c-setsockopt fd SOL_SOCKET 0 err-buf 0) ;; dummy — use getsockopt + ;; Use getsockopt SO_ERROR to check + (let ([getsockopt (foreign-procedure "getsockopt" + (int int int void* void*) int)]) + (getsockopt fd SOL_SOCKET 4 err-buf len-buf) ;; SO_ERROR = 4 + (let ([err (foreign-ref 'int err-buf 0)]) + (foreign-free err-buf) + (foreign-free len-buf) + (when (not (zero? err)) + (c-close fd) + (error 'fiber-tcp-connect "connect() failed" address port err)) + fd)))] + [else + (c-close fd) + (error 'fiber-tcp-connect "connect() failed" address port)]))))) + + ;; ========== Convenience ========== + + (define-syntax with-io-poller + (syntax-rules () + [(_ rt poller-var body ...) + (let ([poller-var (make-io-poller rt)]) + (io-poller-start! poller-var) + (guard (exn [#t (io-poller-stop! poller-var) (raise exn)]) + (let ([result (begin body ...)]) + (io-poller-stop! poller-var) + result)))])) + +) ;; end library --- a/lib/std/os/epoll-native.sls +++ b/lib/std/os/epoll-native.sls @@ -10,7 +10,9 @@ epoll-wait EPOLLIN EPOLLOUT EPOLLERR EPOLLHUP EPOLLET EPOLLONESHOT EPOLLRDHUP EPOLLPRI - EPOLL_CTL_ADD EPOLL_CTL_MOD EPOLL_CTL_DEL) + EPOLL_CTL_ADD EPOLL_CTL_MOD EPOLL_CTL_DEL + ;; eventfd for poller wakeup + eventfd-create eventfd-signal eventfd-drain) (import (chezscheme)) @@ -82,4 +84,26 @@ [events (bytevector-u32-native-ref buf (+ offset 4))]) (loop (+ i 1) (cons (cons fd events) acc))))))))) + ;; --- eventfd for poller wakeup --- + + (define c-eventfd-create + (foreign-procedure "jerboa_eventfd_create" () int)) + (define c-eventfd-signal + (foreign-procedure "jerboa_eventfd_signal" (int) int)) + (define c-eventfd-drain + (foreign-procedure "jerboa_eventfd_drain" (int) int)) + + (define (eventfd-create) + (let ([fd (c-eventfd-create)]) + (when (< fd 0) (error 'eventfd-create "eventfd() failed")) + fd)) + + (define (eventfd-signal fd) + (c-eventfd-signal fd) + (void)) + + (define (eventfd-drain fd) + (c-eventfd-drain fd) + (void)) + ) ;; end library new file mode 100644 --- /dev/null +++ b/tests/test-fiber-io.ss @@ -0,0 +1,161 @@ +;;; Test fiber-aware I/O: echo server with concurrent fiber clients. +;;; Tests Phase 1 of green-wins: epoll integration with fiber scheduler. + +(import (chezscheme)) +(import (std fiber)) +(import (std net io)) + +(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: Poller lifecycle +;; ========================================================================= + +(test "poller creates and stops cleanly" + (let ([rt (make-fiber-runtime 2)]) + (let ([p (make-io-poller rt)]) + (assert-true (io-poller? p) "is poller") + (io-poller-start! p) + (sleep (make-time 'time-duration 50000000 0)) ;; 50ms + (io-poller-stop! p)))) + +;; ========================================================================= +;; Test 2: Fiber-aware TCP listen + accept + read/write +;; ========================================================================= + +(test "echo server with fiber I/O" + (let ([rt (make-fiber-runtime 4)]) + (with-io-poller rt poller + ;; Start server fiber + (let-values ([(listen-fd listen-port) (fiber-tcp-listen "127.0.0.1" 0)]) + (fiber-spawn rt + (lambda () + ;; Accept one client + (let ([client-fd (fiber-tcp-accept listen-fd poller)]) + ;; Echo: read then write back + (let ([buf (make-bytevector 1024)]) + (let ([n (fiber-tcp-read client-fd buf 1024 poller)]) + (when (> n 0) + (fiber-tcp-write client-fd buf n poller)))) + (fiber-tcp-close client-fd))) + "echo-server") + + ;; Start client fiber + (fiber-spawn rt + (lambda () + ;; Small delay to let server fiber start + (fiber-sleep 20) + (let ([fd (fiber-tcp-connect "127.0.0.1" listen-port poller)]) + ;; Send message + (let ([msg (string->bytevector "hello-fiber-io" + (make-transcoder (utf-8-codec)))]) + (fiber-tcp-write fd msg (bytevector-length msg) poller) + ;; Read echo + (let ([buf (make-bytevector 1024)]) + (let ([n (fiber-tcp-read fd buf 1024 poller)]) + (let ([reply (bytevector->string + (let ([b (make-bytevector n)]) + (bytevector-copy! buf 0 b 0 n) b) + (make-transcoder (utf-8-codec)))]) + (assert-equal reply "hello-fiber-io" + "echo reply matches"))))))) + "echo-client") + + ;; Run the runtime (blocks until all fibers done) + (fiber-runtime-run! rt) + (fiber-tcp-close listen-fd))))) + +;; ========================================================================= +;; Test 3: Multiple concurrent connections +;; ========================================================================= + +(test "50 concurrent echo clients" + (let ([rt (make-fiber-runtime 4)] + [num-clients 50] + [results (make-vector 50 #f)]) + (with-io-poller rt poller + (let-values ([(listen-fd listen-port) (fiber-tcp-listen "127.0.0.1" 0)]) + ;; Server: accept and echo in a loop + (fiber-spawn rt + (lambda () + (let loop ([i 0]) + (when (< i num-clients) + (let ([client-fd (fiber-tcp-accept listen-fd poller)]) + (fiber-spawn* + (lambda () + (let ([buf (make-bytevector 256)]) + (let ([n (fiber-tcp-read client-fd buf 256 poller)]) + (when (> n 0) + (fiber-tcp-write client-fd buf n poller)))) + (fiber-tcp-close client-fd)) + "echo-handler")) + (loop (+ i 1))))) + "accept-loop") + + ;; Spawn 50 client fibers + (do ([i 0 (+ i 1)]) + ((= i num-clients)) + (let ([idx i]) + (fiber-spawn rt + (lambda () + (fiber-sleep 10) + (let ([fd (fiber-tcp-connect "127.0.0.1" listen-port poller)]) + (let ([msg (string->bytevector + (string-append "msg-" (number->string idx)) + (make-transcoder (utf-8-codec)))]) + (fiber-tcp-write fd msg (bytevector-length msg) poller) + (let ([buf (make-bytevector 256)]) + (let ([n (fiber-tcp-read fd buf 256 poller)]) + (let ([reply (bytevector->string + (let ([b (make-bytevector n)]) + (bytevector-copy! buf 0 b 0 n) b) + (make-transcoder (utf-8-codec)))]) + (vector-set! results idx + (string=? reply (string-append "msg-" (number->string idx)))))))) + (fiber-tcp-close fd))) + (string-append "client-" (number->string idx))))) + + (fiber-runtime-run! rt) + (fiber-tcp-close listen-fd) + + ;; Check all clients got correct echo + (let ([ok-count (do ([i 0 (+ i 1)] [c 0 (+ c (if (vector-ref results i) 1 0))]) + ((= i num-clients) c))]) + (assert-equal ok-count num-clients "all clients echoed correctly")))))) + +;; ========================================================================= +;; 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))