From 61c3898132d29f8f466c0fab90bfa0ea6068c7e4 Mon Sep 17 00:00:00 2001 From: "kiro-agent[bot]" <245459735+kiro-agent[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 03:20:10 +0000 Subject: [PATCH 1/5] fix(ingress): reject non-finite stimuli and modulators before network step --- src/daemon.rs | 264 ++++++++++++++++++++++++++++++++++++++++++- src/ingress/mod.rs | 53 +++++++++ src/ingress/tests.rs | 33 +++++- 3 files changed, 348 insertions(+), 2 deletions(-) diff --git a/src/daemon.rs b/src/daemon.rs index 078850b..2ba92d7 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -809,6 +809,7 @@ struct TickDiagnostics { receive: OccurrenceLimiter, emit: OccurrenceLimiter, dropped: OccurrenceLimiter, + non_finite: OccurrenceLimiter, } impl TickDiagnostics { @@ -819,12 +820,16 @@ impl TickDiagnostics { receive: OccurrenceLimiter::new(interval, Duration::ZERO), emit: OccurrenceLimiter::new(interval, Duration::ZERO), dropped: OccurrenceLimiter::new(interval, Duration::ZERO), + non_finite: OccurrenceLimiter::new(interval, Duration::ZERO), } } #[cfg(test)] fn suppressed_total(&self) -> u64 { - self.receive.suppressed() + self.emit.suppressed() + self.dropped.suppressed() + self.receive.suppressed() + + self.emit.suppressed() + + self.dropped.suppressed() + + self.non_finite.suppressed() } } @@ -834,6 +839,7 @@ impl Default for TickDiagnostics { receive: OccurrenceLimiter::new(DIAGNOSTIC_INTERVAL, DIAGNOSTIC_MIN_GAP), emit: OccurrenceLimiter::new(DIAGNOSTIC_INTERVAL, DIAGNOSTIC_MIN_GAP), dropped: OccurrenceLimiter::new(DIAGNOSTIC_INTERVAL, DIAGNOSTIC_MIN_GAP), + non_finite: OccurrenceLimiter::new(DIAGNOSTIC_INTERVAL, DIAGNOSTIC_MIN_GAP), } } } @@ -876,6 +882,36 @@ fn run_tick( } }; + // Fail-closed finite gate: a live packet carrying any non-finite stimulus + // or modulator (NaN/+/-Inf) is rejected and counted here, before it can + // reach `network.step` (which errors on non-finite and would be treated as + // fatal). Modulators are validated before any held/cached or network state + // is updated. An already-`rejected` packet is left to the existing + // `rejected_batches` accounting below. Finite inputs are untouched. + let mut backend_packet = backend_packet; + if let Some(packet) = backend_packet.as_ref() + && !packet.rejected + && let Some(offender) = crate::ingress::first_non_finite(packet) + { + stats.rejected_batches += 1; + let msg = format!("Rejected non-finite ingress: {offender}"); + let now = time::Instant::now(); + match diagnostics.non_finite.record(diagnostic_key(&msg), now) { + Some(suppressed) => { + stats.diagnostic_emissions = stats.diagnostic_emissions.saturating_add(1); + warn!(total = stats.rejected_batches, suppressed, "{msg}"); + } + None => { + stats.suppressed_diagnostics = stats.suppressed_diagnostics.saturating_add(1); + } + } + // Drop the packet for this tick: skip IngressObserved, do not + // admit/enqueue it, do not count it accepted, and do not reach + // network.step. decode_inputs then zero-fills (identical to a + // skip / backend `None`). + backend_packet = None; + } + if backend_packet.as_ref().is_some_and(packet_carries_input) { health.apply(HealthEvent::IngressObserved); } @@ -1621,6 +1657,232 @@ channels = 16 assert_eq!(stats.suppressed_diagnostics, 4); } + // ── Non-finite ingress rejection (issue #66) ───────────────────────────── + + /// Source that yields the same (cloned) packet on every tick, for driving + /// the run_tick non-finite gate deterministically. + struct RepeatingSource { + packet: IngressPacket, + } + impl StimulusSource for RepeatingSource { + fn next_ingress(&mut self) -> Result> { + Ok(Some(self.packet.clone())) + } + fn initialize(&mut self, _model_path: Option<&str>) -> Result<()> { + Ok(()) + } + } + + /// Drive `ticks` ticks with the same packet through a fresh 1-channel + /// network, returning the collected sink, stats, diagnostics, and health. + fn drive_repeating( + packet: IngressPacket, + ticks: usize, + interval: u64, + ) -> ( + CollectingSpikeSink, + RuntimeStats, + TickDiagnostics, + HealthHandle, + SpikingNetwork, + ) { + let ingress = BoundedIngress::new(IngressConfig::default()).unwrap(); + let health = HealthHandle::started(HealthLimits::default()); + let mut source = RepeatingSource { packet }; + let mut sink = CollectingSpikeSink::new(); + let mut network = SpikingNetwork::with_dimensions(1, 0, 1); + let mut stimuli = [0.0]; + let mut spike_buf = Vec::new(); + let mut stats = RuntimeStats::default(); + let mut diagnostics = TickDiagnostics::new(interval); + for _ in 0..ticks { + run_tick( + &mut source, + &mut network, + &mut sink, + &mut stimuli, + &mut spike_buf, + &ingress, + &mut TickReport { + health: &health, + stats: &mut stats, + diagnostics: &mut diagnostics, + }, + ); + } + (sink, stats, diagnostics, health, network) + } + + fn assert_non_finite_stimulus_rejected(value: f32) { + let packet = IngressPacket { + stimuli: vec![value], + batch_id: Some(7), + ..Default::default() + }; + let (sink, stats, diagnostics, health, _network) = drive_repeating(packet, 25, 10); + + // Counted as rejected, never accepted, and the daemon keeps ticking + // (25 successful ticks emitted, not fatal). + assert_eq!(stats.rejected_batches, 25, "value {value:?}"); + assert_eq!(stats.accepted_batches, 0, "value {value:?}"); + assert!(stats.last_batch_id.is_none(), "value {value:?}"); + assert_eq!(stats.ticks, 25, "value {value:?}"); + assert_eq!(sink.emitted.len(), 25, "value {value:?}"); + + let snap = health.snapshot(); + assert!(snap.live, "daemon must stay live: value {value:?}"); + assert!( + snap.input_freshness.age_ms.is_none(), + "rejected input must not count as fresh ingress: value {value:?}" + ); + + // Diagnostics are rate-limited exactly like the receive-error path + // (interval 10 over 25 occurrences: emit at 1, 11, 21). + assert_eq!(stats.diagnostic_emissions, 3, "value {value:?}"); + assert_eq!(stats.suppressed_diagnostics, 22, "value {value:?}"); + assert_eq!(diagnostics.suppressed_total(), 4, "value {value:?}"); + } + + #[test] + fn non_finite_stimulus_is_rejected_and_bounded() { + assert_non_finite_stimulus_rejected(f32::NAN); + assert_non_finite_stimulus_rejected(f32::INFINITY); + assert_non_finite_stimulus_rejected(f32::NEG_INFINITY); + } + + #[test] + fn non_finite_modulator_is_rejected_and_state_untouched() { + let packet = IngressPacket { + stimuli: vec![0.0], + modulators: Some(vec![f32::NAN, 0.0, 0.0, 0.0]), + batch_id: Some(9), + ..Default::default() + }; + let (sink, stats, _diag, health, network) = drive_repeating(packet, 3, 1_000); + + assert_eq!(stats.rejected_batches, 3); + assert_eq!(stats.accepted_batches, 0); + assert_eq!(sink.emitted.len(), 3); + assert!(health.snapshot().live); + // Modulators are validated before any network/held state update, so the + // network's modulator snapshot stays at its default. + assert_eq!(network.modulators, NeuroModulators::default()); + } + + #[test] + fn finite_wrong_width_packet_is_accepted_not_rejected() { + // stimuli length (3) != channels (1): still ACCEPTED; decode_inputs + // truncates. accepted_batches increments only because batch_id is Some. + let packet = IngressPacket { + stimuli: vec![0.1, 0.2, 0.3], + batch_id: Some(11), + ..Default::default() + }; + let (sink, stats, _diag, _health, _network) = drive_repeating(packet, 1, 1_000); + assert_eq!(stats.rejected_batches, 0); + assert_eq!(stats.accepted_batches, 1); + assert_eq!(stats.last_batch_id, Some(11)); + assert_eq!(sink.emitted.len(), 1); + + // Same shape but no batch_id: accepted (not rejected) but uncounted, + // exactly as today. + let packet = IngressPacket { + stimuli: vec![0.1, 0.2, 0.3], + ..Default::default() + }; + let (_sink, stats, _diag, _health, _network) = drive_repeating(packet, 1, 1_000); + assert_eq!(stats.rejected_batches, 0); + assert_eq!(stats.accepted_batches, 0); + } + + #[test] + fn non_finite_in_masked_slot_is_still_rejected() { + // Fail-closed: a non-finite value in a valid_mask==false slot is still + // rejected. Masking does not neutralize the finite check. + let packet = IngressPacket { + stimuli: vec![f32::NAN], + valid_mask: Some(vec![false]), + batch_id: Some(13), + ..Default::default() + }; + let (sink, stats, _diag, _health, _network) = drive_repeating(packet, 2, 1_000); + assert_eq!(stats.rejected_batches, 2); + assert_eq!(stats.accepted_batches, 0); + assert!(stats.last_valid_mask.is_none()); + assert_eq!(sink.emitted.len(), 2); + } + + #[test] + fn empty_input_is_accepted_and_ticks() { + // Empty stimuli + None modulators (the stub default) is accepted and + // ticks with zeroed stimuli; nothing is rejected. + let (sink, stats, _diag, health, _network) = + drive_repeating(IngressPacket::default(), 1, 1_000); + assert_eq!(stats.rejected_batches, 0); + assert_eq!(stats.accepted_batches, 0); + assert_eq!(sink.emitted.len(), 1); + assert!(health.snapshot().input_freshness.age_ms.is_none()); + } + + #[test] + fn first_non_finite_classifies_ingress_packets() { + use crate::ingress::{NonFiniteInput, first_non_finite}; + + // Empty packet and finite packet: accepted (None). + assert_eq!(first_non_finite(&IngressPacket::default()), None); + let finite = IngressPacket { + stimuli: vec![0.0, 1.0, -2.5], + modulators: Some(vec![0.1, 0.2, 0.3, 0.4]), + ..Default::default() + }; + assert_eq!(first_non_finite(&finite), None); + + // NaN/Inf/-Inf in stimuli report the first offending index. + for value in [f32::NAN, f32::INFINITY, f32::NEG_INFINITY] { + let packet = IngressPacket { + stimuli: vec![0.0, value, 3.0], + ..Default::default() + }; + assert_eq!( + first_non_finite(&packet), + Some(NonFiniteInput::Stimulus { index: 1 }), + "value {value:?}" + ); + } + + // Non-finite only in modulators reports a Modulator offender. + let packet = IngressPacket { + stimuli: vec![0.0, 1.0], + modulators: Some(vec![0.0, f32::NAN, 0.0, 0.0]), + ..Default::default() + }; + assert_eq!( + first_non_finite(&packet), + Some(NonFiniteInput::Modulator { index: 1 }) + ); + + // A non-finite value in a masked-out slot is still flagged (fail-closed). + let packet = IngressPacket { + stimuli: vec![f32::INFINITY], + valid_mask: Some(vec![false]), + ..Default::default() + }; + assert_eq!( + first_non_finite(&packet), + Some(NonFiniteInput::Stimulus { index: 0 }) + ); + + // A short modulator tail is still fully scanned for present elements. + let packet = IngressPacket { + modulators: Some(vec![0.0, f32::NAN]), + ..Default::default() + }; + assert_eq!( + first_non_finite(&packet), + Some(NonFiniteInput::Modulator { index: 1 }) + ); + } + struct ScriptedStimulusSource { packet: IngressPacket, } diff --git a/src/ingress/mod.rs b/src/ingress/mod.rs index 3c895e6..1015b10 100644 --- a/src/ingress/mod.rs +++ b/src/ingress/mod.rs @@ -38,6 +38,59 @@ use crate::backend::IngressPacket; use queue::{BoundedQueue, WaitMode}; +/// First non-finite (`NaN`, `+Inf`, or `-Inf`) value found in an [`IngressPacket`]. +/// +/// Returned by [`first_non_finite`] so the tick loop can reject and count a +/// live invalid packet before it reaches `SpikingNetwork::step` (which would +/// otherwise be a fatal error). The index identifies the offending element +/// within its respective vector. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum NonFiniteInput { + /// A stimulus value at `packet.stimuli[index]` is not finite. + Stimulus { + /// Index into `packet.stimuli`. + index: usize, + }, + /// A modulator value at `packet.modulators[index]` is not finite. + Modulator { + /// Index into `packet.modulators`. + index: usize, + }, +} + +impl std::fmt::Display for NonFiniteInput { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Stimulus { index } => write!(f, "non-finite stimulus at index {index}"), + Self::Modulator { index } => write!(f, "non-finite modulator at index {index}"), + } + } +} + +/// Return the first non-finite (`NaN`, `+Inf`, or `-Inf`) value in `packet`. +/// +/// Scans `packet.stimuli` first, then `packet.modulators` (when `Some`), +/// returning the first offender or `None` when every present value is finite. +/// +/// This is a fail-closed check: stimuli are scanned regardless of +/// `packet.valid_mask`, so a non-finite value is flagged even in a masked-out +/// (`false`) slot. Masking is a downstream decode concern and is not a reason +/// to admit `NaN`/`Inf` into the network. Empty stimuli with `None` modulators +/// returns `None` (an empty packet is accepted). A modulator vector shorter than +/// [`crate::backend::NEUROMODULATOR_COUNT`](crate::backend::NEUROMODULATOR_COUNT) +/// is still fully scanned for the elements that are present. +pub fn first_non_finite(packet: &IngressPacket) -> Option { + if let Some(index) = packet.stimuli.iter().position(|v| !v.is_finite()) { + return Some(NonFiniteInput::Stimulus { index }); + } + if let Some(mods) = packet.modulators.as_ref() + && let Some(index) = mods.iter().position(|v| !v.is_finite()) + { + return Some(NonFiniteInput::Modulator { index }); + } + None +} + /// Hard cap so a TOML typo cannot request an enormous `VecDeque`. pub const MAX_QUEUE_CAPACITY: usize = 16_384; diff --git a/src/ingress/tests.rs b/src/ingress/tests.rs index 3216b69..8f87b44 100644 --- a/src/ingress/tests.rs +++ b/src/ingress/tests.rs @@ -9,7 +9,8 @@ use crate::backend::IngressPacket; use super::{ BoundedIngress, ClassMetrics, EnqueueOutcome, IngressConfig, MAX_BLOCK_TIMEOUT_MS, - MAX_PAYLOAD_LEN, MAX_QUEUE_CAPACITY, MessageClass, OverflowPolicy, + MAX_PAYLOAD_LEN, MAX_QUEUE_CAPACITY, MessageClass, NonFiniteInput, OverflowPolicy, + first_non_finite, }; fn wait_until(timeout: Duration, mut pred: impl FnMut() -> bool) -> bool { @@ -426,6 +427,36 @@ fn admit_skips_empty_backend_placeholder() { assert_eq!(ingress.pop(MessageClass::Sensory).unwrap().stimuli[0], 7.0); } +#[test] +fn first_non_finite_is_the_choke_point_for_direct_ingress() { + // The direct enqueue / admit_backend_packet path drains through run_tick, + // which uses `first_non_finite` as the single rejection choke point. This + // documents that contract against packets shaped like the direct helpers. + + // Existing finite direct-ingress helpers are accepted (None). + assert_eq!(first_non_finite(&pkt(1.0)), None); + assert_eq!(first_non_finite(&reward_pkt(0.5)), None); + + // A stimulus packet with a non-finite value is flagged. + let mut bad = pkt(f32::NAN); + assert_eq!( + first_non_finite(&bad), + Some(NonFiniteInput::Stimulus { index: 0 }) + ); + bad.stimuli = vec![f32::INFINITY]; + assert_eq!( + first_non_finite(&bad), + Some(NonFiniteInput::Stimulus { index: 0 }) + ); + + // A reward/modulator packet with a non-finite value is flagged. + let bad_reward = reward_pkt(f32::NEG_INFINITY); + assert_eq!( + first_non_finite(&bad_reward), + Some(NonFiniteInput::Modulator { index: 0 }) + ); +} + #[test] fn admit_empty_stimuli_still_enqueues_modulators() { let ingress = BoundedIngress::new(IngressConfig::tiny_fixture()).unwrap(); From a9b4c80a65bc463eb825f3e9162ae76f8a5cf174 Mon Sep 17 00:00:00 2001 From: "kiro-agent[bot]" <245459735+kiro-agent[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 03:24:39 +0000 Subject: [PATCH 2/5] fix(ingress): guard held modulators and cover typed IPC non-finite rejection --- src/backend.rs | 88 +++++++++++++++++++++++++++++++++++++++++-- src/ingress/corpus.rs | 63 +++++++++++++++++++++++++++++++ 2 files changed, 147 insertions(+), 4 deletions(-) diff --git a/src/backend.rs b/src/backend.rs index bcf6916..f9cb2a7 100644 --- a/src/backend.rs +++ b/src/backend.rs @@ -242,10 +242,22 @@ mod zmq_impl { .as_nanos() as u64 } - fn hold_modulators(&mut self, mods: &[f32]) { + /// Cache `mods` as the held modulator vector, replayed on idle ticks. + /// + /// Returns `false` (and caches nothing) when any value is non-finite + /// (NaN or ±Inf). This is the belt-and-suspenders guard required by + /// issue #66: modulators must be validated before any held/cached state + /// is updated, so a bad vector can never be replayed. The typed wire + /// path already rejects non-finite modulators at deserialization, but + /// this guard also protects any future non-typed producer. + fn hold_modulators(&mut self, mods: &[f32]) -> bool { + if mods.iter().any(|v| !v.is_finite()) { + return false; + } let dst = self.last_modulators.get_or_insert_with(Vec::new); dst.clear(); dst.extend_from_slice(mods); + true } fn skip_ingress(&self, rejected: bool) -> Option { @@ -281,8 +293,16 @@ mod zmq_impl { 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); + if let Some(mods) = packet.modulators.as_ref() + && !self.hold_modulators(mods) + { + // Non-finite modulators must never be cached or + // replayed (issue #66). Reject this frame like + // any other invalid frame instead of holding it. + tracing::warn!( + "Rejected ingress frame: non-finite modulator value" + ); + return Ok(self.skip_ingress(true)); } // Modulation-only: keep draining so a Stimuli frame // in the same tick is not delayed by one period. @@ -290,7 +310,13 @@ mod zmq_impl { } Ok(mut packet) => { if let Some(mods) = packet.modulators.as_ref() { - self.hold_modulators(mods); + if !self.hold_modulators(mods) { + // See above: never cache non-finite modulators. + tracing::warn!( + "Rejected ingress frame: non-finite modulator value" + ); + return Ok(self.skip_ingress(true)); + } } else { self.attach_held_modulators(&mut packet); } @@ -543,6 +569,60 @@ mod zmq_impl { assert!(!idle.rejected); assert_eq!(idle.modulators.as_deref(), Some(&[0.4, 0.3, 0.2, 1.0][..])); } + + #[test] + fn non_finite_modulators_are_never_held_or_replayed() { + // Invariant (issue #66): a non-finite modulator frame must never be + // cached into `last_modulators` and replayed on idle ticks. + // + // Layer that rejects it here: the corpus-ipc typed DESERIALIZE layer. + // A `NeuromodulatorSnapshot` with a non-finite field cannot even be + // built into a wire frame with valid JSON — `serde_json` encodes NaN + // as `null`, and `NeuromodulatorSnapshot`'s validating deserialize + // (`check_range` -> `check_finite`) rejects it. So `accept_ipc_json` + // returns an error and `next_ingress` takes the invalid-frame arm + // (`warn!("Rejected ingress frame: ..")` + `skip_ingress(true)`) + // before `hold_modulators` is ever called. The `hold_modulators` + // finite guard is the belt-and-suspenders backstop for any future + // non-typed path. + let (publisher, mut source) = bind_loopback(4); + let snapshot = NeuromodulatorSnapshot { + tick: 1, + dopamine: f32::NAN, + 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(); + + // The frame is rejected (skipped), never held. Because nothing is + // cached, subsequent ticks return `None` (no held modulators, no + // stimuli), and no idle tick ever replays the bad vector. + let mut saw_rejected = false; + for _ in 0..50 { + match source.next_ingress().expect("ingress must not hard-fail") { + Some(packet) => { + // A rejected/skip packet may carry the (still None) held + // modulators, but must never carry the NaN vector. + assert!( + packet.modulators.is_none(), + "non-finite modulators must never be held/replayed, got {:?}", + packet.modulators + ); + if packet.rejected { + saw_rejected = true; + } + } + None => std::thread::sleep(Duration::from_millis(10)), + } + } + assert!( + saw_rejected, + "the non-finite modulator frame must be surfaced as a rejected/skip outcome" + ); + } } } diff --git a/src/ingress/corpus.rs b/src/ingress/corpus.rs index 67cc17d..06ccb19 100644 --- a/src/ingress/corpus.rs +++ b/src/ingress/corpus.rs @@ -317,4 +317,67 @@ mod tests { let err = accept_ipc_message(IpcMessage::Ping, &policy(), 1_000).unwrap_err(); assert_eq!(err, IngressError::UnexpectedVariant("Ping")); } + + // ── Non-finite typed-IPC rejection (issue #66) ─────────────────────────── + // + // Observed layering (verified against corpus-ipc 0.1.0): + // * The direct decoded-struct path (`accept_ipc_message`) rejects a + // non-finite `values` entry inside `accept_stimulus_batch` via + // `batch.validate()` -> `check_finite_slice("values", ..)`, surfaced as + // `IngressError::Stimulus(_)`. + // * The JSON wire path (`accept_ipc_json`) never even reaches that check: + // `serde_json` serializes NaN/±Inf as `null`, so a non-finite value + // cannot round-trip and `de_values` fails to parse `null` as `f32`, + // surfaced as `IngressError::Deserialize(_)`. + // Both layers reject the frame; the invariant is that a non-finite stimulus + // value is never admitted, regardless of its `valid_mask` slot. + + #[test] + fn non_finite_stimulus_value_is_rejected_by_validate_on_decoded_struct() { + for bad in [f32::NAN, f32::INFINITY, f32::NEG_INFINITY] { + let mut batch = sample_batch(); + batch.values = vec![bad, 0.0, 0.25, 0.5]; + let err = accept_ipc_message(IpcMessage::Stimuli(batch), &policy(), 1_000).unwrap_err(); + // batch.validate() -> check_finite_slice produces Stimulus(_). + assert!( + matches!(err, IngressError::Stimulus(_)), + "expected Stimulus(_) for {bad:?}, got {err:?}" + ); + } + } + + #[test] + fn non_finite_stimulus_in_masked_slot_is_still_rejected() { + // Fail-closed: the typed layer's `check_finite_slice` scans every value + // regardless of `valid_mask`, so a non-finite value in a masked-out + // (false) slot is still rejected rather than admitted as a placeholder. + let mut batch = sample_batch(); + batch.values = vec![0.0, f32::NAN, 0.25, 0.5]; + batch.valid_mask = Some(vec![true, false, true, true]); + let err = accept_ipc_message(IpcMessage::Stimuli(batch), &policy(), 1_000).unwrap_err(); + match err { + IngressError::Stimulus(msg) => assert!( + msg.contains("values[1]"), + "expected the masked slot index in the error, got {msg}" + ), + other => panic!("expected Stimulus(_) for masked non-finite value, got {other}"), + } + } + + #[test] + fn non_finite_stimulus_value_cannot_round_trip_json() { + // On the wire, NaN/±Inf serialize to JSON `null`, so the frame is + // rejected at the deserialize layer (`de_values` cannot parse `null` + // as f32) before `accept_stimulus_batch`/`validate()` is reached. + for bad in [f32::NAN, f32::INFINITY, f32::NEG_INFINITY] { + let mut batch = sample_batch(); + batch.values = vec![bad, 0.0, 0.25, 0.5]; + let bytes = serde_json::to_vec(&IpcMessage::Stimuli(batch)).unwrap(); + let err = accept_ipc_json(&bytes, &policy(), 1_000).unwrap_err(); + assert!( + matches!(err, IngressError::Deserialize(_)), + "expected Deserialize(_) for wire {bad:?}, got {err:?}" + ); + } + } } From 5eb6c54022552e04ca93dff3dcb01fe44cba67a7 Mon Sep 17 00:00:00 2001 From: "kiro-agent[bot]" <245459735+kiro-agent[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 03:30:17 +0000 Subject: [PATCH 3/5] docs(ingress): clarify short non-finite modulator tail is rejected --- src/ingress/mod.rs | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/src/ingress/mod.rs b/src/ingress/mod.rs index 1015b10..ee7ba96 100644 --- a/src/ingress/mod.rs +++ b/src/ingress/mod.rs @@ -79,6 +79,15 @@ impl std::fmt::Display for NonFiniteInput { /// returns `None` (an empty packet is accepted). A modulator vector shorter than /// [`crate::backend::NEUROMODULATOR_COUNT`](crate::backend::NEUROMODULATOR_COUNT) /// is still fully scanned for the elements that are present. +/// +/// This is intentional and stricter than downstream decoding: `decode_inputs` +/// only reads modulators when the vector has at least +/// [`NEUROMODULATOR_COUNT`](crate::backend::NEUROMODULATOR_COUNT) elements and +/// otherwise falls back to default modulators for a short (finite) tail. A +/// non-finite value in a short tail (1..`NEUROMODULATOR_COUNT` elements) is +/// still rejected here rather than silently dropped, so this is not an +/// inconsistency with `decode_inputs`: the fail-closed check runs before decode +/// and never lets `NaN`/`Inf` reach it. pub fn first_non_finite(packet: &IngressPacket) -> Option { if let Some(index) = packet.stimuli.iter().position(|v| !v.is_finite()) { return Some(NonFiniteInput::Stimulus { index }); From 05eb09907dc0f69f2beb06d204a1911064d67e44 Mon Sep 17 00:00:00 2001 From: "kiro-agent[bot]" <245459735+kiro-agent[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 03:39:20 +0000 Subject: [PATCH 4/5] fix(ingress): reject non-finite direct bounded-ingress packets The finite-value gate only validated backend packets from source.next_ingress(), so packets injected through the public BrainstemDaemon::ingress() -> BoundedIngress::enqueue handle could still carry NaN/Inf through drain_for_tick() into network.step, tripping the fatal shutdown the fix intends to prevent (issue #66, direct bounded-ingress path). Add a second fail-closed non-finite gate on the drained/merged packet, the single choke point through which both the admitted-backend and direct-enqueue paths flow. Reuses first_non_finite, the diagnostics.non_finite occurrence limiter, and the same rejected_batches/diagnostic counters and message shape as the backend gate. Rejected ticks stay live (not fatal), do not count a successful tick, and emit no spike batch. Finite inputs are unchanged. Adds regression tests covering a non-finite stimulus (Sensory) and modulator (Reward) enqueued directly, over NaN/+Inf/-Inf. --- src/daemon.rs | 121 ++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 121 insertions(+) diff --git a/src/daemon.rs b/src/daemon.rs index 2ba92d7..ac341b2 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -935,6 +935,37 @@ fn run_tick( let packet = drained.into_packet(); apply_queue_pressure(health, ingress); + // Second fail-closed finite gate, on the drained/merged packet. The + // backend-packet gate above only sees `source.next_ingress()`; packets + // injected through the public `BoundedIngress` handle + // (`BrainstemDaemon::ingress()` -> `enqueue`/`try_enqueue`) never pass it + // and surface only here, in the drained packet. This drained packet is the + // single choke point through which both the admitted-backend path and the + // direct-enqueue path flow, so validating it before `decode_inputs` closes + // the direct-ingress path required by issue #66. Modulators are validated + // here too, before any held/cached or network state is updated. + if let Some(offender) = crate::ingress::first_non_finite(&packet) { + stats.rejected_batches += 1; + let msg = format!("Rejected non-finite ingress: {offender}"); + let now = time::Instant::now(); + match diagnostics.non_finite.record(diagnostic_key(&msg), now) { + Some(suppressed) => { + stats.diagnostic_emissions = stats.diagnostic_emissions.saturating_add(1); + warn!(total = stats.rejected_batches, suppressed, "{msg}"); + } + None => { + stats.suppressed_diagnostics = stats.suppressed_diagnostics.saturating_add(1); + } + } + // Skip the step for this tick without touching the network: do not call + // `network.step`, do not count a successful tick (`stats.ticks`), and do + // not emit a spike batch for the rejected input. The daemon stays live + // and keeps ticking (not fatal), mirroring the safe no-step exit used + // when there is nothing valid to publish. + health.apply(HealthEvent::TickSucceeded); + return; + } + let modulators = decode_inputs(&packet, stimuli); // Note: decode_inputs already zero-fills any remaining channels when packet.stimuli is shorter. @@ -1824,6 +1855,96 @@ channels = 16 assert!(health.snapshot().input_freshness.age_ms.is_none()); } + /// Reproduces the reported #66 gap: a non-finite packet injected through the + /// public `BoundedIngress` handle bypasses the backend-packet gate (the + /// backend source here yields `None`) and would previously reach + /// `network.step` via the drained packet, tripping the fatal shutdown. The + /// drained-packet gate must reject it, keep the daemon live, count it as a + /// rejected batch, and never advance a successful tick for it. + fn assert_direct_enqueue_non_finite_rejected(class: MessageClass, packet: IngressPacket) { + // Backend source contributes nothing this tick; the only non-finite + // data comes from the direct `enqueue`, exercising the path that never + // passed the backend-packet gate. + struct NoneSource; + impl StimulusSource for NoneSource { + fn next_ingress(&mut self) -> Result> { + Ok(None) + } + fn initialize(&mut self, _model_path: Option<&str>) -> Result<()> { + Ok(()) + } + } + + let ingress = BoundedIngress::new(IngressConfig::default()).unwrap(); + let health = HealthHandle::started(HealthLimits::default()); + let mut source = NoneSource; + let mut sink = CollectingSpikeSink::new(); + let mut network = SpikingNetwork::with_dimensions(2, 0, 2); + let mut stimuli = vec![0.0; 2]; + let mut spike_buf = Vec::new(); + let mut stats = RuntimeStats::default(); + let mut diagnostics = TickDiagnostics::new(1_000); + + // Inject directly through the SAME ingress handle passed to run_tick. + assert!(ingress.enqueue(class, packet).accepted()); + + run_tick( + &mut source, + &mut network, + &mut sink, + &mut stimuli, + &mut spike_buf, + &ingress, + &mut TickReport { + health: &health, + stats: &mut stats, + diagnostics: &mut diagnostics, + }, + ); + + // Daemon stays live (not fatal) and the rejected input never drove + // network.step: no successful tick advanced and no spike batch emitted. + let snap = health.snapshot(); + assert!( + snap.live, + "daemon must stay live after direct non-finite enqueue" + ); + assert_eq!(stats.rejected_batches, 1); + assert_eq!( + stats.ticks, 0, + "rejected input must not advance a network tick" + ); + assert!(sink.emitted.is_empty(), "no spike batch for rejected input"); + // Held network state is untouched (modulators validated before decode). + assert_eq!(network.modulators, NeuroModulators::default()); + } + + #[test] + fn direct_enqueue_non_finite_stimulus_is_rejected_not_fatal() { + for value in [f32::NAN, f32::INFINITY, f32::NEG_INFINITY] { + assert_direct_enqueue_non_finite_rejected( + MessageClass::Sensory, + IngressPacket { + stimuli: vec![value, 0.0], + ..Default::default() + }, + ); + } + } + + #[test] + fn direct_enqueue_non_finite_modulator_is_rejected_not_fatal() { + for value in [f32::NAN, f32::INFINITY, f32::NEG_INFINITY] { + assert_direct_enqueue_non_finite_rejected( + MessageClass::Reward, + IngressPacket { + modulators: Some(vec![value, 0.0, 0.0, 0.0]), + ..Default::default() + }, + ); + } + } + #[test] fn first_non_finite_classifies_ingress_packets() { use crate::ingress::{NonFiniteInput, first_non_finite}; From af275893d7db92e140419876501f2e6ee46a07e3 Mon Sep 17 00:00:00 2001 From: "kiro-agent[bot]" <245459735+kiro-agent[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 03:54:11 +0000 Subject: [PATCH 5/5] refactor(ingress): de-duplicate non-finite rejection paths Extract the shared non-finite rejection body from run_tick's two fail-closed gates into record_non_finite_rejection, and flatten the backend-packet gate's let-chain via a backend_non_finite helper. In ZmqStimulusSource::next_ingress, collapse the two duplicated hold_modulators reject branches into a hold_or_reject helper. Pure code-health cleanup: no behavior change. Counters, diagnostic message, rate limiter, warn fields, packet-drop and TickSucceeded early-return semantics are all preserved. Reduces the CodeScene Complex Method / Complex Conditional (run_tick) and Bumpy Road (next_ingress) findings on #66. --- src/backend.rs | 39 ++++++++++++++++------------- src/daemon.rs | 67 +++++++++++++++++++++++++++++--------------------- 2 files changed, 61 insertions(+), 45 deletions(-) diff --git a/src/backend.rs b/src/backend.rs index f9cb2a7..d15b754 100644 --- a/src/backend.rs +++ b/src/backend.rs @@ -277,6 +277,20 @@ mod zmq_impl { packet.modulators.clone_from(&self.last_modulators); } } + + /// Cache `mods` for idle replay, or reject the frame when they are + /// non-finite (issue #66: never cache or replay non-finite modulators). + /// Returns `Some(skip_packet)` for the caller to return on rejection, or + /// `None` when the modulators were held successfully. Shared by both + /// held-modulator branches in `next_ingress` so the warn message and + /// `skip_ingress(true)` reject path stay identical. + fn hold_or_reject(&mut self, mods: &[f32]) -> Option> { + if self.hold_modulators(mods) { + return None; + } + tracing::warn!("Rejected ingress frame: non-finite modulator value"); + Some(self.skip_ingress(true)) + } } impl StimulusSource for ZmqStimulusSource { @@ -294,31 +308,22 @@ mod zmq_impl { 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) + && let Some(skip) = self.hold_or_reject(mods) { - // Non-finite modulators must never be cached or - // replayed (issue #66). Reject this frame like - // any other invalid frame instead of holding it. - tracing::warn!( - "Rejected ingress frame: non-finite modulator value" - ); - return Ok(self.skip_ingress(true)); + return Ok(skip); } // 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() { - if !self.hold_modulators(mods) { - // See above: never cache non-finite modulators. - tracing::warn!( - "Rejected ingress frame: non-finite modulator value" - ); - return Ok(self.skip_ingress(true)); + match packet.modulators.as_ref() { + Some(mods) => { + if let Some(skip) = self.hold_or_reject(mods) { + return Ok(skip); + } } - } else { - self.attach_held_modulators(&mut packet); + None => self.attach_held_modulators(&mut packet), } return Ok(Some(packet)); } diff --git a/src/daemon.rs b/src/daemon.rs index ac341b2..0c49ac4 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -804,6 +804,42 @@ fn diagnostic_key(message: &str) -> u64 { }) } +/// Non-finite offender in a live backend packet, or `None` when the packet is +/// already `rejected` (left to the existing `rejected_batches` accounting) or +/// carries only finite values. Flattens the backend-gate condition so the check +/// is a single `Option` combinator instead of a multi-branch `let`-chain. +fn backend_non_finite(packet: &IngressPacket) -> Option { + if packet.rejected { + return None; + } + crate::ingress::first_non_finite(packet) +} + +/// Record one non-finite ingress rejection: bump `rejected_batches`, format the +/// shared diagnostic message, and emit-or-suppress through the bounded +/// `non_finite` rate limiter. Shared by both fail-closed gates in `run_tick` +/// (backend packet and drained packet) so the counting/logging is byte-for-byte +/// identical. Callers own the surrounding control flow (dropping the backend +/// packet, or the `TickSucceeded` + early return on the drained gate). +fn record_non_finite_rejection( + offender: crate::ingress::NonFiniteInput, + stats: &mut RuntimeStats, + diagnostics: &mut TickDiagnostics, +) { + stats.rejected_batches += 1; + let msg = format!("Rejected non-finite ingress: {offender}"); + let now = time::Instant::now(); + match diagnostics.non_finite.record(diagnostic_key(&msg), now) { + Some(suppressed) => { + stats.diagnostic_emissions = stats.diagnostic_emissions.saturating_add(1); + warn!(total = stats.rejected_batches, suppressed, "{msg}"); + } + None => { + stats.suppressed_diagnostics = stats.suppressed_diagnostics.saturating_add(1); + } + } +} + #[derive(Debug)] struct TickDiagnostics { receive: OccurrenceLimiter, @@ -889,22 +925,8 @@ fn run_tick( // is updated. An already-`rejected` packet is left to the existing // `rejected_batches` accounting below. Finite inputs are untouched. let mut backend_packet = backend_packet; - if let Some(packet) = backend_packet.as_ref() - && !packet.rejected - && let Some(offender) = crate::ingress::first_non_finite(packet) - { - stats.rejected_batches += 1; - let msg = format!("Rejected non-finite ingress: {offender}"); - let now = time::Instant::now(); - match diagnostics.non_finite.record(diagnostic_key(&msg), now) { - Some(suppressed) => { - stats.diagnostic_emissions = stats.diagnostic_emissions.saturating_add(1); - warn!(total = stats.rejected_batches, suppressed, "{msg}"); - } - None => { - stats.suppressed_diagnostics = stats.suppressed_diagnostics.saturating_add(1); - } - } + if let Some(offender) = backend_packet.as_ref().and_then(backend_non_finite) { + record_non_finite_rejection(offender, stats, diagnostics); // Drop the packet for this tick: skip IngressObserved, do not // admit/enqueue it, do not count it accepted, and do not reach // network.step. decode_inputs then zero-fills (identical to a @@ -945,18 +967,7 @@ fn run_tick( // the direct-ingress path required by issue #66. Modulators are validated // here too, before any held/cached or network state is updated. if let Some(offender) = crate::ingress::first_non_finite(&packet) { - stats.rejected_batches += 1; - let msg = format!("Rejected non-finite ingress: {offender}"); - let now = time::Instant::now(); - match diagnostics.non_finite.record(diagnostic_key(&msg), now) { - Some(suppressed) => { - stats.diagnostic_emissions = stats.diagnostic_emissions.saturating_add(1); - warn!(total = stats.rejected_batches, suppressed, "{msg}"); - } - None => { - stats.suppressed_diagnostics = stats.suppressed_diagnostics.saturating_add(1); - } - } + record_non_finite_rejection(offender, stats, diagnostics); // Skip the step for this tick without touching the network: do not call // `network.step`, do not count a successful tick (`stats.ticks`), and do // not emit a spike batch for the rejected input. The daemon stays live