Repository navigation
fix(OMN-15232): fail-close the deserialization-path DLQ publish in event_bus_kafka - #2497
Conversation
EventBusKafka._consume_loop awaited _publish_raw_to_dlq() on the deserialization-failure path and DISCARDED the returned bool, then continued. Consumers built here run with enable_auto_commit defaulting to True, so the client committed the fetch position regardless of whether the DLQ write ever landed -- a failed DLQ publish on a poison message was a silent, committed, unrecoverable drop. Same defect class as OMN-14936, whose gate landed only in runtime/event_bus_subcontract_wiring.py. The module audit this ticket requires found a second ungated site: _dispatch_to_subscriber discarded the DLQ result on the retries-exhausted branch, and MixinKafkaDlq._publish_to_dlq did not even return a persistence signal (always None), so no caller could have gated on it. Changes: - MixinKafkaDlq._publish_to_dlq now returns bool (the `success` it already computed), mirroring the _publish_raw_to_dlq contract from OMN-14936. - _dispatch_to_subscriber returns "safe to advance offset"; False only when retries were exhausted AND the DLQ write was not confirmed. - _consume_loop gates both sites and, when persistence is unconfirmed, rewinds the fetch position via consumer.seek(tp, msg.offset) -- the same "does NOT advance the committed offset" idiom KafkaTransport.nack uses. Withholding a commit is a no-op in this loop (it never commits; the client does), so the rewind is what makes it fail-closed under both auto-commit and manual-commit models. - Bounded backoff between rewinds so an unreachable DLQ cannot hot-spin. - Static AST test ratchets that no DLQ-publish call site in event_bus_kafka.py discards its persistence result. RED-first: 4 of 5 new tests fail against dev, all 5 pass with the fix.
📝 WalkthroughWalkthroughKafka DLQ publishing now returns a persistence result. The Kafka consumer uses that result to allow offset advancement only after confirmed DLQ persistence, otherwise rewinding to the failed record for deserialization and exhausted-handler failures. ChangesKafka DLQ fail-closed flow
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant KafkaConsumer
participant EventBusKafka
participant MixinKafkaDlq
KafkaConsumer->>EventBusKafka: deliver failed record
EventBusKafka->>MixinKafkaDlq: publish record to DLQ
MixinKafkaDlq-->>EventBusKafka: return persistence result
alt result is true
EventBusKafka-->>KafkaConsumer: allow offset advancement
else result is false
EventBusKafka->>KafkaConsumer: seek to failed partition and offset
end
Possibly related PRs
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
🧹 Nitpick comments (1)
src/omnibase_infra/event_bus/event_bus_kafka.py (1)
2139-2242: 🩺 Stability & Availability | 🔵 TrivialRewind mechanism is correct; consider a rewind counter for observability.
Verified against aiokafka's own documented pattern (
consumer.seek(tp, msg.offset)immediately followed bygetone()returning the same message) — theseek-based rewind here is a supported, immediate-effect API, not a buffered/delayed one.One operational gap: if the DLQ stays unreachable, this partition will loop
seek → sleep(1s) → refetch → failindefinitely (by design, to avoid silent loss), but there's currently no counter/metric distinguishing this from normal traffic — only ERROR-level logs. Consider incrementing a metric (e.g. via the existing DLQ metrics mechanism) on eachdlq_unpersisted_offset_rewoundso an operator can alert on a stuck partition rather than relying on log-grepping.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/omnibase_infra/event_bus/event_bus_kafka.py` around lines 2139 - 2242, The _rewind_after_unpersisted_dlq method lacks a metric for repeated failed DLQ persistence retries. Using the existing DLQ metrics mechanism, increment a counter each time the successful rewind path emits dlq_unpersisted_offset_rewound, including suitable topic, group, partition, and failure-stage dimensions if supported, while leaving the rewind and backoff behavior unchanged.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Nitpick comments:
In `@src/omnibase_infra/event_bus/event_bus_kafka.py`:
- Around line 2139-2242: The _rewind_after_unpersisted_dlq method lacks a metric
for repeated failed DLQ persistence retries. Using the existing DLQ metrics
mechanism, increment a counter each time the successful rewind path emits
dlq_unpersisted_offset_rewound, including suitable topic, group, partition, and
failure-stage dimensions if supported, while leaving the rewind and backoff
behavior unchanged.
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: c0698441-42be-4213-a6f7-2179915ef2df
📒 Files selected for processing (3)
src/omnibase_infra/event_bus/event_bus_kafka.pysrc/omnibase_infra/event_bus/mixin_kafka_dlq.pytests/unit/event_bus/test_deser_dlq_fail_closed_omn15232.py
…ibase_infra#2497 (#5118) * evidence(OMN-15232): author OCC companion for OmniNode-ai/omnibase_infra#2497 OCC companion by node_pr_lifecycle_fix_effect (OMN-13317 F1 / OMN-13990 / OMN-14285). Product PR head ceb33d3415fbc5c2a80615bc3e6929674bfe5147. * evidence(OMN-15232): self-bind OCC#5118 + rebind contract_sha256 --------- Co-authored-by: omnimarket-bot <bot@omninode.ai>
|
| Verdict | Meaning | Blocks merge? |
|---|---|---|
passed |
No critical findings | No |
blocked |
CRITICAL findings found | Yes |
degraded |
All models unavailable (infra) | No (pilot) |
Powered by omniintelligence.review_pairing.cli_review — node-based adversarial review via HandlerLlmCliSubprocess (OMN-8468/OMN-8524)
Defect
src/omnibase_infra/event_bus/event_bus_kafka.pyawaited_publish_raw_to_dlq(...)on the deserialization-failure path, discarded the returnedbool, thencontinued:Consumers built by this class pass
enable_auto_commit=self._config.enable_auto_commit, which defaults toTrue, and this loop never commits at all — the client commits the fetch position on its own cadence. So a failed DLQ publish on a poison message was a silent, committed, unrecoverable drop.Same defect class as OMN-14936, at a call site that fix did not cover: the OMN-14936 gate landed only in
runtime/event_bus_subcontract_wiring.pyas 5xif not dlq_persisted: return.Second ungated site found by the required module audit
_dispatch_to_subscriberdiscarded the DLQ result on the retries-exhausted branch — andMixinKafkaDlq._publish_to_dlqreturnedNoneunconditionally, so no caller could have gated on it. The static ratchet added here flagged both sites (lines 2075 and 2233) on the unfixed tree.Fix — and why a rewind, not "withhold the commit"
The OMN-14936 gate works by not calling commit, because that path is a manual-commit model. That mechanism is a no-op in this loop: it has no commit call to withhold; the aiokafka client commits for it. Declining to commit here would leave the fail-open behaviour exactly as-is.
The fail-closed action that works under both commit models is a rewind of the fetch position:
This is the same "does NOT advance the committed offset" idiom
KafkaTransport.nackalready uses inevent_bus/kafka_transport.py. Under auto-commit the committer then commits the rewound position, which cannot be past the un-persisted message; under manual commit the message is simply refetched. Either way Kafka redelivers and the DLQ write is retried instead of the record being lost. Only the failed message's own partition is rewound — siblings untouched, matching the OMN-14757 per-partition discipline.Changes
mixin_kafka_dlq.py_publish_to_dlqreturnsbool(thesuccessit already computed), mirroring the_publish_raw_to_dlqcontract from OMN-14936event_bus_kafka.py_dispatch_to_subscriberreturns "safe to advance offset";Falseonly when retries exhausted AND DLQ unconfirmedevent_bus_kafka.py_consume_loopgates both sites; new_rewind_after_unpersisted_dlqhelperevent_bus_kafka.pyDLQ_UNPERSISTED_REWIND_BACKOFF_SECONDSbounds the retry rate so an unreachable DLQ cannot hot-spinOnly an explicit
Falsecounts as a confirmed non-persist — duck-typed hosts/test doubles returningNonekeep prior behaviour, matching the allowance the OMN-14936 gate makes atevent_bus_subcontract_wiring.py:607.Accepted trade-off: with the DLQ down, a poison message now stalls its partition instead of being dropped. That is the intended fail-closed semantics and is what OMN-14936 already chose ("Kafka will redeliver this message and we retry the DLQ write rather than silently losing it"). On the fan-out path, redelivery may duplicate for subscribers that already succeeded — at-least-once, which is the delivery contract here, and strictly preferable to losing the record.
Seams
Two private return types widen
None -> bool. Both are private, have no external callers, andmypy src/omnibase_infra/is clean across 2615 source files.MixinKafkaDlq._publish_to_dlq:None→boolEventBusKafka._dispatch_to_subscriber:None→boolEventBusSubcontractWiring._publish_to_dlqis a different method on a different class (alreadyboolunder OMN-14936) and is untouched.ProtocolKafkaDlqHostdoes not declare either method, so no protocol surface moved.RED-first proof
tests/unit/event_bus/test_deser_dlq_fail_closed_omn15232.pydrives the real_consume_loopagainst a fake consumer — the artifact that runs, not a surrogate.Against unfixed
dev(4 failed, 1 passed):The static ratchet named both defects precisely:
With the fix:
5 passed.The one pre-passing test is the deliberate control — it proves the fix is a gate, not a blanket rewind (a confirmed DLQ persist must still let the offset advance, or poison messages replay forever).
Verification (rule 11a — patch-transferred to
.200, gates + push from.200)Diff was generated locally,
scp'd,git apply'd on.200, and sha256-verified identical on both hosts for all three files before any gate ran (no vacuous green). Both worktrees on basef7fb7cde. Commit and push executed on.200.Full-suite run (
pytest tests/ -q -n auto):26351 passed, 25 failed, 42 errors. Every failure/error is intests/integration,tests/performance, ortests/replayand needs live infra absent on.200— Kafka onlocalhost:19092, Postgres, and the LLM SLO endpoints. Zero failures intests/unit/when the unit tree is run on its own (22058 passed). The 4 unit failures that appeared inside the combined xdist session are cross-suite env-var pollution from the integration tests sharing the worker pool: those same 3 files pass 65/65 both in this patched worktree and in the cleandevcheckout on.200.No gate was bypassed; no
--no-verify, no skip tokens.OCC
Companion (net-new-file-only): OmniNode-ai/onex_change_control#5115 — opens and merges before this PR lands, per strict OCC ordering (OMN-15222). Its three evidence probes were verified present at
ceb33d341and absent ondev.Closes OMN-15232
Evidence-Ticket: OMN-15232
Evidence-Source: OCC#5118
Summary by CodeRabbit