feat: Phase 6.3 — TCP/TLS transport adapter for Raft clusters
ober
a850cafc2ebcda61e46e6bb19364e1c9d2e4b726
--- a/Makefile +++ b/Makefile @@ -7,7 +7,7 @@ CHEZ_EXT_DIR ?= $(HOME)/src CHEZ_EXT_LIBDIRS = $(CHEZ_EXT_DIR)/chez-lmdb:$(CHEZ_EXT_DIR)/chez-duckdb FULL_LIBDIRS = $(LIBDIRS):$(CHEZ_EXT_LIBDIRS) -.PHONY: test test-cluster build clean check bench bench-quick mbrainz mbrainz-quick +.PHONY: test test-cluster test-transport build clean check bench bench-quick mbrainz mbrainz-quick # Run the core test suite (in-memory, no FFI deps) test: @@ -17,6 +17,10 @@ test: test-cluster: $(SCHEME) --libdirs "$(LIBDIRS)" --script tests/test-cluster.ss +# Run TCP transport tests (two nodes connected via loopback) +test-transport: + $(SCHEME) --libdirs "$(LIBDIRS)" --script tests/test-transport.ss + # Run tests including LMDB backend test-lmdb: $(SCHEME) --libdirs "$(FULL_LIBDIRS)" --script tests/test-lmdb.ss --- a/lib/jerboa-db/replication.ss +++ b/lib/jerboa-db/replication.ss @@ -34,6 +34,7 @@ ;; Lifecycle start-replication stop-replication + start-replication-from-node! ;; wrap a pre-configured raft-node (for transport adapter) ;; Multi-node local cluster convenience start-local-cluster @@ -161,6 +162,17 @@ (raft-stop! (replication-state-raft-node state)) (void)) + ;; start-replication-from-node! : raft-node replication-config -> replication-state + ;; + ;; Wraps a caller-provided, already-started raft-node in a replication-state. + ;; Use this when the raft-node has been configured externally before being wrapped + ;; — for example, by (jerboa-db transport), which wires TCP proxy channels into + ;; the node's peers list before calling raft-start!. + ;; + ;; The caller is responsible for having called raft-start! on the node. + (def (start-replication-from-node! node config) + (make-replication-state config node #f 0 #t)) + ;; ========================================================================= ;; Status / introspection ;; ========================================================================= new file mode 100644 --- /dev/null +++ b/lib/jerboa-db/transport.ss @@ -0,0 +1,338 @@ +#!chezscheme +;;; (jerboa-db transport) — TCP transport adapter for (std raft) +;;; +;;; Bridges each raft-node's in-process channels to TCP sockets so nodes +;;; can run in separate OS processes or on separate machines. +;;; +;;; Architecture: +;;; Each Raft node has: +;;; - An inbox channel: the Raft message loop reads incoming messages here. +;;; - A peers list: ((id . channel) pairs) the node writes outgoing messages here. +;;; +;;; For each peer, the transport installs a "proxy channel" in the peers list. +;;; - Outbound fiber: reads proxy-ch → serializes → writes length-framed FASL to TCP +;;; - Inbound fiber: reads length-framed FASL from TCP → puts to node inbox +;;; +;;; Connections are uni-directional per pair: +;;; Node A dials B for A→B traffic (A writes via out-port of A's outbound connection) +;;; Node B dials A for B→A traffic (B writes via out-port of B's outbound connection) +;;; Each node also runs an accept loop; accepted connections feed the inbox. +;;; +;;; Wire protocol: +;;; [4 bytes big-endian uint32: body length][body: FASL-encoded Raft message] +;;; +;;; client-propose messages (#(client-propose cmd reply-ch)) are local-only and +;;; are dropped by the outbound fiber (they carry live reply channels). +;;; +;;; Apply semantics: +;;; Unlike start-local-db-cluster (which reads the cluster leader's full log), +;;; start-transport-db-node! uses each node's own committed log via +;;; replication-for-each-committed!. This is correct for multi-process clusters. + +(library (jerboa-db transport) + (export + ;; Core transport lifecycle + start-transport-node! + stop-transport-node! + ;; Add a peer after startup (needed for staged test setup) + transport-node-add-peer! + ;; Accessors + transport-node-replication-state + transport-node-listen-port + transport-node? + ;; Convenience: transport + DB connection in one call + start-transport-db-node!) + + (import (except (chezscheme) + make-hash-table hash-table? + sort sort! + printf fprintf + path-extension path-absolute? + with-input-from-string with-output-to-string + iota 1+ 1- + partition + make-date make-time + log ;; avoid conflict with (jerboa-db core)'s log + atom? meta) + (rename (only (chezscheme) make-time) (make-time chez-make-time)) + (jerboa prelude) + (std net tcp) + (std fasl) + (std misc channel) + (std raft) + (jerboa-db replication) + (jerboa-db core) + (jerboa-db cluster)) + + ;; ========================================================================= + ;; Transport node record + ;; ========================================================================= + + ;; defstruct generates: make-transport-node transport-node? transport-node-<field> + ;; and transport-node-<field>-set! for each field. + (defstruct transport-node + (raft-node ;; (std raft) raft-node — holds inbox + peers + replication-state ;; replication-state wrapping the raft-node + listen-port ;; actual TCP listen port (important when 0 was requested) + server ;; tcp-server handle for the listen socket + running?)) ;; mutable: set to #f by stop-transport-node! + + ;; ========================================================================= + ;; Stop sentinel + ;; ========================================================================= + + ;; A unique value placed on proxy channels to signal the outbound fiber to exit. + ;; Cannot appear as a real Raft message (which are vectors). + (define %stop% (list 'transport-stop)) + + ;; ========================================================================= + ;; Wire protocol + ;; ========================================================================= + + ;; write-frame : binary-output-port bytevector → void + ;; Writes a 4-byte big-endian uint32 length header followed by body bytes. + (define (write-frame out-port body) + (let* ([len (bytevector-length body)] + [header (make-bytevector 4)]) + (bytevector-u8-set! header 0 (bitwise-and (ash len -24) #xff)) + (bytevector-u8-set! header 1 (bitwise-and (ash len -16) #xff)) + (bytevector-u8-set! header 2 (bitwise-and (ash len -8) #xff)) + (bytevector-u8-set! header 3 (bitwise-and len #xff)) + (put-bytevector out-port header) + (put-bytevector out-port body) + (flush-output-port out-port))) + + ;; read-frame : binary-input-port → bytevector or #f + ;; Reads a length-prefixed frame. Returns #f on EOF or truncated read. + (define (read-frame in-port) + (let ([header (get-bytevector-n in-port 4)]) + (if (and (bytevector? header) (= (bytevector-length header) 4)) + (let ([len (+ (ash (bytevector-u8-ref header 0) 24) + (ash (bytevector-u8-ref header 1) 16) + (ash (bytevector-u8-ref header 2) 8) + (bytevector-u8-ref header 3))]) + (let ([body (get-bytevector-n in-port len)]) + (and (bytevector? body) (= (bytevector-length body) len) body))) + #f))) + + ;; ========================================================================= + ;; Inbound listener + ;; ========================================================================= + + ;; start-inbound-listener! : binary-input-port channel → thread + ;; + ;; Reads FASL-framed Raft messages from in-port and delivers them to + ;; node-inbox. One fiber per accepted inbound TCP connection. + ;; Exits silently on EOF or deserialization error. + (define (start-inbound-listener! in-port node-inbox) + (fork-thread + (lambda () + (let loop () + (let ([frame (guard (exn [#t #f]) (read-frame in-port))]) + (when frame + (let ([msg (guard (exn [#t #f]) (bytevector->fasl frame))]) + (when msg + (guard (exn [#t (void)]) + (channel-put node-inbox msg)))) + (loop))))))) + + ;; ========================================================================= + ;; Accept loop + ;; ========================================================================= + + ;; start-accept-loop! : tcp-server channel → thread + ;; + ;; Accepts inbound TCP connections and starts an inbound listener for each. + ;; The out-port of accepted connections is closed immediately — inbound + ;; connections are receive-only in this scheme (peers dial us for their + ;; outbound traffic, we dial them for ours). + (define (start-accept-loop! server node-inbox) + (fork-thread + (lambda () + (let loop () + (let ([conn (guard (exn [#t #f]) + (let-values ([(in out) (tcp-accept-binary server)]) + (cons in out)))]) + (when conn + (guard (exn [#t (void)]) (close-port (cdr conn))) + (start-inbound-listener! (car conn) node-inbox))) + (loop))))) + + ;; ========================================================================= + ;; Outbound peer connector + ;; ========================================================================= + + ;; start-peer-connector! : channel string integer → thread + ;; + ;; Maintains a persistent outbound TCP connection to one peer. + ;; On connection failure, reconnects with exponential backoff (100ms → 5s). + ;; Stops when the %stop% sentinel is received on proxy-ch. + ;; + ;; client-propose messages (#(client-propose cmd reply-ch)) are dropped + ;; — they carry live reply channels and must never cross process boundaries. + (define (start-peer-connector! proxy-ch peer-host peer-port) + (fork-thread + (lambda () + (let retry ([backoff-ms 100]) + (let ([result (guard (exn [#t #f]) + (let-values ([(in out) (tcp-connect-binary peer-host peer-port)]) + (cons in out)))]) + (if result + (let ([out-port (cdr result)]) + ;; in-port (car result) is intentionally ignored: both ports share + ;; the same fd/closed? flag in fd->binary-ports, so closing in-port + ;; would destroy out-port too. The fd is GC'd when the conn drops. + ;; Run the proxy loop; returns #t if disconnected, #f if stopped + (let ([disconnected? + (let loop () + (let ([msg (channel-get proxy-ch)]) + (cond + ;; Stop sentinel — exit cleanly + [(eq? msg %stop%) #f] + ;; client-propose carries a live channel — skip + [(and (vector? msg) + (> (vector-length msg) 0) + (eq? (vector-ref msg 0) 'client-propose)) + (loop)] + [else + (let ([ok? (guard (exn [#t #f]) + (write-frame out-port (fasl->bytevector msg)) + #t)]) + (if ok? + (loop) + #t))])))]) ;; #t = disconnected, retry + (guard (exn [#t (void)]) (close-port out-port)) + (when disconnected? + (sleep (chez-make-time 'time-duration (* backoff-ms 1000000) 0)) + (retry (min (* backoff-ms 2) 5000))))) + ;; Connection failed — sleep and retry + (begin + (sleep (chez-make-time 'time-duration (* backoff-ms 1000000) 0)) + (retry (min (* backoff-ms 2) 5000))))))))) + + ;; ========================================================================= + ;; Internal helpers + ;; ========================================================================= + + ;; Proxy channel capacity. Large enough that a 5-second reconnect window + ;; (worst-case backoff) never fills the channel at 50 ms heartbeat rate. + (define proxy-ch-capacity 1024) + + ;; wire-peer! : (id host port) → (id . proxy-channel) + ;; + ;; Creates a proxy channel for one peer and starts the outbound connector fiber. + (define (wire-peer! spec) + (let ([peer-id (car spec)] + [peer-host (cadr spec)] + [peer-port (caddr spec)]) + (let ([proxy-ch (make-channel proxy-ch-capacity)]) + (start-peer-connector! proxy-ch peer-host peer-port) + (cons peer-id proxy-ch)))) + + ;; ========================================================================= + ;; Lifecycle + ;; ========================================================================= + + ;; start-transport-node! : + ;; node-id symbol or integer + ;; peer-specs list of (id host port) + ;; data-path string (":memory:" or file path — informational, passed to config) + ;; listen-port integer (0 = OS-assigned; use transport-node-listen-port to read back) + ;; → transport-node + ;; + ;; Creates a raft-node, wires TCP proxy channels for each peer, starts the + ;; TCP listener, and starts the Raft consensus engine. + ;; + ;; To get a full DB-backed node with the cluster API, use start-transport-db-node!. + (def (start-transport-node! node-id peer-specs data-path listen-port) + (let* ([node (make-raft-node node-id)] + [node-inbox (raft-node-inbox node)] + [peer-chans (map wire-peer! peer-specs)]) + ;; Install proxy channels as the node's peers list + (raft-node-peers-set! node peer-chans) + ;; Bind TCP listen socket + (let ([server (tcp-listen "0.0.0.0" listen-port 16)]) + (let ([actual-port (tcp-server-port server)]) + ;; Accept loop feeds the node's inbox from inbound connections + (start-accept-loop! server node-inbox) + ;; Start Raft consensus engine + (raft-start! node) + ;; Wrap in replication-state for cluster API compatibility + (let* ([config (new-replication-config node-id #f data-path)] + [state (start-replication-from-node! node config)]) + (display (str "transport: node " node-id + " listening on port " actual-port "\n")) + (make-transport-node node state actual-port server #t)))))) + + ;; stop-transport-node! : transport-node → void + ;; + ;; Stops the Raft node, signals outbound proxy fibers to exit, and closes + ;; the TCP listen socket. + (def (stop-transport-node! tnode) + (transport-node-running?-set! tnode #f) + ;; Stop Raft (sends stop-signal to inbox) + (stop-replication (transport-node-replication-state tnode)) + ;; Signal each outbound proxy to exit + (for-each (lambda (p) + (guard (exn [#t (void)]) + (channel-put (cdr p) %stop%))) + (raft-node-peers (transport-node-raft-node tnode))) + ;; Close listen socket + (guard (exn [#t (void)]) + (tcp-close (transport-node-server tnode)))) + + ;; transport-node-add-peer! : + ;; transport-node peer-id string integer → void + ;; + ;; Adds a new peer to an already-running transport node and starts an + ;; outbound connector to it. This is the key primitive for staged test + ;; setup, where port numbers are only known after start: + ;; + ;; (def a (start-transport-node! 0 '() ":memory:" 0)) + ;; (def b (start-transport-node! 1 `((0 "127.0.0.1" ,(transport-node-listen-port a))) + ;; ":memory:" 0)) + ;; (transport-node-add-peer! a 1 "127.0.0.1" (transport-node-listen-port b)) + (def (transport-node-add-peer! tnode peer-id peer-host peer-port) + (let* ([node (transport-node-raft-node tnode)] + [new-entry (wire-peer! (list peer-id peer-host peer-port))]) + ;; raft-node-add-peer! is thread-safe and also initialises next-index / + ;; match-index if the node is already a leader, preventing send-heartbeats! + ;; from crashing with a (cdr #f) on the new peer's missing entry. + (raft-node-add-peer! node peer-id (cdr new-entry)))) + + ;; ========================================================================= + ;; Convenience: transport + DB + replicated-conn + ;; ========================================================================= + + ;; start-transport-db-node! : + ;; node-id symbol or integer + ;; peer-specs list of (id host port) + ;; data-path string (":memory:" or file path for persistent storage) + ;; listen-port integer (0 = OS-assigned) + ;; → (values transport-node replicated-conn) + ;; + ;; Creates the transport node, opens a DB connection at data-path, and + ;; starts an apply fiber that replicates committed Raft entries into the + ;; local DB. Returns a replicated-conn compatible with the full cluster API: + ;; cluster-transact!, cluster-db, cluster-status, cluster-leader? + (def (start-transport-db-node! node-id peer-specs data-path listen-port) + (let* ([tnode (start-transport-node! node-id peer-specs data-path listen-port)] + [state (transport-node-replication-state tnode)] + [conn (connect data-path)] + [fiber (fork-thread + (lambda () + (let loop () + (sleep (chez-make-time 'time-duration 50000000 0)) ;; 50 ms + (when (replication-running? state) + (replication-for-each-committed! state + (lambda (entry-index tx-ops) + (guard (exn [#t + (display + (str "transport apply: skip failed tx at index " + entry-index "\n"))]) + (transact! conn tx-ops)))) + (loop)))))] + [rconn (make-replicated-conn conn state fiber)]) + (values tnode rconn))) + +) ;; end library new file mode 100644 --- /dev/null +++ b/tests/test-transport.ss @@ -0,0 +1,192 @@ +(import (jerboa prelude) + (rename (only (chezscheme) make-time) (make-time chez-make-time)) + (jerboa-db core) + (jerboa-db replication) + (jerboa-db cluster) + (jerboa-db transport)) + +;; ---- Test harness (same as test-cluster.ss) ---- + +(def test-count 0) +(def pass-count 0) +(def fail-count 0) + +(defrule (test name body ...) + (begin + (set! test-count (+ test-count 1)) + (guard (exn [#t (set! fail-count (+ fail-count 1)) + (displayln "FAIL: " name) + (displayln " Error: " + (if (message-condition? exn) + (condition-message exn) + exn))]) + body ... + (set! pass-count (+ pass-count 1)) + (displayln "PASS: " name)))) + +(defrule (assert-equal actual expected) + (let ([a actual] [e expected]) + (unless (equal? a e) + (error 'assert-equal (format "Expected ~s but got ~s" e a))))) + +(defrule (assert-true expr) + (unless expr (error 'assert-true "Expected true"))) + +(defrule (assert-false expr) + (when expr (error 'assert-false "Expected false"))) + +;; Wait up to timeout-ms for (pred) to return truthy, polling every 30 ms. +(def (wait-until pred timeout-ms) + (let loop ([elapsed 0]) + (cond + [(pred) #t] + [(>= elapsed timeout-ms) #f] + [else + (sleep (chez-make-time 'time-duration 30000000 0)) + (loop (+ elapsed 30))]))) + +;; ============================================================ +(displayln "") +(displayln "=== Jerboa-DB Transport Tests (TCP) ===") +;; ============================================================ + +;; Global state for the multi-test cluster. +;; Nodes are started once and shared across tests; stopped at the end. +(def tnode-a #f) (def rconn-a #f) (def port-a #f) +(def tnode-b #f) (def rconn-b #f) (def port-b #f) + +;; ---- 1. Node startup ---- + +(test "start node A (no initial peers, OS-assigned port)" + (let-values ([(t r) (start-transport-db-node! 'node-a '() ":memory:" 0)]) + (set! tnode-a t) + (set! rconn-a r) + (set! port-a (transport-node-listen-port t))) + (assert-true (transport-node? tnode-a)) + (assert-true (> port-a 0))) + +(test "start node B (peer = A), then add B as peer to A" + (let-values ([(t r) (start-transport-db-node! + 'node-b + `((node-a "127.0.0.1" ,port-a)) + ":memory:" 0)]) + (set! tnode-b t) + (set! rconn-b r) + (set! port-b (transport-node-listen-port t))) + ;; Now that we know B's port, tell A about B + (transport-node-add-peer! tnode-a 'node-b "127.0.0.1" port-b) + (assert-true (transport-node? tnode-b)) + (assert-true (> port-b 0))) + +;; ---- 2. Leader election ---- + +(test "a leader is elected within 1500ms" + ;; Give time for TCP connections to establish + Raft election (~150–300 ms) + (let ([got-leader + (wait-until + (lambda () (or (cluster-leader? rconn-a) (cluster-leader? rconn-b))) + 1500)]) + (assert-true got-leader))) + +;; ---- 3. Identify leader/follower ---- + +(def leader #f) +(def follower #f) + +(test "identify leader and follower" + (set! leader (if (cluster-leader? rconn-a) rconn-a rconn-b)) + (set! follower (if (cluster-leader? rconn-a) rconn-b rconn-a)) + (assert-true (cluster-leader? leader)) + (assert-false (cluster-leader? follower))) + +;; ---- 4. Transact through leader ---- + +(test "schema transact through leader" + (cluster-transact! leader + (list + '((db/ident . person/name) + (db/valueType . db.type/string) + (db/cardinality . db.cardinality/one)) + '((db/ident . person/age) + (db/valueType . db.type/long) + (db/cardinality . db.cardinality/one)))) + (assert-true #t)) ;; no exception = pass + +(test "data transact through leader" + (cluster-transact! leader + (list + '((person/name . "Alice") (person/age . 30)) + '((person/name . "Bob") (person/age . 25)))) + (assert-true #t)) + +;; ---- 5. Query from leader ---- + +(test "query from leader returns data" + (let ([results + (q '[(find ?n ?a) + (where (?e person/name ?n) + (?e person/age ?a))] + (cluster-db leader))]) + (assert-equal (length results) 2) + (assert-true (member '("Alice" 30) results)) + (assert-true (member '("Bob" 25) results)))) + +;; ---- 6. Replication to follower ---- + +(test "follower applies transactions within 500ms" + ;; The apply fiber polls every 50 ms; committed entries should appear quickly. + (let ([replicated + (wait-until + (lambda () + (let ([db (cluster-db follower)]) + (not (null? (q '[(find ?n) (where (?e person/name ?n))] + db))))) + 500)]) + (assert-true replicated))) + +(test "query from follower returns same data as leader" + (let ([results + (q '[(find ?n ?a) + (where (?e person/name ?n) + (?e person/age ?a))] + (cluster-db follower))]) + (assert-equal (length results) 2) + (assert-true (member '("Alice" 30) results)) + (assert-true (member '("Bob" 25) results)))) + +;; ---- 7. Additional transact verifies continued replication ---- + +(test "second transact replicates to follower" + (cluster-transact! leader + (list '((person/name . "Carol") (person/age . 35)))) + (let ([replicated + (wait-until + (lambda () + (not (null? (q '[(find ?n) (where (?e person/name "Carol") + (?e person/name ?n))] + (cluster-db follower))))) + 500)]) + (assert-true replicated))) + +;; ---- 8. Cluster status ---- + +(test "cluster-status includes Raft fields and basis-tx" + (let ([status (cluster-status leader)]) + (assert-true (assoc 'role status)) + (assert-true (assoc 'term status)) + (assert-true (assoc 'basis-tx status)) + (assert-true (assoc 'last-applied status)))) + +;; ---- 9. Stop both nodes ---- + +(test "stop both nodes cleanly" + (stop-transport-node! tnode-a) + (stop-transport-node! tnode-b) + (assert-true #t)) + +;; ---- Summary ---- + +(displayln "") +(displayln (str "Results: " pass-count "/" test-count " passed, " + fail-count " failed")) +(when (> fail-count 0) (exit 1))