fix(ipc): bound PUB lifecycle, loss semantics and safe bind defaults - #74
Conversation
Make the built-in corpus-ipc ZeroMQ PUB path explicit and bounded without overstating PUB reliability (issue #68 / LIM-1320, v0.3.0 blocker). - Add serde-default config keys spine_pub_bind_host (loopback default), spine_pub_sndhwm (finite, 1000), spine_pub_linger_ms (finite, 0), and spine_pub_send_empty_batches (true). Existing TOML stays backward compatible. - Set explicit finite SNDHWM and finite LINGER on the PUB socket before bind, and bind to tcp://<spine_pub_bind_host>:<spine_pub_port> (loopback by default instead of all interfaces). - Add empty-batch policy plus application-side send accounting (attempted/suppressed/failed) on ZmqSpikeSink, documenting that subscriber-side PUB loss is inherently unobservable and is not measured. Wire payload unchanged. - Add feature-gated tests: no-subscriber, empty-batch both ways, slow/disconnected subscriber, bounded teardown, and wire-payload round-trip. - Document the new keys and semantics in README and AGENTS.md. Stub build/tests unchanged; corpus-ipc gate green.
|
Skipping PR review because a bot author is detected. If you want to trigger CodeAnt AI, comment |
|
Important Review skippedBot user detected. To trigger a single review, invoke the ⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Advanced Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
Comment |
Not up to standards ⛔🔴 Issues
|
| Category | Results |
|---|---|
| Comprehensibility | 1 minor |
🟢 Metrics 41 complexity · 12 duplication
Metric Results Complexity 41 Duplication 12
NEW Get contextual insights on your PRs based on Codacy's metrics, along with PR and Jira context, without leaving GitHub. Enable AI reviewer
TIP This summary will be updated as you push new changes.
| #[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, |
There was a problem hiding this comment.
These fields need range validation before they reach ZeroMQ. spine_pub_sndhwm = 0 means no HWM limit, and spine_pub_linger_ms = -1 means infinite linger; both deserialize successfully and are passed straight to set_sndhwm/set_linger, so a valid config can defeat this PR’s central bounded-lifecycle guarantee and potentially make shutdown block indefinitely. Please reject sndhwm <= 0 and linger_ms < 0 during config validation (with tests), or use types that cannot represent those sentinel values.
Severity 10/10 · View on dashboard · PR Review Settings
There was a problem hiding this comment.
Fixed in df081f8 — validate_pub_socket_opts rejects spine_pub_sndhwm <= 0 and spine_pub_linger_ms < 0 at config load, with unit tests.
|
@coderabbitai review |
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
|
|
@CodeAnt-AI review |
🤖 CodeAnt AI — Review Status
|
Thanks for using CodeAnt! 🎉We're free for open-source projects. if you're enjoying it, help us grow by sharing. Share on X · |
| let pub_endpoint = format!("tcp://{}:{}", cfg.spine_pub_bind_host, cfg.spine_pub_port); | ||
| pub_socket | ||
| .bind(&pub_endpoint) |
There was a problem hiding this comment.
Suggestion: An IPv6 bind host such as ::1 produces tcp://::1:5556, which is not valid ZeroMQ TCP endpoint syntax and makes valid interface configuration fail.
Assessment: 🟠 Major · 🔁 Occurrence: Sometimes · 🏷️ Type error
Prompt for AI Agent 🤖
This is a comment left during a code review.
**Path:** src/bin/brainstem_daemon.rs
**Line:** 113:115
**Comment:**
*Type Error: An IPv6 bind host such as `::1` produces `tcp://::1:5556`, which is not valid ZeroMQ TCP endpoint syntax and makes valid interface configuration fail.
Validate the correctness of the flagged issue. If correct, How can I resolve this? If you propose a fix, implement it and please make it concise.
Once fix is implemented, also check other comments on the same PR, and ask user if the user wants to fix the rest of the comments as well. if said yes, then fetch all the comments validate the correctness and implement a minimal fixThere was a problem hiding this comment.
Fixed in df081f8 — pub_endpoint brackets IPv6 literals so ::1 becomes tcp://[::1]:5556.
| 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, |
There was a problem hiding this comment.
Suggestion: These i32 settings are accepted without validation, so -1 enables infinite linger and 0 enables unlimited SNDHWM, breaking the documented finite lifecycle guarantees.
Assessment: 🔴 Critical · 🔁 Occurrence: Rarely · 🏷️ Api mismatch
Prompt for AI Agent 🤖
This is a comment left during a code review.
**Path:** src/daemon.rs
**Line:** 72:79
**Comment:**
*Api Mismatch: These `i32` settings are accepted without validation, so `-1` enables infinite linger and `0` enables unlimited SNDHWM, breaking the documented finite lifecycle guarantees.
Validate the correctness of the flagged issue. If correct, How can I resolve this? If you propose a fix, implement it and please make it concise.
Once fix is implemented, also check other comments on the same PR, and ask user if the user wants to fix the rest of the comments as well. if said yes, then fetch all the comments validate the correctness and implement a minimal fixThere was a problem hiding this comment.
Fixed in df081f8 — config load rejects nonpositive SNDHWM and negative LINGER before they reach ZeroMQ.
| 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}"))?; |
There was a problem hiding this comment.
🔴 Zero send limit permits unbounded buffering
With spine_pub_sndhwm = 0, the PUB socket has no send queue limit. ZeroMQ treats zero SNDHWM as unlimited, so slow subscribers can exhaust publisher memory.
Learn more
The PUB socket's send high-water mark limits how many messages it queues for each subscriber. ZeroMQ interprets SNDHWM zero as unlimited, and the new public configuration field accepts zero without validation. The binary passes that value directly to set_sndhwm, so a slow subscriber defeats the intended buffer bound.
Example: Configure spine_pub_sndhwm = 0 and connect a subscriber that stops reading. The publisher keeps emitting frames without the configured finite queue limit, instead of dropping frames when the queue fills.
Recommended fix: Validate DaemonConfig::spine_pub_sndhwm as strictly positive before creating or binding the PUB socket. Cover zero and negative TOML values with tests, while preserving the default of 1000.
Was this helpful? React with 👍 or 👎 to provide feedback.
There was a problem hiding this comment.
Fixed in df081f8 — spine_pub_sndhwm <= 0 fails closed at config load.
| pub_socket | ||
| .set_linger(cfg.spine_pub_linger_ms) | ||
| .map_err(|e| anyhow::anyhow!("ZMQ linger: {e}"))?; |
There was a problem hiding this comment.
🔴 Negative linger can stall daemon shutdown
With spine_pub_linger_ms = -1, shutdown can wait indefinitely for queued PUB messages. ZeroMQ treats negative-one LINGER as infinite, and the binary passes that value through unchanged.
Learn more
LINGER controls how long ZeroMQ waits for pending messages when a socket or its context closes. ZeroMQ interprets -1 as an infinite wait. The new public configuration field accepts -1, which the binary passes to set_linger. If a connected subscriber stops reading while messages remain queued, signal-driven teardown can hang instead of completing within a finite time.
Example: Configure spine_pub_linger_ms = -1, fill the PUB queue for a stalled subscriber, and send SIGTERM. The daemon closes its socket with infinite linger and can wait indefinitely rather than exiting promptly.
Recommended fix: Validate DaemonConfig::spine_pub_linger_ms as nonnegative before applying the socket option; cover -1 in a config or startup test. Retain zero as the bounded default.
Was this helpful? React with 👍 or 👎 to provide feedback.
There was a problem hiding this comment.
Fixed in df081f8 — spine_pub_linger_ms < 0 fails closed at config load.
| 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"); |
There was a problem hiding this comment.
Fixed in df081f8 — corpus-ipc test now covers a connected subscriber that stops reading while the publisher emits past SNDHWM.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f6b6e16b2b
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| .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}"))?; |
There was a problem hiding this comment.
Reject unbounded PUB option sentinels
When a corpus-ipc configuration sets spine_pub_sndhwm = 0 or spine_pub_linger_ms = -1, these libzmq-valid sentinel values respectively remove the queue limit or request infinite linger. Because DaemonConfig::load does not validate either field, such a configuration starts successfully but defeats the advertised bounded lifecycle, and infinite linger can block shutdown while messages are pending. Reject a nonpositive SNDHWM and negative linger before applying these options.
AGENTS.md reference: AGENTS.md:L92-L94
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Fixed in df081f8 — DaemonConfig::load rejects nonpositive SNDHWM and negative LINGER via validate_pub_socket_opts.
Address converged PR review findings on the bounded PUB lifecycle change: - Fail closed at config load on libzmq-valid sentinel values that defeat the bounded-lifecycle guarantee. `DaemonConfig::load` now rejects `spine_pub_sndhwm <= 0` (0 disables the send high-water mark, enabling unbounded buffering) and `spine_pub_linger_ms < 0` (negative requests infinite LINGER, which can block shutdown indefinitely), via a new `validate_pub_socket_opts` using the existing `with_context` wrapping. Valid configs and the omitted-key defaults (1000 / 0) still load. - Bracket IPv6 literal bind hosts. A new `pub_endpoint` helper wraps IPv6 literals so `::1` becomes `tcp://[::1]:5556` instead of the invalid `tcp://::1:5556`; IPv4 literals and hostnames are unchanged. The binary uses it for the PUB bind. - Strengthen slow-subscriber coverage in the corpus-ipc backend tests: a subscriber that connects but stops reading while the publisher emits far beyond SNDHWM. Every emit still returns Ok and only `attempted` advances (failed == 0, suppressed == 0), documenting the silent over-HWM drop. Adds stub-build unit tests for the config validation and IPv6 bracketing.
| 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. |
There was a problem hiding this comment.
Dismissed — daemon.rs is a pre-existing multi-responsibility hotspot; this PR only adds bounded PUB config validation helpers and does not expand the cohesion surface enough to justify a module split here. Out of scope for the PUB lifecycle fix.
There was a problem hiding this comment.
Gates Failed
Prevent hotspot decline
(1 hotspot with Low Cohesion, Code Duplication)
Enforce critical code health rules
(1 file with Low Cohesion)
Enforce advisory code health rules
(2 files with Large Assertion Blocks, Code Duplication)
Our agent can fix these. Install it.
Gates Passed
3 Quality Gates Passed
Reason for failure
| Prevent hotspot decline | Violations | Code Health Impact | |
|---|---|---|---|
| daemon.rs | 2 rules in this hotspot | 6.04 → 5.22 | Suppress |
| Enforce critical code health rules | Violations | Code Health Impact | |
|---|---|---|---|
| daemon.rs | 1 critical rule | 6.04 → 5.22 | Suppress |
| Enforce advisory code health rules | Violations | Code Health Impact | |
|---|---|---|---|
| daemon.rs | 1 advisory rule | 6.04 → 5.22 | Suppress |
| backend.rs | 1 advisory rule | 8.96 → 8.41 | Suppress |
Quality Gate Profile: Pay Down Tech Debt
Install CodeScene MCP: safeguard and uplift AI-generated code. Catch issues early with our IDE extension and CLI tool.
| 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:?}"), | ||
| } | ||
| } |
There was a problem hiding this comment.
❌ New issue: Large Assertion Blocks
The test suite contains 4 assertion blocks with at least 4 assertions, threshold = 4
| 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. |
There was a problem hiding this comment.
❌ New issue: Low Cohesion
This module has at least 7 different responsibilities amongst its 48 functions, threshold = 3
| 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. |
There was a problem hiding this comment.
❌ Getting worse: Code Duplication
introduced similar code in: tests.config_load_rejects_negative_linger,tests.config_load_rejects_negative_sndhwm,tests.config_load_rejects_zero_sndhwm
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: df081f8de6
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| pub_socket | ||
| .bind(&pub_endpoint) |
There was a problem hiding this comment.
Enable IPv6 before binding IPv6 hosts
When spine_pub_bind_host is an IPv6 literal such as ::1, pub_endpoint now produces a valid bracketed endpoint, but the PUB socket still retains libzmq's default ZMQ_IPV6 = 0. Consequently, binding tcp://[::1]:<port> fails at startup instead of providing the IPv6 support covered by the new helper and tests; set the socket's IPv6 option before bind when an IPv6 host is configured.
Useful? React with 👍 / 👎.
#76) * refactor(daemon): extract config into src/config.rs to reduce module cohesion Move the CONFIG responsibility (RuntimeMode, DaemonConfig + serde defaults, DaemonConfig::load, validate_pub_socket_opts, pub_endpoint) out of daemon.rs into a new src/config.rs and re-export it from daemon.rs so every public and crate-internal path resolves unchanged. Pure code movement + re-exports: no behavior, signature, error-message, serde-attribute, or wire-payload changes. Fixes the CodeScene 'Prevent hotspot decline' / Low Cohesion gate that failed on merged PR #74 by lowering daemon.rs's distinct-responsibility count. * test(config): consolidate duplicated config-load tests Reduce Code Duplication in src/config.rs by factoring the shared valid TOML body into config_toml_with_log_level/base_config_toml helpers and a load_base_config helper, and by folding the four reject cases (unknown log level, zero sndhwm, negative sndhwm, negative linger) into a single table-driven test. All prior coverage and error-substring assertions are preserved; only the #[cfg(test)] module changed. --------- Co-authored-by: kiro-agent[bot] <245459735+kiro-agent[bot]@users.noreply.github.com>
User description
Problem
The built-in
corpus-ipcZeroMQ PUB egress was under-specified for v0.3.0 release qualification (issue #68 / LIM-1320):tcp://*:<port>— all interfaces — by default, exposing the readout beyond loopback without explicit opt-in.Resulting behavior
SNDHWMand a finiteLINGERbeforebind, mirroring the SUB-side style. Finite LINGER keeps teardown/shutdown from blocking on pending messages.tcp://{spine_pub_bind_host}:{spine_pub_port}, defaulting to loopback127.0.0.1. Broader exposure now requires explicit configuration.spine_pub_send_empty_batches(defaulttrue) makes the policy explicit and configurable while preserving today's per-tick always-send behavior; whenfalse, empty batches are suppressed (not sent).ZmqSpikeSinknow tracksSinkSendStats { attempted, suppressed, failed }— the only outcomes the publisher can observe. Code and docs are explicit that subscriber-side PUB loss (slow/absent/over-SNDHWM) is inherently unobservable by the publisher and is not measured. Asend()returningOkmeans the frame was accepted into ZeroMQ's egress, not that any subscriber received it.#[serde(default)], backward compatible — existing TOML without them still loads):spine_pub_bind_host(127.0.0.1),spine_pub_sndhwm(1000),spine_pub_linger_ms(0),spine_pub_send_empty_batches(true). Documented in README truth table + AGENTS.md.Preserved
corpus-ipc.{"Spikes":{...}}withbatch_id/timestamp/spikes) is byte-for-byte unchanged, proven by a round-trip deserialize test.run_tickerror handling, SUB ingress behavior, and the MSRV pin (1.98.1) are untouched.The PUB listener now defaults to loopback (127.0.0.1) instead of all interfaces. Deployments that rely on remote subscribers must set
spine_pub_bind_host = "0.0.0.0"(or a specific interface address). This is the intended safe default per the acceptance criteria.Validation
Independently re-run on this branch:
cargo fmt --checksilent;cargo clippy --locked --all-targets -- -D warningsclean;cargo build --lockedok;cargo test --locked146 passed / 0 failed.CC=gcc CXX=g++, libzmq built via zeromq-src):clippy --all-targets --all-features -D warningsclean;build --all-featuresok;test --all-features167 lib + 9 smoke passed / 0 failed, including 6 new feature-gated PUB tests (no-subscriber, empty-batch both ways, slow/disconnected-subscriber, teardown bounded by finite LINGER, wire-payload round-trip). The corpus-ipc lifecycle tests run in the Linuxcorpus-ipcCI job.Limitations
attempted/suppressed/failedcounters are exposed viastats()but not yet surfaced to health/metrics at runtime; wiring them into observability would broaden scope beyond this issue.Tracking
Part of #68 (implementation counterpart of LIM-1320). This is one blocker on an open v0.3.0 qualification issue, so this is a non-closing reference — the issue stays open.
🤖 Delivered by the Kiro agent.
Summary by cubic
Brings the built-in
corpus-ipcZeroMQ PUB egress to v0.3.0 qualification (LIM-1320) by bounding its lifecycle and shipping safe defaults: finiteSNDHWMandLINGERare set beforebind, teardown no longer blocks, and the listener now binds to loopback (127.0.0.1) instead of all interfaces.Behavior change to note
tcp://<spine_pub_bind_host>:<spine_pub_port>with a loopback default. Deployments relying on remote subscribers must setspine_pub_bind_host = "0.0.0.0"(or a specific interface). IPv6 literal hosts are auto-bracketed (::1→tcp://[::1]:port).spine_pub_bind_host(127.0.0.1),spine_pub_sndhwm(1000),spine_pub_linger_ms(0), andspine_pub_send_empty_batches(true). Config load now rejects values that defeat the bounded lifecycle:spine_pub_sndhwm <= 0andspine_pub_linger_ms < 0.ZmqSpikeSinktracks application-sideattempted/suppressed/failedcounts. Subscriber-side PUB loss (slow, absent, or over-SNDHWMsubscribers) is inherently unobservable by the publisher and documented as best-effort, not measured.SNDHWMslow subscribers, bounded teardown, and payload round-trip.Written for commit df081f8. Summary will update on new commits.
CodeAnt-AI Description
Bound and secure the ZeroMQ spike publisher lifecycle
What Changed
Impact
✅ Safer default network exposure✅ Bounded shutdown time✅ Configurable empty-batch traffic💡 Usage Guide
Checking Your Pull Request
Every time you make a pull request, our system automatically looks through it. We check for security issues, mistakes in how you're setting up your infrastructure, and common code problems. We do this to make sure your changes are solid and won't cause any trouble later.
Talking to CodeAnt AI
Got a question or need a hand with something in your pull request? You can easily get in touch with CodeAnt AI right here. Just type the following in a comment on your pull request, and replace "Your question here" with whatever you want to ask:
This lets you have a chat with CodeAnt AI about your pull request, making it easier to understand and improve your code.
Example
Preserve Org Learnings with CodeAnt
You can record team preferences so CodeAnt AI applies them in future reviews. Reply directly to the specific CodeAnt AI suggestion (in the same thread) and replace "Your feedback here" with your input:
This helps CodeAnt AI learn and adapt to your team's coding style and standards.
Example
Retrigger review
Ask CodeAnt AI to review the PR again, by typing:
Check Your Repository Health
To analyze the health of your code repository, visit our dashboard at https://app.codeant.ai. This tool helps you identify potential issues and areas for improvement in your codebase, ensuring your repository maintains high standards of code health.