Step 8 complete: Distributed Computing (Steps 28-30)
ober
1c3d0352ea11d213025f55afe689661c94a97fe8
--- a/Makefile +++ b/Makefile @@ -90,6 +90,7 @@ test-features: @$(SCHEME) --libdirs $(LIBDIRS) --script tests/test-ffi-bind.ss @$(SCHEME) --libdirs $(LIBDIRS) --script tests/test-match2.ss @$(SCHEME) --libdirs $(LIBDIRS) --script tests/test-staging.ss + @$(SCHEME) --libdirs $(LIBDIRS) --script tests/test-cluster.ss test-all: test test-features test-wrappers new file mode 100644 --- /dev/null +++ b/lib/std/actor/cluster.sls @@ -0,0 +1,290 @@ +#!chezscheme +;;; (std actor cluster) — Node Discovery and Distributed Supervision +;;; +;;; Step 28: Node Discovery and Clustering +;;; start-node! — start a cluster node (in-process simulation) +;;; stop-node! — stop a cluster node +;;; node? — predicate +;;; node-name — node name +;;; node-alive? — is the node alive? +;;; cluster-join! — join a cluster +;;; cluster-nodes — list all known nodes +;;; whereis — look up actor in registry, optionally on remote node +;;; +;;; Step 29: Distributed Supervision +;;; make-distributed-supervisor — supervisor managing actors across nodes +;;; dsupervisor-start-child! — start child with placement strategy +;;; dsupervisor-which-children — list children and their node placements +;;; dsupervisor-restart-on-failure — called when a node fails, restarts children +;;; +;;; Note: This is an in-process simulation. Real distributed clustering +;;; would require TCP/UDP sockets and serialization. + +(library (std actor cluster) + (export + ;; Step 28: Node management + start-node! + stop-node! + node? + node-name + node-id + node-alive? + current-node + + ;; Step 28: Cluster operations + cluster-join! + cluster-leave! + cluster-nodes + cluster-node-by-name + on-node-join + on-node-leave + + ;; Step 28: Remote actor registry + remote-register! + remote-unregister! + remote-whereis + whereis/any + + ;; Step 29: Distributed supervision + make-distributed-supervisor + distributed-supervisor? + dsupervisor-start-child! + dsupervisor-which-children + dsupervisor-stop-child! + dsupervisor-handle-node-failure! + + ;; Step 29: Placement strategies + strategy/round-robin + strategy/least-loaded + strategy/local-first) + + (import (chezscheme)) + + ;; ========== Node ========== + + (define-record-type (node make-node-raw node?) + (fields + (immutable name node-name) + (immutable id node-id) + (mutable alive? node-alive? node-set-alive!) + (immutable registry node-registry) ;; eq-hashtable: name → pid + (immutable metadata node-metadata) ;; string-hashtable: key → val + (immutable mutex node-mutex))) + + (define (start-node! name . opts) + (let ([id (string->symbol (format "node-~a-~a" name (time-second (current-time))))] + [meta (make-hashtable equal-hash equal?)]) + ;; Process keyword options + (let loop ([opts opts]) + (unless (null? opts) + (cond + [(eq? (car opts) '#:listen) + (hashtable-set! meta 'listen (cadr opts)) + (loop (cddr opts))] + [(eq? (car opts) '#:cookie) + (hashtable-set! meta 'cookie (cadr opts)) + (loop (cddr opts))] + [(eq? (car opts) '#:seeds) + (hashtable-set! meta 'seeds (cadr opts)) + (loop (cddr opts))] + [else (loop (cdr opts))]))) + (let ([n (make-node-raw name id #t + (make-eq-hashtable) meta (make-mutex))]) + ;; Register globally + (with-mutex *cluster-mutex* + (hashtable-set! *cluster-nodes* id n)) + ;; Notify join hooks + (for-each (lambda (hook) (hook n)) *join-hooks*) + n))) + + (define (stop-node! n) + (with-mutex (node-mutex n) + (node-set-alive! n #f)) + ;; Notify leave hooks + (for-each (lambda (hook) (hook n)) *leave-hooks*) + ;; Remove from global cluster + (with-mutex *cluster-mutex* + (hashtable-delete! *cluster-nodes* (node-id n)))) + + ;; Thread-local current node (simulated via parameter) + (define *current-node* (make-parameter #f)) + (define (current-node) (*current-node*)) + + ;; ========== Cluster ========== + + ;; Global cluster registry + (define *cluster-mutex* (make-mutex)) + (define *cluster-nodes* (make-eq-hashtable)) ;; id → node + (define *join-hooks* '()) + (define *leave-hooks* '()) + + (define (cluster-join! node1 node2) + ;; Simulate joining: node1 discovers node2 (and vice versa) + ;; In real implementation this would exchange membership lists + (with-mutex *cluster-mutex* + (hashtable-set! *cluster-nodes* (node-id node1) node1) + (hashtable-set! *cluster-nodes* (node-id node2) node2))) + + (define (cluster-leave! node) + (stop-node! node)) + + (define (cluster-nodes) + (with-mutex *cluster-mutex* + (let-values ([(ids nodes) (hashtable-entries *cluster-nodes*)]) + (filter node-alive? (vector->list nodes))))) + + (define (cluster-node-by-name name) + (let ([nodes (cluster-nodes)]) + (let loop ([ns nodes]) + (cond + [(null? ns) #f] + [(equal? (node-name (car ns)) name) (car ns)] + [else (loop (cdr ns))])))) + + (define (on-node-join hook) + (set! *join-hooks* (cons hook *join-hooks*))) + + (define (on-node-leave hook) + (set! *leave-hooks* (cons hook *leave-hooks*))) + + ;; ========== Remote Actor Registry ========== + + (define (remote-register! node name pid) + (when (node-alive? node) + (with-mutex (node-mutex node) + (hashtable-set! (node-registry node) name pid)))) + + (define (remote-unregister! node name) + (with-mutex (node-mutex node) + (hashtable-delete! (node-registry node) name))) + + (define (remote-whereis node name) + ;; Look up a named actor on a specific node + (and (node-alive? node) + (with-mutex (node-mutex node) + (hashtable-ref (node-registry node) name #f)))) + + (define (whereis/any name) + ;; Find a named actor on any alive node (returns first found) + (let loop ([nodes (cluster-nodes)]) + (cond + [(null? nodes) #f] + [(remote-whereis (car nodes) name) => values] + [else (loop (cdr nodes))]))) + + ;; ========== Step 29: Distributed Supervisor ========== + + ;; Child spec: (id proc node . restart-type) + ;; restart-type: 'permanent | 'transient | 'temporary + + (define-record-type (distributed-supervisor make-dsup-raw distributed-supervisor?) + (fields + (immutable name dsup-name) + (mutable children dsup-children dsup-set-children!) ;; list of child-info + (immutable strategy dsup-strategy) ;; placement strategy proc + (immutable mutex dsup-mutex))) + + ;; child-info: (id proc assigned-node pid restart-type) + (define (child-info-id c) (list-ref c 0)) + (define (child-info-proc c) (list-ref c 1)) + (define (child-info-node c) (list-ref c 2)) + (define (child-info-pid c) (list-ref c 3)) + (define (child-info-restart-type c) (list-ref c 4)) + + (define (make-distributed-supervisor name strategy) + (make-dsup-raw name '() strategy (make-mutex))) + + (define (dsupervisor-start-child! dsup id proc . opts) + (let* ([restart-type (if (null? opts) 'permanent (car opts))] + [nodes (cluster-nodes)] + [target ((dsup-strategy dsup) nodes dsup id)] + [pid (and target + (node-alive? target) + ;; Simulate starting the actor on target node + (let ([p (fork-thread proc)]) + (remote-register! target id p) + p))] + [child (list id proc target pid restart-type)]) + (with-mutex (dsup-mutex dsup) + (dsup-set-children! dsup + (cons child + (filter (lambda (c) (not (eq? (child-info-id c) id))) + (dsup-children dsup))))) + pid)) + + (define (dsupervisor-which-children dsup) + (with-mutex (dsup-mutex dsup) + (map (lambda (c) + (list (child-info-id c) + (and (child-info-node c) (node-name (child-info-node c))) + (child-info-restart-type c))) + (dsup-children dsup)))) + + (define (dsupervisor-stop-child! dsup id) + (with-mutex (dsup-mutex dsup) + (let ([child (find (lambda (c) (eq? (child-info-id c) id)) + (dsup-children dsup))]) + (when child + (let ([node (child-info-node child)]) + (when node (remote-unregister! node id))) + (dsup-set-children! dsup + (filter (lambda (c) (not (eq? (child-info-id c) id))) + (dsup-children dsup))))))) + + (define (dsupervisor-handle-node-failure! dsup failed-node) + ;; Find all children on the failed node and restart them on survivors + (let ([affected + (with-mutex (dsup-mutex dsup) + (filter (lambda (c) + (and (child-info-node c) + (eq? (node-id (child-info-node c)) + (node-id failed-node)))) + (dsup-children dsup)))]) + ;; Restart permanent and transient children + (for-each + (lambda (child) + (let ([id (child-info-id child)] + [proc (child-info-proc child)] + [rt (child-info-restart-type child)]) + (when (memq rt '(permanent transient)) + (dsupervisor-start-child! dsup id proc rt)))) + affected))) + + ;; ========== Placement Strategies ========== + + ;; strategy/round-robin: place children in rotating order across nodes + (define *rr-counter* 0) + (define *rr-mutex* (make-mutex)) + + (define (strategy/round-robin nodes dsup id) + (if (null? nodes) #f + (with-mutex *rr-mutex* + (let ([idx (remainder *rr-counter* (length nodes))]) + (set! *rr-counter* (+ *rr-counter* 1)) + (list-ref nodes idx))))) + + ;; strategy/least-loaded: pick node with fewest registered actors + (define (strategy/least-loaded nodes dsup id) + (if (null? nodes) #f + (let loop ([best (car nodes)] [best-count (node-load (car nodes))] + [rest (cdr nodes)]) + (if (null? rest) + best + (let ([c (node-load (car rest))]) + (if (< c best-count) + (loop (car rest) c (cdr rest)) + (loop best best-count (cdr rest)))))))) + + (define (node-load n) + (with-mutex (node-mutex n) + (let-values ([(ks _) (hashtable-entries (node-registry n))]) + (vector-length ks)))) + + ;; strategy/local-first: prefer the current node, fall back to round-robin + (define (strategy/local-first nodes dsup id) + (let ([current (current-node)]) + (if (and current (memq current nodes)) + current + (strategy/round-robin nodes dsup id)))) + + ) ;; end library new file mode 100644 --- /dev/null +++ b/lib/std/actor/crdt.sls @@ -0,0 +1,471 @@ +#!chezscheme +;;; (std actor crdt) — Conflict-Free Replicated Data Types +;;; +;;; Step 30: CRDT-Based Distributed State +;;; +;;; CRDTs are data types that can be concurrently updated by multiple nodes +;;; and then merged without conflict. The merge operation (join) is: +;;; - Commutative: merge(a, b) = merge(b, a) +;;; - Associative: merge(a, merge(b, c)) = merge(merge(a, b), c) +;;; - Idempotent: merge(a, a) = a +;;; +;;; Implemented types: +;;; G-Counter — grow-only counter (increment only) +;;; PN-Counter — positive-negative counter (increment and decrement) +;;; OR-Set — observed-remove set (add/remove without tombstones) +;;; LWW-Register — last-write-wins register (timestamp-based) +;;; MV-Register — multi-value register (vector-clock-based, concurrent writes preserved) +;;; G-Set — grow-only set (add only) + +(library (std actor crdt) + (export + ;; G-Counter + make-gcounter + gcounter? + gcounter-increment! + gcounter-value + gcounter-merge! + gcounter-state + + ;; PN-Counter + make-pncounter + pncounter? + pncounter-increment! + pncounter-decrement! + pncounter-value + pncounter-merge! + + ;; G-Set + make-gset + gset? + gset-add! + gset-member? + gset-value + gset-merge! + + ;; OR-Set + make-orset + orset? + orset-add! + orset-remove! + orset-member? + orset-value + orset-merge! + + ;; LWW-Register + make-lww-register + lww-register? + lww-register-set! + lww-register-value + lww-register-timestamp + lww-register-merge! + + ;; MV-Register + make-mv-register + mv-register? + mv-register-set! + mv-register-values ;; returns list of concurrent values + mv-register-merge! + + ;; Vector clock utilities + make-vclock + vclock? + vclock-increment! + vclock-get + vclock-merge! + vclock-happens-before? + vclock-concurrent? + vclock->alist) + + (import (chezscheme)) + + ;; ========== Utilities ========== + + ;; new-uuid is replaced by new-tag below — not used directly + + ;; Better unique tag generator using random + time + (define *tag-counter* 0) + (define *tag-mutex* (make-mutex)) + (define (new-tag) + (with-mutex *tag-mutex* + (set! *tag-counter* (+ *tag-counter* 1)) + (format "~a:~a" (time-second (current-time)) *tag-counter*))) + + (define (alist-merge a b merge-val) + ;; Merge two alists, applying merge-val to values for shared keys + (let loop ([b b] [result a]) + (if (null? b) + result + (let* ([kv (car b)] + [k (car kv)] + [v (cdr kv)] + [existing (assoc k result)]) + (if existing + (loop (cdr b) + (cons (cons k (merge-val (cdr existing) v)) + (filter (lambda (x) (not (equal? (car x) k))) result))) + (loop (cdr b) (cons kv result))))))) + + ;; ========== Vector Clock ========== + + ;; Vector clock: eq-hashtable mapping node-id → integer count + (define-record-type (vclock make-vclock-raw vclock?) + (fields (immutable clock vclock-clock))) + + (define (make-vclock) + (make-vclock-raw (make-eq-hashtable))) + + (define (vclock-increment! vc node-id) + (let ([c (vclock-clock vc)]) + (hashtable-set! c node-id + (+ 1 (hashtable-ref c node-id 0))))) + + (define (vclock-get vc node-id) + (hashtable-ref (vclock-clock vc) node-id 0)) + + (define (vclock-merge! target other) + ;; Merge other into target by taking max of each component + (let-values ([(keys vals) (hashtable-entries (vclock-clock other))]) + (vector-for-each + (lambda (k v) + (let ([current (hashtable-ref (vclock-clock target) k 0)]) + (when (> v current) + (hashtable-set! (vclock-clock target) k v)))) + keys vals))) + + (define (vclock-happens-before? a b) + ;; a happens-before b: every entry in a ≤ corresponding in b, and at least one < + (let-values ([(keys-a vals-a) (hashtable-entries (vclock-clock a))]) + (let ([has-strict #f] + [all-leq #t]) + (vector-for-each + (lambda (k va) + (let ([vb (vclock-get b k)]) + (when (> va vb) (set! all-leq #f)) + (when (< va vb) (set! has-strict #t)))) + keys-a vals-a) + ;; Also check any keys only in b + (and all-leq + (or has-strict + (let-values ([(keys-b _) (hashtable-entries (vclock-clock b))]) + (vector-any + (lambda (k) + (let ([va (vclock-get a k)] + [vb (vclock-get b k)]) + (> vb va))) + keys-b))))))) + + (define (vector-any pred vec) + (let loop ([i 0]) + (cond + [(= i (vector-length vec)) #f] + [(pred (vector-ref vec i)) #t] + [else (loop (+ i 1))]))) + + (define (vclock-concurrent? a b) + ;; Concurrent: neither happens-before the other + (not (or (vclock-happens-before? a b) + (vclock-happens-before? b a) + (vclock-equal? a b)))) + + (define (vclock-equal? a b) + (let-values ([(keys-a vals-a) (hashtable-entries (vclock-clock a))] + [(keys-b vals-b) (hashtable-entries (vclock-clock b))]) + (and (= (vector-length keys-a) (vector-length keys-b)) + (vector-for-all + (lambda (k v) + (= v (vclock-get b k))) + keys-a vals-a)))) + + (define (vector-for-all pred . vecs) + (let ([len (vector-length (car vecs))]) + (let loop ([i 0]) + (cond + [(= i len) #t] + [(apply pred (map (lambda (v) (vector-ref v i)) vecs)) + (loop (+ i 1))] + [else #f])))) + + (define (vclock->alist vc) + (let-values ([(keys vals) (hashtable-entries (vclock-clock vc))]) + (map cons (vector->list keys) (vector->list vals)))) + + ;; ========== G-Counter ========== + ;; + ;; Grow-only counter. Each node has its own counter. + ;; Value = sum of all node counters. + ;; Merge = pairwise max. + + (define-record-type (gcounter make-gcounter-raw gcounter?) + (fields (immutable node-id gcounter-node-id) + (immutable counts gcounter-counts) ;; eq-hashtable: node-id → count + (immutable mutex gcounter-mutex))) + + (define (make-gcounter node-id) + (let ([ht (make-eq-hashtable)]) + (hashtable-set! ht node-id 0) + (make-gcounter-raw node-id ht (make-mutex)))) + + (define (gcounter-increment! gc . args) + (let ([amount (if (null? args) 1 (car args))]) + (with-mutex (gcounter-mutex gc) + (let ([node (gcounter-node-id gc)]) + (hashtable-set! (gcounter-counts gc) node + (+ amount (hashtable-ref (gcounter-counts gc) node 0))))))) + + (define (gcounter-value gc) + (with-mutex (gcounter-mutex gc) + (let-values ([(keys vals) (hashtable-entries (gcounter-counts gc))]) + (vector-fold-right + 0 vals)))) + + (define (vector-fold-right f init vec) + (let loop ([i 0] [acc init]) + (if (= i (vector-length vec)) + acc + (loop (+ i 1) (f (vector-ref vec i) acc))))) + + (define (gcounter-state gc) + (with-mutex (gcounter-mutex gc) + (let-values ([(keys vals) (hashtable-entries (gcounter-counts gc))]) + (map cons (vector->list keys) (vector->list vals))))) + + (define (gcounter-merge! target other) + ;; Merge other's counts into target by taking pairwise max + (with-mutex (gcounter-mutex target) + (let-values ([(keys vals) (hashtable-entries (gcounter-counts other))]) + (vector-for-each + (lambda (k v) + (let ([current (hashtable-ref (gcounter-counts target) k 0)]) + (when (> v current) + (hashtable-set! (gcounter-counts target) k v)))) + keys vals)))) + + ;; ========== PN-Counter ========== + ;; + ;; Positive-Negative counter: two G-Counters (positive, negative). + ;; Value = pos.value - neg.value + ;; Merge = merge both G-Counters separately + + (define-record-type (pncounter make-pncounter-raw pncounter?) + (fields (immutable node-id pncounter-node-id) + (immutable positive pncounter-positive) + (immutable negative pncounter-negative))) + + (define (make-pncounter node-id) + (make-pncounter-raw node-id + (make-gcounter node-id) + (make-gcounter node-id))) + + (define (pncounter-increment! pnc . args) + (gcounter-increment! (pncounter-positive pnc) + (if (null? args) 1 (car args)))) + + (define (pncounter-decrement! pnc . args) + (gcounter-increment! (pncounter-negative pnc) + (if (null? args) 1 (car args)))) + + (define (pncounter-value pnc) + (- (gcounter-value (pncounter-positive pnc)) + (gcounter-value (pncounter-negative pnc)))) + + (define (pncounter-merge! target other) + (gcounter-merge! (pncounter-positive target) (pncounter-positive other)) + (gcounter-merge! (pncounter-negative target) (pncounter-negative other))) + + ;; ========== G-Set ========== + ;; + ;; Grow-only set. Elements can only be added, never removed. + ;; Merge = set union. + + (define-record-type (gset make-gset-raw gset?) + (fields (immutable elements gset-elements) ;; hashtable: elem → #t + (immutable mutex gset-mutex))) + + (define (make-gset) + (make-gset-raw (make-hashtable equal-hash equal?) (make-mutex))) + + (define (gset-add! gs elem) + (with-mutex (gset-mutex gs) + (hashtable-set! (gset-elements gs) elem #t))) + + (define (gset-member? gs elem) + (with-mutex (gset-mutex gs) + (hashtable-ref (gset-elements gs) elem #f))) + + (define (gset-value gs) + (with-mutex (gset-mutex gs) + (let-values ([(keys _) (hashtable-entries (gset-elements gs))]) + (vector->list keys)))) + + (define (gset-merge! target other) + (with-mutex (gset-mutex target) + (let-values ([(keys vals) (hashtable-entries (gset-elements other))]) + (vector-for-each + (lambda (k v) + (hashtable-set! (gset-elements target) k v)) + keys vals)))) + + ;; ========== OR-Set ========== + ;; + ;; Observed-Remove Set. Each element is tagged with unique IDs. + ;; Add: add (elem, tag) to 'added' set. + ;; Remove: remove all tags for elem from 'added' set. + ;; Member: elem has at least one tag in 'added' not in 'removed'. + ;; Merge: union of added sets, union of removed sets. + + (define-record-type (orset make-orset-raw orset?) + (fields (immutable added orset-added) ;; hashtable: (elem . tag) → #t + (immutable removed orset-removed) ;; hashtable: (elem . tag) → #t + (immutable mutex orset-mutex))) + + (define (make-orset) + (make-orset-raw (make-hashtable equal-hash equal?) (make-hashtable equal-hash equal?) (make-mutex))) + + (define (orset-add! os elem) + (with-mutex (orset-mutex os) + (let ([tag (new-tag)]) + (hashtable-set! (orset-added os) (cons elem tag) #t)))) + + (define (orset-remove! os elem) + (with-mutex (orset-mutex os) + ;; Move all tags for this element from added to removed + (let-values ([(pairs _) (hashtable-entries (orset-added os))]) + (vector-for-each + (lambda (pair) + (when (equal? (car pair) elem) + (hashtable-delete! (orset-added os) pair) + (hashtable-set! (orset-removed os) pair #t))) + pairs)))) + + (define (orset-member? os elem) + (with-mutex (orset-mutex os) + (let-values ([(pairs _) (hashtable-entries (orset-added os))]) + (vector-any + (lambda (pair) (equal? (car pair) elem)) + pairs)))) + + (define (orset-value os) + (with-mutex (orset-mutex os) + (let ([seen (make-hashtable equal-hash equal?)]) + (let-values ([(pairs _) (hashtable-entries (orset-added os))]) + (vector-for-each + (lambda (pair) + (hashtable-set! seen (car pair) #t)) + pairs)) + (let-values ([(elems _) (hashtable-entries seen)]) + (vector->list elems))))) + + (define (orset-merge! target other) + ;; Union of added, union of removed, then subtract removed from added + (with-mutex (orset-mutex target) + ;; Union removed + (let-values ([(pairs _) (hashtable-entries (orset-removed other))]) + (vector-for-each + (lambda (pair) + (hashtable-set! (orset-removed target) pair #t)) + pairs)) + ;; Union added (not already removed) + (let-values ([(pairs _) (hashtable-entries (orset-added other))]) + (vector-for-each + (lambda (pair) + (unless (hashtable-ref (orset-removed target) pair #f) + (hashtable-set! (orset-added target) pair #t))) + pairs)) + ;; Remove any added entries that are in removed + (let-values ([(pairs _) (hashtable-entries (orset-removed target))]) + (vector-for-each + (lambda (pair) + (hashtable-delete! (orset-added target) pair)) + pairs)))) + + ;; ========== LWW-Register ========== + ;; + ;; Last-Write-Wins Register. Stores a single value with a timestamp. + ;; On merge, the higher-timestamp value wins. + + (define-record-type (lww-register make-lww-raw lww-register?) + (fields (mutable value lww-value lww-set-value!) + (mutable timestamp lww-timestamp lww-set-timestamp!) + (immutable mutex lww-mutex))) + + (define (make-lww-register) + (make-lww-raw #f -inf.0 (make-mutex))) + + (define (lww-register-value r) + (with-mutex (lww-mutex r) (lww-value r))) + + (define (lww-register-timestamp r) + (with-mutex (lww-mutex r) (lww-timestamp r))) + + (define (lww-register-set! r val . args) + (let ([ts (if (null? args) + (inexact (time-second (current-time))) + (car args))]) + (with-mutex (lww-mutex r) + (when (> ts (lww-timestamp r)) + (lww-set-value! r val) + (lww-set-timestamp! r ts))))) + + (define (lww-register-merge! target other) + (with-mutex (lww-mutex target) + (let ([other-ts (lww-timestamp other)] + [other-val (lww-value other)]) + (when (> other-ts (lww-timestamp target)) + (lww-set-value! target other-val) + (lww-set-timestamp! target other-ts))))) + + ;; ========== MV-Register ========== + ;; + ;; Multi-Value Register. Uses vector clocks to track causality. + ;; Concurrent writes are both preserved (unlike LWW which picks one). + ;; Reading returns a list of all concurrent values. + + (define-record-type (mv-register make-mv-raw mv-register?) + (fields (mutable entries mv-entries mv-set-entries!) ;; list of (vclock . value) + (immutable mutex mv-mutex))) + + (define (make-mv-register) + (make-mv-raw '() (make-mutex))) + + (define (mv-register-values r) + (with-mutex (mv-mutex r) + (map cdr (mv-entries r)))) + + (define (mv-register-set! r node-id val) + (with-mutex (mv-mutex r) + ;; Create new vector clock by merging all current clocks and incrementing + (let ([new-vc (make-vclock)]) + ;; Merge all current vector clocks + (for-each + (lambda (entry) + (vclock-merge! new-vc (car entry))) + (mv-entries r)) + ;; Increment for this node + (vclock-increment! new-vc node-id) + ;; Replace all entries dominated by new-vc with single new entry + (let ([surviving + (filter + (lambda (entry) + ;; Keep if not dominated by new-vc + (not (vclock-happens-before? (car entry) new-vc))) + (mv-entries r))]) + (mv-set-entries! r (cons (cons new-vc val) surviving)))))) + + (define (mv-register-merge! target other) + (with-mutex (mv-mutex target) + (let ([target-entries (mv-entries target)] + [other-entries (mv-entries other)]) + ;; Merge: keep entries not dominated by any entry in the other set + (let* ([all (append target-entries other-entries)] + [merged + (filter + (lambda (e1) + (not (exists + (lambda (e2) + (and (not (eq? e1 e2)) + (vclock-happens-before? (car e1) (car e2)))) + all))) + all)]) + (mv-set-entries! target merged))))) + + ) ;; end library new file mode 100644 --- /dev/null +++ b/tests/test-cluster.ss @@ -0,0 +1,402 @@ +#!chezscheme +;;; Tests for (std actor cluster) and (std actor crdt) + +(import (chezscheme) (std actor crdt) (std actor cluster)) + +(define pass 0) +(define fail 0) + +(define-syntax test + (syntax-rules () + [(_ name expr expected) + (guard (exn [#t (set! fail (+ fail 1)) + (printf "FAIL ~a: ~a~%" name + (if (message-condition? exn) (condition-message exn) exn))]) + (let ([got expr]) + (if (equal? got expected) + (begin (set! pass (+ pass 1)) (printf " ok ~a~%" name)) + (begin (set! fail (+ fail 1)) + (printf "FAIL ~a: got ~s expected ~s~%" name got expected)))))])) + +(printf "--- (std actor crdt) + (std actor cluster) tests ---~%") + +;;; ======== G-Counter ======== + +(printf "~%-- G-Counter --~%") + +(let ([gc (make-gcounter 'node1)]) + (test "gcounter: initial value" + (gcounter-value gc) + 0) + + (gcounter-increment! gc) + (gcounter-increment! gc) + (test "gcounter: after 2 increments" + (gcounter-value gc) + 2) + + (gcounter-increment! gc 5) + (test "gcounter: increment by 5" + (gcounter-value gc) + 7) + + ;; Merge two G-Counters + (let ([gc2 (make-gcounter 'node2)]) + (gcounter-increment! gc2 10) + (gcounter-merge! gc gc2) + (test "gcounter: merge from another node" + (gcounter-value gc) + 17) + + ;; Idempotent merge + (gcounter-merge! gc gc2) + (test "gcounter: merge is idempotent" + (gcounter-value gc) + 17) + + ;; Commutative merge + (let ([gc3 (make-gcounter 'node1)] + [gc4 (make-gcounter 'node2)]) + (gcounter-increment! gc3 3) + (gcounter-increment! gc4 7) + (gcounter-merge! gc3 gc4) + (let ([merged-3-then-4 (gcounter-value gc3)]) + (let ([gc5 (make-gcounter 'node2)] + [gc6 (make-gcounter 'node1)]) + (gcounter-increment! gc5 7) + (gcounter-increment! gc6 3) + (gcounter-merge! gc5 gc6) + (test "gcounter: merge is commutative" + merged-3-then-4 + (gcounter-value gc5)))))) + + (test "gcounter: state is alist" + (pair? (gcounter-state gc)) + #t)) + +;;; ======== PN-Counter ======== + +(printf "~%-- PN-Counter --~%") + +(let ([pnc (make-pncounter 'node1)]) + (test "pncounter: initial value" + (pncounter-value pnc) + 0) + + (pncounter-increment! pnc 5) + (pncounter-decrement! pnc 2) + (test "pncounter: increment 5, decrement 2" + (pncounter-value pnc) + 3) + + (let ([pnc2 (make-pncounter 'node2)]) + (pncounter-increment! pnc2 10) + (pncounter-merge! pnc pnc2) + (test "pncounter: merge" + (pncounter-value pnc) + 13))) + +;;; ======== G-Set ======== + +(printf "~%-- G-Set --~%") + +(let ([gs (make-gset)]) + (test "gset: initial empty" + (gset-value gs) + '()) + + (gset-add! gs 'apple) + (gset-add! gs 'banana) + (test "gset: member after add" + (gset-member? gs 'apple) + #t) + + (test "gset: non-member" + (gset-member? gs 'cherry) + #f) + + (let ([gs2 (make-gset)]) + (gset-add! gs2 'cherry) + (gset-merge! gs gs2) + (test "gset: after merge has cherry" + (gset-member? gs 'cherry) + #t) + + (test "gset: merge is idempotent" + (begin (gset-merge! gs gs2) (gset-member? gs 'cherry)) + #t)) + + (test "gset: all values present" + (list-sort (lambda (a b) (string<? (symbol->string a) (symbol->string b))) + (gset-value gs)) + '(apple banana cherry))) + +;;; ======== OR-Set ======== + +(printf "~%-- OR-Set --~%") + +(let ([os (make-orset)]) + (test "orset: initial not member" + (orset-member? os 'x) + #f) + + (orset-add! os 'x) + (test "orset: member after add" + (orset-member? os 'x) + #t) + + (orset-remove! os 'x) + (test "orset: not member after remove" + (orset-member? os 'x) + #f) + + ;; Add-wins: concurrent add and remove — add survives in OR-Set + ;; Simulate: node A removes x, node B adds x concurrently + (let ([os-a (make-orset)] + [os-b (make-orset)]) + (orset-add! os-a 'x) + (orset-add! os-b 'x) + ;; os-a removes x, os-b doesn't know about it + (orset-remove! os-a 'x) + ;; Merge: os-b's add (which os-a didn't know about) survives + (orset-merge! os-a os-b) + (test "orset: add-wins on concurrent add/remove" + (orset-member? os-a 'x) + #t))) + +;;; ======== LWW-Register ======== + +(printf "~%-- LWW-Register --~%") + +(let ([r (make-lww-register)]) + (test "lww: initial value is #f" + (lww-register-value r) + #f) + + (lww-register-set! r 'hello 100.0) + (test "lww: value after set" + (lww-register-value r) + 'hello) + + (lww-register-set! r 'world 200.0) + (test "lww: newer value wins" + (lww-register-value r) + 'world) + + (lww-register-set! r 'old 50.0) + (test "lww: older value ignored" + (lww-register-value r) + 'world) + + (let ([r2 (make-lww-register)]) + (lww-register-set! r2 'newest 999.0) + (lww-register-merge! r r2) + (test "lww: merge takes newer" + (lww-register-value r) + 'newest) + + ;; Idempotent + (lww-register-merge! r r2) + (test "lww: merge idempotent" + (lww-register-value r) + 'newest))) + +;;; ======== MV-Register ======== + +(printf "~%-- MV-Register --~%") + +(let ([r (make-mv-register)]) + (test "mv: initial values empty" + (mv-register-values r) + '()) + + (mv-register-set! r 'node1 42) + (test "mv: single value" + (mv-register-values r) + '(42)) + + ;; Sequential update: overwrites + (mv-register-set! r 'node1 99) + (test "mv: sequential update replaces" + (mv-register-values r) + '(99)) +