From 08446ef7b6c865fd1d985d9a4d9961bc7e971229 Mon Sep 17 00:00:00 2001 From: "russell@unturf.com" Date: Fri, 5 Jun 2026 18:18:29 -0400 Subject: [PATCH] examples/cuda-fanout: bend multi-worker fan-out MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds *bend-workers* list with round-robin dispatch, BEND_WORKERS env loader, and helpers (bend-set-workers!, bend-pick-worker, bend-parse-workers-env, bend-load-workers-from-env!). Single-host legacy callers unaffected — when *bend-workers* is empty the dispatcher falls back to *bend-worker-host* / *bend-worker-port*. Validated cross-tier (Python + C lumbda) against a live two-host cluster (3090-ai.foxhop.net:9091, ai.foxhop.net:9092) — round-robin distributes evenly; cgbn mod-mul results byte-identical to gmpy2 reference on both hosts. --- examples/cuda-fanout/bend.lsp | 103 ++++++++++++++++++++++++++++++---- 1 file changed, 92 insertions(+), 11 deletions(-) diff --git a/examples/cuda-fanout/bend.lsp b/examples/cuda-fanout/bend.lsp index d52cc17..afa8d03 100644 --- a/examples/cuda-fanout/bend.lsp +++ b/examples/cuda-fanout/bend.lsp @@ -73,32 +73,97 @@ (define *bend-worker-host* "127.0.0.1") (define *bend-worker-port* 9091) -;; Override the default localhost:9091 endpoint. +;; Override our default localhost:9091 endpoint. (define (bend-set-worker! host port) (set! *bend-worker-host* host) (set! *bend-worker-port* port)) -;; Probe TCP connect; return #t/#f without raising. +;;; -- multi-worker fleet ---------------------------------------- +;;; +;;; A list of (host . port) pairs. When non-empty, bend-dispatch-to-gpu +;;; rotates round-robin across our cluster; a failed pick falls forward +;;; to a peer worker on a tcp-connect error. When empty, dispatcher +;;; falls back to a single *bend-worker-host* / *bend-worker-port* pair +;;; (full back-compat with single-host callers). +;;; +;;; Set via (bend-set-workers! '(("3090-ai.foxhop.net" . 9091) +;;; ("ai.foxhop.net" . 9092))) +;;; or environment variable BEND_WORKERS="host:port,host:port". + +(define *bend-workers* '()) +(define *bend-rr-idx* 0) + +(define (bend-set-workers! lst) + (set! *bend-workers* lst) + (set! *bend-rr-idx* 0)) + +;; Round-robin pick a worker from *bend-workers*. Returns (host . port). +;; Caller asserts *bend-workers* non-empty. +(define (bend-pick-worker) + (let* ((n (length *bend-workers*)) + (idx (if (= n 0) 0 (remainder *bend-rr-idx* n))) + (pick (list-ref *bend-workers* idx))) + (set! *bend-rr-idx* (+ *bend-rr-idx* 1)) + pick)) + +;; Parse "host:port,host:port" into a list of (host . port) pairs. Skips +;; malformed entries silently rather than raising — keeps boot path resilient. +(define (bend-parse-workers-env s) + (let loop ((rest s) (acc '()) (cur "")) + (cond + ((= (string-length rest) 0) + (let ((parsed (bend-parse-one-worker cur))) + (reverse (if parsed (cons parsed acc) acc)))) + ((string=? (substring rest 0 1) ",") + (let ((parsed (bend-parse-one-worker cur))) + (loop (substring rest 1 (string-length rest)) + (if parsed (cons parsed acc) acc) ""))) + (else + (loop (substring rest 1 (string-length rest)) acc + (string-append cur (substring rest 0 1))))))) + +(define (bend-parse-one-worker s) + ;; Split on first ':'; return (host . port-number) or #f if malformed. + (let loop ((i 0)) + (cond + ((>= i (string-length s)) #f) + ((string=? (substring s i (+ i 1)) ":") + (let ((host (substring s 0 i)) + (port (string->number (substring s (+ i 1) (string-length s))))) + (cond + ((or (= (string-length host) 0) (eq? port #f)) #f) + (else (cons host port))))) + (else (loop (+ i 1)))))) + +;; Probe TCP connect to the next pick OR the legacy single host. +;; Returns #t/#f without raising. (define (bend-worker-available?) - (let ((sock (tcp-connect *bend-worker-host* *bend-worker-port*))) - (cond - ((eq? sock #f) #f) - (else (tcp-close sock) #t)))) + (let ((target (cond + ((null? *bend-workers*) + (cons *bend-worker-host* *bend-worker-port*)) + (else (car *bend-workers*))))) ; cheap reachability probe + (let ((sock (tcp-connect (car target) (cdr target)))) + (cond + ((eq? sock #f) #f) + (else (tcp-close sock) #t))))) -;; Open a fresh connection, send the form framed, read framed reply, +;; Open a fresh connection, send our form framed, read framed reply, -;; close. Returns the result from the worker, or raises if the worker +;; close. Returns our result from the worker, or raises if the worker ;; responded with (bend-error...). (define (bend-dispatch-to-gpu quoted-form) - (let ((sock (tcp-connect *bend-worker-host* *bend-worker-port*))) + (let* ((target (cond + ((null? *bend-workers*) + (cons *bend-worker-host* *bend-worker-port*)) + (else (bend-pick-worker)))) + (sock (tcp-connect (car target) (cdr target)))) (cond ((eq? sock #f) - (bend-error "bend: tcp-connect failed to" - (list *bend-worker-host* *bend-worker-port*))) + (bend-error "bend: tcp-connect failed to" target)) (else (wire-send sock quoted-form) (let ((reply (wire-recv sock))) @@ -111,6 +176,22 @@ (bend-error "bend worker error:" (cdr reply))) (else (bend-error "bend: unexpected reply" reply)))))))) +;; Optional boot-time hook: if BEND_WORKERS is set in env, parse it now. +;; Safe to call repeatedly; a missing var is a no-op. +(define (bend-load-workers-from-env!) + (let ((s (get-environment-variable "BEND_WORKERS"))) + (cond + ((or (eq? s #f) (= (string-length s) 0)) #f) + (else + (let ((lst (bend-parse-workers-env s))) + (cond + ((null? lst) #f) + (else + (bend-set-workers! lst) + (display ";;; bend: loaded ") (display (length lst)) + (display " workers from BEND_WORKERS") (newline) + #t))))))) + ;;; -- core dispatcher ------------------------------------------- ;; thunk: 0-arg lambda that evaluates the form in its original scope.