feat: Phase 3 production hardening (metrics, logging, SQLite, TLS, shutdown)
ober
77c302d1ac6eb3c9f2f311957f02ceaeff488c45
--- a/Makefile +++ b/Makefile @@ -1,8 +1,9 @@ JERBOA_HOME ?= $(HOME)/mine/jerboa SCHEME = scheme --libdirs $(JERBOA_HOME)/lib NATIVE = $(JERBOA_HOME)/jerboa-native-rs/target/release +DOCKER_IMAGE = jerboa21/jerboa -.PHONY: run test clean +.PHONY: run run-db run-tls test clean edge-static ## Run the webhook service on port 8080 run: @@ -17,6 +18,38 @@ run-custom: EDGE_SECRET=$(or $(SECRET),) \ $(SCHEME) --script edge.ss +## Run with SQLite persistence (events survive restart) +## Example: make run-db DB=/var/lib/edge/events.db +run-db: + LD_LIBRARY_PATH=$(NATIVE) \ + EDGE_DB_PATH=$(or $(DB),/tmp/edge.db) \ + $(SCHEME) --script edge.ss + +## Run with TLS enabled (requires cert + key PEM files) +## Example: make run-tls CERT=tls/cert.pem KEY=tls/key.pem +run-tls: + LD_LIBRARY_PATH=$(NATIVE) \ + EDGE_TLS_CERT=$(or $(CERT),tls/cert.pem) \ + EDGE_TLS_KEY=$(or $(KEY),tls/key.pem) \ + EDGE_TLS_PORT=$(or $(TLS_PORT),8443) \ + $(SCHEME) --script edge.ss + +## Build a self-contained static Linux binary via Docker (Phase 3.1) +## Produces ./edge — runs on any x86_64 Linux with zero dependencies. +## Binary size target: <25MB +## Requires Docker and the jerboa21/jerboa base image: +## docker pull jerboa21/jerboa (or: make -C $(JERBOA_HOME) docker-build) +edge-static: + @echo "=== Building static edge binary ===" + docker run --rm \ + -v $(JERBOA_HOME):/build/mine/jerboa:ro \ + -v $(PWD):/build/edge \ + -w /build/edge \ + $(DOCKER_IMAGE) \ + /build/mine/jerboa/support/build-static-script.sh edge.ss edge + @echo "=== Built: ./edge ($$(du -sh edge | cut -f1)) ===" + @echo "Run: LD_LIBRARY_PATH= ./edge" + ## Run smoke tests against a running server test: ./test-edge.sh --- a/edge.ss +++ b/edge.ss @@ -8,6 +8,11 @@ ;;; - STM state store (lock-free, snapshot-consistent) ;;; - WebSocket dashboard (live event stream) ;;; - HMAC-SHA256 signature verification +;;; - Prometheus metrics (/metrics endpoint) +;;; - Structured JSON logging (greppable by Loki/Splunk) +;;; - Graceful SIGTERM shutdown (drain, flush SQLite WAL) +;;; - Optional SQLite persistence (EDGE_DB_PATH) +;;; - Optional TLS termination (EDGE_TLS_CERT + EDGE_TLS_KEY) ;;; ;;; Zero external dependencies. Everything is Jerboa stdlib. ;;; @@ -15,6 +20,7 @@ ;;; Test: curl -X POST localhost:8080/hooks/payment \ ;;; -d '{"id":"evt_1","amount":4999}' ;;; curl localhost:8080/api/stats +;;; curl localhost:8080/metrics (import (except (chezscheme) merge make-hash-table hash-table? @@ -35,22 +41,154 @@ (std crypto native) (std pmap) (std security restrict) - (std misc timeout)) + (std misc timeout) + (std metrics) + (std os signal) + (std db sqlite-native) + (std net tls-rustls) + (std net tcp-raw)) ;; ═══════════════════════════════════════════════════════════════ ;; Configuration ;; ═══════════════════════════════════════════════════════════════ -(def *port* (or (and (getenv "EDGE_PORT") (string->number (getenv "EDGE_PORT"))) 8080)) -(def *workers* (or (and (getenv "EDGE_WORKERS") (string->number (getenv "EDGE_WORKERS"))) 4)) -(def *secret* (or (getenv "EDGE_SECRET") "")) +(def *port* (or (and (getenv "EDGE_PORT") (string->number (getenv "EDGE_PORT"))) 8080)) +(def *workers* (or (and (getenv "EDGE_WORKERS") (string->number (getenv "EDGE_WORKERS"))) 4)) +(def *secret* (or (getenv "EDGE_SECRET") "")) +(def *db-path* (getenv "EDGE_DB_PATH")) ;; SQLite file path, or #f for in-memory only +(def *tls-cert* (getenv "EDGE_TLS_CERT")) ;; PEM cert path, or #f to disable TLS +(def *tls-key* (getenv "EDGE_TLS_KEY")) ;; PEM key path, or #f to disable TLS +(def *tls-port* (or (and (getenv "EDGE_TLS_PORT") (string->number (getenv "EDGE_TLS_PORT"))) 8443)) + +;; ═══════════════════════════════════════════════════════════════ +;; Prometheus Metrics (Phase 3.4) +;; +;; All metrics are registered in a single registry exposed at /metrics. +;; Label tracking is not per-type (the simple metrics module does not +;; support label dimensions), so each status gets its own counter. +;; ═══════════════════════════════════════════════════════════════ + +(def *metrics* (make-registry)) +(def *metric-received* (make-counter *metrics* "edge_events_received_total" "Total webhook events received")) +(def *metric-ok* (make-counter *metrics* "edge_events_ok_total" "Events processed successfully")) +(def *metric-error* (make-counter *metrics* "edge_events_error_total" "Events that failed at least once")) +(def *metric-dead-letter* (make-counter *metrics* "edge_events_dead_letter_total" "Events moved to dead-letter queue")) +(def *metric-restarts* (make-counter *metrics* "edge_worker_restarts_total" "Worker actor supervisor restarts")) +(def *metric-duration* (make-histogram *metrics* "edge_processing_duration_seconds" "Event processing wall-clock duration")) +(def *metric-connections* (make-gauge *metrics* "edge_active_connections" "Active WebSocket dashboard connections")) + +;; ═══════════════════════════════════════════════════════════════ +;; Structured JSON Logging (Phase 3.5) +;; +;; All event-path log lines are JSON — greppable by Splunk, Loki, +;; jq, etc. The startup banner stays human-readable. +;; +;; (jlog "worker ready" "worker" 0) +;; (jlog "event processed" "worker" 1 "id" "evt_42" "status" "ok" "ms" 12) +;; ═══════════════════════════════════════════════════════════════ + +(def (jlog msg . pairs) + ;; Emit a JSON log line to stdout. pairs: alternating key value. + (let ([ht (make-hash-table)]) + (hash-put! ht "ts" (format "~a" (time-second (current-time)))) + (hash-put! ht "msg" msg) + (let loop ([ps pairs]) + (unless (or (null? ps) (null? (cdr ps))) + (hash-put! ht (car ps) (cadr ps)) + (loop (cddr ps)))) + (displayln (json-object->string ht)))) + +;; ═══════════════════════════════════════════════════════════════ +;; SQLite Persistence (Phase 3.3) +;; +;; Optional: set EDGE_DB_PATH to a file path to enable persistence. +;; On startup, existing events and dead-letters are loaded into STM. +;; Every store-result! / store-dead-letter! also writes to SQLite. +;; WAL mode is used for concurrent reads with no reader blocking. +;; +;; All DB writes are fire-and-forget: guarded with (guard (e [#t (void)])) +;; so a SQLite error never crashes a worker. +;; ═══════════════════════════════════════════════════════════════ + +(def *db* #f) ;; SQLite handle — #f when EDGE_DB_PATH is not set + +(def (db-init!) + (when *db-path* + (let ([db (sqlite-open *db-path*)]) + (sqlite-exec db "PRAGMA journal_mode=WAL") + (sqlite-exec db + (string-append + "CREATE TABLE IF NOT EXISTS events (" + " id TEXT PRIMARY KEY, type TEXT, status TEXT," + " received TEXT, processed TEXT)")) + (sqlite-exec db + (string-append + "CREATE TABLE IF NOT EXISTS dead_letters (" + " id TEXT PRIMARY KEY, type TEXT, reason TEXT," + " retries INTEGER, ts INTEGER)")) + (set! *db* db) + (jlog "sqlite opened" "path" *db-path*)))) + +(def (db-store-event! event status) + (when *db* + (guard (e [#t (void)]) + (sqlite-execute *db* + "INSERT OR REPLACE INTO events VALUES (?,?,?,?,?)" + (hash-ref event "id" "") + (hash-ref event "type" "unknown") + status + (hash-ref event "received" "") + (format "~a" (time-second (current-time))))))) + +(def (db-store-dead-letter! event reason) + (when *db* + (guard (e [#t (void)]) + (sqlite-execute *db* + "INSERT OR REPLACE INTO dead_letters VALUES (?,?,?,?,?)" + (hash-ref event "id" "") + (hash-ref event "type" "unknown") + reason + (hash-ref event "retry-count" 0) + (time-second (current-time)))))) + +(def (db-load-events!) + ;; Replay persisted events into STM state on startup. + (when *db* + (let ([rows (sqlite-query *db* "SELECT id, type, status, received, processed FROM events")]) + (for-each + (lambda (row) + (let ([id (cdr (assoc "id" row))] + [ht (list->hash-table + `(("type" . ,(cdr (assoc "type" row))) + ("status" . ,(cdr (assoc "status" row))) + ("received" . ,(cdr (assoc "received" row))) + ("processed" . ,(cdr (assoc "processed" row)))))]) + (dosync (alter events-ref (lambda (st) (persistent-map-set st id ht)))))) + rows)) + (let ([dead (sqlite-query *db* "SELECT id, type, reason, retries FROM dead_letters")]) + (for-each + (lambda (row) + (let ([id (cdr (assoc "id" row))] + [ht (list->hash-table + `(("id" . ,(cdr (assoc "id" row))) + ("type" . ,(cdr (assoc "type" row))) + ("reason" . ,(cdr (assoc "reason" row))) + ("retries" . ,(cdr (assoc "retries" row)))))]) + (dosync (alter dead-letter-ref (lambda (st) (persistent-map-set st id ht)))))) + dead)) + (jlog "state restored from sqlite" "path" *db-path*))) + +(def (db-flush-wal!) + (when *db* + (guard (e [#t (void)]) + (sqlite-exec *db* "PRAGMA wal_checkpoint(FULL)")))) ;; ═══════════════════════════════════════════════════════════════ ;; ID Generation ;; ═══════════════════════════════════════════════════════════════ (def id-counter 0) -(def id-mutex (make-mutex)) +(def id-mutex (make-mutex)) (def (gen-id) (with-mutex id-mutex @@ -80,7 +218,8 @@ `(("type" . ,(hash-ref event "type")) ("status" . ,status) ("received" . ,(hash-ref event "received")) - ("processed" . ,(format "~a" (time-second (current-time))))))))))) + ("processed" . ,(format "~a" (time-second (current-time)))))))))) + (db-store-event! event status)) (def (get-event id) (let ([st (ref-deref events-ref)]) @@ -113,7 +252,8 @@ ("type" . ,(hash-ref event "type" "unknown")) ("reason" . ,reason) ("retries" . ,(hash-ref event "retry-count" 0)) - ("time" . ,(format "~a" (time-second (current-time))))))))))) + ("time" . ,(format "~a" (time-second (current-time)))))))))) + (db-store-dead-letter! event reason)) (def (get-dead-letter-list) (let ([result '()]) @@ -184,8 +324,8 @@ (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)) + (jlog "filter eval error" "err" + (if (message-condition? exn) (condition-message exn) (format "~a" exn))) #t]) ;; permissive on error ;; with-timeout uses Chez engines (no fork) — fiber-safe (let ([result (with-timeout 2.0 #t @@ -207,12 +347,14 @@ ;; ═══════════════════════════════════════════════════════════════ (def (register-watcher! ws) - (dosync (alter watchers-ref (lambda (lst) (cons ws lst))))) + (dosync (alter watchers-ref (lambda (lst) (cons ws lst)))) + (gauge-inc! *metric-connections*)) (def (unregister-watcher! ws) (dosync (alter watchers-ref - (lambda (lst) (filter (lambda (w) (not (eq? w ws))) lst))))) + (lambda (lst) (filter (lambda (w) (not (eq? w ws))) lst)))) + (gauge-dec! *metric-connections*)) (def (notify-watchers! event status) (let ([msg (json-object->string @@ -239,16 +381,17 @@ (register-handler! "payment.completed" (lambda (evt) (let ([p (hash-ref evt "payload")]) - (displayln "[payment] " (hash-ref evt "id") - " amount=" (hash-ref p "amount" "?"))))) + (jlog "payment completed" + "id" (hash-ref evt "id") + "amount" (format "~a" (hash-ref p "amount" "?")))))) (register-handler! "user.created" (lambda (evt) - (displayln "[user] " (hash-ref evt "id")))) + (jlog "user created" "id" (hash-ref evt "id")))) (register-handler! "order.shipped" (lambda (evt) - (displayln "[order] " (hash-ref evt "id") " shipped"))) + (jlog "order shipped" "id" (hash-ref evt "id")))) ;; Crash handler — deliberately errors to demonstrate supervisor restart. ;; POST /hooks/test.crash to trigger. The worker that picks this up @@ -283,7 +426,7 @@ ;; Unknown type → log and continue [else (lambda (evt) - (displayln "[" type "] " (hash-ref evt "id")))]))) + (jlog "unhandled type" "type" type "id" (hash-ref evt "id")))]))) ;; ═══════════════════════════════════════════════════════════════ ;; Transducer Pipeline @@ -311,7 +454,7 @@ [(not id) #t] ;; no ID → pass through [(> retry 0) #t] ;; retry → always pass [(hash-key? seen-ids id) - (displayln "[dedup] dropping duplicate " id) #f] + (jlog "dedup drop" "id" id) #f] [else (hash-put! seen-ids id #t) #t])))) ;; Normalize: ensure every event has a timestamp (mapping @@ -340,42 +483,60 @@ (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) ")") + (jlog "re-enqueuing for retry" + "worker" (format "~a" worker-id) + "id" (hash-ref event "id") + "retry" (format "~a" (hash-ref event "retry-count" 0))) (>!! ingest-ch event)))) +(def (time-now-sec) + ;; Fractional seconds for duration measurement + (let ([t (current-time)]) + (+ (time-second t) (/ (time-nanosecond t) 1e9)))) + (def (make-worker id) (lambda (msg) (match msg ['start - (displayln "[worker-" id "] ready") + (jlog "worker ready" "worker" (format "~a" id)) (let loop () (let ([event (<!! ingest-ch)]) (when (and event (not (eof-object? event))) (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)")] + (jlog "filter drop" + "worker" (format "~a" id) + "id" (hash-ref event "id" "?"))] ;; Process the event [else (let* ([type (hash-ref event "type" "unknown")] [retry-count (hash-ref event "retry-count" 0)] [handler (dispatch-handler type)] + [t0 (time-now-sec)] [result (try (ok (handler event)) (catch (e) (err e)))] + [duration (- (time-now-sec) t0)] [status (if (ok? result) "ok" "error")]) (cond ;; Success path [(ok? result) + (counter-inc! *metric-ok*) + (histogram-observe! *metric-duration* duration) (store-result! event status) - (notify-watchers! event status)] + (notify-watchers! event status) + (jlog "event processed" + "worker" (format "~a" id) + "id" (hash-ref event "id" "?") + "type" type + "status" "ok" + "ms" (format "~a" (inexact->exact (round (* duration 1000)))))] ;; Deliberate crash demo — re-raise to prove supervisor restarts [(string=? type "test.crash") + (counter-inc! *metric-restarts*) (store-result! event status) (notify-watchers! event status) - (displayln "[worker-" id "] crashing — supervisor will restart") + (jlog "worker crash demo" "worker" (format "~a" id)) (raise (unwrap-err result))] ;; Retry with exponential backoff: 1s / 2s / 4s [(< retry-count *max-retries*) @@ -383,9 +544,13 @@ [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") + (counter-inc! *metric-error*) + (jlog "event failed retrying" + "worker" (format "~a" id) + "type" type + "retry" (format "~a" (+ retry-count 1)) + "of" (format "~a" *max-retries*) + "err" err-msg) (schedule-retry! next-event delay-secs id) (store-result! event "retrying") (notify-watchers! event "retrying"))] @@ -393,16 +558,19 @@ [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) + (counter-inc! *metric-dead-letter*) + (jlog "event dead-lettered" + "worker" (format "~a" id) + "type" type + "err" 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) + ] ;; close dead-letter [else] + ) ;; close inner (cond) + ) ;; close outer (let*) + ] ;; close outer [else] + ) ;; close outer (cond) (loop) ) ;; close (when) ) ;; close (let ([event ...])) @@ -444,6 +612,7 @@ ("payload" . ,payload) ("received" . ,(format "~a" (time-second (current-time)))) ("status" . "queued")))]) + (counter-inc! *metric-received*) (>!! ingest-ch event) (respond-json 202 (json-object->string @@ -469,6 +638,12 @@ (def (handle-health req) (respond-json 200 "{\"status\":\"ok\"}")) +;; GET /metrics — Prometheus text format exposition +(def (handle-metrics req) + (respond 200 + '(("Content-Type" . "text/plain; version=0.0.4; charset=utf-8")) + (prometheus-format *metrics*))) + ;; GET /dashboard — WebSocket upgrade for live event stream (def (handle-dashboard req) (make-websocket-response @@ -522,7 +697,7 @@ (lambda (st) (persistent-map-set st name (list->hash-table `(("name" . ,name) ("code" . ,code))))))) - (displayln "[filters] registered: " name) + (jlog "filter registered" "name" name) (respond-json 201 (json-object->string (list->hash-table `(("status" . "registered") ("name" . ,name)))))]))) @@ -535,7 +710,7 @@ (dosync (alter filters-ref (lambda (st) (persistent-map-delete st name)))) - (displayln "[filters] removed: " name) + (jlog "filter removed" "name" name) (respond-json 200 (json-object->string (list->hash-table `(("status" . "removed") ("name" . ,name)))))) @@ -560,12 +735,12 @@ [else ;; Store as a code string; dispatch-handler will sandbox-eval it (hash-put! handler-registry type code) - (displayln "[handlers] registered: " type) + (jlog "handler registered" "type" type) (respond-json 201 (json-object->string - (list->hash-table `(("status" . "registered") - ("type" . ,type) - ("sandboxed". #t)))))]))) + (list->hash-table `(("status" . "registered") + ("type" . ,type) + ("sandboxed" . #t)))))]))) ;; GET /api/handlers — list registered handler types (def (handle-list-handlers req) @@ -573,8 +748,8 @@ (hash-for-each (lambda (type h) (set! types (cons - (list->hash-table `(("type" . ,type) - ("kind" . ,(if (procedure? h) "native" "sandboxed")))) + (list->hash-table `(("type" . ,type) + ("kind" . ,(if (procedure? h) "native" "sandboxed")))) types))) handler-registry) (respond-json 200 @@ -616,12 +791,177 @@ (alter dead-letter-ref (lambda (st) (persistent-map-delete st id)))) (>!! ingest-ch replay-event) - (displayln "[dead-letter] replaying " id) + (jlog "dead-letter replaying" "id" id) (respond-json 202 (json-object->string (list->hash-table `(("status" . "replayed") ("id" . ,id))))))]))) ;; ═══════════════════════════════════════════════════════════════ +;; TLS Server (Phase 3.2) +;; +;; Optional: set EDGE_TLS_CERT + EDGE_TLS_KEY for direct HTTPS on +;; EDGE_TLS_PORT (default 8443). Uses the Rust rustls backend. +;; +;; Architecture: one OS thread per TLS connection. The fiber runtime +;; handles the plain-HTTP path; TLS connections use blocking I/O. +;; The same router handles both — TLS connections get a full HTTP +;; parse/dispatch/respond cycle over the encrypted channel. +;; ═══════════════════════════════════════════════════════════════ + +;; Parse an HTTP/1.1 request from a raw string (used by TLS path). +;; Returns a request record or #f on parse failure. +(def (parse-tls-http-request s) + (let* ([sep-pos (string-contains s "\r\n\r\n")] + [header-section (if sep-pos (substring s 0 sep-pos) s)] + [all-lines (tls-split-crlf header-section)] + [req-line (and (pair? all-lines) (car all-lines))] + [parts (and req-line (string-split-spaces* req-line))]) + (and parts (>= (length parts) 2) + (let* ([method (car parts)] + [path (cadr parts)] + [version (if (>= (length parts) 3) (caddr parts) "HTTP/1.1")] + [hdrs (tls-parse-headers (cdr all-lines))] + [cl-str (let ([e (assoc "content-length" hdrs)]) (and e (cdr e)))] + [cl (and cl-str (string->number cl-str))] + [body (and cl sep-pos + (let ([start (+ sep-pos 4)]) + (and (>= (string-length s) (+ start cl)) + (substring s start (+ start cl)))))]) + (make-request method path version hdrs body))))) + +(def (tls-split-crlf s) + ;; Split string on CRLF, returning list of lines + (let loop ([start 0] [acc '()]) + (let ([pos (string-contains (substring s start (string-length s)) "\r\n")]) + (if pos + (let ([line (substring s start (+ start pos))]) + (if (string=? line "") + (reverse acc) ;; blank line = end of headers + (loop (+ start pos 2) (cons line acc)))) + (let ([rest (substring s start (string-length s))]) + (reverse (if (string=? rest "") acc (cons rest acc)))))))) + +(def (string-split-spaces* s) + ;; Split on spaces, dropping empties + (let loop ([i 0] [start 0] [acc '()]) + (cond + [(= i (string-length s)) + (reverse (if (= start i) acc (cons (substring s start i) acc)))] + [(char=? (string-ref s i) #\space) + (loop (+ i 1) (+ i 1) + (if (= start i) acc (cons (substring s start i) acc)))] + [else (loop (+ i 1) start acc)]))) + +(def (tls-parse-headers lines) + ;; Parse header lines into alist — name is lowercased + (let loop ([ls lines] [acc '()]) + (if (null? ls) + (reverse acc) + (let* ([line (car ls)] + [pos (string-contains line ":")]) + (if pos + (loop (cdr ls) + (cons (cons (string-downcase (substring line 0 pos)) + (string-trim (substring line (+ pos 1) (string-length line)))) + acc)) + (loop (cdr ls) acc)))))) + +(def (http-response->bv resp) + ;; Serialize a response record to a bytevector for TLS write + (let* ([status (response-status resp)] + [headers (response-headers resp)] + [body (or (response-body resp) "")] + [body-bv (string->utf8 body)] + [body-len (bytevector-length body-bv)] + [status-text (case status + [(200) "OK"] [(201) "Created"] [(202) "Accepted"] + [(204) "No Content"] + [(400) "Bad Request"] [(401) "Unauthorized"] + [(404) "Not Found"] [(500) "Internal Server Error"] + [else "Unknown"])] + [hdr-str + (string-append + (format "HTTP/1.1 ~a ~a\r\n" status status-text) + (format "Content-Length: ~a\r\n" body-len) + (apply string-append + (map (lambda (h) (format "~a: ~a\r\n" (car h) (cdr h))) + headers)) + "\r\n")] + [hdr-bv (string->utf8 hdr-str)] + [hdr-len (bytevector-length hdr-bv)] + [out (make-bytevector (+ hdr-len body-len))]) + (bytevector-copy! hdr-bv 0 out 0 hdr-len) + (bytevector-copy! body-bv 0 out hdr-len body-len) + out)) + +(def (handle-tls-connection conn-handle router) + ;; Handle one TLS connection: read request, dispatch, write response. + ;; Uses a 16KB buffer — fits typical webhook payloads in one TLS record. + (let* ([buf (make-bytevector 16384)] + [n (rustls-read conn-handle buf 16384)]) + (when (> n 0) + (let* ([raw (let ([out (make-bytevector n)]) + (bytevector-copy! buf 0 out 0 n) out)] + [s (guard (e [#t #f]) (utf8->string raw))] + [req (and s (parse-tls-http-request s))]) + (when req + (let* ([resp (guard (e [#t (respond-json 500 "{\"error\":\"internal\"}")]) + (router-dispatch router req))] + [resp-bv (http-response->bv resp)]) + (rustls-write conn-handle resp-bv (bytevector-length resp-bv)))))))) + +(def (run-tls-server! router tls-port cert-path key-path) + ;; Start TLS server in background OS threads: one thread for accept loop, + ;; one thread per connection. + (let* ([server-ctx (rustls-server-ctx-new cert-path key-path)] + [listen-fd (tcp-listen tls-port)]) + (jlog "tls listening" "port" (format "~a" tls-port) "cert" cert-path) + (fork-thread + (lambda () + (let loop () + (let ([client-fd (guard (e [#t #f]) (tcp-accept listen-fd))]) + (when client-fd + (fork-thread + (lambda () + (let ([conn (guard (e [#t #f]) + (rustls-accept server-ctx client-fd))]) + (when conn + (guard (e [#t (void)]) + (handle-tls-connection conn router)) + (guard (e [#t (void)]) + (rustls-close conn)))))) + (loop)))))))) + +;; ═══════════════════════════════════════════════════════════════ +;; Graceful Shutdown (Phase 3.6) +;; +;; SIGTERM triggers an orderly shutdown: +;; 1. Stop accepting new HTTP connections (fiber-httpd-stop!) +;; 2. Close ingest channel — workers see EOF and exit their loops +;; 3. Brief drain window (3s) for in-flight events to complete +;; 4. Flush SQLite WAL if persistence is enabled +;; 5. Exit 0 +;; +;; add-signal-handler! is safe to call from any thread. +;; ═══════════════════════════════════════════════════════════════ + +(def *httpd-server* #f) ;; set in run-edge after start + +(def (graceful-shutdown!) + (jlog "sigterm received" "msg" "starting graceful shutdown") + ;; Stop accepting new connections + (when *httpd-server* + (guard (e [#t (void)]) (fiber-httpd-stop! *httpd-server*))) + ;; Close ingest channel — workers' (<!! ingest-ch) returns EOF + (guard (e [#t (void)]) (chan-close! ingest-ch)) + ;; Drain window — let in-flight events finish + (sleep (make-time 'time-duration 0 3)) + ;; Flush SQLite WAL + (db-flush-wal!) + (jlog "shutdown complete") + (exit 0)) + +;; ═══════════════════════════════════════════════════════════════ ;; Router & Server ;; ═══════════════════════════════════════════════════════════════ @@ -642,19 +982,31 @@ ;; 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) + ;; Phase 3: Prometheus metrics + (route-get r "/metrics" handle-metrics) ;; Dashboard + health (route-get r "/dashboard" handle-dashboard) (route-get r "/health" handle-health) + ;; Initialize SQLite persistence (Phase 3.3) + (db-init!) + (db-load-events!) + + ;; Register SIGTERM handler (Phase 3.6) + (add-signal-handler! SIGTERM graceful-shutdown!) + ;; Banner first, then start workers (displayln "") - (displayln " Jerboa Edge v0.2.0") + (displayln " Jerboa Edge v0.3.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 " sqlite: ~a~n" (if *db-path* *db-path* "disabled")) + (printf " tls: ~a~n" (if *tls-cert* (format ":~a" *tls-port*) "disabled")) (printf " dashboard: ws://localhost:~a/dashboard~n" *port*) + (printf " metrics: http://localhost:~a/metrics~n" *port*) (displayln "") (displayln " POST /hooks/:type ingest webhook") (displayln " GET /api/events/:id query event") @@ -666,6 +1018,7 @@ (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 /metrics prometheus metrics") (displayln " GET /health liveness") (displayln " GET /dashboard websocket stream") (displayln " ──────────────────────────────────────────────") @@ -674,10 +1027,15 @@ ;; Start supervised worker pool (let ([sup (start-workers *workers*)]) + ;; Start optional TLS server (Phase 3.2) + (when (and *tls-cert* *tls-key*) + (run-tls-server! r *tls-port* *tls-cert* *tls-key*)) + ;; Start fiber HTTP server (let ([srv (fiber-httpd-start *port* (lambda (req) (router-dispatch r req)))]) - (displayln "[edge] listening on :" *port*) + (set! *httpd-server* srv) + (jlog "edge listening" "port" (format "~a" *port*)) ;; Block main thread (let loop () (sleep (make-time 'time-duration 0 3600))