feat: implement Phase 2 — sandboxed filters, hot handlers, retry/dead-letter
ober
3a656ae7015722ed65b7008ee54e86a045bab99a
--- a/edge.ss +++ b/edge.ss @@ -33,7 +33,9 @@ (std transducer) (std text json) (std crypto native) - (std pmap)) + (std pmap) + (std security restrict) + (std misc timeout)) ;; ═══════════════════════════════════════════════════════════════ ;; Configuration @@ -62,8 +64,12 @@ ;; atomic and produces an immutable snapshot — no locks, no races. ;; ═══════════════════════════════════════════════════════════════ -(def events-ref (make-ref (make-persistent-map))) -(def watchers-ref (make-ref '())) +(def events-ref (make-ref (make-persistent-map))) +(def watchers-ref (make-ref '())) +(def filters-ref (make-ref (make-persistent-map))) ;; user-defined event filters +(def dead-letter-ref (make-ref (make-persistent-map))) ;; events exhausted all retries + +(def *max-retries* 3) ;; retries before dead-letter: 1s / 2s / 4s backoff (def (store-result! event status) (dosync @@ -92,7 +98,29 @@ (set! err-n (+ err-n 1)))) st) (list->hash-table - `(("total" . ,total) ("ok" . ,ok-n) ("error" . ,err-n))))) + `(("total" . ,total) + ("ok" . ,ok-n) + ("error" . ,err-n) + ("dead_letter" . ,(persistent-map-size (ref-deref dead-letter-ref))))))) + +(def (store-dead-letter! event reason) + (dosync + (alter dead-letter-ref + (lambda (st) + (persistent-map-set st (hash-ref event "id") + (list->hash-table + `(("id" . ,(hash-ref event "id")) + ("type" . ,(hash-ref event "type" "unknown")) + ("reason" . ,reason) + ("retries" . ,(hash-ref event "retry-count" 0)) + ("time" . ,(format "~a" (time-second (current-time))))))))))) + +(def (get-dead-letter-list) + (let ([result '()]) + (persistent-map-for-each + (lambda (id record) (set! result (cons record result))) + (ref-deref dead-letter-ref)) + result)) ;; ═══════════════════════════════════════════════════════════════ ;; HMAC-SHA256 Signature Verification @@ -120,6 +148,58 @@ (string->utf8 sig)))))) ;; ═══════════════════════════════════════════════════════════════ +;; Sandboxed User Filters (Phase 2.1) +;; +;; Users submit filter expressions via POST /api/filters. +;; Each filter is evaluated per-event using a restricted environment: +;; - Allowlist-only (~60 safe bindings: arithmetic, list ops, predicates) +;; - No I/O, no file access, no network, no FFI +;; - 2-second engine-based timeout (no fork — safe from fibers) +;; +;; Filter code receives `evt` as a nested alist. Top-level event keys: +;; id, type, payload, received, status. Payload is a nested alist. +;; +;; (equal? (cdr (assoc "type" evt)) "payment.completed") +;; (let ([p (cdr (assoc "payload" evt))]) +;; (> (cdr (assoc "amount" p)) 1000)) +;; +;; Errors and timeouts → #t (permissive — broken filter never drops all events). +;; ═══════════════════════════════════════════════════════════════ + +;; Recursively convert hash tables → alists so event data can be +;; passed to the restricted evaluator as plain quoted data. +(def (hash->alist-deep x) + (cond + [(hash-table? x) + (map (lambda (kv) (cons (car kv) (hash->alist-deep (cdr kv)))) + (hash->list x))] + [(pair? x) (cons (hash->alist-deep (car x)) (hash->alist-deep (cdr x)))] + [else x])) + +;; Evaluate one filter expression against an event in a restricted +;; environment (allowlist-only: no I/O, no FFI, no network). +;; Uses Chez engine-based timeout — no fork, safe to call from fibers. +;; Returns #t (pass) or #f (drop). Errors and timeouts → #t (permissive). +(def (eval-user-filter filter-code event) + (let* ([evt-data (hash->alist-deep event)] + [code (format "(let ([evt '~s]) ~a)" evt-data filter-code)]) + (guard (exn [#t + (displayln "[filter] eval error: " + (if (message-condition? exn) (condition-message exn) exn)) + #t]) ;; permissive on error + ;; with-timeout uses Chez engines (no fork) — fiber-safe + (let ([result (with-timeout 2.0 #t + (lambda () (restricted-eval-string code)))]) + (and result #t))))) + +;; Apply all registered user filters to an event. +;; Returns #t only if every filter accepts the event. +(def (apply-user-filters event) + (let ([filters (persistent-map-values (ref-deref filters-ref))]) + (every (lambda (f) (eval-user-filter (hash-ref f "code") event)) + filters))) + +;; ═══════════════════════════════════════════════════════════════ ;; WebSocket Dashboard ;; ;; Connected browsers receive live event-processed notifications. @@ -177,10 +257,33 @@ (lambda (evt) (error 'worker "deliberate crash for demo"))) +;; Restricted environment for hot handlers — adds displayln on top of safe-bindings. +;; Handlers run in worker OS threads so this is fiber-safe. +(def handler-env + (make-restricted-environment + (list (cons 'displayln + (lambda args + (for-each display args) + (newline)))))) + (def (dispatch-handler type) - (or (hash-get handler-registry type) - (lambda (evt) - (displayln "[" type "] " (hash-ref evt "id"))))) + (let ([h (hash-get handler-registry type)]) + (cond + ;; Pre-built native handler (Scheme procedure) + [(procedure? h) h] + ;; Hot-registered code string — run in restricted env per invocation. + ;; Restricted env: safe-bindings + displayln, 5s engine timeout. + ;; No inner guard — exceptions propagate to worker's try/catch for retry. + [(string? h) + (lambda (evt) + (let* ([evt-data (hash->alist-deep evt)] + [code (format "(let ([evt '~s]) ~a)" evt-data h)]) + (with-timeout 5.0 (void) + (lambda () (restricted-eval-string code handler-env)))))] + ;; Unknown type → log and continue + [else + (lambda (evt) + (displayln "[" type "] " (hash-ref evt "id")))]))) ;; ═══════════════════════════════════════════════════════════════ ;; Transducer Pipeline @@ -198,12 +301,15 @@ (lambda (evt) (and (hash-key? evt "type") (hash-key? evt "payload")))) - ;; Deduplicate: drop events with IDs we've already seen + ;; Deduplicate: drop events with IDs we've already seen. + ;; Retried events (retry-count > 0) bypass dedup intentionally. (filtering (lambda (evt) - (let ([id (hash-ref evt "id" #f)]) + (let ([id (hash-ref evt "id" #f)] + [retry (hash-ref evt "retry-count" 0)]) (cond - [(not id) #t] ;; no ID → pass through + [(not id) #t] ;; no ID → pass through + [(> retry 0) #t] ;; retry → always pass [(hash-key? seen-ids id) (displayln "[dedup] dropping duplicate " id) #f] [else (hash-put! seen-ids id #t) #t])))) @@ -228,6 +334,17 @@ ;; crashes, the supervisor restarts it — others continue unaffected. ;; ═══════════════════════════════════════════════════════════════ +;; Schedule a retry: sleep N seconds in a background thread then re-enqueue. +;; Using fork-thread so the calling worker is not blocked. +(def (schedule-retry! event delay-secs worker-id) + (fork-thread + (lambda () + (sleep (make-time 'time-duration 0 delay-secs)) + (displayln "[worker-" worker-id "] re-enqueuing " + (hash-ref event "id") " (retry " + (hash-ref event "retry-count" 0) ")") + (>!! ingest-ch event)))) + (def (make-worker id) (lambda (msg) (match msg @@ -236,24 +353,61 @@ (let loop () (let ([event (<!! ingest-ch)]) (when (and event (not (eof-object? event))) - (let* ([type (hash-ref event "type" "unknown")] - [handler (dispatch-handler type)] - [result (try (ok (handler event)) - (catch (e) (err e)))] - [status (if (ok? result) "ok" "error")]) - (when (err? result) - (let ([e (unwrap-err result)]) - (displayln "[worker-" id "] error: " - (if (message-condition? e) - (condition-message e) e)))) - (store-result! event status) - (notify-watchers! event status) - ;; Deliberate crash: re-raise to kill actor. - ;; Supervisor will restart it automatically. - (when (and (err? result) (string=? type "test.crash")) - (displayln "[worker-" id "] crashing — supervisor will restart") - (raise (unwrap-err result)))) - (loop))))] + (cond + ;; User filter rejected this event — drop and continue + [(not (apply-user-filters event)) + (displayln "[worker-" id "] dropping " (hash-ref event "id") + " (failed user filter)")] + ;; Process the event + [else + (let* ([type (hash-ref event "type" "unknown")] + [retry-count (hash-ref event "retry-count" 0)] + [handler (dispatch-handler type)] + [result (try (ok (handler event)) + (catch (e) (err e)))] + [status (if (ok? result) "ok" "error")]) + (cond + ;; Success path + [(ok? result) + (store-result! event status) + (notify-watchers! event status)] + ;; Deliberate crash demo — re-raise to prove supervisor restarts + [(string=? type "test.crash") + (store-result! event status) + (notify-watchers! event status) + (displayln "[worker-" id "] crashing — supervisor will restart") + (raise (unwrap-err result))] + ;; Retry with exponential backoff: 1s / 2s / 4s + [(< retry-count *max-retries*) + (let* ([e (unwrap-err result)] + [err-msg (if (message-condition? e) (condition-message e) (format "~a" e))] + [delay-secs (expt 2 retry-count)] + [next-event (begin (hash-put! event "retry-count" (+ retry-count 1)) event)]) + (displayln "[worker-" id "] " type " failed (" err-msg + ") — retry " (+ retry-count 1) "/" *max-retries* + " in " delay-secs "s") + (schedule-retry! next-event delay-secs id) + (store-result! event "retrying") + (notify-watchers! event "retrying"))] + ;; All retries exhausted → dead letter + [else + (let* ([e (unwrap-err result)] + [err-msg (if (message-condition? e) (condition-message e) (format "~a" e))]) + (displayln "[worker-" id "] " type " dead-lettered after " + retry-count " retries: " err-msg) + (store-dead-letter! event err-msg) + (store-result! event "dead-letter") + (notify-watchers! event "dead-letter")) ;; close dead-letter let* + ] ;; close dead-letter [else] + ) ;; close inner (cond) + ) ;; close outer (let*) + ] ;; close outer [else] + ) ;; close outer (cond) + (loop) + ) ;; close (when) + ) ;; close (let ([event ...])) + ) ;; close (let loop ()) + ] ;; close ['start ...] clause [_ (void)]))) (def (start-workers n) @@ -332,35 +486,189 @@ (unregister-watcher! ws)))))) ;; ═══════════════════════════════════════════════════════════════ +;; Phase 2 HTTP Handlers +;; ═══════════════════════════════════════════════════════════════ + +;; ── Filters ───────────────────────────────────────────────────── + +;; GET /api/filters — list registered user filters +(def (handle-list-filters req) + (let ([items '()]) + (persistent-map-for-each + (lambda (name f) (set! items (cons f items))) + (ref-deref filters-ref)) + (respond-json 200 + (format "[~a]" + (string-join + (map json-object->string items) + ","))))) + +;; POST /api/filters — register or replace a user filter +;; Body: {"name": "big-payments", "code": "(> (cdr (assoc \"amount\" evt)) 1000)"} +(def (handle-add-filter req) + (let* ([body (or (request-body req) "{}")] + [data (guard (exn [#t (make-hash-table)]) + (string->json-object body))] + [name (hash-get data "name")] + [code (hash-get data "code")]) + (cond + [(not name) + (respond-json 400 "{\"error\":\"missing field: name\"}")] + [(not code) + (respond-json 400 "{\"error\":\"missing field: code\"}")] + [else + (dosync + (alter filters-ref + (lambda (st) + (persistent-map-set st name + (list->hash-table `(("name" . ,name) ("code" . ,code))))))) + (displayln "[filters] registered: " name) + (respond-json 201 + (json-object->string + (list->hash-table `(("status" . "registered") ("name" . ,name)))))]))) + +;; DELETE /api/filters/:name — remove a user filter +(def (handle-delete-filter req) + (let ([name (route-param req "name")]) + (if name + (begin + (dosync + (alter filters-ref + (lambda (st) (persistent-map-delete st name)))) + (displayln "[filters] removed: " name) + (respond-json 200 + (json-object->string + (list->hash-table `(("status" . "removed") ("name" . ,name)))))) + (respond-json 400 "{\"error\":\"missing name\"}")))) + +;; ── Handlers ──────────────────────────────────────────────────── + +;; POST /api/handlers — register a hot handler for a webhook type +;; Body: {"type": "order.shipped", "code": "(displayln \"shipped: \" (cdr (assoc \"id\" evt)))"} +;; Code is evaluated per-event in a sandbox (5s timeout). +(def (handle-register-handler req) + (let* ([body (or (request-body req) "{}")] + [data (guard (exn [#t (make-hash-table)]) + (string->json-object body))] + [type (hash-get data "type")] + [code (hash-get data "code")]) + (cond + [(not type) + (respond-json 400 "{\"error\":\"missing field: type\"}")] + [(not code) + (respond-json 400 "{\"error\":\"missing field: code\"}")] + [else + ;; Store as a code string; dispatch-handler will sandbox-eval it + (hash-put! handler-registry type code) + (displayln "[handlers] registered: " type) + (respond-json 201 + (json-object->string + (list->hash-table `(("status" . "registered") + ("type" . ,type) + ("sandboxed". #t)))))]))) + +;; GET /api/handlers — list registered handler types +(def (handle-list-handlers req) + (let ([types '()]) + (hash-for-each + (lambda (type h) + (set! types (cons + (list->hash-table `(("type" . ,type) + ("kind" . ,(if (procedure? h) "native" "sandboxed")))) + types))) + handler-registry) + (respond-json 200 + (format "[~a]" + (string-join (map json-object->string types) ","))))) + +;; ── Dead Letter ───────────────────────────────────────────────── + +;; GET /api/dead-letter — list events that exhausted all retries +(def (handle-list-dead-letter req) + (let ([items (get-dead-letter-list)]) + (respond-json 200 + (format "[~a]" + (string-join (map json-object->string items) ","))))) + +;; POST /api/dead-letter/:id/replay — re-enqueue a dead-letter event +(def (handle-replay-dead-letter req) + (let* ([id (route-param req "id")] + [record (and id + (persistent-map-has? (ref-deref dead-letter-ref) id) + (persistent-map-ref (ref-deref dead-letter-ref) id))]) + (cond + [(not id) + (respond-json 400 "{\"error\":\"missing id\"}")] + [(not record) + (respond-json 404 "{\"error\":\"not found in dead-letter store\"}")] + [else + ;; Re-enqueue with reset retry count + (let ([replay-event + (list->hash-table + `(("id" . ,id) + ("type" . ,(hash-ref record "type" "unknown")) + ("payload" . ,(make-hash-table)) + ("received" . ,(format "~a" (time-second (current-time)))) + ("retry-count" . 1) ;; >0 so dedup passes it through + ("status" . "replayed")))]) + ;; Remove from dead-letter + (dosync + (alter dead-letter-ref + (lambda (st) (persistent-map-delete st id)))) + (>!! ingest-ch replay-event) + (displayln "[dead-letter] replaying " id) + (respond-json 202 + (json-object->string + (list->hash-table `(("status" . "replayed") ("id" . ,id))))))]))) + +;; ═══════════════════════════════════════════════════════════════ ;; Router & Server ;; ═══════════════════════════════════════════════════════════════ (def (run-edge) (let ([r (make-router)]) ;; Webhook ingestion - (route-post r "/hooks/:type" handle-webhook) + (route-post r "/hooks/:type" handle-webhook) ;; Query API - (route-get r "/api/events/:id" handle-get-event) - (route-get r "/api/stats" handle-stats) + (route-get r "/api/events/:id" handle-get-event) + (route-get r "/api/stats" handle-stats) + ;; Phase 2: User filters + (route-get r "/api/filters" handle-list-filters) + (route-post r "/api/filters" handle-add-filter) + (route-delete r "/api/filters/:name" handle-delete-filter) + ;; Phase 2: Hot handler registration + (route-get r "/api/handlers" handle-list-handlers) + (route-post r "/api/handlers" handle-register-handler) + ;; Phase 2: Dead-letter queue + (route-get r "/api/dead-letter" handle-list-dead-letter) + (route-post r "/api/dead-letter/:id/replay" handle-replay-dead-letter) ;; Dashboard + health - (route-get r "/dashboard" handle-dashboard) - (route-get r "/health" handle-health) + (route-get r "/dashboard" handle-dashboard) + (route-get r "/health" handle-health) ;; Banner first, then start workers (displayln "") - (displayln " Jerboa Edge v0.1.0") - (displayln " ──────────────────────────────────────") + (displayln " Jerboa Edge v0.2.0") + (displayln " ──────────────────────────────────────────────") (printf " port: ~a~n" *port*) (printf " workers: ~a (supervised, one-for-one)~n" *workers*) (printf " hmac: ~a~n" (if (string=? *secret* "") "disabled" "enabled")) + (printf " retries: ~a (1s/2s/4s backoff → dead-letter)~n" *max-retries*) (printf " dashboard: ws://localhost:~a/dashboard~n" *port*) (displayln "") - (displayln " POST /hooks/:type ingest webhook") - (displayln " GET /api/events/:id query event") - (displayln " GET /api/stats statistics") - (displayln " GET /health liveness") - (displayln " GET /dashboard websocket stream") - (displayln " ──────────────────────────────────────") + (displayln " POST /hooks/:type ingest webhook") + (displayln " GET /api/events/:id query event") + (displayln " GET /api/stats statistics") + (displayln " GET /api/filters list user filters") + (displayln " POST /api/filters register filter") + (displayln " DELETE /api/filters/:name remove filter") + (displayln " GET /api/handlers list handlers") + (displayln " POST /api/handlers register hot handler") + (displayln " GET /api/dead-letter dead-letter queue") + (displayln " POST /api/dead-letter/:id/replay replay event") + (displayln " GET /health liveness") + (displayln " GET /dashboard websocket stream") + (displayln " ──────────────────────────────────────────────") (displayln "") ;; Start supervised worker pool