From 22ad7e209983f4ef823fafbfa4ceca591b390284 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 15 Sep 2026 05:21:42 +0000 Subject: [PATCH 1/8] =?UTF-8?q?test:=20add=20Thalamic=20=E2=86=92=20corpus?= =?UTF-8?q?-ipc=20=E2=86=92=20Brainstem=20CPU=20smoke=20harness?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit LIM-1135 / GH#43. Load an explicit JSON checkpoint, decode published IpcMessage::Stimuli frames (valid_mask included), reject schema incompatibility, and keep the Thalamic fixture healthy when Brainstem is unavailable. Pin corpus-ipc 0.1.0 and MSRV 1.98.1. Co-authored-by: Raul Cardenas Montoya --- .devin/blueprint.yaml | 4 +- .github/workflows/ci.yml | 12 +- AGENTS.md | 8 +- CHANGELOG.md | 10 +- Cargo.lock | 308 +---------------------- Cargo.toml | 8 +- Dockerfile | 2 +- README.md | 29 ++- docs/ci.md | 4 +- rust-toolchain.toml | 2 +- src/backend.rs | 211 ++++++++++++---- src/bin/brainstem_daemon.rs | 5 +- src/checkpoint.rs | 363 ++++++++++++++++++++++++++++ src/daemon.rs | 132 +++++++++- src/ingress.rs | 272 +++++++++++++++++++++ src/lib.rs | 12 +- tests/fixtures/thalamic_producer.rs | 104 ++++++++ tests/thalamic_brainstem_smoke.rs | 233 ++++++++++++++++++ 18 files changed, 1327 insertions(+), 392 deletions(-) create mode 100644 src/checkpoint.rs create mode 100644 src/ingress.rs create mode 100644 tests/fixtures/thalamic_producer.rs create mode 100644 tests/thalamic_brainstem_smoke.rs diff --git a/.devin/blueprint.yaml b/.devin/blueprint.yaml index 2f7a992..106ccd2 100644 --- a/.devin/blueprint.yaml +++ b/.devin/blueprint.yaml @@ -2,8 +2,8 @@ initialize: | sudo apt-get update sudo apt-get install -y libzmq3-dev - rustup toolchain install 1.97.1 - rustup default 1.97.1 + rustup toolchain install 1.98.1 + rustup default 1.98.1 rustup component add rustfmt clippy maintenance: | cargo fetch diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index ba24109..78d1409 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -27,9 +27,9 @@ jobs: persist-credentials: false - name: Install Rust toolchain - uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # 1.97.1 + uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # 1.98.1 with: - toolchain: "1.97.1" + toolchain: "1.98.1" components: rustfmt - name: Check formatting @@ -48,9 +48,9 @@ jobs: persist-credentials: false - name: Install Rust toolchain - uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # 1.97.1 + uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # 1.98.1 with: - toolchain: "1.97.1" + toolchain: "1.98.1" components: clippy - name: Cache cargo build @@ -88,9 +88,9 @@ jobs: sudo apt-get install -y libzmq3-dev - name: Install Rust toolchain - uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # 1.97.1 + uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # 1.98.1 with: - toolchain: "1.97.1" + toolchain: "1.98.1" components: clippy - name: Cache cargo build diff --git a/AGENTS.md b/AGENTS.md index cb91cea..7673529 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -45,7 +45,7 @@ If a `--all-features` build fails because the C++ compiler cannot find a standar ## Cursor Cloud setup -This repository is preconfigured on the Cursor Cloud virtual machine. The Rust toolchain is pinned to **1.97.1 only** via `rust-toolchain.toml`. At startup the environment runs `cargo fetch`. +This repository is preconfigured on the Cursor Cloud virtual machine. The Rust toolchain is pinned to **1.98.1 only** via `rust-toolchain.toml`. At startup the environment runs `cargo fetch`. ### Backend features @@ -83,7 +83,7 @@ spine_pub_port = 5556 model_path = "~/models/soma16.mem" ``` -The `model_path` is not used by the stub backend. -With `--features corpus-ipc` the binary passes it literally to `ZmqStimulusSource::initialize`. -`~` is not expanded, and the pinned `ZmqBrainBackend` currently ignores `_model_path`. +The `model_path` is loaded as a JSON checkpoint when the file exists; otherwise the runtime constructs a blank `SpikingNetwork::with_dimensions`. +With `--features corpus-ipc` the binary also connects a ZMQ SUB socket for JSON `IpcMessage` frames (`CORPUS_IPC_ZMQ_READOUT_IPC` / `SPIKENAUT_ZMQ_READOUT_IPC`). +`~` is not expanded. With the stub backend, `brainstem-daemon` runs a headless spiking-neural-network tick loop and logs `🔌 Using stub backend`. diff --git a/CHANGELOG.md b/CHANGELOG.md index 5c98976..88e6ec5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,18 +12,22 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - GitHub Actions CI matrix: stub build/test/clippy on Linux, macOS, and Windows; rustfmt and optional `corpus-ipc` / libzmq jobs on Linux only (`docs/ci.md`). -- Stub vs `corpus-ipc` backend feature truth table in `README.md`: which Cargo flags wire which backend, which TOML keys apply, and which env vars are no-ops under stub. Documents that `model_path` is passed literally (no `~` expansion) but currently ignored by pinned `ZmqBrainBackend`, that `CORPUS_IPC_ZMQ_READOUT_IPC` is binary-set compatibility only, and that `log_level` is binary tracing-init only. +- Stub vs `corpus-ipc` backend feature truth table in `README.md`: which Cargo flags wire which backend, which TOML keys apply, and which env vars are no-ops under stub. Documents that `model_path` is a checkpoint path when the file exists (otherwise a blank network), that `CORPUS_IPC_ZMQ_READOUT_IPC` is the SUB endpoint for typed JSON `IpcMessage` frames, and that `log_level` is binary tracing-init only. - GitHub Actions CI workflow for formatting, clippy, build, and test validation. - Config-driven `ServiceRegistry` and `BrainstemDaemon` in the library. - `DaemonConfig.services` field for registering named, enabled services. - `## Role and boundary matrix` documentation in `README.md`. - Local `StimulusSource` / `SpikeSink` traits + `IngressPacket` / `SpikeEvent` (owned by this crate). - `BackendPair` + `BackendPair::stub()` for pluggable I/O. -- In-crate stub backend (`StubStimulusSource`, `NoopSpikeSink`, `CollectingSpikeSink` under `#[cfg(test)]` for our own tests; not re-exported for downstream test use). +- In-crate stub backend (`StubStimulusSource`, `NoopSpikeSink`, `CollectingSpikeSink`). - `BrainstemDaemon::with_backend(cfg, pair)` constructor for tests and custom backends. - Test coverage for the non-`corpus-ipc` (stub) path that runs under `--no-default-features`. - Graceful `SIGTERM` handling alongside the existing `SIGINT` (Ctrl-C): the tick loop now breaks, flushes the backend, and exits `0` on either signal. +- Explicit JSON checkpoint loader (`src/checkpoint.rs`) that fails closed on schema, dimension, NaN, and blank-weight fixtures. +- Typed `corpus-ipc` ingress validation (`src/ingress.rs`) for width, freshness, `valid_mask`, and schema token `corpus-ipc.stimulus.v1`. +- CPU-only Thalamic → corpus-ipc → Brainstem integration smoke test (`tests/thalamic_brainstem_smoke.rs`, `--features corpus-ipc`). +- `BrainstemDaemon::run_for_ticks` for bounded, signal-free tick runs. ### Changed @@ -36,6 +40,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Binary now logs the active backend mode (`🔌 stub` / `📡 ZMQ corpus-ipc`). - `decode_inputs` now accepts `&IngressPacket` (with explicit `None` modulator fallback). - All direct `corpus_ipc` / `zmq` usage is now feature-gated (except the compatibility `CORPUS_IPC_READOUT_ENV` const). +- `corpus-ipc` feature pins crate version **0.1.0** at git rev `3ad764a` (`IpcMessage` / `StimulusBatch`; crates.io does not yet resolve `corpus-ipc = "0.1"`). ZMQ ingress decodes JSON `IpcMessage::Stimuli` rather than the legacy binary readout packet. +- MSRV / `rust-toolchain.toml` aligned to **1.98.1** so the 0.1.0 `corpus-ipc` crate compiles. ### Removed diff --git a/Cargo.lock b/Cargo.lock index 667c302..3457d03 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -67,78 +67,6 @@ version = "1.0.102" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f202df86484c868dbad7eaa557ef785d5c66295e41b460ef922eca0723b842c" -[[package]] -name = "async-trait" -version = "0.1.89" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9035ad2d096bed7955a320ee7e2230574d28fd3c3a0f186cbea1ff3c7eed5dbb" -dependencies = [ - "proc-macro2", - "quote", - "syn", -] - -[[package]] -name = "atomic-waker" -version = "1.1.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0" - -[[package]] -name = "axum" -version = "0.7.9" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "edca88bc138befd0323b20752846e6587272d3b03b0343c8ea28a6f819e6e71f" -dependencies = [ - "async-trait", - "axum-core", - "bytes", - "futures-util", - "http", - "http-body", - "http-body-util", - "hyper", - "hyper-util", - "itoa", - "matchit", - "memchr", - "mime", - "percent-encoding", - "pin-project-lite", - "rustversion", - "serde", - "serde_json", - "serde_path_to_error", - "serde_urlencoded", - "sync_wrapper", - "tokio", - "tower 0.5.3", - "tower-layer", - "tower-service", - "tracing", -] - -[[package]] -name = "axum-core" -version = "0.4.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "09f2bd6146b97ae3359fa0cc6d6b376d9539582c7b4220f041a33ec24c226199" -dependencies = [ - "async-trait", - "bytes", - "futures-util", - "http", - "http-body", - "http-body-util", - "mime", - "pin-project-lite", - "rustversion", - "sync_wrapper", - "tower-layer", - "tower-service", - "tracing", -] - [[package]] name = "bitflags" version = "1.3.2" @@ -252,15 +180,10 @@ checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570" [[package]] name = "corpus-ipc" version = "0.1.0" -source = "git+https://github.com/Limen-Neural/corpus-ipc?rev=78220a6413c20f202252016afbf5981b7350dfd0#78220a6413c20f202252016afbf5981b7350dfd0" +source = "git+https://github.com/Limen-Neural/corpus-ipc?rev=3ad764a9765ead27d6120d25ecd963b9338940fb#3ad764a9765ead27d6120d25ecd963b9338940fb" dependencies = [ - "axum", "serde", - "serde_json", "thiserror", - "tokio", - "tower 0.4.13", - "zmq", ] [[package]] @@ -379,48 +302,6 @@ version = "0.1.9" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5baebc0774151f905a1a2cc41989300b1e6fbb29aff0ceffa1064fdd3088d582" -[[package]] -name = "form_urlencoded" -version = "1.2.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cb4cb245038516f5f85277875cdaa4f7d2c9a0fa0468de06ed190163b1581fcf" -dependencies = [ - "percent-encoding", -] - -[[package]] -name = "futures-channel" -version = "0.3.32" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "07bbe89c50d7a535e539b8c17bc0b49bdb77747034daa8087407d655f3f7cc1d" -dependencies = [ - "futures-core", -] - -[[package]] -name = "futures-core" -version = "0.3.32" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7e3450815272ef58cec6d564423f6e755e25379b217b0bc688e295ba24df6b1d" - -[[package]] -name = "futures-task" -version = "0.3.32" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "037711b3d59c33004d3856fbdc83b99d4ff37a24768fa1be9ce3538a1cde4393" - -[[package]] -name = "futures-util" -version = "0.3.32" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6" -dependencies = [ - "futures-core", - "futures-task", - "pin-project-lite", - "slab", -] - [[package]] name = "getrandom" version = "0.2.17" @@ -456,86 +337,6 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" -[[package]] -name = "http" -version = "1.4.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6970f50e31d6fc17d3fa27329444bfa74e196cf62e95052a3f6fee181dba6425" -dependencies = [ - "bytes", - "itoa", -] - -[[package]] -name = "http-body" -version = "1.0.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1efedce1fb8e6913f23e0c92de8e62cd5b772a67e7b3946df930a62566c93184" -dependencies = [ - "bytes", - "http", -] - -[[package]] -name = "http-body-util" -version = "0.1.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b021d93e26becf5dc7e1b75b1bed1fd93124b374ceb73f43d4d4eafec896a64a" -dependencies = [ - "bytes", - "futures-core", - "http", - "http-body", - "pin-project-lite", -] - -[[package]] -name = "httparse" -version = "1.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6dbf3de79e51f3d586ab4cb9d5c3e2c14aa28ed23d180cf89b4df0454a69cc87" - -[[package]] -name = "httpdate" -version = "1.0.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df3b46402a9d5adb4c86a0cf463f42e19994e3ee891101b1841f30a545cb49a9" - -[[package]] -name = "hyper" -version = "1.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "55281c53a1894c864990125767da440a4e630446785086f52523b20033b74498" -dependencies = [ - "atomic-waker", - "bytes", - "futures-channel", - "futures-core", - "http", - "http-body", - "httparse", - "httpdate", - "itoa", - "pin-project-lite", - "smallvec", - "tokio", -] - -[[package]] -name = "hyper-util" -version = "0.1.20" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "96547c2556ec9d12fb1578c4eaf448b04993e7fb79cbaad930a656880a6bdfa0" -dependencies = [ - "bytes", - "http", - "http-body", - "hyper", - "pin-project-lite", - "tokio", - "tower-service", -] - [[package]] name = "indexmap" version = "2.14.0" @@ -623,24 +424,12 @@ dependencies = [ "regex-automata", ] -[[package]] -name = "matchit" -version = "0.7.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0e7465ac9959cc2b1404e8e2367b43684a6d13790fe23056cc8c6c5a6b7bcb94" - [[package]] name = "memchr" version = "2.8.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "88904434abc2901f197fe8cc55f0445e7ded921dba5911dad2e2b39b48e663c4" -[[package]] -name = "mime" -version = "0.3.17" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6877bb514081ee2a7ff5ef9de3281f14a4dd4bceac4c09388074a6b5df8a139a" - [[package]] name = "mio" version = "1.2.1" @@ -713,12 +502,6 @@ dependencies = [ "windows-link", ] -[[package]] -name = "percent-encoding" -version = "2.3.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" - [[package]] name = "pin-project-lite" version = "0.2.17" @@ -850,18 +633,6 @@ version = "0.8.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4" -[[package]] -name = "rustversion" -version = "1.0.22" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b39cdef0fa800fc44525c84ccb54a029961a8215f9619753635a9c0d2538d46d" - -[[package]] -name = "ryu" -version = "1.0.23" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" - [[package]] name = "same-file" version = "1.0.6" @@ -920,17 +691,6 @@ dependencies = [ "zmij", ] -[[package]] -name = "serde_path_to_error" -version = "0.1.20" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "10a9ff822e371bb5403e391ecd83e182e0e77ba7f6fe0160b795797109d1b457" -dependencies = [ - "itoa", - "serde", - "serde_core", -] - [[package]] name = "serde_spanned" version = "0.6.9" @@ -940,18 +700,6 @@ dependencies = [ "serde", ] -[[package]] -name = "serde_urlencoded" -version = "0.7.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d3491c14715ca2294c4d6a88f15e84739788c1d030eed8c110436aafdaa2f3fd" -dependencies = [ - "form_urlencoded", - "itoa", - "ryu", - "serde", -] - [[package]] name = "sharded-slab" version = "0.1.7" @@ -977,12 +725,6 @@ dependencies = [ "libc", ] -[[package]] -name = "slab" -version = "0.4.12" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0c790de23124f9ab44544d7ac05d60440adc586479ce501c1d6d7da3cd8c9cf5" - [[package]] name = "smallvec" version = "1.15.2" @@ -1016,12 +758,6 @@ dependencies = [ "unicode-ident", ] -[[package]] -name = "sync_wrapper" -version = "1.0.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0bf256ce5efdfa370213c1dabab5935a12e49f2c58d15e9eac2870d3b4f27263" - [[package]] name = "system-deps" version = "6.2.2" @@ -1139,52 +875,12 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5d99f8c9a7727884afe522e9bd5edbfc91a3312b36a77b5fb8926e4c31a41801" -[[package]] -name = "tower" -version = "0.4.13" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b8fa9be0de6cf49e536ce1851f987bd21a43b771b09473c3549a6c853db37c1c" -dependencies = [ - "tower-layer", - "tower-service", - "tracing", -] - -[[package]] -name = "tower" -version = "0.5.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ebe5ef63511595f1344e2d5cfa636d973292adc0eec1f0ad45fae9f0851ab1d4" -dependencies = [ - "futures-core", - "futures-util", - "pin-project-lite", - "sync_wrapper", - "tokio", - "tower-layer", - "tower-service", - "tracing", -] - -[[package]] -name = "tower-layer" -version = "0.3.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "121c2a6cda46980bb0fcd1647ffaf6cd3fc79a013de288782836f6df9c48780e" - -[[package]] -name = "tower-service" -version = "0.3.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8df9b6e13f2d32c91b9bd719c00d1958837bc7dec474d94952798cc8e69eeec3" - [[package]] name = "tracing" version = "0.1.44" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "63e71662fa4b2a2c3a26f570f037eb95bb1f85397f3cd8076caed2f026a6d100" dependencies = [ - "log", "pin-project-lite", "tracing-attributes", "tracing-core", @@ -1295,7 +991,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.48.0", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index e328f53..caf0657 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -2,7 +2,7 @@ name = "brainstem-daemon" version = "0.1.2" edition = "2024" -rust-version = "1.97.1" +rust-version = "1.98.1" description = "SNN Engine + Inference + Telemetry + AI (Tier 1 — headless)" authors = ["Raul Cardenas Montoya "] license = "MIT OR Apache-2.0" @@ -11,7 +11,7 @@ repository = "https://github.com/Limen-Neural/brainstem-daemon" [dependencies] anyhow = "1" clap = { version = "4", features = ["derive"] } -corpus-ipc = { git = "https://github.com/Limen-Neural/corpus-ipc", rev = "78220a6413c20f202252016afbf5981b7350dfd0", features = ["zmq"], optional = true } +corpus-ipc = { git = "https://github.com/Limen-Neural/corpus-ipc", rev = "3ad764a9765ead27d6120d25ecd963b9338940fb", default-features = false, optional = true } neuromod = "0.4.0" serde = { version = "1", features = ["derive"] } serde_json = "1" @@ -30,6 +30,10 @@ corpus-ipc = ["dep:corpus-ipc", "dep:zmq"] name = "brainstem-daemon" path = "src/bin/brainstem_daemon.rs" +[[test]] +name = "thalamic_brainstem_smoke" +required-features = ["corpus-ipc"] + # Cargo profiles (dev / release / release-with-debug / test / bench). # Keep settings conservative — no maximum-optimization knobs that over-promise # for a headless inference daemon. Aligned with the `neuromod` profile pattern. diff --git a/Dockerfile b/Dockerfile index 90778ff..32f571d 100644 --- a/Dockerfile +++ b/Dockerfile @@ -7,7 +7,7 @@ # CI and contributors can validate: # cargo fmt --check, clippy, build, test inside the image. -FROM rust:1.97.1-bookworm AS base +FROM rust:1.98.1-bookworm AS base WORKDIR /app # Common system deps for the full feature set (libzmq). Core-only builds do not need this. # pkgconf provides /usr/bin/pkg-config on Debian bookworm. diff --git a/README.md b/README.md index b766678..ac93c15 100644 --- a/README.md +++ b/README.md @@ -20,7 +20,7 @@ Headless spiking neural-network runtime written in Rust. ## Building -Requires **Rust 1.97.1 only** (`rust-toolchain.toml`). Do not use other toolchains. +Requires **Rust 1.98.1 only** (`rust-toolchain.toml`). Do not use other toolchains. ```bash # Release build, default stub backend (no libzmq) @@ -102,7 +102,19 @@ Default Cargo features are empty (`default = []` in `Cargo.toml`). That path use Enabling the feature does **not** change `BrainstemDaemon::new()` or `try_new()`. Those always inject `BackendPair::stub()`. Only `src/bin/brainstem_daemon.rs` constructs `ZmqStimulusSource` + `ZmqSpikeSink` when `corpus-ipc` is on. -Library users who want live ZMQ must build that pair themselves under `#[cfg(feature = "corpus-ipc")]` and pass it to `with_backend` / `try_with_backend`. Call `StimulusSource::initialize(...)` on the source first (as the binary does). Neither constructor nor `run` calls `initialize`; skipping it makes ingress fail with `ZmqBrainBackend not initialized`. +Library users who want live ZMQ must build that pair themselves under `#[cfg(feature = "corpus-ipc")]` and pass it to `with_backend` / `try_with_backend`. Call `StimulusSource::initialize(...)` on the source first (as the binary does). Neither constructor nor `run` calls `initialize`; skipping it makes ingress fail with `ZMQ stimulus source not initialized`. + +### Integration smoke (Thalamic → corpus-ipc → Brainstem) + +`tests/thalamic_brainstem_smoke.rs` is a CPU-only cross-contract harness (no GPU). It runs when the `corpus-ipc` feature is enabled: + +```bash +cargo test --locked --features corpus-ipc --test thalamic_brainstem_smoke +``` + +The Thalamic fixture (`tests/fixtures/thalamic_producer.rs`) produces `IpcMessage::Stimuli(StimulusBatch)` from simulated telemetry. It does not import `neuromod` or own a `SpikingNetwork`. Brainstem loads an explicit JSON checkpoint (`write_nonblank_checkpoint`) before ticking, preserves `valid_mask` across the wire, and rejects incompatible schema/JSON loudly. A separate assertion keeps the fixture's safety flag healthy when Brainstem/transport is absent. + +The default stub `cargo test` path does not compile this harness (`required-features = ["corpus-ipc"]`). #### Config keys and env vars @@ -112,20 +124,17 @@ Library users who want live ZMQ must build that pair themselves under `#[cfg(fea | `tick_rate_hz` | used | used | | `log_level` | binary tracing init only; unused by `::new()` / `run` | binary tracing init only; unused by `::new()` / `run` | | `services` | used (`ServiceRegistry`) | used | -| `spine_sub_port` | parsed, **no-op** | sets `SPIKENAUT_ZMQ_READOUT_IPC` to `tcp://127.0.0.1:` (also sets unused `CORPUS_IPC_ZMQ_READOUT_IPC` for compatibility) | +| `spine_sub_port` | parsed, **no-op** | sets `SPIKENAUT_ZMQ_READOUT_IPC` and `CORPUS_IPC_ZMQ_READOUT_IPC` to `tcp://127.0.0.1:` | | `spine_pub_port` | parsed, **no-op** | binds ZMQ PUB `tcp://*:` | -| `model_path` | parsed, **no-op** (`StubStimulusSource::initialize` ignores it) | passed literally to `initialize` (no `~` expansion); pinned `ZmqBrainBackend` currently ignores `_model_path` | +| `model_path` | parsed; if the file exists it is loaded as a JSON checkpoint, otherwise a blank `with_dimensions` network is used | same checkpoint load; ZMQ source `initialize` only connects the SUB socket | **Settings that only take effect with `corpus-ipc`** (the `brainstem-daemon` binary built `--features corpus-ipc`): -- `spine_sub_port` (drives `SPIKENAUT_ZMQ_READOUT_IPC`) +- `spine_sub_port` (drives `SPIKENAUT_ZMQ_READOUT_IPC` and `CORPUS_IPC_ZMQ_READOUT_IPC`) - `spine_pub_port` -- `SPIKENAUT_ZMQ_READOUT_IPC` (const `CORPUS_IPC_READOUT_ENV`; this is what pinned `ZmqBrainBackend::initialize` reads) - -**Passed through / set, but currently unused by the pinned dep:** +- `SPIKENAUT_ZMQ_READOUT_IPC` / `CORPUS_IPC_ZMQ_READOUT_IPC` (JSON `IpcMessage` SUB endpoint) -- `model_path` (literal filesystem path; `~` is not expanded; passed to `initialize`, which names the argument `_model_path` and does not consume it) -- `CORPUS_IPC_ZMQ_READOUT_IPC` (the binary still sets this alongside `SPIKENAUT_ZMQ_READOUT_IPC` for compatibility; pinned `corpus-ipc` does not read it) +The ZMQ source subscribes to **JSON** `corpus_ipc::IpcMessage` frames (`Stimuli` / `Neuromodulators`). It does not use the legacy binary readout packet. Under stub those TOML keys are still parsed. The env vars are unset by the default binary. Nothing in this crate reads them without the `corpus-ipc` feature. diff --git a/docs/ci.md b/docs/ci.md index 4f989aa..b6cb1af 100644 --- a/docs/ci.md +++ b/docs/ci.md @@ -13,7 +13,9 @@ system `libzmq` is available. | `corpus-ipc` | `ubuntu-latest` | `--all-features` (the `corpus-ipc` feature, which enables the optional `zmq` dependency) | clippy, build, test | Default features are empty. Stub jobs do **not** install `libzmq` and do -not pass `--features corpus-ipc`. +not pass `--features corpus-ipc`. The Thalamic → corpus-ipc → Brainstem +integration smoke (`tests/thalamic_brainstem_smoke.rs`) is compiled only in +the Linux `corpus-ipc` job (`required-features = ["corpus-ipc"]`). ## Skips diff --git a/rust-toolchain.toml b/rust-toolchain.toml index 5ec5869..c990140 100644 --- a/rust-toolchain.toml +++ b/rust-toolchain.toml @@ -1,4 +1,4 @@ [toolchain] -channel = "1.97.1" +channel = "1.98.1" components = ["rustfmt", "clippy"] profile = "minimal" diff --git a/src/backend.rs b/src/backend.rs index 02e9ff3..6f9a1a3 100644 --- a/src/backend.rs +++ b/src/backend.rs @@ -23,6 +23,13 @@ pub struct IngressPacket { /// Optional raw modulator values (e.g. [dopamine, cortisol, acetylcholine, tempo, ...]). /// When `None`, the caller should use defaults (see `decode_inputs`). pub modulators: Option>, + /// Per-channel validity mask copied from a typed `StimulusBatch` when present. + /// `false` means the corresponding stimulus is a placeholder, not a real zero. + pub valid_mask: Option>, + /// Optional typed batch id from `corpus-ipc` stimulus ingress. + pub batch_id: Option, + /// Optional stimulus timestamp in nanoseconds from `corpus-ipc`. + pub timestamp_ns: Option, } /// Local spike event type (independent of any external crate). @@ -103,6 +110,9 @@ impl StimulusSource for StubStimulusSource { Ok(Some(IngressPacket { stimuli: Vec::new(), modulators: None, + valid_mask: None, + batch_id: None, + timestamp_ns: None, })) } @@ -120,14 +130,14 @@ impl SpikeSink for NoopSpikeSink { } } -/// Collecting sink for tests. Collects every emitted batch. -#[cfg(test)] +/// Collecting sink for tests and the integration smoke harness. +#[derive(Default)] pub struct CollectingSpikeSink { pub emitted: Vec>, } -#[cfg(test)] impl CollectingSpikeSink { + /// Create an empty collector. pub fn new() -> Self { Self { emitted: Vec::new(), @@ -135,14 +145,6 @@ impl CollectingSpikeSink { } } -#[cfg(test)] -impl Default for CollectingSpikeSink { - fn default() -> Self { - Self::new() - } -} - -#[cfg(test)] impl SpikeSink for CollectingSpikeSink { fn emit(&mut self, spikes: &[SpikeEvent], _batch_time: std::time::Duration) -> Result<()> { self.emitted.push(spikes.to_vec()); @@ -156,15 +158,19 @@ impl SpikeSink for CollectingSpikeSink { #[cfg(feature = "corpus-ipc")] mod zmq_impl { use super::*; - // In the current pinned corpus-ipc revision, the main trait is exported as - // `NeuralBackend` (deprecated alias). Importing it brings the trait methods - // into scope for ZmqBrainBackend. - use corpus_ipc::NeuralBackend as BackendConnector; - use corpus_ipc::{SpikeBatch, SpikeEvent as CorpusSpikeEvent, SpineMessage, ZmqBrainBackend}; + use crate::ingress::{IngressPolicy, accept_ipc_json}; + use corpus_ipc::{IpcMessage, SpikeBatch, SpikeEvent as CorpusSpikeEvent}; + + /// Env var read by this source for the SUB endpoint. The binary also sets + /// the historical `SPIKENAUT_ZMQ_READOUT_IPC` alias. + pub const CORPUS_IPC_ZMQ_READOUT_ENV: &str = "CORPUS_IPC_ZMQ_READOUT_IPC"; + const LEGACY_READOUT_ENV: &str = "SPIKENAUT_ZMQ_READOUT_IPC"; pub struct ZmqStimulusSource { - inner: ZmqBrainBackend, + socket: Option, channels: usize, + last_modulators: Option>, + max_age: std::time::Duration, } impl Default for ZmqStimulusSource { @@ -175,52 +181,86 @@ mod zmq_impl { impl ZmqStimulusSource { pub fn new() -> Self { - Self { - inner: ZmqBrainBackend::new(), - channels: 0, - } + Self::with_channels(0) } - /// Construct with known channel count so `next_ingress` can split - /// stimulus prefix from appended neuromodulator tail (4 floats). - /// - /// The default `new()` uses `channels=0`, which means the entire readout - /// is passed as stimuli and no modulators are extracted. Library users - /// who want automatic modulator extraction must use `with_channels(cfg.channels)`. + /// Construct with the configured ingress width used for width checks. pub fn with_channels(ch: usize) -> Self { Self { - inner: ZmqBrainBackend::new(), + socket: None, channels: ch, + last_modulators: None, + max_age: std::time::Duration::from_secs(1), } } + + /// Override the freshness window applied to typed `StimulusBatch` timestamps. + pub fn with_max_age(mut self, max_age: std::time::Duration) -> Self { + self.max_age = max_age; + self + } + + /// Connect the SUB socket to an explicit endpoint (tests / custom wiring). + pub fn connect(&mut self, endpoint: &str) -> Result<()> { + let context = ::zmq::Context::new(); + let socket = context + .socket(::zmq::SUB) + .map_err(|e| anyhow::anyhow!("failed to create ZMQ SUB socket: {e}"))?; + socket + .set_subscribe(b"") + .map_err(|e| anyhow::anyhow!("ZMQ subscribe: {e}"))?; + socket + .set_rcvhwm(16) + .map_err(|e| anyhow::anyhow!("ZMQ rcvhwm: {e}"))?; + socket + .connect(endpoint) + .map_err(|e| anyhow::anyhow!("ZMQ connect to {endpoint}: {e}"))?; + self.socket = Some(SafeSocket { socket }); + Ok(()) + } + + fn readout_endpoint() -> String { + std::env::var(CORPUS_IPC_ZMQ_READOUT_ENV) + .or_else(|_| std::env::var(LEGACY_READOUT_ENV)) + .unwrap_or_else(|_| "tcp://127.0.0.1:5555".to_string()) + } + + fn now_ns() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos() as u64 + } } impl StimulusSource for ZmqStimulusSource { fn next_ingress(&mut self) -> Result> { - let readout = self.inner.process_signals(&[])?; - let ch = self.channels; - if ch > 0 && readout.len() > ch { - let stimuli = readout[..ch].to_vec(); - let modulators = if readout.len() >= ch + 4 { - Some(readout[ch..ch + 4].to_vec()) - } else { - None - }; - Ok(Some(IngressPacket { - stimuli, - modulators, - })) - } else { - Ok(Some(IngressPacket { - stimuli: readout, - modulators: None, - })) + let socket = self + .socket + .as_ref() + .ok_or_else(|| anyhow::anyhow!("ZMQ stimulus source not initialized"))?; + match socket.socket.recv_bytes(::zmq::DONTWAIT) { + Ok(buf) => { + let policy = IngressPolicy::new(self.channels, Some(self.max_age)); + let mut packet = accept_ipc_json(&buf, &policy, Self::now_ns()) + .map_err(|e| anyhow::anyhow!("{e}"))?; + if let Some(mods) = packet.modulators.as_ref() { + self.last_modulators = Some(mods.clone()); + } else { + packet.modulators = self.last_modulators.clone(); + } + Ok(Some(packet)) + } + Err(::zmq::Error::EAGAIN) => Ok(None), + Err(e) => Err(anyhow::anyhow!("ZMQ recv failed: {e}")), } } - fn initialize(&mut self, model_path: Option<&str>) -> Result<()> { - self.inner.initialize(model_path)?; - Ok(()) + fn initialize(&mut self, _model_path: Option<&str>) -> Result<()> { + if self.socket.is_some() { + return Ok(()); + } + self.connect(&Self::readout_endpoint()) } } @@ -268,7 +308,7 @@ mod zmq_impl { let cap = self.corpus_buf.capacity(); let corpus_spikes = std::mem::replace(&mut self.corpus_buf, Vec::with_capacity(cap)); - let msg = SpineMessage::Spikes(SpikeBatch { + let msg = IpcMessage::Spikes(SpikeBatch { session_id: None, batch_id, timestamp, @@ -285,6 +325,77 @@ mod zmq_impl { Ok(()) } } + + #[cfg(test)] + mod tests { + use super::*; + use crate::ingress::STIMULUS_SCHEMA; + use corpus_ipc::{BatchMetadata, StimulusBatch}; + use std::collections::HashMap; + use std::time::Duration; + + fn sample_frame(channels: usize, batch_id: u64) -> Vec { + let mut values = vec![0.0; channels]; + let mut valid_mask = vec![true; channels]; + values[0] = 1.0; + if channels > 1 { + valid_mask[1] = false; + } + let mut custom = HashMap::new(); + custom.insert("schema".into(), STIMULUS_SCHEMA.to_string()); + let batch = StimulusBatch { + session_id: Some("zmq-loopback".into()), + batch_id, + timestamp: std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos() as u64, + values, + valid_mask: Some(valid_mask), + metadata: Some(BatchMetadata { + processing_latency_ns: None, + source: Some("thalamic-relay-fixture".into()), + custom, + }), + }; + serde_json::to_vec(&IpcMessage::Stimuli(batch)).expect("serialize") + } + + #[test] + fn typed_json_frame_round_trips_over_zmq() { + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let port = listener.local_addr().unwrap().port(); + drop(listener); + let endpoint = format!("tcp://127.0.0.1:{port}"); + + let context = ::zmq::Context::new(); + let publisher = context.socket(::zmq::PUB).unwrap(); + publisher.bind(&endpoint).unwrap(); + + let mut source = + ZmqStimulusSource::with_channels(4).with_max_age(Duration::from_secs(5)); + source.connect(&endpoint).unwrap(); + std::thread::sleep(Duration::from_millis(150)); + + let frame = sample_frame(4, 100); + publisher.send(&frame, 0).unwrap(); + + let packet = (0..50) + .find_map(|_| match source.next_ingress() { + Ok(Some(packet)) => Some(packet), + Ok(None) => { + std::thread::sleep(Duration::from_millis(10)); + None + } + Err(err) => panic!("ingress error: {err}"), + }) + .expect("ZMQ SUB should receive a typed StimulusBatch"); + + assert_eq!(packet.batch_id, Some(100)); + assert_eq!(packet.valid_mask.as_ref().map(|m| m[1]), Some(false)); + assert!((packet.stimuli[0] - 1.0).abs() < f32::EPSILON); + } + } } #[cfg(feature = "corpus-ipc")] diff --git a/src/bin/brainstem_daemon.rs b/src/bin/brainstem_daemon.rs index 91a83a4..40195be 100644 --- a/src/bin/brainstem_daemon.rs +++ b/src/bin/brainstem_daemon.rs @@ -79,8 +79,9 @@ async fn run(cfg: DaemonConfig, config_path: PathBuf) -> anyhow::Result<()> { // is intentionally conservative. let mut source = brainstem_daemon::backend::ZmqStimulusSource::with_channels(cfg.channels); - // Pass the model path through (was dropped before). The pinned - // ZmqBrainBackend::initialize takes `_model_path` and currently ignores it. + // Pass the model path through for initialize's signature. Checkpoint + // loading is performed by `BrainstemDaemon::run` from `cfg.model_path`. + // The ZMQ source only connects the SUB socket here. let model_path = cfg.model_path.to_string_lossy(); source diff --git a/src/checkpoint.rs b/src/checkpoint.rs new file mode 100644 index 0000000..3b5f235 --- /dev/null +++ b/src/checkpoint.rs @@ -0,0 +1,363 @@ +// SPDX-License-Identifier: MIT OR Apache-2.0 +// Copyright 2026 Raul Montoya Cardenas + +//! Explicit Spikenaut/neuromod checkpoint load and validation. +//! +//! Live ticks should restore a recorded network rather than silently +//! constructing a blank `SpikingNetwork::with_dimensions` when a checkpoint +//! file is present. Missing files still yield `None` so the historical stub +//! path (dummy `model_path`) keeps working; a file that exists but is invalid +//! fails closed. + +use std::fs; +use std::path::{Path, PathBuf}; + +use anyhow::{Context, Result, bail}; +use neuromod::SpikingNetwork; +use serde::{Deserialize, Serialize}; + +/// Checkpoint schema version understood by this loader. +pub const CHECKPOINT_SCHEMA_VERSION: u32 = 1; + +/// Dimensions the checkpoint must match. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct NetworkDims { + pub lif_count: usize, + pub izh_count: usize, + pub channels: usize, +} + +/// Provenance recorded after a successful load. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct CheckpointIdentity { + pub model_id: String, + pub schema_version: u32, + pub path: PathBuf, + pub fingerprint: String, + pub dims: NetworkDims, +} + +/// On-disk envelope around a serializable `neuromod::SpikingNetwork`. +#[derive(Serialize, Deserialize)] +pub struct CheckpointFile { + pub schema_version: u32, + pub model_id: String, + pub network: SpikingNetwork, +} + +impl CheckpointFile { + /// Validate schema, dimensions, finite parameters, and non-blank weights. + pub fn validate(&self, expected: NetworkDims) -> Result<()> { + if self.schema_version != CHECKPOINT_SCHEMA_VERSION { + bail!( + "checkpoint schema_version {} is incompatible with runtime {}", + self.schema_version, + CHECKPOINT_SCHEMA_VERSION + ); + } + if self.model_id.trim().is_empty() { + bail!("checkpoint model_id must be non-empty"); + } + + let got = NetworkDims { + lif_count: self.network.neurons.len(), + izh_count: self.network.iz_neurons.len(), + channels: self.network.num_channels, + }; + if got != expected { + bail!( + "checkpoint dimensions lif={}, izh={}, channels={} do not match config lif={}, izh={}, channels={}", + got.lif_count, + got.izh_count, + got.channels, + expected.lif_count, + expected.izh_count, + expected.channels + ); + } + + if self.network.input_spike_times.len() != expected.channels + || self.network.predictive_state.len() != expected.channels + { + bail!("checkpoint input-state width does not match channel count"); + } + + let mut any_weight = false; + for (idx, neuron) in self.network.neurons.iter().enumerate() { + if neuron.weights.len() != expected.channels { + bail!( + "checkpoint LIF[{idx}] weight width {} != channels {}", + neuron.weights.len(), + expected.channels + ); + } + if !neuron.membrane_potential.is_finite() + || !neuron.decay_rate.is_finite() + || !neuron.threshold.is_finite() + || !neuron.base_threshold.is_finite() + { + bail!("checkpoint LIF[{idx}] has a non-finite parameter"); + } + for (ch, weight) in neuron.weights.iter().enumerate() { + if !weight.is_finite() { + bail!("checkpoint LIF[{idx}] weight[{ch}] is not finite"); + } + if weight.abs() > 1e-8 { + any_weight = true; + } + } + } + + if expected.lif_count > 0 && !any_weight { + bail!("checkpoint looks blank: all LIF input weights are zero"); + } + + Ok(()) + } +} + +/// Load `path` when it exists. `Ok(None)` if the path is absent. +pub fn try_load_checkpoint( + path: &Path, + expected: NetworkDims, +) -> Result> { + if !path.exists() { + return Ok(None); + } + Ok(Some(load_checkpoint(path, expected)?)) +} + +/// Load and validate a checkpoint file. Fails if the file is missing or invalid. +pub fn load_checkpoint( + path: &Path, + expected: NetworkDims, +) -> Result<(SpikingNetwork, CheckpointIdentity)> { + let bytes = + fs::read(path).with_context(|| format!("failed to read checkpoint {}", path.display()))?; + let parsed: CheckpointFile = serde_json::from_slice(&bytes) + .with_context(|| format!("failed to parse checkpoint JSON {}", path.display()))?; + parsed.validate(expected)?; + + let identity = CheckpointIdentity { + model_id: parsed.model_id.clone(), + schema_version: parsed.schema_version, + path: path.to_path_buf(), + fingerprint: fnv1a64_hex(&bytes), + dims: expected, + }; + Ok((parsed.network, identity)) +} + +/// Write a non-blank, dimension-checked smoke checkpoint to `path`. +pub fn write_nonblank_checkpoint( + path: &Path, + model_id: &str, + dims: NetworkDims, +) -> Result { + if dims.lif_count == 0 { + bail!("smoke checkpoint requires at least one LIF neuron"); + } + if dims.channels == 0 { + bail!("smoke checkpoint requires at least one input channel"); + } + + let mut network = + SpikingNetwork::with_dimensions(dims.lif_count, dims.izh_count, dims.channels); + for (idx, neuron) in network.neurons.iter_mut().enumerate() { + for (ch, weight) in neuron.weights.iter_mut().enumerate() { + *weight = 0.12 + (idx as f32) * 0.03 + (ch as f32) * 0.01; + } + if let Some(first) = neuron.weights.first_mut() { + // Strong channel-0 coupling so a unit stimulus deterministically exceeds threshold. + if idx == 0 { + *first = 2.0; + } + } + neuron.threshold = 0.02; + neuron.base_threshold = 0.02; + neuron.membrane_potential = 0.0; + } + + let file = CheckpointFile { + schema_version: CHECKPOINT_SCHEMA_VERSION, + model_id: model_id.to_string(), + network, + }; + file.validate(dims)?; + + if let Some(parent) = path.parent() + && !parent.as_os_str().is_empty() + { + fs::create_dir_all(parent) + .with_context(|| format!("failed to create {}", parent.display()))?; + } + let bytes = serde_json::to_vec_pretty(&file).context("failed to serialize checkpoint")?; + fs::write(path, &bytes) + .with_context(|| format!("failed to write checkpoint {}", path.display()))?; + + Ok(CheckpointIdentity { + model_id: model_id.to_string(), + schema_version: CHECKPOINT_SCHEMA_VERSION, + path: path.to_path_buf(), + fingerprint: fnv1a64_hex(&bytes), + dims, + }) +} + +fn fnv1a64_hex(bytes: &[u8]) -> String { + const OFFSET: u64 = 0xcbf29ce484222325; + const PRIME: u64 = 0x100000001b3; + let mut hash = OFFSET; + for byte in bytes { + hash ^= u64::from(*byte); + hash = hash.wrapping_mul(PRIME); + } + format!("{hash:016x}") +} + +#[cfg(test)] +mod tests { + use super::*; + use std::path::Path; + use std::time::{SystemTime, UNIX_EPOCH}; + + fn temp_path(name: &str) -> PathBuf { + let nanos = SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("clock") + .as_nanos(); + std::env::temp_dir().join(format!( + "brainstem-checkpoint-{name}-{}-{nanos}.json", + std::process::id() + )) + } + + fn smoke_dims() -> NetworkDims { + NetworkDims { + lif_count: 4, + izh_count: 0, + channels: 4, + } + } + + #[test] + fn missing_path_is_none() { + let path = PathBuf::from("/tmp/brainstem-daemon-missing-checkpoint-does-not-exist.json"); + let loaded = try_load_checkpoint(&path, smoke_dims()).unwrap(); + assert!(loaded.is_none()); + } + + #[test] + fn round_trip_nonblank_checkpoint() { + let path = temp_path("round-trip"); + let identity = write_nonblank_checkpoint(&path, "smoke-v1", smoke_dims()).unwrap(); + let (network, loaded) = load_checkpoint(&path, smoke_dims()).unwrap(); + assert_eq!(loaded.model_id, "smoke-v1"); + assert_eq!(loaded.fingerprint, identity.fingerprint); + assert_eq!(network.neurons.len(), 4); + assert!(network.neurons[0].weights[0].abs() > 1.0); + let _ = fs::remove_file(&path); + } + + #[test] + fn blank_weights_are_rejected() { + let path = temp_path("blank"); + let network = SpikingNetwork::with_dimensions(2, 0, 2); + let file = CheckpointFile { + schema_version: CHECKPOINT_SCHEMA_VERSION, + model_id: "blank".into(), + network, + }; + fs::write(&path, serde_json::to_vec(&file).unwrap()).unwrap(); + let err = match load_checkpoint( + &path, + NetworkDims { + lif_count: 2, + izh_count: 0, + channels: 2, + }, + ) { + Ok(_) => panic!("expected blank checkpoint to fail"), + Err(err) => err, + }; + let _ = fs::remove_file(&path); + let message = err.to_string(); + assert!(message.contains("blank"), "unexpected error: {message}"); + } + + #[test] + fn schema_mismatch_fails_loudly() { + let path = temp_path("schema"); + let mut file = CheckpointFile { + schema_version: CHECKPOINT_SCHEMA_VERSION, + model_id: "x".into(), + network: SpikingNetwork::with_dimensions(4, 0, 4), + }; + file.network.neurons[0].weights[0] = 1.0; + file.schema_version = 99; + fs::write(&path, serde_json::to_vec(&file).unwrap()).unwrap(); + let err = expect_load_err(&path, smoke_dims()); + let _ = fs::remove_file(&path); + let message = err.to_string(); + assert!( + message.contains("schema_version"), + "unexpected error: {message}" + ); + } + + #[test] + fn dimension_mismatch_fails() { + let path = temp_path("dims"); + write_nonblank_checkpoint(&path, "smoke-v1", smoke_dims()).unwrap(); + let err = expect_load_err( + &path, + NetworkDims { + lif_count: 8, + izh_count: 0, + channels: 4, + }, + ); + let _ = fs::remove_file(&path); + let message = err.to_string(); + assert!( + message.contains("dimensions"), + "unexpected error: {message}" + ); + } + + #[test] + fn non_finite_weight_is_rejected_by_validate() { + let mut network = SpikingNetwork::with_dimensions(4, 0, 4); + network.neurons[0].weights[0] = f32::NAN; + let file = CheckpointFile { + schema_version: CHECKPOINT_SCHEMA_VERSION, + model_id: "nan".into(), + network, + }; + let err = file.validate(smoke_dims()).unwrap_err(); + assert!( + err.to_string().contains("not finite"), + "unexpected error: {err}" + ); + } + + #[test] + fn corrupt_json_fails_loudly() { + let path = temp_path("corrupt"); + fs::write(&path, "{not-json").unwrap(); + let err = expect_load_err(&path, smoke_dims()); + let _ = fs::remove_file(&path); + let message = err.to_string(); + assert!( + message.contains("parse") || message.contains("JSON") || message.contains("expected"), + "unexpected error: {message}" + ); + } + + fn expect_load_err(path: &Path, dims: NetworkDims) -> anyhow::Error { + match load_checkpoint(path, dims) { + Ok(_) => panic!("expected checkpoint load to fail"), + Err(err) => err, + } + } +} diff --git a/src/daemon.rs b/src/daemon.rs index bd8ccb5..d273fad 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -17,6 +17,7 @@ use tracing::{error, info, warn}; use crate::backend::{ BackendPair, IngressPacket, SpikeEvent as LocalSpikeEvent, SpikeSink, StimulusSource, }; +use crate::checkpoint::{CheckpointIdentity, NetworkDims, try_load_checkpoint}; use crate::registry::{ServiceConfig, ServiceRegistry}; // Keep the const for compatibility when the corpus-ipc feature is used. @@ -61,6 +62,17 @@ impl DaemonConfig { } } +/// Observable counters from a bounded tick run (smoke / integration harness). +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct RuntimeStats { + pub ticks: u64, + pub accepted_batches: u64, + pub rejected_batches: u64, + pub last_batch_id: Option, + pub last_valid_mask: Option>, + pub loaded_checkpoint: Option, +} + /// Headless spiking-network daemon. /// /// Owns the tick loop and delegates I/O to pluggable `StimulusSource` + `SpikeSink`. @@ -130,6 +142,8 @@ impl BrainstemDaemon { pub async fn run(self) -> Result<()> { let cfg = self.config; let mut backend = self.backend; + let mut stats = RuntimeStats::default(); + let mut network = instantiate_network(&cfg, &mut stats)?; if cfg.tick_rate_hz == 0 || cfg.tick_rate_hz > 1_000_000 { anyhow::bail!("tick_rate_hz must be in range 1..=1_000_000"); @@ -139,8 +153,6 @@ impl BrainstemDaemon { let mut ticker = time::interval(tick_duration); ticker.set_missed_tick_behavior(time::MissedTickBehavior::Skip); - let mut network = - SpikingNetwork::with_dimensions(cfg.lif_count, cfg.izh_count, cfg.channels); let mut stimuli = vec![0.0; cfg.channels]; let mut spike_buf: Vec = Vec::with_capacity(128); @@ -155,6 +167,7 @@ impl BrainstemDaemon { &mut *backend.sink, &mut stimuli, &mut spike_buf, + &mut stats, ); } _ = &mut shutdown => { @@ -176,6 +189,36 @@ impl BrainstemDaemon { Ok(()) } + + /// Drive a bounded number of ticks without waiting for a termination signal. + /// + /// Used by the CPU-only Thalamic → corpus-ipc → Brainstem smoke harness. + pub fn run_for_ticks(self, ticks: u64) -> Result { + let cfg = self.config; + let mut backend = self.backend; + let mut stats = RuntimeStats::default(); + let mut network = instantiate_network(&cfg, &mut stats)?; + let mut stimuli = vec![0.0; cfg.channels]; + let mut spike_buf: Vec = Vec::with_capacity(128); + + for _ in 0..ticks { + run_tick( + &mut *backend.source, + &mut network, + &mut *backend.sink, + &mut stimuli, + &mut spike_buf, + &mut stats, + ); + } + + backend.sink.flush().context("failed to flush spike sink")?; + backend + .source + .shutdown() + .context("failed to shut down stimulus source")?; + Ok(stats) + } } /// Wait for a termination request: `SIGINT` (Ctrl-C) on every platform, plus @@ -247,6 +290,38 @@ fn validate_neuron_count(config: &DaemonConfig) -> Result<()> { Ok(()) } +fn instantiate_network(cfg: &DaemonConfig, stats: &mut RuntimeStats) -> Result { + let expected = NetworkDims { + lif_count: cfg.lif_count, + izh_count: cfg.izh_count, + channels: cfg.channels, + }; + match try_load_checkpoint(&cfg.model_path, expected)? { + Some((network, identity)) => { + info!( + model_id = %identity.model_id, + schema_version = identity.schema_version, + fingerprint = %identity.fingerprint, + path = %identity.path.display(), + "loaded explicit checkpoint" + ); + stats.loaded_checkpoint = Some(identity); + Ok(network) + } + None => { + info!( + path = %cfg.model_path.display(), + "no checkpoint file present; constructing a blank network from dimensions" + ); + Ok(SpikingNetwork::with_dimensions( + cfg.lif_count, + cfg.izh_count, + cfg.channels, + )) + } + } +} + // Trait-based tick loop (works with or without corpus-ipc feature) fn run_tick( @@ -255,6 +330,7 @@ fn run_tick( sink: &mut dyn SpikeSink, stimuli: &mut [f32], spike_buf: &mut Vec, + stats: &mut RuntimeStats, ) { let packet = match source.next_ingress() { Ok(Some(p)) => p, @@ -265,14 +341,24 @@ fn run_tick( IngressPacket { stimuli: Vec::new(), modulators: None, + valid_mask: None, + batch_id: None, + timestamp_ns: None, } } Err(e) => { - warn!("Failed to receive from stimulus source: {e}"); + error!("Failed to receive from stimulus source: {e}"); + stats.rejected_batches += 1; return; } }; + if packet.batch_id.is_some() { + stats.accepted_batches += 1; + stats.last_batch_id = packet.batch_id; + stats.last_valid_mask = packet.valid_mask.clone(); + } + let modulators = decode_inputs(&packet, stimuli); // Note: decode_inputs already zero-fills any remaining channels when packet.stimuli is shorter. @@ -284,6 +370,7 @@ fn run_tick( return; } }; + stats.ticks += 1; // Single timestamp for both per-spike time and batch metadata (keeps them consistent). let now = SystemTime::now() @@ -341,6 +428,13 @@ fn decode_inputs(packet: &IngressPacket, stimuli: &mut [f32]) -> NeuroModulators if readout.len() < channels { stimuli[upto..].fill(0.0); } + if let Some(mask) = packet.valid_mask.as_ref() { + for (idx, valid) in mask.iter().enumerate().take(channels) { + if !*valid { + stimuli[idx] = 0.0; + } + } + } match packet.modulators.as_ref() { Some(mods) if mods.len() >= 4 => { @@ -369,7 +463,14 @@ pub(crate) fn run_tick_for_test( stimuli: &mut [f32], spike_buf: &mut Vec, ) { - run_tick(source, network, sink, stimuli, spike_buf); + run_tick( + source, + network, + sink, + stimuli, + spike_buf, + &mut RuntimeStats::default(), + ); } #[cfg(test)] @@ -418,6 +519,9 @@ mod tests { let packet = IngressPacket { stimuli: vec![0.1, 0.2, 0.3, 0.4], modulators: None, + valid_mask: None, + batch_id: None, + timestamp_ns: None, }; let mut stimuli = vec![0.0; 4]; let _mods = decode_inputs(&packet, &mut stimuli); @@ -429,6 +533,9 @@ mod tests { let packet = IngressPacket { stimuli: vec![0.0; 4], modulators: Some(vec![0.5, 0.6, 0.7, 0.8]), + valid_mask: None, + batch_id: None, + timestamp_ns: None, }; let mut stimuli = vec![0.0; 4]; let mods = decode_inputs(&packet, &mut stimuli); @@ -443,6 +550,9 @@ mod tests { let packet = IngressPacket { stimuli: vec![0.1, 0.2], modulators: None, + valid_mask: None, + batch_id: None, + timestamp_ns: None, }; let mut stimuli = vec![0.0; 4]; let mods = decode_inputs(&packet, &mut stimuli); @@ -450,6 +560,20 @@ mod tests { assert_eq!(mods, NeuroModulators::default()); } + #[test] + fn decode_inputs_zeros_invalid_channels() { + let packet = IngressPacket { + stimuli: vec![1.0, 0.5, 0.25, 0.1], + modulators: None, + valid_mask: Some(vec![true, false, true, true]), + batch_id: Some(1), + timestamp_ns: Some(1), + }; + let mut stimuli = vec![0.0; 4]; + let _mods = decode_inputs(&packet, &mut stimuli); + assert_eq!(stimuli, vec![1.0, 0.0, 0.25, 0.1]); + } + #[test] fn daemon_allows_u16_max_total_neurons() { let mut cfg = sample_config(); diff --git a/src/ingress.rs b/src/ingress.rs new file mode 100644 index 0000000..7c51ff4 --- /dev/null +++ b/src/ingress.rs @@ -0,0 +1,272 @@ +// SPDX-License-Identifier: MIT OR Apache-2.0 +// Copyright 2026 Raul Montoya Cardenas + +//! Typed `corpus-ipc` ingress validation. +//! +//! Brainstem consumes published `IpcMessage` payloads. It does not copy the +//! wire structs: decoding goes through `corpus_ipc::{IpcMessage, StimulusBatch}`. + +use std::time::Duration; + +use crate::backend::IngressPacket; +use corpus_ipc::{IpcMessage, StimulusBatch}; + +/// Schema token Brainstem requires in `StimulusBatch.metadata.custom["schema"]`. +pub const STIMULUS_SCHEMA: &str = "corpus-ipc.stimulus.v1"; + +/// Ingress acceptance policy. +#[derive(Debug, Clone)] +pub struct IngressPolicy { + pub expected_channels: usize, + pub max_age: Option, +} + +impl IngressPolicy { + /// Policy for a configured channel width and optional freshness window. + pub fn new(expected_channels: usize, max_age: Option) -> Self { + Self { + expected_channels, + max_age, + } + } +} + +/// Why a typed ingress payload was rejected. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum IngressError { + Deserialize(String), + UnexpectedVariant(&'static str), + Width { expected: usize, got: usize }, + Stale { age_ns: u64, max_age_ns: u64 }, + Schema(String), + Stimulus(String), +} + +impl std::fmt::Display for IngressError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Deserialize(msg) => write!(f, "IPC schema deserialize error: {msg}"), + Self::UnexpectedVariant(name) => { + write!( + f, + "unexpected IpcMessage variant on stimulus ingress: {name}" + ) + } + Self::Width { expected, got } => { + write!( + f, + "stimulus width {got} does not match configured channels {expected}" + ) + } + Self::Stale { age_ns, max_age_ns } => { + write!( + f, + "stimulus is stale: age {age_ns} ns exceeds {max_age_ns} ns" + ) + } + Self::Schema(msg) => write!(f, "stimulus schema incompatibility: {msg}"), + Self::Stimulus(msg) => write!(f, "stimulus batch invalid: {msg}"), + } + } +} + +impl std::error::Error for IngressError {} + +/// Decode and validate a JSON `IpcMessage` from the corpus-ipc wire contract. +pub fn accept_ipc_json( + bytes: &[u8], + policy: &IngressPolicy, + now_ns: u64, +) -> Result { + let message: IpcMessage = + serde_json::from_slice(bytes).map_err(|err| IngressError::Deserialize(err.to_string()))?; + accept_ipc_message(message, policy, now_ns) +} + +/// Validate an already-decoded `IpcMessage`. +pub fn accept_ipc_message( + message: IpcMessage, + policy: &IngressPolicy, + now_ns: u64, +) -> Result { + match message { + IpcMessage::Stimuli(batch) => accept_stimulus_batch(batch, policy, now_ns), + IpcMessage::Neuromodulators(snapshot) => Ok(IngressPacket { + stimuli: Vec::new(), + modulators: Some(vec![ + snapshot.dopamine, + snapshot.cortisol, + snapshot.acetylcholine, + snapshot.tempo, + ]), + valid_mask: None, + batch_id: None, + timestamp_ns: None, + }), + IpcMessage::Spikes(_) => Err(IngressError::UnexpectedVariant("Spikes")), + IpcMessage::Embeddings(_) => Err(IngressError::UnexpectedVariant("Embeddings")), + IpcMessage::Loss(_) => Err(IngressError::UnexpectedVariant("Loss")), + IpcMessage::ConfigUpdate(_) => Err(IngressError::UnexpectedVariant("ConfigUpdate")), + IpcMessage::GradientUpdate(_) => Err(IngressError::UnexpectedVariant("GradientUpdate")), + IpcMessage::EligibilityTraces(_) => { + Err(IngressError::UnexpectedVariant("EligibilityTraces")) + } + IpcMessage::TrainingComplete => Err(IngressError::UnexpectedVariant("TrainingComplete")), + IpcMessage::Shutdown => Err(IngressError::UnexpectedVariant("Shutdown")), + IpcMessage::Ping => Err(IngressError::UnexpectedVariant("Ping")), + } +} + +fn accept_stimulus_batch( + batch: StimulusBatch, + policy: &IngressPolicy, + now_ns: u64, +) -> Result { + batch.validate().map_err(IngressError::Stimulus)?; + + match batch.metadata.as_ref() { + Some(meta) => match meta.custom.get("schema") { + Some(schema) if schema == STIMULUS_SCHEMA => {} + Some(schema) => { + return Err(IngressError::Schema(format!( + "unsupported schema '{schema}', expected {STIMULUS_SCHEMA}" + ))); + } + None => { + return Err(IngressError::Schema( + "missing metadata.custom.schema".into(), + )); + } + }, + None => { + return Err(IngressError::Schema("missing batch metadata".into())); + } + } + + if batch.values.len() != policy.expected_channels { + return Err(IngressError::Width { + expected: policy.expected_channels, + got: batch.values.len(), + }); + } + + if let Some(max_age) = policy.max_age { + let max_age_ns = max_age.as_nanos() as u64; + let age_ns = now_ns.saturating_sub(batch.timestamp); + if age_ns > max_age_ns { + return Err(IngressError::Stale { age_ns, max_age_ns }); + } + } + + let mut stimuli = batch.values; + if let Some(mask) = batch.valid_mask.as_ref() { + for (value, valid) in stimuli.iter_mut().zip(mask.iter()) { + if !*valid { + *value = 0.0; + } + } + } + + Ok(IngressPacket { + stimuli, + modulators: None, + valid_mask: batch.valid_mask, + batch_id: Some(batch.batch_id), + timestamp_ns: Some(batch.timestamp), + }) +} + +#[cfg(test)] +mod tests { + use super::*; + use corpus_ipc::BatchMetadata; + use std::collections::HashMap; + + fn sample_batch() -> StimulusBatch { + let mut custom = HashMap::new(); + custom.insert("schema".into(), STIMULUS_SCHEMA.into()); + StimulusBatch { + session_id: Some("smoke".into()), + batch_id: 7, + timestamp: 1_000, + values: vec![1.0, 0.0, 0.25, 0.5], + valid_mask: Some(vec![true, false, true, true]), + metadata: Some(BatchMetadata { + processing_latency_ns: None, + source: Some("thalamic-relay-fixture".into()), + custom, + }), + } + } + + fn policy() -> IngressPolicy { + IngressPolicy::new(4, Some(Duration::from_secs(1))) + } + + #[test] + fn valid_mask_survives_json_round_trip() { + let message = IpcMessage::Stimuli(sample_batch()); + let bytes = serde_json::to_vec(&message).unwrap(); + let packet = accept_ipc_json(&bytes, &policy(), 1_000).unwrap(); + assert_eq!(packet.valid_mask, Some(vec![true, false, true, true])); + assert_eq!(packet.stimuli, vec![1.0, 0.0, 0.25, 0.5]); + assert_eq!(packet.batch_id, Some(7)); + } + + #[test] + fn width_mismatch_fails() { + let mut batch = sample_batch(); + batch.values.push(0.1); + batch.valid_mask = Some(vec![true, false, true, true, true]); + let err = accept_ipc_message(IpcMessage::Stimuli(batch), &policy(), 1_000).unwrap_err(); + assert!(matches!( + err, + IngressError::Width { + expected: 4, + got: 5 + } + )); + } + + #[test] + fn stale_batch_fails() { + let mut batch = sample_batch(); + batch.timestamp = 0; + let err = accept_ipc_message( + IpcMessage::Stimuli(batch), + &policy(), + Duration::from_secs(2).as_nanos() as u64, + ) + .unwrap_err(); + assert!(matches!(err, IngressError::Stale { .. })); + } + + #[test] + fn unknown_schema_fails_loudly() { + let mut batch = sample_batch(); + batch + .metadata + .as_mut() + .unwrap() + .custom + .insert("schema".into(), "corpus-ipc.stimulus.v0".into()); + let err = accept_ipc_message(IpcMessage::Stimuli(batch), &policy(), 1_000).unwrap_err(); + match err { + IngressError::Schema(msg) => assert!(msg.contains("unsupported schema")), + other => panic!("expected schema error, got {other}"), + } + } + + #[test] + fn copied_udp_json_is_not_ipc_message() { + let bytes = br#"{"type":"Stimuli","values":[0.1,0.2]}"#; + let err = accept_ipc_json(bytes, &policy(), 1_000).unwrap_err(); + assert!(matches!(err, IngressError::Deserialize(_))); + } + + #[test] + fn unexpected_variant_fails() { + let err = accept_ipc_message(IpcMessage::Ping, &policy(), 1_000).unwrap_err(); + assert_eq!(err, IngressError::UnexpectedVariant("Ping")); + } +} diff --git a/src/lib.rs b/src/lib.rs index d6b4281..d4b5959 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -4,8 +4,18 @@ //! Brainstem daemon library: config-driven service registry and runtime. pub mod backend; +pub mod checkpoint; pub mod daemon; +#[cfg(feature = "corpus-ipc")] +pub mod ingress; pub mod registry; // Re-export the new pluggable I/O surface (pub from day one). -pub use backend::{BackendPair, IngressPacket, SpikeEvent, SpikeSink, StimulusSource}; +pub use backend::{ + BackendPair, CollectingSpikeSink, IngressPacket, SpikeEvent, SpikeSink, StimulusSource, +}; +pub use checkpoint::{ + CheckpointIdentity, NetworkDims, load_checkpoint, try_load_checkpoint, + write_nonblank_checkpoint, +}; +pub use daemon::{BrainstemDaemon, DaemonConfig, RuntimeStats}; diff --git a/tests/fixtures/thalamic_producer.rs b/tests/fixtures/thalamic_producer.rs new file mode 100644 index 0000000..3eec4ec --- /dev/null +++ b/tests/fixtures/thalamic_producer.rs @@ -0,0 +1,104 @@ +// SPDX-License-Identifier: MIT OR Apache-2.0 +// Copyright 2026 Raul Montoya Cardenas + +//! Thalamic-side fixture for the Brainstem integration smoke test. +//! +//! This module must not depend on `neuromod` or own a `SpikingNetwork`. +//! It produces typed `corpus-ipc` sensory frames from simulated telemetry and +//! keeps a local safety/health flag independent of whether Brainstem is +//! reachable. + +use std::collections::HashMap; +use std::time::{SystemTime, UNIX_EPOCH}; + +use anyhow::{Result, anyhow}; +use brainstem_daemon::ingress::STIMULUS_SCHEMA; +use corpus_ipc::{BatchMetadata, IpcMessage, StimulusBatch}; + +/// Simulated Thalamic producer: sensory + safety only. +#[derive(Debug, Default)] +pub struct ThalamicProducer { + /// Deterministic hardware-safety path. Independent of IPC/Brainstem. + pub safety_healthy: bool, + pub published: u64, + pub publish_errors: u64, +} + +impl ThalamicProducer { + pub fn new() -> Self { + Self { + safety_healthy: true, + published: 0, + publish_errors: 0, + } + } + + /// Build a representative normalized telemetry frame. + pub fn simulate_telemetry(&self, batch_id: u64, channels: usize) -> StimulusBatch { + let mut values = vec![0.0; channels]; + let mut valid_mask = vec![true; channels]; + if channels > 0 { + values[0] = 1.0; + } + if channels > 1 { + // Channel 1 is missing this tick: placeholder 0.0, mask false. + values[1] = 0.0; + valid_mask[1] = false; + } + let mut custom = HashMap::new(); + custom.insert("schema".into(), STIMULUS_SCHEMA.to_string()); + custom.insert("provenance".into(), "simulated-telemetry".into()); + StimulusBatch { + session_id: Some("thalamic-smoke".into()), + batch_id, + timestamp: now_ns(), + values, + valid_mask: Some(valid_mask), + metadata: Some(BatchMetadata { + processing_latency_ns: Some(1_000), + source: Some("thalamic-relay-fixture".into()), + custom, + }), + } + } + + /// Encode through the published `IpcMessage` contract (not a copied struct). + pub fn encode_frame(&self, batch: &StimulusBatch) -> Result> { + batch + .validate() + .map_err(|err| anyhow!("thalamic fixture produced an invalid StimulusBatch: {err}"))?; + Ok(serde_json::to_vec(&IpcMessage::Stimuli(batch.clone()))?) + } + + /// Attempt to publish. Publish failure never clears `safety_healthy`. + pub fn publish(&mut self, transport: Option<&mut dyn StimulusTransport>, frame: &[u8]) { + match transport { + Some(tx) => match tx.send(frame) { + Ok(()) => self.published += 1, + Err(_) => self.publish_errors += 1, + }, + None => { + // Brainstem / transport unavailable. + self.publish_errors += 1; + } + } + self.safety_healthy = true; + } + + /// Evaluate a trivial independent protection predicate. + pub fn safety_tick(&mut self, thermal_ok: bool) { + self.safety_healthy = thermal_ok; + } +} + +/// Narrow send side used by the fixture. Production Thalamic would use ZMQ PUB. +pub trait StimulusTransport { + fn send(&mut self, frame: &[u8]) -> Result<()>; +} + +fn now_ns() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_nanos() as u64 +} diff --git a/tests/thalamic_brainstem_smoke.rs b/tests/thalamic_brainstem_smoke.rs new file mode 100644 index 0000000..c01e668 --- /dev/null +++ b/tests/thalamic_brainstem_smoke.rs @@ -0,0 +1,233 @@ +// SPDX-License-Identifier: MIT OR Apache-2.0 +// Copyright 2026 Raul Montoya Cardenas + +//! End-to-end smoke: Thalamic fixture → corpus-ipc types → Brainstem runtime. +//! +//! Requires `--features corpus-ipc`. CPU-only; no GPU. + +#[path = "fixtures/thalamic_producer.rs"] +mod thalamic; + +use std::collections::VecDeque; +use std::time::Duration; + +use anyhow::Result; +use brainstem_daemon::checkpoint::{NetworkDims, write_nonblank_checkpoint}; +use brainstem_daemon::daemon::{BrainstemDaemon, DaemonConfig}; +use brainstem_daemon::ingress::{IngressPolicy, accept_ipc_json}; +use brainstem_daemon::{BackendPair, CollectingSpikeSink, IngressPacket, StimulusSource}; +use corpus_ipc::{IpcMessage, StimulusBatch}; +use thalamic::{StimulusTransport, ThalamicProducer}; + +const CHANNELS: usize = 4; +const LIF: usize = 4; +const IZH: usize = 0; + +struct QueuedStimulusSource { + packets: VecDeque>, +} + +impl StimulusSource for QueuedStimulusSource { + fn next_ingress(&mut self) -> anyhow::Result> { + match self.packets.pop_front() { + Some(Ok(packet)) => Ok(Some(packet)), + Some(Err(err)) => Err(anyhow::anyhow!(err)), + None => Ok(None), + } + } + + fn initialize(&mut self, _model_path: Option<&str>) -> anyhow::Result<()> { + Ok(()) + } +} + +struct CaptureTransport { + frames: Vec>, +} + +impl StimulusTransport for CaptureTransport { + fn send(&mut self, frame: &[u8]) -> Result<()> { + self.frames.push(frame.to_vec()); + Ok(()) + } +} + +fn smoke_dims() -> NetworkDims { + NetworkDims { + lif_count: LIF, + izh_count: IZH, + channels: CHANNELS, + } +} + +fn smoke_config(model_path: std::path::PathBuf) -> DaemonConfig { + DaemonConfig { + tick_rate_hz: 1000, + log_level: "info".into(), + spine_sub_port: 5555, + spine_pub_port: 5556, + model_path, + lif_count: LIF, + izh_count: IZH, + channels: CHANNELS, + services: Vec::new(), + } +} + +fn now_ns() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos() as u64 +} + +#[test] +fn thalamic_fixture_has_no_spiking_network() { + let src = include_str!("fixtures/thalamic_producer.rs"); + assert!( + !src.contains("use neuromod"), + "Thalamic fixture must not import neuromod" + ); + assert!( + !src.contains("SpikingNetwork::"), + "Thalamic fixture must not construct a SpikingNetwork" + ); +} + +#[test] +fn checkpoint_is_loaded_instead_of_blank_network() { + let dir = std::env::temp_dir().join(format!( + "brainstem-smoke-ckpt-{}-{}", + std::process::id(), + now_ns() + )); + std::fs::create_dir_all(&dir).unwrap(); + let path = dir.join("smoke.json"); + let identity = write_nonblank_checkpoint(&path, "smoke-thalamic-v1", smoke_dims()).unwrap(); + + let source = QueuedStimulusSource { + packets: VecDeque::new(), + }; + let sink = CollectingSpikeSink::new(); + let pair = BackendPair { + source: Box::new(source), + sink: Box::new(sink), + }; + let stats = BrainstemDaemon::try_with_backend(smoke_config(path.clone()), pair) + .unwrap() + .run_for_ticks(1) + .unwrap(); + + let loaded = stats.loaded_checkpoint.expect("checkpoint must be loaded"); + assert_eq!(loaded.model_id, "smoke-thalamic-v1"); + assert_eq!(loaded.fingerprint, identity.fingerprint); + let _ = std::fs::remove_dir_all(dir); +} + +#[test] +fn typed_frame_crosses_ipc_and_produces_runtime_result() { + let dir = std::env::temp_dir().join(format!( + "brainstem-smoke-e2e-{}-{}", + std::process::id(), + now_ns() + )); + std::fs::create_dir_all(&dir).unwrap(); + let path = dir.join("smoke.json"); + write_nonblank_checkpoint(&path, "smoke-thalamic-v1", smoke_dims()).unwrap(); + + let mut thalamic = ThalamicProducer::new(); + let batch = thalamic.simulate_telemetry(42, CHANNELS); + let expected_mask = batch.valid_mask.clone(); + let bytes = thalamic.encode_frame(&batch).unwrap(); + let mut transport = CaptureTransport { frames: Vec::new() }; + thalamic.publish(Some(&mut transport), &bytes); + assert_eq!(thalamic.published, 1); + assert!(thalamic.safety_healthy); + + let policy = IngressPolicy::new(CHANNELS, Some(Duration::from_secs(1))); + let packet = accept_ipc_json(&transport.frames[0], &policy, now_ns()).unwrap(); + assert_eq!(packet.batch_id, Some(42)); + assert_eq!(packet.valid_mask, expected_mask); + + let source = QueuedStimulusSource { + packets: VecDeque::from([Ok(packet)]), + }; + let sink = CollectingSpikeSink::new(); + let pair = BackendPair { + source: Box::new(source), + sink: Box::new(sink), + }; + let stats = BrainstemDaemon::try_with_backend(smoke_config(path), pair) + .unwrap() + .run_for_ticks(1) + .unwrap(); + + assert_eq!(stats.ticks, 1); + assert_eq!(stats.accepted_batches, 1); + assert_eq!(stats.last_batch_id, Some(42)); + assert_eq!(stats.last_valid_mask, expected_mask); + assert!(stats.loaded_checkpoint.is_some()); + let _ = std::fs::remove_dir_all(dir); +} + +#[test] +fn validity_mask_survives_wire_contract() { + let thalamic = ThalamicProducer::new(); + let batch = thalamic.simulate_telemetry(1, CHANNELS); + let wire = thalamic.encode_frame(&batch).unwrap(); + let decoded: IpcMessage = serde_json::from_slice(&wire).unwrap(); + match decoded { + IpcMessage::Stimuli(StimulusBatch { + valid_mask, values, .. + }) => { + assert_eq!(valid_mask, batch.valid_mask); + assert_eq!(values.len(), CHANNELS); + assert!(!valid_mask.as_ref().unwrap()[1]); + } + other => panic!("expected Stimuli, got {other:?}"), + } +} + +#[test] +fn schema_incompatibility_fails_loudly() { + let policy = IngressPolicy::new(CHANNELS, Some(Duration::from_secs(1))); + + let old_udp = br#"{"type":"Stimuli","values":[1.0,0.0,0.0,0.0]}"#; + let err = accept_ipc_json(old_udp, &policy, now_ns()).unwrap_err(); + assert!(err.to_string().contains("deserialize")); + + let thalamic = ThalamicProducer::new(); + let mut batch = thalamic.simulate_telemetry(2, CHANNELS); + batch + .metadata + .as_mut() + .unwrap() + .custom + .insert("schema".into(), "not-a-supported-schema".into()); + let bytes = serde_json::to_vec(&IpcMessage::Stimuli(batch)).unwrap(); + let err = accept_ipc_json(&bytes, &policy, now_ns()).unwrap_err(); + assert!(err.to_string().contains("schema")); + + let ping = serde_json::to_vec(&IpcMessage::Ping).unwrap(); + let err = accept_ipc_json(&ping, &policy, now_ns()).unwrap_err(); + assert!(err.to_string().contains("Ping")); +} + +#[test] +fn thalamic_stays_healthy_when_brainstem_unavailable() { + let mut thalamic = ThalamicProducer::new(); + thalamic.safety_tick(true); + let batch = thalamic.simulate_telemetry(9, CHANNELS); + let bytes = thalamic.encode_frame(&batch).unwrap(); + + // No transport: Brainstem is not running. + thalamic.publish(None, &bytes); + assert!(thalamic.safety_healthy); + assert_eq!(thalamic.publish_errors, 1); + assert_eq!(thalamic.published, 0); + + // Safety evaluation continues after the failed publish. + thalamic.safety_tick(true); + assert!(thalamic.safety_healthy); + let _still_producing = thalamic.simulate_telemetry(10, CHANNELS); +} From 9072863a9c51bdbcf2c0692c81b292aaf80eed93 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 15 Sep 2026 05:53:48 +0000 Subject: [PATCH 2/8] fix: tick on rejected ZMQ frames and hold last neuromodulators Protocol validation failures (Ping, stale, width, other IpcMessage variants) now skip ingress with Ok(None) instead of Err, so run_tick still advances the SNN. Idle EAGAIN ticks return the last neuromodulator snapshot so dopamine/cortisol/acetylcholine/tempo no longer reset. --- src/backend.rs | 126 +++++++++++++++++++++++++++++++++++++++++++++---- 1 file changed, 117 insertions(+), 9 deletions(-) diff --git a/src/backend.rs b/src/backend.rs index 6f9a1a3..98b7f5b 100644 --- a/src/backend.rs +++ b/src/backend.rs @@ -231,6 +231,20 @@ mod zmq_impl { .unwrap_or_default() .as_nanos() as u64 } + + /// Skip ingress this tick while retaining the last neuromodulator snapshot. + /// + /// `None` means there is no held snapshot, matching the `StimulusSource` + /// contract (`Ok(None)` still advances the network on zeroed stimuli). + fn skip_ingress(&self) -> Option { + self.last_modulators.as_ref().map(|mods| IngressPacket { + stimuli: Vec::new(), + modulators: Some(mods.clone()), + valid_mask: None, + batch_id: None, + timestamp_ns: None, + }) + } } impl StimulusSource for ZmqStimulusSource { @@ -242,16 +256,24 @@ mod zmq_impl { match socket.socket.recv_bytes(::zmq::DONTWAIT) { Ok(buf) => { let policy = IngressPolicy::new(self.channels, Some(self.max_age)); - let mut packet = accept_ipc_json(&buf, &policy, Self::now_ns()) - .map_err(|e| anyhow::anyhow!("{e}"))?; - if let Some(mods) = packet.modulators.as_ref() { - self.last_modulators = Some(mods.clone()); - } else { - packet.modulators = self.last_modulators.clone(); + match accept_ipc_json(&buf, &policy, Self::now_ns()) { + Ok(mut packet) => { + if let Some(mods) = packet.modulators.as_ref() { + self.last_modulators = Some(mods.clone()); + } else { + packet.modulators = self.last_modulators.clone(); + } + Ok(Some(packet)) + } + Err(e) => { + // Consumed-but-invalid payload (stale, Ping, width, ...): + // skip ingress this tick but still let the caller tick. + tracing::warn!("Rejected ingress frame: {e}"); + Ok(self.skip_ingress()) + } } - Ok(Some(packet)) } - Err(::zmq::Error::EAGAIN) => Ok(None), + Err(::zmq::Error::EAGAIN) => Ok(self.skip_ingress()), Err(e) => Err(anyhow::anyhow!("ZMQ recv failed: {e}")), } } @@ -330,7 +352,7 @@ mod zmq_impl { mod tests { use super::*; use crate::ingress::STIMULUS_SCHEMA; - use corpus_ipc::{BatchMetadata, StimulusBatch}; + use corpus_ipc::{BatchMetadata, NeuromodulatorSnapshot, StimulusBatch}; use std::collections::HashMap; use std::time::Duration; @@ -395,6 +417,92 @@ mod zmq_impl { assert_eq!(packet.valid_mask.as_ref().map(|m| m[1]), Some(false)); assert!((packet.stimuli[0] - 1.0).abs() < f32::EPSILON); } + + fn bind_loopback(channels: usize) -> (::zmq::Socket, ZmqStimulusSource) { + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let port = listener.local_addr().unwrap().port(); + drop(listener); + let endpoint = format!("tcp://127.0.0.1:{port}"); + + let context = ::zmq::Context::new(); + let publisher = context.socket(::zmq::PUB).unwrap(); + publisher.bind(&endpoint).unwrap(); + + let mut source = + ZmqStimulusSource::with_channels(channels).with_max_age(Duration::from_secs(5)); + source.connect(&endpoint).unwrap(); + std::thread::sleep(Duration::from_millis(150)); + (publisher, source) + } + + fn poll_ingress(source: &mut ZmqStimulusSource) -> Result> { + let mut last = Ok(None); + for _ in 0..50 { + last = source.next_ingress(); + match &last { + Err(err) => panic!("ingress error: {err}"), + Ok(Some(_)) => return last, + Ok(None) => std::thread::sleep(Duration::from_millis(10)), + } + } + last + } + + #[test] + fn rejected_frame_does_not_hard_fail_ingress() { + let (publisher, mut source) = bind_loopback(4); + let ping = serde_json::to_vec(&IpcMessage::Ping).expect("serialize Ping"); + publisher.send(&ping, 0).unwrap(); + + // Ping is consumed-but-invalid: must be Ok(None), not Err, so the + // subsequent valid batch can still be accepted on the next tick. + let ping_result = poll_ingress(&mut source); + assert!( + ping_result.is_ok(), + "rejected Ping must not be a hard receive failure: {ping_result:?}" + ); + assert!( + matches!(ping_result, Ok(None)), + "Ping has no held modulators, so skip-ingress is Ok(None): {ping_result:?}" + ); + + let frame = sample_frame(4, 101); + publisher.send(&frame, 0).unwrap(); + let packet = poll_ingress(&mut source) + .expect("recv") + .expect("ZMQ SUB should receive a typed StimulusBatch after Ping"); + assert_eq!(packet.batch_id, Some(101)); + } + + #[test] + fn idle_ticks_hold_last_modulators() { + let (publisher, mut source) = bind_loopback(4); + let snapshot = NeuromodulatorSnapshot { + tick: 1, + dopamine: 0.4, + cortisol: 0.3, + acetylcholine: 0.2, + tempo: 1.0, + }; + let frame = + serde_json::to_vec(&IpcMessage::Neuromodulators(snapshot)).expect("serialize"); + publisher.send(&frame, 0).unwrap(); + + let packet = poll_ingress(&mut source) + .expect("recv") + .expect("ZMQ SUB should receive Neuromodulators"); + assert_eq!( + packet.modulators.as_deref(), + Some(&[0.4, 0.3, 0.2, 1.0][..]) + ); + + let idle = source + .next_ingress() + .expect("EAGAIN with held modulators must not be a hard failure"); + let idle = idle.expect("idle tick must carry last neuromodulator snapshot"); + assert!(idle.stimuli.is_empty()); + assert_eq!(idle.modulators.as_deref(), Some(&[0.4, 0.3, 0.2, 1.0][..])); + } } } From eface214ef5a5bc95a54ef7ed87ba56b391d4b3a Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 17 Sep 2026 21:19:48 +0000 Subject: [PATCH 3/8] fix: satisfy DeepSource Default/clone_from on smoke fixture and skip_ingress Implement ThalamicProducer::default with explicit healthy fields and have new() call Default so RS-A1008 does not treat default() as a recursive Self constructor. Reuse last_modulators via clone_from in skip_ingress. Co-authored-by: Raul Cardenas Montoya --- src/backend.rs | 8 ++++---- tests/fixtures/thalamic_producer.rs | 16 ++++++++-------- 2 files changed, 12 insertions(+), 12 deletions(-) diff --git a/src/backend.rs b/src/backend.rs index a7466b6..bcf6916 100644 --- a/src/backend.rs +++ b/src/backend.rs @@ -252,12 +252,12 @@ mod zmq_impl { if !rejected && self.last_modulators.is_none() { return None; } - Some(IngressPacket { - stimuli: Vec::new(), - modulators: self.last_modulators.clone(), + let mut packet = IngressPacket { rejected, ..IngressPacket::default() - }) + }; + packet.modulators.clone_from(&self.last_modulators); + Some(packet) } fn attach_held_modulators(&self, packet: &mut IngressPacket) { diff --git a/tests/fixtures/thalamic_producer.rs b/tests/fixtures/thalamic_producer.rs index 3736cf6..1ef4de9 100644 --- a/tests/fixtures/thalamic_producer.rs +++ b/tests/fixtures/thalamic_producer.rs @@ -24,14 +24,20 @@ pub struct ThalamicProducer { pub publish_errors: u64, } -impl ThalamicProducer { - pub fn new() -> Self { +impl Default for ThalamicProducer { + fn default() -> Self { Self { safety_healthy: true, published: 0, publish_errors: 0, } } +} + +impl ThalamicProducer { + pub fn new() -> Self { + Self::default() + } /// Build a representative normalized telemetry frame. pub fn simulate_telemetry(&self, batch_id: u64, channels: usize) -> StimulusBatch { @@ -90,12 +96,6 @@ impl ThalamicProducer { } } -impl Default for ThalamicProducer { - fn default() -> Self { - Self::new() - } -} - /// Narrow send side used by the fixture. Production Thalamic would use ZMQ PUB. pub trait StimulusTransport { fn send(&mut self, frame: &[u8]) -> Result<()>; From ba8240a890396816e20d459426e79784b67c30c6 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 17 Sep 2026 21:25:00 +0000 Subject: [PATCH 4/8] fix: clear Codacy method-length and test temp_dir findings Split stimulus schema and freshness checks out of accept_stimulus_batch so the function stays under the 50-line limit. Write smoke fixtures under CARGO_TARGET_TMPDIR (or crate target/) instead of the shared system temp dir. Co-authored-by: Raul Cardenas Montoya --- src/ingress/corpus.rs | 74 ++++++++++++++++--------------- tests/thalamic_brainstem_smoke.rs | 8 +++- 2 files changed, 45 insertions(+), 37 deletions(-) diff --git a/src/ingress/corpus.rs b/src/ingress/corpus.rs index 0725a9c..67cc17d 100644 --- a/src/ingress/corpus.rs +++ b/src/ingress/corpus.rs @@ -130,6 +130,43 @@ pub fn accept_ipc_message( } } +fn require_stimulus_schema(batch: &StimulusBatch) -> Result<(), IngressError> { + match batch.metadata.as_ref() { + Some(meta) => match meta.custom.get("schema") { + Some(schema) if schema == STIMULUS_SCHEMA => Ok(()), + Some(schema) => Err(IngressError::Schema(format!( + "unsupported schema '{schema}', expected {STIMULUS_SCHEMA}" + ))), + None => Err(IngressError::Schema( + "missing metadata.custom.schema".into(), + )), + }, + None => Err(IngressError::Schema("missing batch metadata".into())), + } +} + +fn reject_stale_or_future( + timestamp: u64, + now_ns: u64, + max_age: Option, +) -> Result<(), IngressError> { + let Some(max_age) = max_age else { + return Ok(()); + }; + if timestamp > now_ns { + return Err(IngressError::Future { + timestamp_ns: timestamp, + now_ns, + }); + } + let max_age_ns = max_age.as_nanos() as u64; + let age_ns = now_ns - timestamp; + if age_ns > max_age_ns { + return Err(IngressError::Stale { age_ns, max_age_ns }); + } + Ok(()) +} + fn accept_stimulus_batch( batch: StimulusBatch, policy: &IngressPolicy, @@ -138,47 +175,14 @@ fn accept_stimulus_batch( batch .validate() .map_err(|err| IngressError::Stimulus(err.to_string()))?; - - match batch.metadata.as_ref() { - Some(meta) => match meta.custom.get("schema") { - Some(schema) if schema == STIMULUS_SCHEMA => {} - Some(schema) => { - return Err(IngressError::Schema(format!( - "unsupported schema '{schema}', expected {STIMULUS_SCHEMA}" - ))); - } - None => { - return Err(IngressError::Schema( - "missing metadata.custom.schema".into(), - )); - } - }, - None => { - return Err(IngressError::Schema("missing batch metadata".into())); - } - } - + require_stimulus_schema(&batch)?; if policy.expected_channels != 0 && batch.values.len() != policy.expected_channels { return Err(IngressError::Width { expected: policy.expected_channels, got: batch.values.len(), }); } - - if let Some(max_age) = policy.max_age { - if batch.timestamp > now_ns { - return Err(IngressError::Future { - timestamp_ns: batch.timestamp, - now_ns, - }); - } - let max_age_ns = max_age.as_nanos() as u64; - let age_ns = now_ns - batch.timestamp; - if age_ns > max_age_ns { - return Err(IngressError::Stale { age_ns, max_age_ns }); - } - } - + reject_stale_or_future(batch.timestamp, now_ns, policy.max_age)?; // Do not zero masked slots here. `valid_mask` is preserved so downstream // consumers can tell a producer 0.0 from a masked placeholder; the tick // loop applies the mask once in `decode_inputs`. diff --git a/tests/thalamic_brainstem_smoke.rs b/tests/thalamic_brainstem_smoke.rs index 6905609..656e88c 100644 --- a/tests/thalamic_brainstem_smoke.rs +++ b/tests/thalamic_brainstem_smoke.rs @@ -27,8 +27,12 @@ struct TempDir(PathBuf); impl TempDir { fn new(prefix: &str) -> Self { - let dir = - std::env::temp_dir().join(format!("{prefix}-{}-{}", std::process::id(), now_ns())); + // Cargo sets CARGO_TARGET_TMPDIR for tests; fall back to this crate's + // target/ so fixtures never use the shared system temp directory. + let root = std::env::var_os("CARGO_TARGET_TMPDIR") + .map(PathBuf::from) + .unwrap_or_else(|| PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("target")); + let dir = root.join(format!("{prefix}-{}-{}", std::process::id(), now_ns())); std::fs::create_dir_all(&dir).unwrap(); Self(dir) } From 892370065c877fbda5433242a763ed30e887cabe Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 17 Sep 2026 21:30:27 +0000 Subject: [PATCH 5/8] fix: use Option::map_or_else for smoke fixture scratch root DeepSource RS antipattern: replace map + unwrap_or_else with map_or_else. Co-authored-by: Raul Cardenas Montoya --- tests/thalamic_brainstem_smoke.rs | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/tests/thalamic_brainstem_smoke.rs b/tests/thalamic_brainstem_smoke.rs index 656e88c..3770f52 100644 --- a/tests/thalamic_brainstem_smoke.rs +++ b/tests/thalamic_brainstem_smoke.rs @@ -29,9 +29,10 @@ impl TempDir { fn new(prefix: &str) -> Self { // Cargo sets CARGO_TARGET_TMPDIR for tests; fall back to this crate's // target/ so fixtures never use the shared system temp directory. - let root = std::env::var_os("CARGO_TARGET_TMPDIR") - .map(PathBuf::from) - .unwrap_or_else(|| PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("target")); + let root = std::env::var_os("CARGO_TARGET_TMPDIR").map_or_else( + || PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("target"), + PathBuf::from, + ); let dir = root.join(format!("{prefix}-{}-{}", std::process::id(), now_ns())); std::fs::create_dir_all(&dir).unwrap(); Self(dir) From 84bdd80b421104a7c8ec748e93c5418a17cc504f Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 17 Sep 2026 21:32:20 +0000 Subject: [PATCH 6/8] fix: group tick health and stats to satisfy clippy argument limit run_tick would exceed clippy::too-many-arguments after merging #53 health reporting with RuntimeStats. Bundle those observers in TickReport. Co-authored-by: Raul Cardenas Montoya --- src/daemon.rs | 46 ++++++++++++++++++++++++++++++++-------------- 1 file changed, 32 insertions(+), 14 deletions(-) diff --git a/src/daemon.rs b/src/daemon.rs index a30a3cc..ed5e857 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -291,8 +291,10 @@ impl BrainstemDaemon { &mut stimuli, &mut spike_buf, &ingress, - &health, - &mut stats, + &mut TickReport { + health: &health, + stats: &mut stats, + }, ); } _ = &mut shutdown => { @@ -357,8 +359,10 @@ impl BrainstemDaemon { &mut stimuli, &mut spike_buf, &ingress, - &health, - &mut stats, + &mut TickReport { + health: &health, + stats: &mut stats, + }, ); } @@ -666,6 +670,12 @@ fn validate_restored_pair( // Trait-based tick loop (works with or without corpus-ipc feature) +/// Health reporter plus smoke counters updated together on every tick. +struct TickReport<'a> { + health: &'a HealthHandle, + stats: &'a mut RuntimeStats, +} + fn run_tick( source: &mut dyn StimulusSource, network: &mut SpikingNetwork, @@ -673,9 +683,9 @@ fn run_tick( stimuli: &mut [f32], spike_buf: &mut Vec, ingress: &BoundedIngress, - health: &HealthHandle, - stats: &mut RuntimeStats, + report: &mut TickReport<'_>, ) { + let TickReport { health, stats } = report; let backend_packet = match source.next_ingress() { Ok(Some(p)) => Some(p), Ok(None) => None, @@ -850,8 +860,10 @@ pub(crate) fn run_tick_for_test( stimuli, spike_buf, &ingress, - &health, - &mut RuntimeStats::default(), + &mut TickReport { + health: &health, + stats: &mut RuntimeStats::default(), + }, ); } @@ -1249,8 +1261,10 @@ channels = 16 &mut stimuli, &mut spike_buf, &ingress, - &health, - &mut RuntimeStats::default(), + &mut TickReport { + health: &health, + stats: &mut RuntimeStats::default(), + }, ); assert_eq!(sink.emitted.len(), 1); @@ -1552,8 +1566,10 @@ block_timeout_ms = 0 &mut stimuli, &mut spike_buf, &ingress, - &health, - &mut RuntimeStats::default(), + &mut super::TickReport { + health: &health, + stats: &mut RuntimeStats::default(), + }, ); assert_eq!(stimuli, vec![0.1, 0.2]); @@ -1602,8 +1618,10 @@ block_timeout_ms = 0 &mut stimuli, &mut spike_buf, &ingress, - &health, - &mut RuntimeStats::default(), + &mut super::TickReport { + health: &health, + stats: &mut RuntimeStats::default(), + }, ); assert_eq!(stimuli, vec![0.3, 0.4]); From be808bc9294fbfb8dd29bf1aab17c02d775e7482 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 17 Sep 2026 21:37:01 +0000 Subject: [PATCH 7/8] refactor: share initialize and checkpoint boot between run paths run and run_for_ticks both connect the stimulus source then validate the restored network. Extract boot_network so health events stay in lockstep and Codacy duplication from the merge stays down. Co-authored-by: Raul Cardenas Montoya --- src/daemon.rs | 70 ++++++++++++++++++++++++--------------------------- 1 file changed, 33 insertions(+), 37 deletions(-) diff --git a/src/daemon.rs b/src/daemon.rs index ed5e857..f8006d9 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -249,27 +249,13 @@ impl BrainstemDaemon { } }; - if let Err(e) = initialize_source(&mut *backend.source, &cfg, &health) { - abort_startup(&ingress, &mut backend, control_stop, control_task).await; - return Err(e); - } - - let restored = take_restored_network(&cfg, restored); - let (mut network, provenance) = match restored { + let (mut network, _) = match boot_network(&mut *backend.source, &cfg, &health, restored) { Ok(pair) => pair, Err(err) => { - health.apply(HealthEvent::CheckpointRejected { - detail: err.to_string(), - }); abort_startup(&ingress, &mut backend, control_stop, control_task).await; return Err(err); } }; - health.apply(HealthEvent::CheckpointValidated { - identity: checkpoint_identity(&provenance), - }); - - log_model_provenance(&cfg, &provenance); let tick_duration = Duration::from_nanos(1_000_000_000 / u64::from(cfg.tick_rate_hz)); let mut ticker = time::interval(tick_duration); @@ -324,28 +310,15 @@ impl BrainstemDaemon { let health = self.health; let mut stats = RuntimeStats::default(); - if let Err(e) = initialize_source(&mut *backend.source, &cfg, &health) { - ingress.shutdown(); - shutdown_backend(&mut backend); - return Err(e); - } - - let restored = match take_restored_network(&cfg, None) { - Ok(pair) => pair, - Err(err) => { - health.apply(HealthEvent::CheckpointRejected { - detail: err.to_string(), - }); - ingress.shutdown(); - shutdown_backend(&mut backend); - return Err(err); - } - }; - let (mut network, provenance) = restored; - health.apply(HealthEvent::CheckpointValidated { - identity: checkpoint_identity(&provenance), - }); - log_model_provenance(&cfg, &provenance); + let (mut network, provenance) = + match boot_network(&mut *backend.source, &cfg, &health, None) { + Ok(pair) => pair, + Err(err) => { + ingress.shutdown(); + shutdown_backend(&mut backend); + return Err(err); + } + }; stats.loaded_checkpoint = Some(provenance); let mut stimuli = vec![0.0; cfg.channels]; @@ -523,6 +496,29 @@ fn initialize_source( Ok(()) } +fn boot_network( + source: &mut dyn StimulusSource, + cfg: &DaemonConfig, + health: &HealthHandle, + restored: Option<(SpikingNetwork, ModelProvenance)>, +) -> Result<(SpikingNetwork, ModelProvenance)> { + initialize_source(source, cfg, health)?; + let (network, provenance) = match take_restored_network(cfg, restored) { + Ok(pair) => pair, + Err(err) => { + health.apply(HealthEvent::CheckpointRejected { + detail: err.to_string(), + }); + return Err(err); + } + }; + health.apply(HealthEvent::CheckpointValidated { + identity: checkpoint_identity(&provenance), + }); + log_model_provenance(cfg, &provenance); + Ok((network, provenance)) +} + fn checkpoint_identity(provenance: &ModelProvenance) -> CheckpointIdentity { let digest = if provenance.content_sha256 == "none" { None From 25f22e42fc43b2706f2ae90fea5de7d70d92441e Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 17 Sep 2026 21:39:12 +0000 Subject: [PATCH 8/8] fix: split run_for_ticks helpers to stay under Codacy line limit Extract drive_ticks and finish_bounded_run so the bounded smoke entry stays under 50 lines without changing shutdown-on-flush-error behavior. Co-authored-by: Raul Cardenas Montoya --- src/daemon.rs | 91 +++++++++++++++++++++++++++++++-------------------- 1 file changed, 55 insertions(+), 36 deletions(-) diff --git a/src/daemon.rs b/src/daemon.rs index f8006d9..bf63bb1 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -320,42 +320,16 @@ impl BrainstemDaemon { } }; stats.loaded_checkpoint = Some(provenance); - - let mut stimuli = vec![0.0; cfg.channels]; - let mut spike_buf: Vec = Vec::with_capacity(128); - - for _ in 0..ticks { - run_tick( - &mut *backend.source, - &mut network, - &mut *backend.sink, - &mut stimuli, - &mut spike_buf, - &ingress, - &mut TickReport { - health: &health, - stats: &mut stats, - }, - ); - } - - let flush_err = backend - .sink - .flush() - .context("failed to flush spike sink") - .err(); - let shutdown_err = backend - .source - .shutdown() - .context("failed to shut down stimulus source") - .err(); - ingress.shutdown(); - if let Some(err) = flush_err { - return Err(err); - } - if let Some(err) = shutdown_err { - return Err(err); - } + drive_ticks( + ticks, + &mut backend, + &mut network, + &ingress, + &health, + &mut stats, + cfg.channels, + ); + finish_bounded_run(&mut backend, &ingress)?; Ok(stats) } } @@ -372,6 +346,51 @@ fn shutdown_backend(backend: &mut BackendPair) { } } +fn drive_ticks( + ticks: u64, + backend: &mut BackendPair, + network: &mut SpikingNetwork, + ingress: &BoundedIngress, + health: &HealthHandle, + stats: &mut RuntimeStats, + channels: usize, +) { + let mut stimuli = vec![0.0; channels]; + let mut spike_buf: Vec = Vec::with_capacity(128); + for _ in 0..ticks { + run_tick( + &mut *backend.source, + network, + &mut *backend.sink, + &mut stimuli, + &mut spike_buf, + ingress, + &mut TickReport { health, stats }, + ); + } +} + +fn finish_bounded_run(backend: &mut BackendPair, ingress: &BoundedIngress) -> Result<()> { + let flush_err = backend + .sink + .flush() + .context("failed to flush spike sink") + .err(); + let shutdown_err = backend + .source + .shutdown() + .context("failed to shut down stimulus source") + .err(); + ingress.shutdown(); + if let Some(err) = flush_err { + return Err(err); + } + if let Some(err) = shutdown_err { + return Err(err); + } + Ok(()) +} + /// Wait for a termination request: `SIGINT` (Ctrl-C) on every platform, plus /// `SIGTERM` (the default `systemctl stop` / `kill` signal) on Unix. #[cfg(unix)]