feat: TCP transport adapter support in (std raft) and (std net tcp)
ober
51e8663b9c4e9477faa5af8d2ad0be00b0490476
--- a/lib/std/net/tcp.sls +++ b/lib/std/net/tcp.sls @@ -20,7 +20,12 @@ tcp-listen tcp-accept tcp-close tcp-connect tcp-server? tcp-server-port - with-tcp-server) + with-tcp-server + ;; Binary-port variants — for protocols that need raw bytevector I/O + ;; (e.g. FASL framing). Same semantics as tcp-accept/tcp-connect but + ;; returns (values binary-input-port binary-output-port) with no + ;; UTF-8 transcoder. + tcp-accept-binary tcp-connect-binary) (import (chezscheme)) @@ -297,4 +302,89 @@ (eol-style none) (error-handling-mode replace))))))) + ;; ========== Binary-Port Variants ========== + + (define (fd->binary-ports fd name) + ;; Like fd->ports but returns raw binary ports (no UTF-8 transcoder). + ;; Use for protocols that need put-bytevector / get-bytevector-n, e.g. + ;; FASL-framed Raft transport messages. + (set-nonblocking! fd) + (let ([closed? #f]) + (let ([in (make-custom-binary-input-port + (string-append name "-bin-in") + (lambda (bv start count) + (if closed? 0 + (let ([buf (make-bytevector count)]) + (let retry () + (let ([n (c-read fd buf count)]) + (cond + [(> n 0) + (bytevector-copy! buf 0 bv start n) + n] + [(and (< n 0) + (let ([e (get-errno)]) + (or (= e EINTR) (= e EAGAIN)))) + (sleep *retry-delay*) + (retry)] + [else 0])))))) + #f #f + (lambda () + (unless closed? + (set! closed? #t) + (c-close fd))))] + [out (make-custom-binary-output-port + (string-append name "-bin-out") + (lambda (bv start count) + (if closed? 0 + (let ([buf (make-bytevector count)]) + (bytevector-copy! bv start buf 0 count) + (let lp ([written 0]) + (if (= written count) + count + (let ([n (c-write fd + (let ([tmp (make-bytevector (- count written))]) + (bytevector-copy! buf written tmp 0 (- count written)) + tmp) + (- count written))]) + (cond + [(> n 0) (lp (+ written n))] + [(and (< n 0) + (let ([e (get-errno)]) + (or (= e EINTR) (= e EAGAIN)))) + (sleep *retry-delay*) + (lp written)] + [else written]))))))) + #f #f #f)]) + (values in out)))) + + (define (tcp-accept-binary srv) + ;; Like tcp-accept but returns (values binary-input-port binary-output-port). + (let ([fd (tcp-server-fd srv)]) + (let loop () + (let ([client-fd (c-accept fd 0 0)]) + (cond + [(>= client-fd 0) (fd->binary-ports client-fd "tcp-client")] + [(let ([e (get-errno)]) (or (= e EINTR) (= e EAGAIN))) + (sleep *retry-delay*) + (loop)] + [else (error 'tcp-accept-binary "accept() failed")]))))) + + (define (tcp-connect-binary address port) + ;; Like tcp-connect but returns (values binary-input-port binary-output-port). + (let ([fd (c-socket AF_INET SOCK_STREAM 0)]) + (when (< fd 0) + (error 'tcp-connect-binary "socket() failed")) + (let ([addr (make-sockaddr-in address port)]) + (let loop () + (let ([rc (c-connect fd addr SOCKADDR_IN_SIZE)]) + (cond + [(>= rc 0) + (foreign-free addr) + (fd->binary-ports fd "tcp-connection")] + [(= (get-errno) EINTR) (loop)] + [else + (foreign-free addr) + (c-close fd) + (error 'tcp-connect-binary "connect() failed" address port)])))))) + ) ;; end library --- a/lib/std/raft.sls +++ b/lib/std/raft.sls @@ -26,7 +26,12 @@ raft-commit-index make-raft-cluster raft-cluster-nodes - raft-cluster-leader) + raft-cluster-leader + ;; Transport adapter hooks — needed by (jerboa-db transport) + raft-node-inbox ;; (raft-node → channel) read-only: inject inbound messages + raft-node-peers ;; (raft-node → list) read-only: current peers list + raft-node-peers-set! + raft-node-add-peer!) ;; safe runtime peer addition (updates next/match-index for leaders) (import (chezscheme) (std misc channel)) @@ -159,9 +164,30 @@ (define (majority count) (+ (quotient count 2) 1)) - ;; Random election timeout: 150–300ms - (define (election-timeout-ms) - (+ 150 (random 150))) + ;; Random election timeout: 150–300ms per node. + ;; A node-specific constant offset (derived from the node's id) is added to + ;; the random component so that two nodes starting simultaneously cannot end + ;; up with identical timeouts even if Chez's PRNG returns the same value from + ;; concurrent threads (which it can, since random is not thread-safe). + (define (node-id-hash id) + ;; Simple integer hash of a symbol or integer node identifier. + (cond + [(integer? id) id] + [(symbol? id) + (let* ([s (symbol->string id)] + [n (string-length s)]) + (let loop ([i 0] [h 0]) + (if (= i n) + h + (loop (+ i 1) (+ (* h 31) (char->integer (string-ref s i)))))))] + [else 0])) + + (define (election-timeout-ms node) + ;; Node-specific base: 0–74ms determined by node id, random: 0–74ms. + ;; Total range: 150–298ms; guaranteed unique initial value per node id. + (+ 150 + (modulo (node-id-hash (raft-node-id node)) 75) + (random 75))) ;; Heartbeat interval: 50ms (define heartbeat-interval-ms 50) @@ -228,7 +254,7 @@ (when (raft-node-running? node) (let ([t (fork-thread (lambda () - (let ([timeout-ms (election-timeout-ms)]) + (let ([timeout-ms (election-timeout-ms node)]) (sleep (make-time 'time-duration (* timeout-ms 1000000) 0)) (when (and (raft-node-running? node) @@ -267,18 +293,22 @@ [commit (raft-node-commit-index node)]) (for-each (lambda (peer-entry) - (let* ([peer-id (car peer-entry)] - [ni (cdr (assv peer-id (raft-node-next-index node)))] - [prev-index (- ni 1)] - [prev-term (let ([e (log-entry-at node prev-index)]) - (if e (log-entry-term e) 0))] - ;; entries from ni onwards - [entries (filter (lambda (e) - (>= (log-entry-index e) ni)) - (raft-node-log node))]) - (send-to-peer node peer-id - (make-append-entries term id prev-index prev-term - entries commit)))) + (guard (exn [#t (void)]) ;; never let a bad peer crash the message loop + (let* ([peer-id (car peer-entry)] + ;; Safe lookup: peers added after become-leader! may not be in next-index yet; + ;; fall back to log-last-index+1 so the peer gets a complete resync. + [ni-entry (assv peer-id (raft-node-next-index node))] + [ni (if ni-entry (cdr ni-entry) (+ (log-last-index node) 1))] + [prev-index (- ni 1)] + [prev-term (let ([e (log-entry-at node prev-index)]) + (if e (log-entry-term e) 0))] + ;; entries from ni onwards + [entries (filter (lambda (e) + (>= (log-entry-index e) ni)) + (raft-node-log node))]) + (send-to-peer node peer-id + (make-append-entries term id prev-index prev-term + entries commit))))) (raft-node-peers node)))) ;; ========== Election Check ========== @@ -414,7 +444,7 @@ (let* ([existing (raft-node-log node)] ;; Keep entries before first new entry [keep (filter (lambda (e) - (< (log-entry-index e) prev-idx)) + (<= (log-entry-index e) prev-idx)) existing)]) ;; Append entries from leader (raft-node-log-set! node (append keep entries)))) @@ -481,6 +511,26 @@ ;; ========== Node Lifecycle ========== + ;; raft-node-add-peer! : raft-node peer-id channel -> void + ;; + ;; Safely adds a new peer to a running node. Unlike raft-node-peers-set!, + ;; this function also initialises next-index and match-index for the new peer + ;; when the node is currently a leader, so the next call to send-heartbeats! + ;; does not crash with a missing assv entry. + ;; + ;; Thread-safe: acquires the node's mutex. + (define (raft-node-add-peer! node peer-id peer-chan) + (with-mutex (raft-node-mutex node) + (unless (assv peer-id (raft-node-peers node)) + (raft-node-peers-set! node + (cons (cons peer-id peer-chan) (raft-node-peers node))) + (when (eq? (raft-node-state node) 'leader) + (let ([ni (+ (log-last-index node) 1)]) + (raft-node-next-index-set! node + (cons (cons peer-id ni) (raft-node-next-index node))) + (raft-node-match-index-set! node + (cons (cons peer-id 0) (raft-node-match-index node)))))))) + (define (raft-start! node) (raft-node-running?-set! node #t) ;; Start main message loop in a thread @@ -490,8 +540,10 @@ (let loop () (when (raft-node-running? node) (let ([msg (channel-get (raft-node-inbox node))]) - (with-mutex (raft-node-mutex node) - (handle-message! node msg)) + ;; Guard prevents any handler error from crashing the loop thread. + (guard (exn [#t (void)]) + (with-mutex (raft-node-mutex node) + (handle-message! node msg))) (loop)))))) node)