diff --git a/AGENTS.md b/AGENTS.md index c0ec33c..a099626 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -89,6 +89,10 @@ tick_rate_hz = 1000 log_level = "info" spine_sub_port = 5555 spine_pub_port = 5556 +spine_pub_bind_host = "127.0.0.1" # PUB bind host; loopback default, "0.0.0.0" to expose broadly +spine_pub_sndhwm = 1000 # explicit finite send high-water mark +spine_pub_linger_ms = 0 # finite LINGER: drop pending on close (bounded teardown) +spine_pub_send_empty_batches = true # send a frame every tick even with zero spikes model_path = "/var/lib/soma/snn_model.json" # Distill sidecar JSON; `~` is not expanded ``` diff --git a/README.md b/README.md index 5a81772..4d682be 100644 --- a/README.md +++ b/README.md @@ -158,6 +158,11 @@ log_level = "info" # error|warn|info|debug|trace # ZMQ (still required in TOML; no-ops under the default stub backend) spine_sub_port = 5555 # stimuli in spine_pub_port = 5556 # spikes out +# Built-in PUB lifecycle (all optional; corpus-ipc only, no-ops under stub). +# spine_pub_bind_host = "127.0.0.1" # loopback default; use "0.0.0.0" to expose broadly +# spine_pub_sndhwm = 1000 # explicit finite send high-water mark +# spine_pub_linger_ms = 0 # finite LINGER: drop pending on close (bounded teardown) +# spine_pub_send_empty_batches = true # send a frame every tick even with zero spikes # Service registry (optional; empty by default) # Trading/mining-specific adapters are intentionally excluded from defaults. @@ -201,7 +206,7 @@ Default Cargo features are empty (`default = []` in `Cargo.toml`). That path use | Cargo flags | Wired backend | `libzmq` | Binary (`brainstem-daemon`) | Library `BrainstemDaemon::new()` / `try_new()` | |---|---|---|---|---| | default / `--no-default-features` | stub | not required | no backend sockets by default; `control_bind` opens the control listener; logs `🔌 Using stub backend` | stub | -| `--features corpus-ipc` | ZMQ / `corpus-ipc` | required | SUB via env, PUB on `spine_pub_port`; logs `📡 Using ZMQ corpus-ipc backend` | **still stub** | +| `--features corpus-ipc` | ZMQ / `corpus-ipc` | required | SUB via env, PUB on `spine_pub_bind_host:spine_pub_port` (loopback by default); logs `📡 Using ZMQ corpus-ipc backend` | **still stub** | | `--all-features` | same as `corpus-ipc` | required | same as `--features corpus-ipc` | **still stub** | 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. @@ -222,7 +227,11 @@ Health snapshots, probe paths, and the transition table live in [`docs/health.md | `control_bind` | optional HTTP control surface; unset = no listener | same | | `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://*:` | +| `spine_pub_port` | parsed, **no-op** | binds ZMQ PUB `tcp://:` (loopback by default) | +| `spine_pub_bind_host` | parsed, **no-op** | PUB bind host; default `127.0.0.1` (loopback). Set `0.0.0.0` (or a specific interface) to opt in to broader exposure | +| `spine_pub_sndhwm` | parsed, **no-op** | explicit finite send high-water mark applied to the PUB socket before bind (default `1000`) | +| `spine_pub_linger_ms` | parsed, **no-op** | finite LINGER (ms) applied before bind so teardown/shutdown cannot block (default `0` = drop pending on close) | +| `spine_pub_send_empty_batches` | parsed, **no-op** | empty-batch policy; `true` (default) sends a frame every tick even with zero spikes, `false` suppresses empty batches at the sink | | `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`): @@ -238,6 +247,14 @@ Health snapshots, probe paths, and the transition table live in [`docs/health.md 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. +##### PUB lifecycle, bind default, and loss semantics + +The built-in `corpus-ipc` PUB egress is explicitly bounded. Before bind it sets a finite send high-water mark (`spine_pub_sndhwm`) and a finite LINGER (`spine_pub_linger_ms`, default `0` so close/teardown drops pending messages instead of blocking), then binds to `tcp://:` — loopback (`127.0.0.1`) by default, so broader exposure such as `0.0.0.0` requires explicit configuration. + +Empty-batch policy: `spine_pub_send_empty_batches` (default `true`) makes the sink publish a frame every tick even when the batch has zero spikes, preserving the current behavior for subscribers that rely on per-tick / heartbeat frames. Set it to `false` to suppress empty batches at the sink. + +The sink accounts for three **application-side** send outcomes it can directly observe: `attempted` (a frame was handed to ZeroMQ), `suppressed` (an empty batch withheld under the policy), and `failed` (`send()` returned an error). It does **not** and cannot measure subscriber-side delivery: ZeroMQ PUB is best-effort and silently drops messages for subscribers that are slow, absent, or over `SNDHWM`. That subscriber-side loss is inherently not observable by the publisher, so it is documented as best-effort and never reported as a measured count. + 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 diff --git a/src/backend.rs b/src/backend.rs index d15b754..ab50bbb 100644 --- a/src/backend.rs +++ b/src/backend.rs @@ -363,23 +363,93 @@ mod zmq_impl { } unsafe impl Send for SafeSocket {} + /// Application-side send counters observed by [`ZmqSpikeSink`]. + /// + /// These count only what the publisher itself can observe at the moment it + /// hands a batch to (or withholds it from) ZeroMQ. They do NOT and cannot + /// measure subscriber-side PUB loss (see [`ZmqSpikeSink`] docs). + #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] + pub struct SinkSendStats { + /// `send()` was actually called on the socket (a frame was handed to ZeroMQ). + pub attempted: u64, + /// An empty batch was withheld because the empty-batch policy suppresses empties. + pub suppressed: u64, + /// `send()` returned an error. + pub failed: u64, + } + + /// ZeroMQ PUB spike egress. + /// + /// # Loss semantics (read carefully) + /// + /// This sink accounts for exactly three **application-side** outcomes that + /// the publisher can directly observe, exposed via [`Self::stats`]: + /// + /// - `attempted`: `send()` was called (a frame was handed to ZeroMQ); + /// - `suppressed`: an empty batch was withheld under the empty-batch policy + /// (`send_empty_batches == false`), so no frame was handed to ZeroMQ; + /// - `failed`: `send()` returned an error. + /// + /// It does **not** measure subscriber-side delivery. ZeroMQ PUB is + /// best-effort: it silently drops messages for subscribers that are slow, + /// absent, or over the send high-water mark (SNDHWM). That subscriber-side + /// loss is **inherently not observable by the publisher** — an `attempted` + /// send that returns `Ok` says the frame was accepted into ZeroMQ's egress, + /// not that any subscriber received it. Do not read these counters as a + /// delivery guarantee or a measure of dropped-at-subscriber messages. pub struct ZmqSpikeSink { socket: std::sync::Mutex, /// Reusable buffer to convert to corpus-ipc event type without allocating every tick. corpus_buf: Vec, + /// When `false`, empty spike batches are suppressed (not sent) and + /// counted as `suppressed`. Defaults to `true` (send empty batches), + /// preserving the historical always-send behavior. + send_empty_batches: bool, + /// Application-side send accounting (see [`SinkSendStats`]). + stats: SinkSendStats, } impl ZmqSpikeSink { + /// Construct a sink that always sends, including empty batches. + /// + /// Backward-compatible with existing callers/tests that relied on the + /// original always-send behavior. pub fn new(socket: ::zmq::Socket) -> Self { + Self::with_policy(socket, true) + } + + /// Construct a sink with an explicit empty-batch policy. + /// + /// When `send_empty_batches` is `false`, `emit(&[])` is suppressed (no + /// frame is sent) and counted as `suppressed`. + pub fn with_policy(socket: ::zmq::Socket, send_empty_batches: bool) -> Self { Self { socket: std::sync::Mutex::new(SafeSocket { socket }), corpus_buf: Vec::new(), + send_empty_batches, + stats: SinkSendStats::default(), } } + + /// Snapshot of the application-side send counters. + /// + /// See the type-level docs: these reflect only attempted / suppressed / + /// failed sends the publisher can observe, never subscriber-side PUB loss. + pub fn stats(&self) -> SinkSendStats { + self.stats + } } impl SpikeSink for ZmqSpikeSink { fn emit(&mut self, spikes: &[SpikeEvent], batch_time: std::time::Duration) -> Result<()> { + // Documented empty-batch suppression policy: when configured to + // suppress empties, an empty batch is not handed to ZeroMQ. This is + // an application-side decision the publisher fully observes. + if spikes.is_empty() && !self.send_empty_batches { + self.stats.suppressed += 1; + return Ok(()); + } + // Use the tick-level timestamp passed by run_tick so batch metadata // stays aligned with the SpikeEvent.time values in this batch. let batch_id = batch_time.as_millis() as u64; @@ -413,7 +483,16 @@ mod zmq_impl { .socket .lock() .map_err(|_| anyhow::anyhow!("ZMQ socket mutex poisoned"))?; - guard.socket.send(payload, 0)?; + // `send()` returning Ok means the frame was accepted into ZeroMQ's + // egress, NOT that any subscriber received it. Subscriber-side drops + // (slow/absent/over-SNDHWM) are silent and unobservable here. + self.stats.attempted += 1; + if let Err(e) = guard.socket.send(payload, 0) { + self.stats.failed += 1; + // Preserve existing behavior: propagate the error so run_tick's + // emit_errors path and fatal handling still apply. + return Err(anyhow::anyhow!("ZMQ PUB send failed: {e}")); + } Ok(()) } } @@ -628,8 +707,247 @@ mod zmq_impl { "the non-finite modulator frame must be surfaced as a rejected/skip outcome" ); } + + // ───────────────────────────────────────────────────────────────── + // PUB / ZmqSpikeSink egress tests (issue #68 / LIM-1320). + // + // These mirror the SUB-side `bind_loopback`/`get_last_endpoint` + // pattern but for the sink side: the sink owns a PUB socket bound on + // `tcp://127.0.0.1:*`, and a SUB receiver connects to that endpoint. + + /// Build a `ZmqSpikeSink` over a PUB socket bound on loopback with the + /// given empty-batch policy, and return the sink plus the resolved + /// endpoint so tests can attach a subscriber. + fn bind_sink(send_empty_batches: bool) -> (ZmqSpikeSink, String) { + bind_sink_with_hwm(send_empty_batches, 1000) + } + + /// Like [`bind_sink`] but with an explicit send high-water mark, so + /// tests can force the over-SNDHWM drop path with a small bound. + fn bind_sink_with_hwm(send_empty_batches: bool, sndhwm: i32) -> (ZmqSpikeSink, String) { + let context = ::zmq::Context::new(); + let pub_socket = context.socket(::zmq::PUB).unwrap(); + // Apply the same bounded-lifecycle options the binary sets. + pub_socket.set_sndhwm(sndhwm).unwrap(); + pub_socket.set_linger(0).unwrap(); + pub_socket.bind("tcp://127.0.0.1:*").unwrap(); + let endpoint = match pub_socket.get_last_endpoint().unwrap() { + Ok(ep) => ep, + Err(bytes) => String::from_utf8(bytes).expect("endpoint utf8"), + }; + let sink = ZmqSpikeSink::with_policy(pub_socket, send_empty_batches); + (sink, endpoint) + } + + fn connect_subscriber(endpoint: &str) -> ::zmq::Socket { + let context = ::zmq::Context::new(); + let sub = context.socket(::zmq::SUB).unwrap(); + sub.set_subscribe(b"").unwrap(); + sub.set_rcvhwm(16).unwrap(); + sub.connect(endpoint).unwrap(); + // Allow the PUB/SUB connection to settle before publishing. + std::thread::sleep(Duration::from_millis(150)); + sub + } + + fn sample_spikes() -> Vec { + vec![ + SpikeEvent { + channel: 0, + time: 7, + strength: 1.0, + }, + SpikeEvent { + channel: 3, + time: 9, + strength: 0.5, + }, + ] + } + + /// Poll a SUB socket for one frame, tolerating startup slop. + fn recv_frame(sub: &::zmq::Socket) -> Option> { + for _ in 0..50 { + match sub.recv_bytes(::zmq::DONTWAIT) { + Ok(buf) => return Some(buf), + Err(::zmq::Error::EAGAIN) => { + std::thread::sleep(Duration::from_millis(10)); + } + Err(e) => panic!("SUB recv failed: {e}"), + } + } + None + } + + #[test] + fn no_subscriber_emit_is_ok_and_counts_attempted() { + // With no connected subscriber, PUB does not error; the frame is + // silently dropped by ZeroMQ. That drop is NOT observable here — we + // can only assert the application-side `attempted` count. + let (mut sink, _endpoint) = bind_sink(true); + for _ in 0..3 { + sink.emit(&sample_spikes(), Duration::from_millis(1)) + .expect("emit with no subscriber must be Ok (loss is silent/unobservable)"); + } + let stats = sink.stats(); + assert_eq!(stats.attempted, 3); + assert_eq!(stats.failed, 0); + assert_eq!(stats.suppressed, 0); + } + + #[test] + fn empty_batch_policy_sends_when_enabled() { + let (mut sink, endpoint) = bind_sink(true); + let sub = connect_subscriber(&endpoint); + + sink.emit(&[], Duration::from_millis(2)) + .expect("empty emit must be Ok when send_empty_batches=true"); + let stats = sink.stats(); + assert_eq!(stats.attempted, 1); + assert_eq!(stats.suppressed, 0); + + let buf = recv_frame(&sub).expect("subscriber must receive the empty-batch frame"); + let msg: IpcMessage = serde_json::from_slice(&buf).expect("deserialize IpcMessage"); + match msg { + IpcMessage::Spikes(batch) => assert!(batch.spikes.is_empty()), + other => panic!("expected Spikes, got {other:?}"), + } + } + + #[test] + fn empty_batch_policy_suppresses_when_disabled() { + let (mut sink, endpoint) = bind_sink(false); + let sub = connect_subscriber(&endpoint); + + sink.emit(&[], Duration::from_millis(2)) + .expect("suppressed empty emit still returns Ok"); + let stats = sink.stats(); + assert_eq!(stats.suppressed, 1); + assert_eq!(stats.attempted, 0); + assert_eq!(stats.failed, 0); + + assert!( + recv_frame(&sub).is_none(), + "no frame must be delivered when the empty batch is suppressed" + ); + + // A non-empty batch still sends under the suppress-empties policy. + sink.emit(&sample_spikes(), Duration::from_millis(3)) + .expect("non-empty emit must send under suppress-empties policy"); + assert_eq!(sink.stats().attempted, 1); + assert!( + recv_frame(&sub).is_some(), + "the non-empty batch must still be delivered" + ); + } + + #[test] + fn slow_or_disconnected_subscriber_drops_are_not_observable() { + // Connect a subscriber, then drop it, then keep emitting. The + // publisher cannot observe the subscriber-side drops: every emit + // still returns Ok and only `attempted` advances. + let (mut sink, endpoint) = bind_sink(true); + let sub = connect_subscriber(&endpoint); + drop(sub); + + for _ in 0..5 { + sink.emit(&sample_spikes(), Duration::from_millis(1)) + .expect("emit to a disconnected subscriber must remain Ok"); + } + let stats = sink.stats(); + assert_eq!( + stats.attempted, 5, + "only application-side attempts are counted" + ); + assert_eq!(stats.failed, 0); + assert_eq!(stats.suppressed, 0); + } + + #[test] + fn slow_connected_subscriber_over_sndhwm_drops_are_not_observable() { + // A subscriber CONNECTS but then STOPS READING while the publisher + // emits far more than SNDHWM batches. Once the per-subscriber send + // queue fills, ZeroMQ silently drops on the PUB side. That loss is + // NOT observable to the publisher: every emit still returns Ok and + // only `attempted` advances (failed == 0, suppressed == 0). + let (mut sink, endpoint) = bind_sink_with_hwm(true, 2); + // Connect and subscribe, then never call recv on it. + let _slow_sub = connect_subscriber(&endpoint); + + // Emit well beyond SNDHWM (2) so the send queue is guaranteed to + // overflow. Bounded, deterministic loop; no timing dependence. + let emit_count = 50; + for _ in 0..emit_count { + sink.emit(&sample_spikes(), Duration::from_millis(1)) + .expect("emit to a slow (connected, non-reading) subscriber must remain Ok"); + } + + let stats = sink.stats(); + assert_eq!( + stats.attempted, emit_count, + "every application-side attempt is counted even when ZeroMQ drops over SNDHWM" + ); + assert_eq!( + stats.failed, 0, + "over-HWM drops are silent on the PUB side; emit must not report failure" + ); + assert_eq!(stats.suppressed, 0); + } + + #[test] + fn sink_teardown_is_bounded_by_finite_linger() { + // With a finite LINGER (0), constructing, emitting with no + // subscriber, then dropping the sink must complete promptly and not + // hang on pending PUB messages. This is the observable, testable + // proxy for signal-driven teardown not blocking on egress. + let (mut sink, _endpoint) = bind_sink(true); + for _ in 0..10 { + sink.emit(&sample_spikes(), Duration::from_millis(1)) + .expect("emit before teardown"); + } + let start = std::time::Instant::now(); + drop(sink); + assert!( + start.elapsed() < Duration::from_secs(2), + "finite LINGER must bound sink teardown; took {:?}", + start.elapsed() + ); + } + + #[test] + fn wire_payload_is_preserved_round_trip() { + // Prove the unversioned tagged-JSON payload is unchanged: emit one + // non-empty batch and deserialize the received bytes back into + // corpus_ipc::IpcMessage, asserting Spikes with the expected + // batch_id/timestamp derived from batch_time and the same spikes. + let (mut sink, endpoint) = bind_sink(true); + let sub = connect_subscriber(&endpoint); + + let batch_time = Duration::from_millis(1234); + let spikes = sample_spikes(); + sink.emit(&spikes, batch_time) + .expect("emit non-empty batch"); + + let buf = recv_frame(&sub).expect("subscriber must receive the frame"); + let msg: IpcMessage = serde_json::from_slice(&buf).expect("deserialize IpcMessage"); + match msg { + IpcMessage::Spikes(batch) => { + assert_eq!(batch.session_id, None); + assert_eq!(batch.batch_id, batch_time.as_millis() as u64); + assert_eq!(batch.timestamp, batch_time.as_nanos() as u64); + assert!(batch.metadata.is_none()); + assert_eq!(batch.spikes.len(), spikes.len()); + for (got, want) in batch.spikes.iter().zip(spikes.iter()) { + assert_eq!(got.channel, want.channel); + assert_eq!(got.time, want.time); + assert!((got.strength - want.strength).abs() < f32::EPSILON); + } + } + other => panic!("expected Spikes, got {other:?}"), + } + } } } #[cfg(feature = "corpus-ipc")] -pub use zmq_impl::{ZmqSpikeSink, ZmqStimulusSource}; +pub use zmq_impl::{SinkSendStats, ZmqSpikeSink, ZmqStimulusSource}; diff --git a/src/bin/brainstem_daemon.rs b/src/bin/brainstem_daemon.rs index 83524f5..70fa8b6 100644 --- a/src/bin/brainstem_daemon.rs +++ b/src/bin/brainstem_daemon.rs @@ -98,17 +98,32 @@ async fn run(cfg: DaemonConfig, config_path: PathBuf) -> anyhow::Result<()> { let pub_socket = zmq_context .socket(zmq::PUB) .map_err(|e| anyhow::anyhow!("failed to create ZMQ PUB socket: {e}"))?; + // Bounded PUB lifecycle: set an explicit finite send high-water mark and + // a finite LINGER before bind (mirrors the SUB-side style in + // ZmqStimulusSource::connect). The finite LINGER keeps teardown and + // signal shutdown from blocking indefinitely on pending messages. pub_socket - .bind(&format!("tcp://*:{}", cfg.spine_pub_port)) - .map_err(|e| { - anyhow::anyhow!("failed to bind ZMQ PUB on {}: {e}", cfg.spine_pub_port) - })?; + .set_sndhwm(cfg.spine_pub_sndhwm) + .map_err(|e| anyhow::anyhow!("ZMQ sndhwm: {e}"))?; + pub_socket + .set_linger(cfg.spine_pub_linger_ms) + .map_err(|e| anyhow::anyhow!("ZMQ linger: {e}"))?; + // Loopback by default (`spine_pub_bind_host` = 127.0.0.1); broader + // exposure (e.g. 0.0.0.0) requires explicit configuration. + let pub_endpoint = + brainstem_daemon::daemon::pub_endpoint(&cfg.spine_pub_bind_host, cfg.spine_pub_port); + pub_socket + .bind(&pub_endpoint) + .map_err(|e| anyhow::anyhow!("failed to bind ZMQ PUB on {pub_endpoint}: {e}"))?; - info!("📡 Using ZMQ corpus-ipc backend (spine ports active)"); + info!("📡 Using ZMQ corpus-ipc backend (spine ports active, PUB bound on {pub_endpoint})"); BackendPair { source: Box::new(source), - sink: Box::new(brainstem_daemon::backend::ZmqSpikeSink::new(pub_socket)), + sink: Box::new(brainstem_daemon::backend::ZmqSpikeSink::with_policy( + pub_socket, + cfg.spine_pub_send_empty_batches, + )), } }; diff --git a/src/checkpoint/tests.rs b/src/checkpoint/tests.rs index 0151951..7137027 100644 --- a/src/checkpoint/tests.rs +++ b/src/checkpoint/tests.rs @@ -29,6 +29,10 @@ fn live_config(model_path: PathBuf, lif: usize, channels: usize) -> DaemonConfig log_level: "info".to_string(), spine_sub_port: 5555, spine_pub_port: 5556, + spine_pub_bind_host: "127.0.0.1".to_string(), + spine_pub_sndhwm: 1000, + spine_pub_linger_ms: 0, + spine_pub_send_empty_batches: true, model_path, lif_count: lif, izh_count: 0, @@ -389,6 +393,10 @@ fn simulation_uses_blank_with_dimensions_and_does_not_claim_spikenaut() { log_level: "info".to_string(), spine_sub_port: 5555, spine_pub_port: 5556, + spine_pub_bind_host: "127.0.0.1".to_string(), + spine_pub_sndhwm: 1000, + spine_pub_linger_ms: 0, + spine_pub_send_empty_batches: true, model_path: PathBuf::from("/no/such/model.json"), lif_count: 4, izh_count: 1, diff --git a/src/daemon.rs b/src/daemon.rs index fa6b794..f940d01 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -55,6 +55,37 @@ pub struct DaemonConfig { pub log_level: String, pub spine_sub_port: u16, pub spine_pub_port: u16, + /// Bind host for the built-in `corpus-ipc` ZMQ PUB listener. + /// + /// Defaults to loopback (`127.0.0.1`) so the daemon does not expose the + /// spike egress socket on all interfaces implicitly. Set to `0.0.0.0` (or a + /// specific interface address) to opt in to broader exposure. Parsed but a + /// no-op under the stub backend (mirrors `spine_pub_port`); only the + /// `corpus-ipc` PUB path consumes it. + #[serde(default = "default_spine_pub_bind_host")] + pub spine_pub_bind_host: String, + /// Explicit finite send high-water mark (SNDHWM) applied to the PUB socket + /// before bind. Bounds how many outbound messages ZeroMQ queues per + /// subscriber before it silently drops. Parsed but a no-op under the stub + /// backend; only the `corpus-ipc` PUB path consumes it. + #[serde(default = "default_spine_pub_sndhwm")] + pub spine_pub_sndhwm: i32, + /// Finite LINGER (milliseconds) applied to the PUB socket before bind so + /// teardown and signal shutdown cannot block indefinitely on pending + /// messages. `0` (the default) drops any pending messages on close, the + /// safe bounded default for a best-effort PUB. Parsed but a no-op under the + /// stub backend; only the `corpus-ipc` PUB path consumes it. + #[serde(default = "default_spine_pub_linger_ms")] + pub spine_pub_linger_ms: i32, + /// Empty-spike-batch compatibility policy for the PUB egress. + /// + /// `true` (the default) preserves the current behavior where the tick loop + /// emits a frame every tick even when the batch has zero spikes, so + /// subscribers relying on per-tick / heartbeat frames keep working. Set to + /// `false` to suppress empty batches at the sink. Parsed but a no-op under + /// the stub backend; only the `corpus-ipc` PUB path consumes it. + #[serde(default = "default_spine_pub_send_empty_batches")] + pub spine_pub_send_empty_batches: bool, pub model_path: PathBuf, pub lif_count: usize, pub izh_count: usize, @@ -78,6 +109,26 @@ pub struct DaemonConfig { pub control_bind: Option, } +/// Default PUB bind host: loopback, requiring explicit opt-in for broader exposure. +fn default_spine_pub_bind_host() -> String { + "127.0.0.1".to_string() +} + +/// Default finite send high-water mark for the PUB socket. +fn default_spine_pub_sndhwm() -> i32 { + 1000 +} + +/// Default finite LINGER (ms) for the PUB socket: drop pending on close. +fn default_spine_pub_linger_ms() -> i32 { + 0 +} + +/// Default empty-batch policy: send empty batches (preserves current behavior). +fn default_spine_pub_send_empty_batches() -> bool { + true +} + impl DaemonConfig { /// Load daemon configuration from a TOML file. pub fn load(path: &std::path::Path) -> Result { @@ -100,6 +151,8 @@ impl DaemonConfig { .with_context(|| format!("failed to parse config from {}", path.display()))?; validate_log_level(&cfg.log_level) .with_context(|| format!("invalid config from {}", path.display()))?; + validate_pub_socket_opts(&cfg) + .with_context(|| format!("invalid config from {}", path.display()))?; Ok(cfg) } } @@ -656,6 +709,41 @@ fn log_model_provenance(config: &DaemonConfig, provenance: &ModelProvenance) { } } +/// Fail closed on PUB socket options whose libzmq-valid sentinel values would +/// defeat the bounded PUB lifecycle guarantee. +/// +/// `spine_pub_sndhwm` must be `> 0`: `0` disables the send high-water mark +/// (unbounded buffering) and negatives are nonsensical. `spine_pub_linger_ms` +/// must be `>= 0`: a negative value requests infinite LINGER, which can block +/// shutdown indefinitely; `0` (the default) is the safe drop-on-close default. +pub(crate) fn validate_pub_socket_opts(config: &DaemonConfig) -> Result<()> { + if config.spine_pub_sndhwm <= 0 { + bail!( + "spine_pub_sndhwm must be > 0 (0 disables the send high-water mark, which defeats the bounded PUB lifecycle)" + ); + } + if config.spine_pub_linger_ms < 0 { + bail!( + "spine_pub_linger_ms must be >= 0 (a negative value requests infinite LINGER, which can block shutdown indefinitely)" + ); + } + Ok(()) +} + +/// Build a ZeroMQ TCP endpoint string, bracketing IPv6 literal hosts. +/// +/// ZeroMQ TCP endpoints require IPv6 literals to be bracketed +/// (`tcp://[::1]:5556`); an unbracketed `format!("tcp://{host}:{port}")` on +/// `::1` yields the invalid `tcp://::1:5556`. IPv4 literals and hostnames have +/// no `:` and are emitted unchanged; already-bracketed hosts are left as-is. +pub fn pub_endpoint(host: &str, port: u16) -> String { + if host.contains(':') && !host.starts_with('[') { + format!("tcp://[{host}]:{port}") + } else { + format!("tcp://{host}:{port}") + } +} + pub(crate) fn validate_neuron_count(config: &DaemonConfig) -> Result<()> { let total = config .lif_count @@ -1203,6 +1291,10 @@ mod tests { log_level: "info".to_string(), spine_sub_port: 5555, spine_pub_port: 5556, + spine_pub_bind_host: "127.0.0.1".to_string(), + spine_pub_sndhwm: 1000, + spine_pub_linger_ms: 0, + spine_pub_send_empty_batches: true, model_path: PathBuf::from("/tmp/model.mem"), lif_count: 16, izh_count: 5, @@ -2514,6 +2606,152 @@ block_timeout_ms = 0 assert_eq!(sink.emitted.len(), 1); } + #[test] + fn config_load_rejects_zero_sndhwm() { + let path = write_config_toml( + "sndhwm-zero", + r#" +tick_rate_hz = 1000 +log_level = "info" +spine_sub_port = 5555 +spine_pub_port = 5556 +model_path = "/tmp/model.mem" +lif_count = 1 +izh_count = 0 +channels = 1 +spine_pub_sndhwm = 0 +"#, + ); + let err = DaemonConfig::load(&path).unwrap_err(); + let _ = std::fs::remove_file(&path); + let message = format!("{err:#}"); + assert!( + message.contains("spine_pub_sndhwm must be > 0"), + "{message}" + ); + } + + #[test] + fn config_load_rejects_negative_sndhwm() { + let path = write_config_toml( + "sndhwm-negative", + r#" +tick_rate_hz = 1000 +log_level = "info" +spine_sub_port = 5555 +spine_pub_port = 5556 +model_path = "/tmp/model.mem" +lif_count = 1 +izh_count = 0 +channels = 1 +spine_pub_sndhwm = -1 +"#, + ); + let err = DaemonConfig::load(&path).unwrap_err(); + let _ = std::fs::remove_file(&path); + let message = format!("{err:#}"); + assert!( + message.contains("spine_pub_sndhwm must be > 0"), + "{message}" + ); + } + + #[test] + fn config_load_rejects_negative_linger() { + let path = write_config_toml( + "linger-negative", + r#" +tick_rate_hz = 1000 +log_level = "info" +spine_sub_port = 5555 +spine_pub_port = 5556 +model_path = "/tmp/model.mem" +lif_count = 1 +izh_count = 0 +channels = 1 +spine_pub_linger_ms = -1 +"#, + ); + let err = DaemonConfig::load(&path).unwrap_err(); + let _ = std::fs::remove_file(&path); + let message = format!("{err:#}"); + assert!( + message.contains("spine_pub_linger_ms must be >= 0"), + "{message}" + ); + } + + #[test] + fn config_load_accepts_defaults_without_pub_socket_keys() { + // Backward compat: a minimal TOML omitting the PUB socket keys keeps the + // safe defaults (sndhwm 1000, linger 0) and loads cleanly. + let path = write_config_toml( + "pub-defaults", + r#" +tick_rate_hz = 1000 +log_level = "info" +spine_sub_port = 5555 +spine_pub_port = 5556 +model_path = "/tmp/model.mem" +lif_count = 1 +izh_count = 0 +channels = 1 +"#, + ); + let cfg = DaemonConfig::load(&path).unwrap(); + let _ = std::fs::remove_file(&path); + assert_eq!(cfg.spine_pub_sndhwm, 1000); + assert_eq!(cfg.spine_pub_linger_ms, 0); + assert_eq!(cfg.spine_pub_bind_host, "127.0.0.1"); + } + + #[test] + fn config_load_accepts_zero_linger_explicitly() { + let path = write_config_toml( + "linger-zero", + r#" +tick_rate_hz = 1000 +log_level = "info" +spine_sub_port = 5555 +spine_pub_port = 5556 +model_path = "/tmp/model.mem" +lif_count = 1 +izh_count = 0 +channels = 1 +spine_pub_sndhwm = 4 +spine_pub_linger_ms = 0 +"#, + ); + let cfg = DaemonConfig::load(&path).unwrap(); + let _ = std::fs::remove_file(&path); + assert_eq!(cfg.spine_pub_sndhwm, 4); + assert_eq!(cfg.spine_pub_linger_ms, 0); + } + + #[test] + fn pub_endpoint_brackets_ipv6_literals() { + assert_eq!(super::pub_endpoint("::1", 5556), "tcp://[::1]:5556"); + assert_eq!(super::pub_endpoint("fe80::1", 5556), "tcp://[fe80::1]:5556"); + } + + #[test] + fn pub_endpoint_leaves_ipv4_and_hostnames_unbracketed() { + assert_eq!( + super::pub_endpoint("127.0.0.1", 5556), + "tcp://127.0.0.1:5556" + ); + assert_eq!(super::pub_endpoint("0.0.0.0", 5556), "tcp://0.0.0.0:5556"); + assert_eq!( + super::pub_endpoint("localhost", 5556), + "tcp://localhost:5556" + ); + } + + #[test] + fn pub_endpoint_leaves_already_bracketed_hosts_untouched() { + assert_eq!(super::pub_endpoint("[::1]", 5556), "tcp://[::1]:5556"); + } + // Sends a real SIGTERM to this test process, so it's `#[ignore]`d by default: // `cargo test` runs the whole crate's tests in parallel threads of one process, // and a process-wide SIGTERM delivered before the handler below finishes diff --git a/tests/thalamic_brainstem_smoke.rs b/tests/thalamic_brainstem_smoke.rs index 6c14e90..b7a26d1 100644 --- a/tests/thalamic_brainstem_smoke.rs +++ b/tests/thalamic_brainstem_smoke.rs @@ -133,6 +133,10 @@ fn smoke_config(model_path: PathBuf) -> DaemonConfig { log_level: "info".into(), spine_sub_port: 5555, spine_pub_port: 5556, + spine_pub_bind_host: "127.0.0.1".into(), + spine_pub_sndhwm: 1000, + spine_pub_linger_ms: 0, + spine_pub_send_empty_batches: true, model_path, lif_count: LIF, izh_count: IZH,