Harden listener limits and toolchain verification
ober
1d3dbf1d5676d994231a04abd3ca7954b00e2109
--- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -17,7 +17,7 @@ jobs: verify: runs-on: ubuntu-latest steps: - - uses: actions/checkout@v4 + - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4.3.1 - name: Install system tools run: | --- a/.github/workflows/security-baseline.yml +++ b/.github/workflows/security-baseline.yml @@ -13,7 +13,7 @@ jobs: baseline: runs-on: ubuntu-latest steps: - - uses: actions/checkout@v4 + - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4.3.1 - name: Required release files run: | --- a/Makefile +++ b/Makefile @@ -2,20 +2,16 @@ # crate, so building secmon-agent needs only `jerbuild`, a C compiler, and # cargo — no jerboa source checkout and no separately-built Chez/native lib. JERBUILD ?= jerbuild -JH := $(shell $(JERBUILD) --jerboa-home 2>/dev/null) -ifeq ($(JH),) -$(error jerbuild not found on PATH (or '$(JERBUILD) --jerboa-home' failed). Install jerbuild, or set JERBUILD=/path/to/jerbuild) -endif +JH = $(shell $(JERBUILD) --jerboa-home 2>/dev/null) -NATIVE_A := $(JH)/jerboa-native-rs/target/release/libjerboa_native.a +NATIVE_A = $(JH)/jerboa-native-rs/target/release/libjerboa_native.a BIN := secmon-agent BIN_DIR := $(HOME)/.local/bin VENDOR := $(CURDIR)/vendor JSQLITE_REPO ?= $(VENDOR)/jsqlite -JSQLITE_URL ?= https://git.sr.ht/~lisp/jsqlite JSQLITE_SRC ?= $(JSQLITE_REPO)/src -LIBDIRS := --libdirs lib:$(JSQLITE_SRC):$(JH)/lib -JEXEC := $(JERBUILD) exec $(LIBDIRS) +LIBDIRS = --libdirs lib:$(JSQLITE_SRC):$(JH)/lib +JEXEC = $(JERBUILD) exec $(LIBDIRS) CARGO_AUDIT ?= $(shell command -v cargo-audit 2>/dev/null || printf '%s/.cargo/bin/cargo-audit' "$$HOME") DIST_DIR ?= dist/release-evidence SBOM_DIR ?= dist/sbom @@ -23,7 +19,7 @@ REPRO_DIR ?= dist/reproducibility SOAK_DIR ?= dist/soak BINARY_SMOKE_DIR ?= dist/binary-smoke -.PHONY: all help build binary keygen agent collector analyze test security audit verify sbom reproducibility-report soak-evidence binary-smoke sanitize-evidence release-evidence install clean ensure-jsqlite +.PHONY: all help build binary keygen agent collector analyze test security audit verify sbom reproducibility-report soak-evidence binary-smoke sanitize-evidence release-evidence install clean toolchain-check ensure-jsqlite .DEFAULT_GOAL := help all: binary @@ -54,11 +50,12 @@ help: # jerbuild's cache; we then regenerate the FFI symbol list from that archive # (platform-correct) and relink. A static main.c sets JERBOA_STATIC=1 so the # std modules use the registered symbols instead of dlopen. -ensure-jsqlite: - @if [ ! -f "$(JSQLITE_SRC)/jsqlite/api.ss" ]; then \ - mkdir -p "$(VENDOR)"; \ - git clone --depth 1 "$(JSQLITE_URL)" "$(JSQLITE_REPO)"; \ - fi +toolchain-check: + support/verify-toolchain.sh "$(JERBUILD)" + +ensure-jsqlite: toolchain-check + @test "$(abspath $(JSQLITE_REPO))" = "$(abspath $(VENDOR)/jsqlite)" || { echo "JSQLITE override is not permitted by the locked release build" >&2; exit 1; } + support/fetch-locked-deps.sh jsqlite @test -f "$(JSQLITE_SRC)/jsqlite/api.ss" # The final pass stabilizes the WPO payload after the generated FFI symbol set @@ -92,28 +89,29 @@ test: ensure-jsqlite fi @echo "All tests passed." -security: - scripts/security-check.sh +security: toolchain-check + JERBUILD="$(JERBUILD)" scripts/security-check.sh -audit: +audit: toolchain-check @if ! [ -x "$(CARGO_AUDIT)" ]; then \ echo "cargo-audit is required. Install with: cargo install cargo-audit --locked"; \ exit 1; \ fi cd patches/jerboa-native-rs && "$(CARGO_AUDIT)" audit + cd "$(JH)/jerboa-native-rs" && "$(CARGO_AUDIT)" audit verify: security test audit binary -sbom: +sbom: toolchain-check JERBUILD="$(JERBUILD)" JSECMONLIB_SBOM_DIR="$(SBOM_DIR)" scripts/sbom.sh -reproducibility-report: +reproducibility-report: toolchain-check JERBUILD="$(JERBUILD)" JSECMONLIB_REPRO_DIR="$(REPRO_DIR)" scripts/reproducibility-report.sh -soak-evidence: +soak-evidence: toolchain-check JSECMONLIB_SOAK_DIR="$(SOAK_DIR)" scripts/soak-evidence.sh -binary-smoke: +binary-smoke: toolchain-check JSECMONLIB_BINARY_SMOKE_DIR="$(BINARY_SMOKE_DIR)" scripts/binary-smoke.sh sanitize-evidence: @@ -127,12 +125,16 @@ release-evidence: verify sbom reproducibility-report soak-evidence binary-smoke uname -srm > "$(DIST_DIR)/build-environment.txt" printf 'jerbuild_version=' >> "$(DIST_DIR)/build-environment.txt" $(JERBUILD) --version >> "$(DIST_DIR)/build-environment.txt" - printf 'jerboa_home_status=present\n' >> "$(DIST_DIR)/build-environment.txt" + support/verify-toolchain.sh "$(JERBUILD)" >> "$(DIST_DIR)/build-environment.txt" cargo metadata --manifest-path patches/jerboa-native-rs/Cargo.toml --locked --format-version 1 > "$(DIST_DIR)/cargo-metadata-native.json" + cargo metadata --manifest-path "$(JH)/jerboa-native-rs/Cargo.toml" --locked --format-version 1 > "$(DIST_DIR)/cargo-metadata-bundled-native.json" cd patches/jerboa-native-rs && "$(CARGO_AUDIT)" audit > "$(CURDIR)/$(DIST_DIR)/rustsec-native.txt" + cd "$(JH)/jerboa-native-rs" && "$(CARGO_AUDIT)" audit > "$(CURDIR)/$(DIST_DIR)/rustsec-bundled-native.txt" shasum -a 256 secmon-agent > "$(DIST_DIR)/secmon-agent-sha256.txt" shasum -a 256 \ .jerbuild \ + dependency-lock.tsv \ + toolchain-lock.tsv \ Makefile \ SECURITY.md \ .jerboa/security.json \ @@ -148,6 +150,8 @@ release-evidence: verify sbom reproducibility-report soak-evidence binary-smoke support/static-main.c \ support/gen-ffi-symbols.sh \ support/ensure-jerboa.sh \ + support/fetch-locked-deps.sh \ + support/verify-toolchain.sh \ > "$(DIST_DIR)/source-inputs-sha256.txt" cp support/ffi-symbols.gen "$(DIST_DIR)/ffi-symbols.gen" { if command -v otool >/dev/null 2>&1; then otool -L secmon-agent; elif command -v ldd >/dev/null 2>&1; then ldd secmon-agent; else echo "no dynamic-linkage inspector found"; fi; } > "$(DIST_DIR)/native-linkage.txt" --- a/README.md +++ b/README.md @@ -3,6 +3,13 @@ Static Jerboa security-monitor agent, collector/analyzer source, and native FFI registration support. +The agent listener defaults to `127.0.0.1:31337`, binds the requested numeric +IPv4 address through the full socket path, and reports the OS-confirmed bound +address. A fixed worker pool, bounded queue, global/per-source active and rate +budgets, strict frames, and absolute authentication/idle read-and-write deadlines +protect the pre-auth boundary from Slowloris/thread exhaustion. `make test` +includes real non-loopback refusal, byte-drip, and non-reading-peer regressions. + This repo is not public-production-ready until `make release-evidence` is captured on the target Linux or FreeBSD release host and the blockers in `docs/release-evidence.md` are closed. --- a/bin/agent.ss +++ b/bin/agent.ss @@ -102,13 +102,15 @@ (spawn-webshell-monitor proc-provider emit! 2000 hostname) (spawn-revshell-monitor proc-provider net-provider emit! 2000 hostname) - (when debug-mode - (displayln (format "All monitors started. Listening on ~a" - (agent-config-listen-addr config)))) - - ;; Start TCP server (blocks) + ;; Start the listener, then report the OS-confirmed bound address. (let ([server (make-poll-server psk-auth buffer (agent-config-listen-addr config))]) - (poll-server-run! server))))))) + (let ([runtime (poll-server-start! server)]) + (when debug-mode + (displayln + (format "All monitors started. Listening on ~a:~a" + (poll-server-listen-host runtime) + (poll-server-listen-port runtime)))) + (poll-server-wait! runtime)))))))) (apply main (cdr (command-line))) new file mode 100644 --- /dev/null +++ b/dependency-lock.tsv @@ -0,0 +1,2 @@ +# name url commit tree path +jsqlite https://git.sr.ht/~lisp/jsqlite 63d737101b3b74bff3b9db75cc35f34bff6df0b7 79bfcf7a52d233076301c5cba9cd56d1b3e07ebf vendor/jsqlite --- a/docs/deployment-security.md +++ b/docs/deployment-security.md @@ -17,11 +17,20 @@ on the target Linux or FreeBSD host and archived with the release evidence. ## Network -- Bind the agent listener to a management interface or loopback tunnel, not a - public interface. +- The default `SECMON_LISTEN` is `127.0.0.1:31337`. Remote operation must name + an explicit numeric IPv4 management address; hostnames and ambiguous values + fail closed. +- The listener reports the address returned by `getsockname`, not merely the + requested configuration. Bind to a management interface or loopback tunnel, + never a public interface. - Firewall the listener to collector hosts. - Rotate PSKs when a collector host or monitored host is rebuilt, reassigned, or suspected compromised. +- Preserve the fixed worker pool and bounded queue, global/per-source active and + rate caps, five-second authentication deadline, 30-second authenticated + whole-frame/idle read-and-write deadline, and maximum encrypted-frame size. A + per-syscall timeout alone is not sufficient because a peer can drip bytes or + slowly drain a response before it expires. ## Sandbox Plan @@ -36,7 +45,9 @@ on the target Linux or FreeBSD host and archived with the release evidence. ## Release Evidence Required -- `make release-evidence` from a clean checkout with a pinned `JERBUILD`. +- `make release-evidence` from a clean checkout with the exact Jerbuild binary + and bundle hashes in `toolchain-lock.tsv`. `support/verify-toolchain.sh` must + pass before the tool is executed by a release target. - `dist/release-evidence/reproducibility/result.txt` showing `prewarm_status=present` and matching repeated binary, FFI symbol, release-input, and cargo metadata status. @@ -45,6 +56,9 @@ on the target Linux or FreeBSD host and archived with the release evidence. - `dist/release-evidence/binary-smoke/status.txt` showing target-platform startup smoke, or an explicit platform block for non-target local evidence. - Current RustSec, source scanner, secret scan, and external review notes. +- `tests/listener-bind-test.ss` showing an OS-confirmed loopback bind, refusal + through the discovered non-loopback address, byte-drip deadline closure, and + bounded response writes to a peer that never reads. For release candidates, run target proof capture with fail-closed controls: --- a/docs/release-evidence.md +++ b/docs/release-evidence.md @@ -18,8 +18,8 @@ artifacts under `dist/release-evidence/`: - `git-commit.txt` and `git-status.txt`. - `build-environment.txt` with host-neutral OS and Jerboa toolchain identity. -- Cargo dependency metadata and RustSec audit output for the patched native - Rust crate. +- Cargo dependency metadata and RustSec audit output for both the reviewed + native overlay lock and the exact native crate bundled by `JERBUILD`. - SHA-256 hash for `secmon-agent`. - Static-binary source input hashes, generated FFI symbol list, and dynamic linkage output from `otool` or `ldd`. @@ -45,6 +45,17 @@ Reproducibility is claimed only when `reproducibility/result.txt` records `prewarm_status=present`, matching `binary_status`, `ffi_symbols_status`, `release_inputs_status`, `cargo_metadata_status`, and overall `status`. +`dependency-lock.tsv` pins jsqlite to one commit and tree; the release build +fails on mutable branch-tip clones, dirty vendor state, or a different origin. +`patches/jerboa-native-rs` is a reviewed overlay/provenance snapshot rather than +a standalone crate (unmodified source modules intentionally live in the locked +Jerboa bundle), so its lock is audited directly while compilation is proven by +the full static Jerbuild build. +`toolchain-lock.tsv` pins the Jerbuild executable SHA-256 and its embedded +bundle SHA-256. The former network bootstrap has been replaced by a local-only +installer that accepts only a separately acquired executable matching both +locks; a checksum fetched beside an archive is not accepted as authenticity. + Production readiness still requires target-host runtime evidence. The external review packet must include this file, `SECURITY.md`, `docs/threat-model.md`, `docs/deployment-security.md`, `.jerboa/security.json`, the generated evidence --- a/docs/threat-model.md +++ b/docs/threat-model.md @@ -29,13 +29,23 @@ binary with a fixed native FFI symbol set. - The agent never receives the collector private key. It can encrypt events but must not decrypt historical telemetry. - Collector connections must complete PSK authentication before event access. -- Frame parsing must enforce size limits and exact reads. +- Listener configuration defaults to loopback and is honored through numeric + IPv4 bytes in `sockaddr_in`; the OS-confirmed address is reported. +- Frame I/O enforces size limits, exact reads/writes, and an absolute + whole-frame deadline. Fixed workers plus bounded global/per-source connection + and rate budgets prevent unauthenticated thread exhaustion and slow-reader + response stalls. - Static binaries must use the exact FFI symbol list generated from the pinned Jerboa/native toolchain. +- Debugger detection uses the native non-destructive tracer query; it never + resumes Scheme execution in a raw-fork child. - Production deployments must run under a dedicated service account with the smallest file and network access compatible with the selected monitors. - Release candidates must include SBOM, repeated-build reproducibility, binary-smoke, and soak/load status evidence. +- jsqlite is fetched only at the exact commit/tree in `dependency-lock.tsv`. + The reviewed native overlay lock and the exact bundled native lock both pass + current RustSec audit. ## Current Non-Production Gaps --- a/lib/secmon/config.sls +++ b/lib/secmon/config.sls @@ -42,7 +42,7 @@ (define (load-agent-config) (make-agent-config - (or (getenv (obfstr "SECMON_LISTEN")) (obfstr "0.0.0.0:31337")) + (or (getenv (obfstr "SECMON_LISTEN")) (obfstr "127.0.0.1:31337")) (or (and (getenv (obfstr "SECMON_POLL_MS")) (string->number (getenv (obfstr "SECMON_POLL_MS")))) 100) new file mode 100644 --- /dev/null +++ b/lib/secmon/server/limits.sls @@ -0,0 +1,206 @@ +(library (secmon server limits) + (export + make-listener-policy listener-policy? default-listener-policy + listener-policy-workers listener-policy-max-active + listener-policy-handshake-ms listener-policy-idle-ms + make-admission admission-try! admission-release! + admission-active admission-peer-active + make-fixed-worker-pool fixed-worker-pool-start! + fixed-worker-pool-submit! fixed-worker-pool-stop!) + (import (chezscheme) (jerboa prelude clean)) + + ;; #(listener-policy workers max-active per-peer-active global-rate + ;; per-peer-rate rate-window-ms handshake-ms idle-ms) + (define (make-listener-policy workers max-active per-peer-active + global-rate per-peer-rate rate-window-ms + handshake-ms idle-ms) + (for-each + (lambda (entry) + (unless (and (integer? (cdr entry)) (> (cdr entry) 0)) + (error 'make-listener-policy "bounds must be positive integers" entry))) + `((workers . ,workers) (max-active . ,max-active) + (per-peer-active . ,per-peer-active) (global-rate . ,global-rate) + (per-peer-rate . ,per-peer-rate) (rate-window-ms . ,rate-window-ms) + (handshake-ms . ,handshake-ms) (idle-ms . ,idle-ms))) + (when (> max-active workers) + (error 'make-listener-policy + "max-active must not exceed the fixed worker count")) + (when (> per-peer-active max-active) + (error 'make-listener-policy + "per-peer cap must not exceed the global active cap")) + (vector 'listener-policy workers max-active per-peer-active + global-rate per-peer-rate rate-window-ms handshake-ms idle-ms)) + + (define (listener-policy? value) + (and (vector? value) (= (vector-length value) 9) + (eq? (vector-ref value 0) 'listener-policy))) + (define (listener-policy-workers p) (vector-ref p 1)) + (define (listener-policy-max-active p) (vector-ref p 2)) + (define (policy-peer-active p) (vector-ref p 3)) + (define (policy-global-rate p) (vector-ref p 4)) + (define (policy-peer-rate p) (vector-ref p 5)) + (define (policy-window-ms p) (vector-ref p 6)) + (define (listener-policy-handshake-ms p) (vector-ref p 7)) + (define (listener-policy-idle-ms p) (vector-ref p 8)) + + (define default-listener-policy + (make-listener-policy 32 32 8 128 32 10000 5000 30000)) + + (define (wall-ms) + (let ([t (current-time)]) + (+ (* (time-second t) 1000) + (quotient (time-nanosecond t) 1000000)))) + + ;; #(admission policy mutex rows active global-start global-attempts now) + ;; source row: #(peer active window-start attempts) + (define (make-admission policy . maybe-now) + (unless (listener-policy? policy) + (error 'make-admission "invalid listener policy")) + (vector 'admission policy (make-mutex) '() 0 0 0 + (if (pair? maybe-now) (car maybe-now) wall-ms))) + (define (a-policy a) (vector-ref a 1)) + (define (a-mutex a) (vector-ref a 2)) + (define (a-rows a) (vector-ref a 3)) + (define (a-rows-set! a x) (vector-set! a 3 x)) + (define (a-active a) (vector-ref a 4)) + (define (a-active-set! a x) (vector-set! a 4 x)) + (define (a-global-start a) (vector-ref a 5)) + (define (a-global-start-set! a x) (vector-set! a 5 x)) + (define (a-global-attempts a) (vector-ref a 6)) + (define (a-global-attempts-set! a x) (vector-set! a 6 x)) + (define (a-now a) (vector-ref a 7)) + + (define (with-lock mutex thunk) + (dynamic-wind + (lambda () (mutex-acquire mutex)) + thunk + (lambda () (mutex-release mutex)))) + + (define (find-row rows peer) + (let loop ([xs rows]) + (cond [(null? xs) #f] + [(string=? (vector-ref (car xs) 0) peer) (car xs)] + [else (loop (cdr xs))]))) + + (define (admission-try! admission peer) + (with-lock + (a-mutex admission) + (lambda () + (let* ([policy (a-policy admission)] + [now ((a-now admission))] + [window (policy-window-ms policy)]) + (when (or (= (a-global-start admission) 0) + (>= (- now (a-global-start admission)) window)) + (a-global-start-set! admission now) + (a-global-attempts-set! admission 0)) + (let* ([rows (filter + (lambda (row) + (or (> (vector-ref row 1) 0) + (< (- now (vector-ref row 2)) window))) + (a-rows admission))] + [row (find-row rows peer)]) + (a-rows-set! admission rows) + (when (and (not row) (< (length rows) 1024)) + (set! row (vector peer 0 now 0)) + (a-rows-set! admission (cons row rows))) + (a-global-attempts-set! admission (+ 1 (a-global-attempts admission))) + (if (not row) + #f + (begin + (when (>= (- now (vector-ref row 2)) window) + (vector-set! row 2 now) + (vector-set! row 3 0)) + (vector-set! row 3 (+ 1 (vector-ref row 3))) + (if (and (<= (a-global-attempts admission) + (policy-global-rate policy)) + (<= (vector-ref row 3) (policy-peer-rate policy)) + (< (a-active admission) + (listener-policy-max-active policy)) + (< (vector-ref row 1) (policy-peer-active policy))) + (begin + (a-active-set! admission (+ 1 (a-active admission))) + (vector-set! row 1 (+ 1 (vector-ref row 1))) + #t) + #f)))))))) + + (define (admission-release! admission peer) + (with-lock + (a-mutex admission) + (lambda () + (let ([row (find-row (a-rows admission) peer)]) + (when (and row (> (vector-ref row 1) 0)) + (vector-set! row 1 (- (vector-ref row 1) 1)) + (a-active-set! admission (max 0 (- (a-active admission) 1)))))))) + + (define (admission-active admission) + (with-lock (a-mutex admission) (lambda () (a-active admission)))) + (define (admission-peer-active admission peer) + (with-lock + (a-mutex admission) + (lambda () + (let ([row (find-row (a-rows admission) peer)]) + (if row (vector-ref row 1) 0))))) + + ;; #(fixed-pool workers max-queue mutex condition queue closed threads) + (define (make-fixed-worker-pool workers max-queue) + (unless (and (integer? workers) (> workers 0) + (integer? max-queue) (> max-queue 0)) + (error 'make-fixed-worker-pool "invalid worker/queue bounds")) + (vector 'fixed-pool workers max-queue (make-mutex) (make-condition) + '() #f '())) + (define (p-workers p) (vector-ref p 1)) + (define (p-max-queue p) (vector-ref p 2)) + (define (p-mutex p) (vector-ref p 3)) + (define (p-condition p) (vector-ref p 4)) + (define (p-queue p) (vector-ref p 5)) + (define (p-queue-set! p x) (vector-set! p 5 x)) + (define (p-closed? p) (vector-ref p 6)) + (define (p-closed-set! p x) (vector-set! p 6 x)) + + (define (pool-take! pool) + (mutex-acquire (p-mutex pool)) + (let loop () + (cond + [(pair? (p-queue pool)) + (let ([job (car (p-queue pool))]) + (p-queue-set! pool (cdr (p-queue pool))) + (mutex-release (p-mutex pool)) + job)] + [(p-closed? pool) + (mutex-release (p-mutex pool)) + #f] + [else + (condition-wait (p-condition pool) (p-mutex pool)) + (loop)]))) + + (define (fixed-worker-pool-start! pool) + (let loop ([n (p-workers pool)] [threads '()]) + (if (= n 0) + (begin (vector-set! pool 7 threads) pool) + (loop (- n 1) + (cons (fork-thread + (lambda () + (let worker-loop () + (let ([job (pool-take! pool)]) + (when job + (guard (e [#t (void)]) (job)) + (worker-loop)))))) + threads))))) + + (define (fixed-worker-pool-submit! pool job) + (mutex-acquire (p-mutex pool)) + (let ([accepted (and (not (p-closed? pool)) + (< (length (p-queue pool)) (p-max-queue pool)))]) + (when accepted + (p-queue-set! pool (append (p-queue pool) (list job))) + (condition-signal (p-condition pool))) + (mutex-release (p-mutex pool)) + accepted)) + + (define (fixed-worker-pool-stop! pool) + (mutex-acquire (p-mutex pool)) + (p-closed-set! pool #t) + (condition-broadcast (p-condition pool)) + (mutex-release (p-mutex pool))) + + ) --- a/lib/secmon/server/listener.sls +++ b/lib/secmon/server/listener.sls @@ -1,11 +1,19 @@ (library (secmon server listener) - (export make-poll-server poll-server-run!) + (export make-poll-server poll-server-run! + poll-server-start! poll-server-wait! poll-server-stop! + poll-server-listen-host poll-server-listen-port + poll-server-active-connections + open-tcp-listener tcp-listener? tcp-listener-host tcp-listener-port + close-tcp-listener!) (import (chezscheme) (jerboa prelude clean) (secmon crypto psk) (secmon buffer ring) - (secmon server protocol)) + (secmon server protocol) + (secmon server limits) + (prefix (secmon server tcp) tcp:) + (only (std os posix) posix-errno posix-strerror)) (define-record-type poll-server (fields psk-auth buffer listen-addr start-time max-challenge-age) @@ -15,62 +23,199 @@ (lambda (psk-auth buffer listen-addr) (new psk-auth buffer listen-addr (current-time) 60))))) + (define-record-type tcp-listener + (fields fd host port) + (nongenerative secmon-tcp-listener-type)) + + (define-record-type accepted-client + (fields fd peer accepted-ms) + (nongenerative secmon-accepted-client-type)) + + (define-record-type server-runtime + (fields server listener stopped admission pool policy) + (nongenerative secmon-server-runtime-type)) + + (define *darwin?* + (memq (machine-type) '(a6osx ta6osx i3osx ti3osx arm64osx tarm64osx))) + (define *freebsd?* + (memq (machine-type) '(a6fb ta6fb i3fb ti3fb arm64fb tarm64fb))) + (define *bsd?* (or *darwin?* *freebsd?*)) + + ;; Resolve libc FFI symbols from the host/static process. Do not call + ;; load-shared-object while a library initializes: static images cannot catch + ;; that loader failure. + (define c-socket (foreign-procedure "socket" (int int int) int)) + (define c-bind (foreign-procedure "bind" (int u8* int) int)) + (define c-listen (foreign-procedure "listen" (int int) int)) + (define c-accept + (foreign-procedure __collect_safe "accept" (int u8* u8*) int)) + (define c-close (foreign-procedure "close" (int) int)) + (define c-dup (foreign-procedure "dup" (int) int)) + (define c-setsockopt + (foreign-procedure "setsockopt" (int int int u8* int) int)) + (define c-inet-pton (foreign-procedure "inet_pton" (int string u8*) int)) + (define c-getsockname + (foreign-procedure "getsockname" (int u8* u8*) int)) + (define c-fcntl2 (foreign-procedure "fcntl" (int int) int)) + (define c-fcntl3 + (foreign-procedure (__varargs_after 2) "fcntl" (int int int) int)) + + (define AF-INET 2) + (define SOCK-STREAM 1) + (define SOL-SOCKET (if *bsd?* #xffff 1)) + (define SO-REUSEADDR (if *bsd?* 4 2)) + (define SO-SNDTIMEO (if *bsd?* #x1005 21)) + (define SO-RCVTIMEO (if *bsd?* #x1006 20)) + (define F-GETFD 1) + (define F-SETFD 2) + (define FD-CLOEXEC 1) + (define F-GETFL 3) + (define F-SETFL 4) + (define O-NONBLOCK (if *bsd?* 4 #x800)) + (define EINTR 4) + (define EAGAIN (if *bsd?* 35 11)) + (define SOCKADDR-IN-SIZE 16) + (define (current-time-ms) (let ([t (current-time)]) (+ (* (time-second t) 1000) (quotient (time-nanosecond t) 1000000)))) + (define (socket-error who message . irritants) + (let ([e (posix-errno)]) + (apply error who + (string-append message ": " (posix-strerror e)) + (append irritants (list e))))) + + (define (set-cloexec! fd who) + (let ([flags (c-fcntl2 fd F-GETFD)]) + (when (< flags 0) (socket-error who "fcntl(F_GETFD) failed" fd)) + (when (< (c-fcntl3 fd F-SETFD (bitwise-ior flags FD-CLOEXEC)) 0) + (socket-error who "fcntl(F_SETFD) failed" fd)))) + + (define (set-nonblocking! fd) + (let ([flags (c-fcntl2 fd F-GETFL)]) + (when (< flags 0) (socket-error 'open-tcp-listener "fcntl(F_GETFL) failed")) + (when (< (c-fcntl3 fd F-SETFL (bitwise-ior flags O-NONBLOCK)) 0) + (socket-error 'open-tcp-listener "fcntl(F_SETFL) failed")))) + + (define (set-int-sockopt! fd level option value) + (let ([bytes (make-bytevector 4 0)]) + (bytevector-s32-native-set! bytes 0 value) + (c-setsockopt fd level option bytes 4))) + + (define (set-socket-timeout! fd timeout-ms) + (let ([tv (make-bytevector 16 0)]) + (bytevector-s64-native-set! tv 0 (quotient timeout-ms 1000)) + (bytevector-s64-native-set! tv 8 (* (mod timeout-ms 1000) 1000)) + (when (< (c-setsockopt fd SOL-SOCKET SO-RCVTIMEO tv 16) 0) + (socket-error 'set-socket-timeout! "SO_RCVTIMEO failed" fd)) + (when (< (c-setsockopt fd SOL-SOCKET SO-SNDTIMEO tv 16) 0) + (socket-error 'set-socket-timeout! "SO_SNDTIMEO failed" fd)))) + + (define (make-sockaddr-in host port) + (unless (and (string? host) (> (string-length host) 0)) + (error 'make-sockaddr-in "host must be a numeric IPv4 address" host)) + (unless (and (integer? port) (>= port 0) (<= port 65535)) + (error 'make-sockaddr-in "port must be in 0..65535" port)) + (let ([addr (make-bytevector SOCKADDR-IN-SIZE 0)] + [ip (make-bytevector 4 0)]) + (unless (= (c-inet-pton AF-INET host ip) 1) + (error 'make-sockaddr-in + "unsupported or ambiguous listen address; use numeric IPv4" host)) + (if *bsd?* + (begin + (bytevector-u8-set! addr 0 SOCKADDR-IN-SIZE) + (bytevector-u8-set! addr 1 AF-INET)) + (bytevector-u16-native-set! addr 0 AF-INET)) + (bytevector-u16-set! addr 2 port (endianness big)) + (bytevector-copy! ip 0 addr 4 4) + addr)) + + (define (sockaddr-host addr) + (format "~a.~a.~a.~a" + (bytevector-u8-ref addr 4) (bytevector-u8-ref addr 5) + (bytevector-u8-ref addr 6) (bytevector-u8-ref addr 7))) + (define (sockaddr-port addr) + (bytevector-u16-ref addr 2 (endianness big))) + + (define (actual-sockaddr fd who) + (let ([addr (make-bytevector SOCKADDR-IN-SIZE 0)] + [len (make-bytevector 4 0)]) + (bytevector-u32-native-set! len 0 SOCKADDR-IN-SIZE) + (when (< (c-getsockname fd addr len) 0) + (socket-error who "getsockname failed" fd)) + addr)) + + (define (open-tcp-listener host port) + (let ([fd (c-socket AF-INET SOCK-STREAM 0)]) + (when (< fd 0) (socket-error 'open-tcp-listener "socket failed")) + (guard (e [#t (c-close fd) (raise e)]) + (set-cloexec! fd 'open-tcp-listener) + (when (< (set-int-sockopt! fd SOL-SOCKET SO-REUSEADDR 1) 0) + (socket-error 'open-tcp-listener "SO_REUSEADDR failed")) + (let ([addr (make-sockaddr-in host port)]) + (when (< (c-bind fd addr SOCKADDR-IN-SIZE) 0) + (socket-error 'open-tcp-listener "bind failed" host port))) + (when (< (c-listen fd 128) 0) + (socket-error 'open-tcp-listener "listen failed" host port)) + (set-nonblocking! fd) + (let ([actual (actual-sockaddr fd 'open-tcp-listener)]) + (make-tcp-listener fd (sockaddr-host actual) (sockaddr-port actual)))))) + + (define (close-tcp-listener! listener) + (guard (e [#t #f]) (c-close (tcp-listener-fd listener)))) + + (define (accept-client listener) + (let ([peer (make-bytevector SOCKADDR-IN-SIZE 0)] + [len (make-bytevector 4 0)]) + (bytevector-u32-native-set! len 0 SOCKADDR-IN-SIZE) + (let ([fd (c-accept (tcp-listener-fd listener) peer len)]) + (if (< fd 0) + (let ([e (posix-errno)]) + (if (or (= e EINTR) (= e EAGAIN)) + #f + (socket-error 'accept-client "accept failed"))) + (begin + (set-cloexec! fd 'accept-client) + (make-accepted-client fd (sockaddr-host peer) (current-time-ms))))))) + (define (uptime-secs server) - (let ([start (poll-server-start-time server)] - [now (current-time)]) - (- (time-second now) (time-second start)))) - - ;; --- TCP I/O helpers --- - - (define (tcp-read-exact port n) - (let ([buf (make-bytevector n)]) - (let loop ([offset 0]) - (if (>= offset n) buf - (let ([b (get-u8 port)]) - (when (eof-object? b) - (error 'tcp-read-exact "connection closed")) - (bytevector-u8-set! buf offset b) - (loop (+ offset 1))))))) - - (define (tcp-write-all port bv) - (do ([i 0 (+ i 1)]) - ((= i (bytevector-length bv))) - (put-u8 port (bytevector-u8-ref bv i))) - (flush-output-port port)) - - ;; --- Transport-encrypted message I/O --- - - (define (send-message! port server payload) + (- (time-second (current-time)) + (time-second (poll-server-start-time server)))) + + ;; The deadline is absolute for the whole field/frame. SO_RCVTIMEO only + ;; bounds one blocking read; checking the absolute clock after each chunk is + ;; what prevents a peer from extending its lease by dripping bytes. + (define (tcp-read-exact fd port n deadline-ms) + (tcp:read-exact/deadline port n deadline-ms fd)) + + (define (tcp-write-all fd bv deadline-ms) + (tcp:write-all/deadline fd bv deadline-ms)) + + (define (send-message! fd server payload deadline-ms) (let* ([encrypted (psk-encrypt-transport (poll-server-psk-auth server) payload)] [len-bytes (pack-u32-le (bytevector-length encrypted))]) - (tcp-write-all port len-bytes) - (tcp-write-all port encrypted))) + (tcp-write-all fd len-bytes deadline-ms) + (tcp-write-all fd encrypted deadline-ms))) - (define (recv-message! port server) - (guard (e [#t #f]) - (let* ([len-bytes (tcp-read-exact port 4)] - [len (unpack-u32-le len-bytes 0)]) - (when (> len MAX-MESSAGE-SIZE) - (error 'recv-message! "message too large" len)) - (let ([encrypted (tcp-read-exact port len)]) - (psk-decrypt-transport (poll-server-psk-auth server) encrypted))))) - - ;; --- Request handling --- + (define (recv-message! fd port server deadline-ms) + (let* ([len-bytes (tcp-read-exact fd port 4 deadline-ms)] + [len (unpack-u32-le len-bytes 0)]) + (when (or (= len 0) (> len MAX-MESSAGE-SIZE)) + (error 'recv-message! "invalid encrypted message length" len)) + (psk-decrypt-transport + (poll-server-psk-auth server) + (tcp-read-exact fd port len deadline-ms)))) (define (handle-request server req) (case (car req) [(get-events-after) - (let ([events (buffer-get-after (poll-server-buffer server) (cadr req))]) - (pack-events-response events))] + (pack-events-response + (buffer-get-after (poll-server-buffer server) (cadr req)))] [(get-events-range) - (let ([events (buffer-get-in-range (poll-server-buffer server) - (cadr req) (caddr req))]) - (pack-events-response events))] + (pack-events-response + (buffer-get-in-range (poll-server-buffer server) (cadr req) (caddr req)))] [(status) (pack-status-response (buffer-count (poll-server-buffer server)) @@ -79,114 +224,137 @@ [(acknowledge) (buffer-clear-before! (poll-server-buffer server) (cadr req)) (pack-acked-response (cadr req))] - [(ping) - (pack-pong-response (current-time-ms))] - [else - (pack-error-msg "unknown request")])) - - ;; --- Connection handler --- - - (define (handle-connection server in-port out-port) - (guard (e [#t (void)]) ;; Silently close on any error (stealth) - ;; 1. Send challenge - (let ([challenge (psk-create-challenge (poll-server-psk-auth server))]) - (send-message! out-port server (pack-challenge challenge)) - ;; 2. Receive response - (let ([resp-bytes (recv-message! in-port server)]) - (unless resp-bytes (error #f "no response")) - (let ([msg-tag (bytevector-u8-ref resp-bytes 0)]) - (unless (= msg-tag MSG-CHALLENGE-RESPONSE) - (error #f "unexpected message")) - ;; 3. Parse and verify - (let* ([resp-data (make-bytevector (- (bytevector-length resp-bytes) 1))] - [_ (bytevector-copy! resp-bytes 1 resp-data 0 - (bytevector-length resp-data))] - [response (unpack-challenge-response resp-data)]) - (if (not (psk-verify-response - (poll-server-psk-auth server) - challenge response - (poll-server-max-challenge-age server))) - ;; Auth failed - (send-message! out-port server (pack-auth-failed)) - ;; 4. Authenticated request loop - (let loop () - (let ([req-bytes (recv-message! in-port server)]) - (when req-bytes - (let ([msg-tag (bytevector-u8-ref req-bytes 0)]) - (when (= msg-tag MSG-REQUEST) - (let* ([req-data (make-bytevector (- (bytevector-length req-bytes) 1))] - [_ (bytevector-copy! req-bytes 1 req-data 0 - (bytevector-length req-data))] - [req (unpack-request req-data)] - [resp (handle-request server req)]) - (send-message! out-port server resp) - (loop)))))))))))))) - - - ;; --- Main server loop --- + [(ping) (pack-pong-response (current-time-ms))] + [else (pack-error-msg "unknown request")])) + + (define (handle-connection server client policy) + (let* ([fd (tcp:accepted-client-fd client)] + [in-port + (guard (e [#t (tcp:close-accepted-client! client) (raise e)]) + (open-fd-input-port fd))]) + (dynamic-wind + (lambda () (void)) + (lambda () + (guard (e [#t (void)]) + (let* ([handshake-ms (listener-policy-handshake-ms policy)] + [deadline (+ (tcp:accepted-client-accepted-ms client) handshake-ms)] + [challenge (psk-create-challenge (poll-server-psk-auth server))]) + (tcp:set-socket-timeout! fd handshake-ms) + (send-message! fd server (pack-challenge challenge) deadline) + (let ([resp-bytes (recv-message! fd in-port server deadline)]) + (unless (and (bytevector? resp-bytes) + (> (bytevector-length resp-bytes) 0) + (= (bytevector-u8-ref resp-bytes 0) + MSG-CHALLENGE-RESPONSE)) + (error 'handle-connection "invalid challenge response")) + (let* ([resp-data (make-bytevector + (- (bytevector-length resp-bytes) 1))] + [_ (bytevector-copy! resp-bytes 1 resp-data 0 + (bytevector-length resp-data))] + [response (unpack-challenge-response resp-data)]) + (if (not (psk-verify-response + (poll-server-psk-auth server) challenge response + (poll-server-max-challenge-age server))) + (send-message! fd server (pack-auth-failed) deadline) + (let request-loop () + (let* ([idle-ms (listener-policy-idle-ms policy)] + [_ (tcp:set-socket-timeout! fd idle-ms)] + [request-deadline (+ (current-time-ms) idle-ms)] + [req-bytes + (recv-message! fd in-port server request-deadline)]) + (when (and (bytevector? req-bytes) + (> (bytevector-length req-bytes) 0) + (= (bytevector-u8-ref req-bytes 0) MSG-REQUEST)) + (let* ([req-data (make-bytevector + (- (bytevector-length req-bytes) 1))] + [_ (bytevector-copy! req-bytes 1 req-data 0 + (bytevector-length req-data))] + [req (unpack-request req-data)]) + (send-message! fd server + (handle-request server req) request-deadline) + (request-loop))))))))))) + (lambda () + (guard (e [#t (void)]) (close-port in-port)))))) + + (define (parse-listen-address value) + (let ([parts (string-split value #\:)]) + (unless (= (length parts) 2) + (error 'poll-server-start! "listen address must be numeric IPv4:port" value)) + (let ([host (car parts)] [port (string->number (cadr parts))]) + (unless (and (integer? port) (> port 0) (<= port 65535)) + (error 'poll-server-start! "invalid listen port" value)) + ;; Constructing the sockaddr performs strict numeric IPv4 validation. + (make-sockaddr-in host port) + (values host port)))) + + (define (accept-loop runtime) + (let loop () + (unless (unbox (server-runtime-stopped runtime)) + (guard (e (#t (sleep (make-time 'time-duration 10000000 0)))) + (let ((client (tcp:accept-client (server-runtime-listener runtime)))) + (cond + ((not client) + (sleep (make-time 'time-duration 10000000 0))) + (else + (let* ((peer (tcp:accepted-client-peer client)) + (admission (server-runtime-admission runtime)) + (pool (server-runtime-pool runtime)) + (policy (server-runtime-policy runtime))) + (cond + ((not (admission-try! admission peer)) + (tcp:close-accepted-client! client)) + (else + (let ((job + (lambda () + (dynamic-wind + (lambda () (void)) + (lambda () + (handle-connection + (server-runtime-server runtime) client policy)) + (lambda () + (admission-release! admission peer)))))) + (unless (fixed-worker-pool-submit! pool job) + (admission-release! admission peer) + (tcp:close-accepted-client! client)))))))))) + (loop)))) + + (define (poll-server-start! server . maybe-policy) + (let ([policy (if (pair? maybe-policy) + (car maybe-policy) default-listener-policy)]) + (unless (listener-policy? policy) + (error 'poll-server-start! "invalid listener policy")) + (let-values ([(host port) + (parse-listen-address (poll-server-listen-addr server))]) + (let* ([listener (tcp:open-tcp-listener host port)] + [admission (make-admission policy)] + [pool (make-fixed-worker-pool + (listener-policy-workers policy) + (listener-policy-max-active policy))] + [runtime (make-server-runtime server listener (box #f) + admission pool policy)]) + (fixed-worker-pool-start! pool) + (fork-thread (lambda () (accept-loop runtime))) + runtime)))) + + (define (poll-server-listen-host runtime) + (tcp:tcp-listener-host (server-runtime-listener runtime))) + (define (poll-server-listen-port runtime) + (tcp:tcp-listener-port (server-runtime-listener runtime))) + (define (poll-server-active-connections runtime) + (admission-active (server-runtime-admission runtime))) + + (define (poll-server-stop! runtime) + (set-box! (server-runtime-stopped runtime) #t) + (tcp:close-tcp-listener! (server-runtime-listener runtime)) + (fixed-worker-pool-stop! (server-runtime-pool runtime))) + + (define (poll-server-wait! runtime) + (let loop () + (unless (unbox (server-runtime-stopped runtime)) + (sleep (make-time 'time-duration 0 3600)) + (loop)))) (define (poll-server-run! server) - (let* ([addr-parts (string-split (poll-server-listen-addr server) #\:)] - [host (if (>= (length addr-parts) 2) (car addr-parts) "0.0.0.0")] - [port (string->number (if (>= (length addr-parts) 2) - (cadr addr-parts) - (car addr-parts)))] - [listener (listen-tcp-port port)]) - (let loop () - (guard (e [#t (loop)]) ;; Silently restart on errors - (let-values ([(in out) (accept-tcp listener)]) - ;; Spawn handler thread per connection - (fork-thread - (lambda () - (guard (e [#t (void)]) - (handle-connection server in out) - (close-port in) - (close-port out)))) - (loop)))))) - - ;; --- TCP server helpers using Chez primitives --- - - (define (listen-tcp-port port) - (let ([sock (foreign-procedure "socket" (int int int) int)]) - (let ([fd (sock 2 1 0)]) ;; AF_INET=2, SOCK_STREAM=1 - (when (< fd 0) - (error 'listen-tcp-port "socket() failed")) - ;; SO_REUSEADDR - (let ([setsockopt (foreign-procedure "setsockopt" (int int int u8* int) int)] - [optval (make-bytevector 4)]) - (bytevector-u32-set! optval 0 1 (endianness little)) - (setsockopt fd 1 2 optval 4)) ;; SOL_SOCKET=1, SO_REUSEADDR=2 - ;; Bind - (let ([bind-fn (foreign-procedure "bind" (int u8* int) int)] - [addr (make-sockaddr-in port)]) - (let ([rc (bind-fn fd addr (bytevector-length addr))]) - (when (< rc 0) - (error 'listen-tcp-port "bind() failed" port)))) - ;; Listen - (let ([listen-fn (foreign-procedure "listen" (int int) int)]) - (let ([rc (listen-fn fd 128)]) - (when (< rc 0) - (error 'listen-tcp-port "listen() failed")))) - fd))) - - (define (make-sockaddr-in port) - ;; struct sockaddr_in: family(2) port(2) addr(4) zero(8) = 16 bytes - (let ([addr (make-bytevector 16 0)]) - (bytevector-u16-set! addr 0 2 (endianness little)) ;; AF_INET - (bytevector-u16-set! addr 2 port (endianness big)) ;; port in network order - ;; addr = 0.0.0.0 (already zeroed) - addr)) + (poll-server-wait! (poll-server-start! server))) - (define (accept-tcp listener-fd) - (let ([accept-fn (foreign-procedure "accept" (int u8* u8*) int)] - [peer-addr (make-bytevector 16 0)] - [addr-len (make-bytevector 4)])