diff --git a/AGENTS.md b/AGENTS.md index a099626..c0ec33c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -89,10 +89,6 @@ 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 4d682be..5a81772 100644 --- a/README.md +++ b/README.md @@ -158,11 +158,6 @@ 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. @@ -206,7 +201,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_bind_host:spine_pub_port` (loopback by default); logs `📡 Using ZMQ corpus-ipc backend` | **still stub** | +| `--features corpus-ipc` | ZMQ / `corpus-ipc` | required | SUB via env, PUB on `spine_pub_port`; 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. @@ -227,11 +222,7 @@ 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://:` (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 | +| `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); 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`): @@ -247,14 +238,6 @@ 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 ab50bbb..d15b754 100644 --- a/src/backend.rs +++ b/src/backend.rs @@ -363,93 +363,23 @@ 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; @@ -483,16 +413,7 @@ mod zmq_impl { .socket .lock() .map_err(|_| anyhow::anyhow!("ZMQ socket mutex poisoned"))?; - // `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}")); - } + guard.socket.send(payload, 0)?; Ok(()) } } @@ -707,247 +628,8 @@ 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::{SinkSendStats, ZmqSpikeSink, ZmqStimulusSource}; +pub use zmq_impl::{ZmqSpikeSink, ZmqStimulusSource}; diff --git a/src/bin/brainstem_daemon.rs b/src/bin/brainstem_daemon.rs index 70fa8b6..83524f5 100644 --- a/src/bin/brainstem_daemon.rs +++ b/src/bin/brainstem_daemon.rs @@ -98,32 +98,17 @@ 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 - .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}"))?; + .bind(&format!("tcp://*:{}", cfg.spine_pub_port)) + .map_err(|e| { + anyhow::anyhow!("failed to bind ZMQ PUB on {}: {e}", cfg.spine_pub_port) + })?; - info!("📡 Using ZMQ corpus-ipc backend (spine ports active, PUB bound on {pub_endpoint})"); + info!("📡 Using ZMQ corpus-ipc backend (spine ports active)"); BackendPair { source: Box::new(source), - sink: Box::new(brainstem_daemon::backend::ZmqSpikeSink::with_policy( - pub_socket, - cfg.spine_pub_send_empty_batches, - )), + sink: Box::new(brainstem_daemon::backend::ZmqSpikeSink::new(pub_socket)), } }; diff --git a/src/checkpoint/tests.rs b/src/checkpoint/tests.rs index 7137027..0151951 100644 --- a/src/checkpoint/tests.rs +++ b/src/checkpoint/tests.rs @@ -29,10 +29,6 @@ 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, @@ -393,10 +389,6 @@ 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 f940d01..fa6b794 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -55,37 +55,6 @@ 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, @@ -109,26 +78,6 @@ 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 { @@ -151,8 +100,6 @@ 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) } } @@ -709,41 +656,6 @@ 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 @@ -1291,10 +1203,6 @@ 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, @@ -2606,152 +2514,6 @@ 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 b7a26d1..6c14e90 100644 --- a/tests/thalamic_brainstem_smoke.rs +++ b/tests/thalamic_brainstem_smoke.rs @@ -133,10 +133,6 @@ 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,