diff --git a/CHANGELOG.md b/CHANGELOG.md index c4d2c1e..f262152 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -44,7 +44,18 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - `## 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`) plus a public + `CollectingSpikeSink` for tests and the Thalamic integration smoke harness. +- CPU-only Thalamic → `corpus-ipc` → Brainstem smoke (`tests/thalamic_brainstem_smoke.rs`, + `required-features = ["corpus-ipc"]`). A Thalamic fixture with no `SpikingNetwork` + publishes typed `IpcMessage::Stimuli(StimulusBatch)` JSON; Brainstem loads a Distill + sidecar, validates schema/width/freshness/`valid_mask`, and ticks through + `BrainstemDaemon::run_for_ticks`. +- Typed `corpus-ipc` ingress (`accept_ipc_json`): schema token + `corpus-ipc.stimulus.v1`, channel-width check (or unspecified-width when + `expected_channels = 0`), future-timestamp and max-age rejection, session_id + passthrough, and `RuntimeStats` accepted/rejected counters. Rejected frames + still advance the network. Modulation-only frames are drained in the same tick. - `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 @@ -61,11 +72,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Changed - Switch optional `corpus-ipc` from a git pin to crates.io `0.1.0` - (`features = ["zmq"]`). ZMQ ingress uses published `ZmqIpcBackend` / - `IpcBackend::process_batch`; egress publishes unversioned - `IpcMessage::Spikes` JSON. The binary sets `CORPUS_IPC_ZMQ_READOUT_IPC` - (what 0.1 reads) and still sets `SPIKENAUT_ZMQ_READOUT_IPC` for older - tooling. + (`features = ["zmq"]`). ZMQ SUB ingress decodes unversioned JSON + `IpcMessage` (`Stimuli` / `Neuromodulators`) with crates.io types; egress + publishes unversioned `IpcMessage::Spikes` JSON. The binary sets + `CORPUS_IPC_ZMQ_READOUT_IPC` and still sets `SPIKENAUT_ZMQ_READOUT_IPC` for + older tooling. - Upgrade `neuromod` from 0.4.0 to crates.io **0.6.0** (pre-1.0 range `>=0.6.0, <0.7.0`). Ingress modulators map to dopamine / serotonin / acetylcholine / norepinephrine; `cortisol`, `tempo`, and `aux_dopamine` diff --git a/Cargo.toml b/Cargo.toml index 7156579..c183ec6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -53,6 +53,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/README.md b/README.md index 17d9434..e922a2b 100644 --- a/README.md +++ b/README.md @@ -149,7 +149,7 @@ 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`. `BrainstemDaemon::run` / `run_with_restored_network` is the sole caller of `StimulusSource::initialize` (the binary no longer initializes first, so the pinned ZMQ backend is not reconnected). A failing `initialize` marks health **fatal** and never becomes ready. +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`. `BrainstemDaemon::run` / `run_with_restored_network` / `run_for_ticks` call `StimulusSource::initialize` (the binary no longer initializes first, so the pinned ZMQ backend is not reconnected). A failing `initialize` marks health **fatal** and never becomes ready. Health snapshots, probe paths, and the transition table live in [`docs/health.md`](docs/health.md). @@ -166,21 +166,23 @@ Health snapshots, probe paths, and the transition table live in [`docs/health.md | `ingress` | used (bounded class queues in the tick loop; health reports aggregate fill) | used (same queues wrap backend packets before the network step) | | `spine_sub_port` | parsed, **no-op** | sets `CORPUS_IPC_ZMQ_READOUT_IPC` to `tcp://127.0.0.1:` (also sets legacy `SPIKENAUT_ZMQ_READOUT_IPC` for compatibility) | | `spine_pub_port` | parsed, **no-op** | binds ZMQ PUB `tcp://*:` | -| `model_path` | used in **live** mode (sidecar JSON); ignored in **simulation** (`StubStimulusSource::initialize` still ignores it) | same live/simulation gate, then passed literally to `initialize` (no `~` expansion); published `ZmqIpcBackend` currently ignores `_model_path` | +| `model_path` | used in **live** mode (sidecar JSON); ignored in **simulation** (`StubStimulusSource::initialize` still ignores it) | same live/simulation gate, then passed literally to `initialize` (no `~` expansion); the ZMQ SUB source connects and ignores `_model_path` | **Settings that only take effect with `corpus-ipc`** (the `brainstem-daemon` binary built `--features corpus-ipc`): - `spine_sub_port` (drives `CORPUS_IPC_ZMQ_READOUT_IPC`) - `spine_pub_port` -- `CORPUS_IPC_ZMQ_READOUT_IPC` (const `CORPUS_IPC_READOUT_ENV`; this is what published `ZmqIpcBackend::initialize` reads) +- `CORPUS_IPC_ZMQ_READOUT_IPC` (const `CORPUS_IPC_READOUT_ENV`; this is what `ZmqStimulusSource::initialize` reads when no explicit `connect` endpoint is set) -**Passed through / set, but currently unused by published `corpus-ipc` 0.1:** +**Passed through / set, but unused by the ZMQ SUB source after connect:** - `model_path` (literal filesystem path; `~` is not expanded; passed to `initialize`, which names the argument `_model_path` and does not consume it; live-mode restore consumes it before that handshake) -- `SPIKENAUT_ZMQ_READOUT_IPC` (const `LEGACY_SPIKENAUT_READOUT_ENV`; the binary still sets this alongside `CORPUS_IPC_ZMQ_READOUT_IPC` for older tooling; published `corpus-ipc` does not read it) +- `SPIKENAUT_ZMQ_READOUT_IPC` (const `LEGACY_SPIKENAUT_READOUT_ENV`; the binary still sets this alongside `CORPUS_IPC_ZMQ_READOUT_IPC` for older tooling) Under stub those ZMQ 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. +ZMQ SUB ingress decodes unversioned JSON `IpcMessage` frames (`Stimuli` / `Neuromodulators`) through crates.io `corpus-ipc` 0.1 types. Width, schema token `corpus-ipc.stimulus.v1`, freshness, and future timestamps are rejected without stopping the tick loop. Modulation-only frames are drained in the same tick so they do not consume a sensory period. + ### Runtime modes: simulation vs loaded Spikenaut `runtime_mode` is independent of the stub vs ZMQ **backend**. Backends move stimuli and spikes. The network itself is restored **before** the tick loop: @@ -216,6 +218,16 @@ The allowed software artifact is the Distill sidecar JSON published as Hugging F Live restore copies LIF weights, membrane, `last_spike`, `decay_rate`, and `threshold` (also seeding `base_threshold`) and sets `RmStdpConfig.reward_lr = 0` so dopamine-gated R-STDP cannot retrain the Distill matrix. `neuromod` 0.6.0 `SpikingNetwork::step` still assigns `decay_rate` from acetylcholine, blends `threshold` toward `0.05..=0.50`, and L1-renormalizes rows whose weights already sum above `1e-6`. Those are engine contracts; this crate does not fork `step`. +### Thalamic → corpus-ipc → Brainstem smoke + +CPU-only integration coverage (no GPU) lives in `tests/thalamic_brainstem_smoke.rs` and is gated on `--features corpus-ipc` so default stub tests never need `libzmq`. + +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 restores a Distill sidecar JSON checkpoint before ticking, preserves `valid_mask` and `session_id` across the wire, and rejects incompatible schema/JSON loudly. A separate assertion keeps the fixture's safety flag healthy when Brainstem/transport is absent, and a later publish cannot clobber a thermal fault. + +```bash +CC=gcc CXX=g++ cargo test --locked --features corpus-ipc --test thalamic_brainstem_smoke +``` + Simulation is the deliberate test/dev path for a blank network. Do not use it as a stand-in for production Spikenaut. ```toml diff --git a/docs/ci.md b/docs/ci.md index e50c4cf..7df7099 100644 --- a/docs/ci.md +++ b/docs/ci.md @@ -17,7 +17,9 @@ All jobs install **Rust 1.98.1** — the same string as `Cargo.toml` | `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/src/backend.rs b/src/backend.rs index bbc6e2b..bcf6916 100644 --- a/src/backend.rs +++ b/src/backend.rs @@ -30,6 +30,18 @@ pub struct IngressPacket { /// (dopamine, serotonin, acetylcholine, norepinephrine). /// 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, + /// Optional producer session from `StimulusBatch.session_id`. + pub session_id: Option, + /// True when the source consumed an invalid/incompatible frame this tick. + /// The tick loop still advances; `RuntimeStats.rejected_batches` counts these. + pub rejected: bool, } /// Local spike event type (independent of any external crate). @@ -107,10 +119,7 @@ pub struct StubStimulusSource; impl StimulusSource for StubStimulusSource { fn next_ingress(&mut self) -> Result> { - Ok(Some(IngressPacket { - stimuli: Vec::new(), - modulators: None, - })) + Ok(Some(IngressPacket::default())) } fn initialize(&mut self, _model_path: Option<&str>) -> Result<()> { @@ -127,14 +136,14 @@ impl SpikeSink for NoopSpikeSink { } } -/// Collecting sink for tests. Collects every emitted batch. -#[cfg(test)] +/// Collecting sink for tests and the Thalamic 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(), @@ -142,14 +151,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()); @@ -163,14 +164,21 @@ impl SpikeSink for CollectingSpikeSink { #[cfg(feature = "corpus-ipc")] mod zmq_impl { use super::*; - // Import `IpcBackend` so `initialize` / `process_batch` are in scope for - // crates.io `corpus-ipc` 0.1 (`ZmqIpcBackend`). - use corpus_ipc::IpcBackend; - use corpus_ipc::{IpcMessage, SpikeBatch, SpikeEvent as CorpusSpikeEvent, ZmqIpcBackend}; + 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"; + /// Cap on modulator-only frames drained in one tick so a flood cannot stall 1 kHz. + const MAX_DRAIN_PER_TICK: usize = 32; pub struct ZmqStimulusSource { - inner: ZmqIpcBackend, + socket: Option, channels: usize, + last_modulators: Option>, + max_age: std::time::Duration, } impl Default for ZmqStimulusSource { @@ -180,54 +188,136 @@ mod zmq_impl { } impl ZmqStimulusSource { + /// Unspecified-width constructor: typed `StimulusBatch` width is not checked. + /// Prefer [`Self::with_channels`] when the daemon config width is known. pub fn new() -> Self { - Self { - inner: ZmqIpcBackend::new(), - channels: 0, - } + Self::with_channels(0) } - /// Construct with known channel count so `next_ingress` can split - /// stimulus prefix from appended neuromodulator tail - /// ([`NEUROMODULATOR_COUNT`] floats: DA / 5-HT / ACh / NE). - /// - /// 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: ZmqIpcBackend::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 + } + + fn hold_modulators(&mut self, mods: &[f32]) { + let dst = self.last_modulators.get_or_insert_with(Vec::new); + dst.clear(); + dst.extend_from_slice(mods); + } + + fn skip_ingress(&self, rejected: bool) -> Option { + if !rejected && self.last_modulators.is_none() { + return None; + } + let mut packet = IngressPacket { + rejected, + ..IngressPacket::default() + }; + packet.modulators.clone_from(&self.last_modulators); + Some(packet) + } + + fn attach_held_modulators(&self, packet: &mut IngressPacket) { + if packet.modulators.is_none() { + packet.modulators.clone_from(&self.last_modulators); } } } impl StimulusSource for ZmqStimulusSource { fn next_ingress(&mut self) -> Result> { - let readout = self.inner.process_batch(&[])?; - let ch = self.channels; - if ch > 0 && readout.len() > ch { - let stimuli = readout[..ch].to_vec(); - let modulators = if readout.len() >= ch + NEUROMODULATOR_COUNT { - Some(readout[ch..ch + NEUROMODULATOR_COUNT].to_vec()) - } else { - None - }; - Ok(Some(IngressPacket { - stimuli, - modulators, - })) - } else { - Ok(Some(IngressPacket { - stimuli: readout, - modulators: None, - })) + let policy = IngressPolicy::new(self.channels, Some(self.max_age)); + + for _ in 0..MAX_DRAIN_PER_TICK { + let recvd = self + .socket + .as_ref() + .ok_or_else(|| anyhow::anyhow!("ZMQ stimulus source not initialized"))? + .socket + .recv_bytes(::zmq::DONTWAIT); + match recvd { + Ok(buf) => match accept_ipc_json(&buf, &policy, Self::now_ns()) { + Ok(packet) if packet.stimuli.is_empty() && packet.modulators.is_some() => { + if let Some(mods) = packet.modulators.as_ref() { + self.hold_modulators(mods); + } + // Modulation-only: keep draining so a Stimuli frame + // in the same tick is not delayed by one period. + continue; + } + Ok(mut packet) => { + if let Some(mods) = packet.modulators.as_ref() { + self.hold_modulators(mods); + } else { + self.attach_held_modulators(&mut packet); + } + return Ok(Some(packet)); + } + Err(e) => { + tracing::warn!("Rejected ingress frame: {e}"); + // Consumed-but-invalid: stop draining this tick so a + // later valid frame is still available next period. + return Ok(self.skip_ingress(true)); + } + }, + Err(::zmq::Error::EAGAIN) => { + return Ok(self.skip_ingress(false)); + } + Err(e) => return Err(anyhow::anyhow!("ZMQ recv failed: {e}")), + } } + + Ok(self.skip_ingress(false)) } - 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()) } } @@ -296,6 +386,164 @@ mod zmq_impl { Ok(()) } } + + #[cfg(test)] + mod tests { + use super::*; + use crate::ingress::STIMULUS_SCHEMA; + use corpus_ipc::{BatchMetadata, NeuromodulatorSnapshot, 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") + } + + fn bind_loopback(channels: usize) -> (::zmq::Socket, ZmqStimulusSource) { + let context = ::zmq::Context::new(); + let publisher = context.socket(::zmq::PUB).unwrap(); + publisher.bind("tcp://127.0.0.1:*").unwrap(); + let endpoint = match publisher.get_last_endpoint().unwrap() { + Ok(ep) => ep, + Err(bytes) => String::from_utf8(bytes).expect("endpoint utf8"), + }; + + 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(packet)) if packet.rejected => { + return last; + } + Ok(Some(_)) => return last, + Ok(None) => std::thread::sleep(Duration::from_millis(10)), + } + } + last + } + + #[test] + fn typed_json_frame_round_trips_over_zmq() { + let (publisher, mut source) = bind_loopback(4); + let frame = sample_frame(4, 100); + publisher.send(&frame, 0).unwrap(); + + let packet = poll_ingress(&mut source) + .expect("recv") + .expect("ZMQ SUB should receive a typed StimulusBatch"); + assert!(!packet.rejected); + 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); + assert_eq!(packet.session_id.as_deref(), Some("zmq-loopback")); + } + + #[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(); + + let ping_result = poll_ingress(&mut source); + assert!( + ping_result.is_ok(), + "rejected Ping must not be a hard receive failure: {ping_result:?}" + ); + let ping_packet = ping_result.expect("ok").expect("rejection is signaled"); + assert!( + ping_packet.rejected, + "Ping must set rejected so RuntimeStats can count it" + ); + assert!(ping_packet.stimuli.is_empty()); + + let frame = sample_frame(4, 101); + publisher.send(&frame, 0).unwrap(); + let packet = (0..50) + .find_map(|_| match source.next_ingress() { + Ok(Some(packet)) if !packet.rejected => Some(packet), + Ok(_) => { + std::thread::sleep(Duration::from_millis(10)); + None + } + Err(err) => panic!("ingress error: {err}"), + }) + .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(); + + // Modulation-only frames are drained; after EAGAIN the held snapshot + // is returned so idle ticks keep DA/cortisol/ACh/tempo. + let packet = (0..50) + .find_map(|_| match source.next_ingress() { + Ok(Some(packet)) if packet.modulators.is_some() => Some(packet), + Ok(_) => { + std::thread::sleep(Duration::from_millis(10)); + None + } + Err(err) => panic!("ingress error: {err}"), + }) + .expect("ZMQ SUB should hold Neuromodulators after drain"); + assert_eq!( + packet.modulators.as_deref(), + Some(&[0.4, 0.3, 0.2, 1.0][..]) + ); + assert!(packet.stimuli.is_empty()); + + 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!(!idle.rejected); + assert_eq!(idle.modulators.as_deref(), Some(&[0.4, 0.3, 0.2, 1.0][..])); + } + } } #[cfg(feature = "corpus-ipc")] diff --git a/src/checkpoint/mod.rs b/src/checkpoint/mod.rs index 14969e8..592ebb1 100644 --- a/src/checkpoint/mod.rs +++ b/src/checkpoint/mod.rs @@ -113,11 +113,16 @@ fn require_live_izh_count(config: &DaemonConfig) -> Result<()> { /// Accepts a JSON file, a Hugging Face hub `config.json`, or a directory /// containing `snn_model.json` or `dataset/merged_v2/snn_model.json`. pub fn resolve_sidecar_path(model_path: &Path) -> Result { - if !model_path.exists() { - bail!( + match model_path.try_exists() { + Ok(true) => {} + Ok(false) => bail!( "live mode cannot enter the tick loop: checkpoint not found at {}", model_path.display() - ); + ), + Err(err) => bail!( + "live mode cannot enter the tick loop: failed to query checkpoint {}: {err}", + model_path.display() + ), } if model_path.is_dir() { diff --git a/src/daemon.rs b/src/daemon.rs index 0910107..bf63bb1 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -101,6 +101,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`. @@ -238,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); @@ -266,6 +263,7 @@ impl BrainstemDaemon { let mut stimuli = vec![0.0; cfg.channels]; let mut spike_buf: Vec = Vec::with_capacity(128); + let mut stats = RuntimeStats::default(); let mut shutdown = std::pin::pin!(shutdown_signal()); @@ -279,7 +277,10 @@ impl BrainstemDaemon { &mut stimuli, &mut spike_buf, &ingress, - &health, + &mut TickReport { + health: &health, + stats: &mut stats, + }, ); } _ = &mut shutdown => { @@ -295,6 +296,42 @@ impl BrainstemDaemon { shutdown_backend(&mut backend); Ok(()) } + + /// Drive a bounded number of ticks without waiting for a termination signal. + /// + /// Used by the CPU-only Thalamic → corpus-ipc → Brainstem smoke harness. + /// Calls [`StimulusSource::initialize`] (same as [`Self::run`]) so a ZMQ + /// source is connected once. Source shutdown always runs, even when sink + /// flush fails. Does not bind `control_bind`. + pub fn run_for_ticks(self, ticks: u64) -> Result { + let cfg = self.config; + let mut backend = self.backend; + let ingress = self.ingress; + let health = self.health; + let mut stats = RuntimeStats::default(); + + 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); + drive_ticks( + ticks, + &mut backend, + &mut network, + &ingress, + &health, + &mut stats, + cfg.channels, + ); + finish_bounded_run(&mut backend, &ingress)?; + Ok(stats) + } } fn shutdown_backend(backend: &mut BackendPair) { @@ -309,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)] @@ -433,6 +515,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 @@ -580,6 +685,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, @@ -587,13 +698,15 @@ fn run_tick( stimuli: &mut [f32], spike_buf: &mut Vec, ingress: &BoundedIngress, - health: &HealthHandle, + report: &mut TickReport<'_>, ) { + let TickReport { health, stats } = report; let backend_packet = match source.next_ingress() { Ok(Some(p)) => Some(p), Ok(None) => None, Err(e) => { warn!("Failed to receive from stimulus source: {e}"); + stats.rejected_batches += 1; None } }; @@ -607,6 +720,13 @@ fn run_tick( // `None` (skip or error) does not enqueue a placeholder that could evict // in-process sensory. decode_inputs zero-fills when drain yields no stimuli. if let Some(packet) = backend_packet { + if packet.rejected { + stats.rejected_batches += 1; + } else if packet.batch_id.is_some() { + stats.accepted_batches += 1; + stats.last_batch_id = packet.batch_id; + stats.last_valid_mask.clone_from(&packet.valid_mask); + } ingress.admit_backend_packet(packet); } let drained = ingress.drain_for_tick(); @@ -631,6 +751,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() @@ -711,6 +832,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() >= NEUROMODULATOR_COUNT => { @@ -740,7 +868,18 @@ pub(crate) fn run_tick_for_test( ) { let ingress = BoundedIngress::new(IngressConfig::default()).expect("default ingress"); let health = HealthHandle::started(HealthLimits::default()); - run_tick(source, network, sink, stimuli, spike_buf, &ingress, &health); + run_tick( + source, + network, + sink, + stimuli, + spike_buf, + &ingress, + &mut TickReport { + health: &health, + stats: &mut RuntimeStats::default(), + }, + ); } #[cfg(test)] @@ -773,6 +912,7 @@ mod tests { IngressPacket { stimuli: vec![tag], modulators: None, + ..Default::default() } } @@ -843,6 +983,7 @@ control_bind = "127.0.0.1:9464" let packet = IngressPacket { stimuli: vec![0.1, 0.2, 0.3, 0.4], modulators: None, + ..Default::default() }; let mut stimuli = vec![0.0; 4]; let _mods = decode_inputs(&packet, &mut stimuli); @@ -854,6 +995,7 @@ control_bind = "127.0.0.1:9464" let packet = IngressPacket { stimuli: vec![0.0; 4], modulators: Some(vec![0.5, 0.6, 0.7, 0.8]), + ..Default::default() }; let mut stimuli = vec![0.0; 4]; let mods = decode_inputs(&packet, &mut stimuli); @@ -868,6 +1010,7 @@ control_bind = "127.0.0.1:9464" let packet = IngressPacket { stimuli: vec![0.0; 2], modulators: Some(vec![0.1, 0.2, 0.3, 0.4, 0.9]), + ..Default::default() }; let mut stimuli = vec![0.0; 2]; let mods = decode_inputs(&packet, &mut stimuli); @@ -887,6 +1030,7 @@ control_bind = "127.0.0.1:9464" let packet = IngressPacket { stimuli: vec![0.1, 0.2], modulators: None, + ..Default::default() }; let mut stimuli = vec![0.0; 4]; let mods = decode_inputs(&packet, &mut stimuli); @@ -894,6 +1038,21 @@ control_bind = "127.0.0.1:9464" 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), + ..Default::default() + }; + 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(); @@ -1117,7 +1276,10 @@ channels = 16 &mut stimuli, &mut spike_buf, &ingress, - &health, + &mut TickReport { + health: &health, + stats: &mut RuntimeStats::default(), + }, ); assert_eq!(sink.emitted.len(), 1); @@ -1172,6 +1334,7 @@ channels = 16 let packet = IngressPacket { stimuli: vec![0.0; 2], modulators: Some(vec![0.5, 0.25, 0.8, 0.1]), + ..Default::default() }; let sink = tick_once(packet, &mut network); @@ -1202,6 +1365,7 @@ channels = 16 let packet = IngressPacket { stimuli: vec![0.2, 0.3], modulators: None, + ..Default::default() }; let sink = tick_once(packet, &mut network); @@ -1222,6 +1386,7 @@ channels = 16 let packet = IngressPacket { stimuli: vec![0.4, 0.1], modulators: Some(vec![0.0, 0.0, 0.0, 0.0]), + ..Default::default() }; let mut source = ScriptedStimulusSource { packet }; @@ -1384,6 +1549,7 @@ block_timeout_ms = 0 IngressPacket { stimuli: vec![9.0], modulators: None, + ..Default::default() }, ); @@ -1393,6 +1559,7 @@ block_timeout_ms = 0 Ok(Some(IngressPacket { stimuli: vec![0.1, 0.2], modulators: Some(vec![0.5, 0.0, 0.0, 0.0]), + ..Default::default() })) } @@ -1414,7 +1581,10 @@ block_timeout_ms = 0 &mut stimuli, &mut spike_buf, &ingress, - &health, + &mut super::TickReport { + health: &health, + stats: &mut RuntimeStats::default(), + }, ); assert_eq!(stimuli, vec![0.1, 0.2]); @@ -1435,6 +1605,7 @@ block_timeout_ms = 0 IngressPacket { stimuli: vec![0.3, 0.4], modulators: None, + ..Default::default() }, ); @@ -1462,7 +1633,10 @@ block_timeout_ms = 0 &mut stimuli, &mut spike_buf, &ingress, - &health, + &mut super::TickReport { + health: &health, + stats: &mut RuntimeStats::default(), + }, ); assert_eq!(stimuli, vec![0.3, 0.4]); diff --git a/src/ingress/corpus.rs b/src/ingress/corpus.rs new file mode 100644 index 0000000..67cc17d --- /dev/null +++ b/src/ingress/corpus.rs @@ -0,0 +1,320 @@ +// 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, Validate}; + +/// 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 { + /// Expected stimulus width. `0` means unspecified: width is not checked. + pub expected_channels: usize, + pub max_age: Option, +} + +impl IngressPolicy { + /// Policy for a configured channel width and optional freshness window. + /// + /// Pass `expected_channels = 0` to accept any width (unspecified-width mode). + 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 }, + Future { timestamp_ns: u64, now_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::Future { + timestamp_ns, + now_ns, + } => { + write!( + f, + "stimulus timestamp {timestamp_ns} is in the future (now {now_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(), + // Positional map from corpus-ipc's DA/cortisol/ACh/tempo snapshot + // onto neuromod 0.6's DA/5-HT/ACh/NE slots. + modulators: Some(vec![ + snapshot.dopamine, + snapshot.cortisol, + snapshot.acetylcholine, + snapshot.tempo, + ]), + ..IngressPacket::default() + }), + 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 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, + now_ns: u64, +) -> Result { + batch + .validate() + .map_err(|err| IngressError::Stimulus(err.to_string()))?; + 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(), + }); + } + 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`. + Ok(IngressPacket { + stimuli: batch.values, + modulators: None, + valid_mask: batch.valid_mask, + batch_id: Some(batch.batch_id), + timestamp_ns: Some(batch.timestamp), + session_id: batch.session_id, + rejected: false, + }) +} + +#[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])); + // Masked channel 1 keeps the producer placeholder; decode_inputs zeros it. + assert_eq!(packet.stimuli, vec![1.0, 0.0, 0.25, 0.5]); + assert_eq!(packet.batch_id, Some(7)); + assert_eq!(packet.session_id.as_deref(), Some("smoke")); + } + + #[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 unspecified_width_accepts_any_channel_count() { + let packet = accept_ipc_message( + IpcMessage::Stimuli(sample_batch()), + &IngressPolicy::new(0, Some(Duration::from_secs(1))), + 1_000, + ) + .unwrap(); + assert_eq!(packet.stimuli.len(), 4); + } + + #[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 future_timestamp_fails() { + let mut batch = sample_batch(); + batch.timestamp = 5_000; + let err = accept_ipc_message(IpcMessage::Stimuli(batch), &policy(), 1_000).unwrap_err(); + assert!(matches!( + err, + IngressError::Future { + timestamp_ns: 5_000, + now_ns: 1_000 + } + )); + } + + #[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/ingress/mod.rs b/src/ingress/mod.rs index 8421b31..3c895e6 100644 --- a/src/ingress/mod.rs +++ b/src/ingress/mod.rs @@ -17,9 +17,17 @@ mod queue; +#[cfg(feature = "corpus-ipc")] +mod corpus; + #[cfg(test)] mod tests; +#[cfg(feature = "corpus-ipc")] +pub use corpus::{ + IngressError, IngressPolicy, STIMULUS_SCHEMA, accept_ipc_json, accept_ipc_message, +}; + use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; @@ -343,9 +351,15 @@ impl DrainedTick { /// control actuator. Callers must inspect `control` (the tick loop logs it) /// so those events are observed rather than starved behind bulk. pub fn into_packet(self) -> IngressPacket { + let sensory = self.sensory.unwrap_or_default(); IngressPacket { - stimuli: self.sensory.map(|p| p.stimuli).unwrap_or_default(), + stimuli: sensory.stimuli, modulators: self.reward.and_then(|p| p.modulators), + valid_mask: sensory.valid_mask, + batch_id: sensory.batch_id, + timestamp_ns: sensory.timestamp_ns, + session_id: sensory.session_id, + rejected: sensory.rejected, } } } @@ -427,6 +441,11 @@ impl BoundedIngress { let IngressPacket { stimuli, modulators, + valid_mask, + batch_id, + timestamp_ns, + session_id, + rejected, } = packet; if !stimuli.is_empty() { let _ = self.try_enqueue( @@ -434,6 +453,11 @@ impl BoundedIngress { IngressPacket { stimuli, modulators: None, + valid_mask, + batch_id, + timestamp_ns, + session_id, + rejected, }, ); } @@ -443,6 +467,7 @@ impl BoundedIngress { IngressPacket { stimuli: Vec::new(), modulators: Some(modulators), + ..IngressPacket::default() }, ); } diff --git a/src/ingress/tests.rs b/src/ingress/tests.rs index fb9fff5..3216b69 100644 --- a/src/ingress/tests.rs +++ b/src/ingress/tests.rs @@ -29,6 +29,7 @@ fn pkt(tag: f32) -> IngressPacket { IngressPacket { stimuli: vec![tag], modulators: None, + ..Default::default() } } @@ -36,6 +37,7 @@ fn reward_pkt(tag: f32) -> IngressPacket { IngressPacket { stimuli: Vec::new(), modulators: Some(vec![tag, 0.0, 0.0, 0.0]), + ..Default::default() } } @@ -390,6 +392,7 @@ fn oversize_payload_is_rejected() { let big = IngressPacket { stimuli: vec![1.0, 2.0, 3.0], modulators: None, + ..Default::default() }; assert_eq!( ingress.enqueue(MessageClass::Sensory, big), @@ -405,6 +408,7 @@ fn admit_backend_packet_splits_stimuli_and_modulators() { ingress.admit_backend_packet(IngressPacket { stimuli: vec![0.5, 0.25], modulators: Some(vec![1.0, 2.0, 3.0, 4.0]), + ..Default::default() }); let drained = ingress.drain_for_tick(); let packet = drained.into_packet(); @@ -428,6 +432,7 @@ fn admit_empty_stimuli_still_enqueues_modulators() { ingress.admit_backend_packet(IngressPacket { stimuli: Vec::new(), modulators: Some(vec![1.0, 2.0, 3.0, 4.0]), + ..Default::default() }); assert_eq!(ingress.metrics().sensory.depth, 0); assert_eq!(ingress.metrics().reward.depth, 1); diff --git a/src/lib.rs b/src/lib.rs index 149ba29..5b8f5bf 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -14,10 +14,11 @@ pub mod runtime; // Re-export the new pluggable I/O surface (pub from day one). pub use backend::{ - BackendPair, IngressPacket, NEUROMODULATOR_COUNT, SpikeEvent, SpikeSink, StimulusSource, + BackendPair, CollectingSpikeSink, IngressPacket, NEUROMODULATOR_COUNT, SpikeEvent, SpikeSink, + StimulusSource, }; pub use checkpoint::{ModelProvenance, restore_network}; -pub use daemon::RuntimeMode; +pub use daemon::{BrainstemDaemon, DaemonConfig, RuntimeMode, RuntimeStats}; pub use health::{ CheckpointIdentity, FakeClock, FatalCode, HealthEvent, HealthHandle, HealthLimits, HealthMachine, HealthPhase, HealthSnapshot, ReasonCode, SystemClock, diff --git a/tests/fixtures/thalamic_producer.rs b/tests/fixtures/thalamic_producer.rs new file mode 100644 index 0000000..1ef4de9 --- /dev/null +++ b/tests/fixtures/thalamic_producer.rs @@ -0,0 +1,109 @@ +// 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, Validate}; + +/// Simulated Thalamic producer: sensory + safety only. +#[derive(Debug)] +pub struct ThalamicProducer { + /// Deterministic hardware-safety path. Independent of IPC/Brainstem. + pub safety_healthy: bool, + pub published: u64, + pub publish_errors: u64, +} + +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 { + 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; + } + } + } + + /// 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..c13f95a --- /dev/null +++ b/tests/thalamic_brainstem_smoke.rs @@ -0,0 +1,317 @@ +// 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::path::{Path, PathBuf}; +use std::time::Duration; + +use anyhow::Result; +use brainstem_daemon::daemon::{BrainstemDaemon, DaemonConfig, RuntimeMode}; +use brainstem_daemon::ingress::{IngressConfig, 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 TempDir(PathBuf); + +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_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) + } + + fn path(&self) -> &Path { + &self.0 + } +} + +impl Drop for TempDir { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.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_sidecar_json() -> String { + r#"{ + "source": "spikenaut_julia", + "encoder": "v3_state_telemetry", + "q88": "signed", + "frozen_lineage": "thalamic-smoke", + "neurons": [ + { + "decay_rate": 0.85, + "membrane_potential": 0.0, + "threshold": 0.02, + "last_spike": false, + "weights": [2.0, 0.12, 0.13, 0.14] + }, + { + "decay_rate": 0.85, + "membrane_potential": 0.0, + "threshold": 0.02, + "last_spike": false, + "weights": [0.15, 0.16, 0.17, 0.18] + }, + { + "decay_rate": 0.85, + "membrane_potential": 0.0, + "threshold": 0.02, + "last_spike": false, + "weights": [0.19, 0.20, 0.21, 0.22] + }, + { + "decay_rate": 0.85, + "membrane_potential": 0.0, + "threshold": 0.02, + "last_spike": false, + "weights": [0.23, 0.24, 0.25, 0.26] + } + ] + }"# + .to_string() +} + +fn write_smoke_sidecar(dir: &Path) -> PathBuf { + let path = dir.join("snn_model.json"); + std::fs::write(&path, smoke_sidecar_json()).expect("write sidecar"); + path +} + +fn smoke_config(model_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, + runtime_mode: RuntimeMode::Live, + services: Vec::new(), + ingress: IngressConfig::default(), + control_bind: None, + } +} + +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 = TempDir::new("brainstem-smoke-ckpt"); + let path = write_smoke_sidecar(dir.path()); + + 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), pair) + .unwrap() + .run_for_ticks(1) + .unwrap(); + + let loaded = stats.loaded_checkpoint.expect("checkpoint must be loaded"); + assert_eq!(loaded.model_id, "spikenaut-snn:thalamic-smoke"); + assert_eq!(loaded.schema_id, "spikenaut-sidecar-json/v1"); +} + +#[test] +fn typed_frame_crosses_ipc_and_produces_runtime_result() { + let dir = TempDir::new("brainstem-smoke-e2e"); + let path = write_smoke_sidecar(dir.path()); + + 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); + assert_eq!(packet.session_id.as_deref(), Some("thalamic-smoke")); + + 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()); +} + +#[test] +fn rejected_frame_is_counted_without_stopping_the_loop() { + let dir = TempDir::new("brainstem-smoke-reject"); + let path = write_smoke_sidecar(dir.path()); + let source = QueuedStimulusSource { + packets: VecDeque::from([Ok(IngressPacket { + rejected: true, + ..IngressPacket::default() + })]), + }; + 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.rejected_batches, 1); + assert_eq!(stats.accepted_batches, 0); +} + +#[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); + + // A later successful-looking publish must not clobber a thermal fault. + thalamic.safety_tick(false); + let mut transport = CaptureTransport { frames: Vec::new() }; + thalamic.publish(Some(&mut transport), &bytes); + assert!(!thalamic.safety_healthy); + assert_eq!(thalamic.published, 1); + + thalamic.safety_tick(true); + assert!(thalamic.safety_healthy); + let _still_producing = thalamic.simulate_telemetry(10, CHANNELS); +}