Repository navigation
feat(OMN-1740): Add operational semantics to EventBusSubcontractWiring - #219
Conversation
Implement idempotency, error classification, and offset commit policies for contract-driven Kafka consumers per ARCH-002. Changes: - Add config models: ModelIdempotencyConfig, ModelOffsetPolicyConfig, ModelConsumerRetryConfig, ModelDlqConfig - Add dispatch_with_transaction() to ProtocolDispatchEngine - Implement idempotency gate with envelope_id deduplication - Add error classification: content errors → DLQ, infra errors → fail-fast - Add offset commit policies (commit_after_handler default) - Add retry tracking with configurable max attempts Test coverage: - 20 new unit tests for idempotency, error classification, offset policy - 6 new integration tests for AC#7 (deduplication) and AC#8 (offset commit) Related: OMN-1740
📝 WalkthroughWalkthroughAdds four event-bus Pydantic models and re-exports them, adds transaction-scoped dispatch API, and significantly extends EventBusSubcontractWiring with idempotency, retry/DLQ classification and handling, offset-commit controls, and associated tests; also bumps omnibase-core and passes node_name into runtime wiring. Changes
Sequence Diagram(s)sequenceDiagram
participant Consumer
participant Wiring as EventBusSubcontractWiring
participant Idempotency as IdempotencyStore
participant Dispatcher as ProtocolDispatchEngine
participant DLQ as DLQPublisher
participant Offset as OffsetManager
Consumer->>Wiring: deliver_message(message)
Wiring->>Idempotency: is_processed(envelope_id)?
alt duplicate
Idempotency-->>Wiring: processed
Wiring->>Offset: commit_offset(correlation_id)
else new_message
Idempotency-->>Wiring: not_processed
Wiring->>Dispatcher: dispatch_with_transaction(envelope, tx?)
alt success
Dispatcher-->>Wiring: success
Wiring->>Idempotency: record_processed(envelope_id)
Wiring->>Offset: commit_offset(correlation_id)
else content_error
Dispatcher-->>Wiring: content_error
alt dlq_and_commit
Wiring->>DLQ: publish_to_dlq(message, error, consumer_group)
Wiring->>Offset: commit_offset(correlation_id)
else fail_fast
Wiring-->>Consumer: raise ProtocolConfigurationError
end
else infra_error
Dispatcher-->>Wiring: infra_error
alt retry_available
Wiring->>Wiring: increment_retry_count(correlation_id)
else retries_exhausted
Wiring->>DLQ: publish_to_dlq(message, error, consumer_group)
opt dlq_and_commit
Wiring->>Offset: commit_offset(correlation_id)
end
end
end
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing touches
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Actionable comments posted: 5
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
src/omnibase_infra/runtime/event_bus_subcontract_wiring.py (1)
332-352:⚠️ Potential issue | 🟡 MinorKeep DLQ consumer_group aligned with the node_name used for subscriptions.
wire_subscriptions()can accept a differentnode_namethan the constructor, but_publish_to_dlq()usesself._node_name. This can produce mismatched consumer_group metadata in DLQ entries.🔧 Suggested fix
for topic_suffix in subcontract.subscribe_topics: + if node_name != self._node_name: + self._logger.warning( + "node_name override in wire_subscriptions: %s -> %s", + self._node_name, + node_name, + ) + self._node_name = node_name full_topic = self.resolve_topic(topic_suffix) group_id = f"{self._environment}.{node_name}"
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/models/event_bus/model_idempotency_config.py`:
- Around line 6-19: Update the module/class docstrings to consistently refer to
envelope_id instead of event_id: find the top-level module docstring and the
"Idempotency Overview" section plus any docstrings in the IdempotencyConfig
model (or any docstring mentioning the idempotency store) and replace the
incorrect `event_id` mentions with `envelope_id` (specifically update the lines
referenced in the review where `event_id` appears). Ensure any examples or
bullets describing the deduplication behavior mention `envelope.envelope_id` as
the key used for INSERT ... ON CONFLICT DO NOTHING and update the
retention/pruning description to reference `envelope_id` as well.
In `@src/omnibase_infra/protocols/protocol_dispatch_engine.py`:
- Around line 168-215: The docstring example for dispatch_with_transaction
currently shows using isinstance(tx, asyncpg.Connection) and raising TypeError;
update it to demonstrate duck-typing and raising the project-specific OnexError
instead: change the example in the dispatch_with_transaction docstring to
perform a minimal, protocol-style duck-typing check (e.g., verify required
method/attribute on tx) and raise OnexError (not TypeError) when the check
fails, and briefly mention that implementations should prefer duck-typing and
raise OnexError for transaction-type mismatches.
In `@src/omnibase_infra/runtime/event_bus_subcontract_wiring.py`:
- Around line 562-577: The duplicate branch currently returns before committing,
causing infinite redelivery; update the block after
_idempotency_store.check_and_record (when not is_new) to still
acknowledge/commit the message before exiting—e.g., invoke the same commit logic
used by commit_after_handler (or call the consumer commit/ack method your wiring
uses) and then log via self._logger.info and return; ensure you reference
envelope_id, topic, correlation_id, self._node_name and reuse the existing
commit routine so duplicates are acknowledged and not redelivered.
In `@src/omnibase_infra/runtime/service_message_dispatch_engine.py`:
- Around line 1494-1511: dispatch_with_transaction currently dereferences
envelope (e.g., envelope.correlation_id, envelope.trace_id) and assumes topic is
truthy before calling dispatch, which can raise AttributeError instead of the
expected ModelOnexError; update dispatch_with_transaction to validate inputs
up-front: check that envelope is not None and that topic is non-empty (or meets
whatever contract) and if invalid raise/return the appropriate ModelOnexError,
only then access envelope.correlation_id/trace_id, and finally call
self.dispatch(topic=topic, envelope=envelope); keep tx unused for now but
preserve the existing tx acknowledgement.
- Add node_name="runtime-host" to EventBusSubcontractWiring call - Update omnibase-core dependency to ^0.10.0 - Fixes mypy error: Missing positional argument "node_name" Related: OMN-1740
Includes Kafka import lint guard validator (OMN-1745). Related: OMN-1740
MAJOR fixes: - Fix retry-count off-by-one in calculate_delay_ms and get_all_delays_ms - Replace ValueError with OnexError in retry config (OnexError-only rule) - Commit offsets for deduped messages to prevent infinite redelivery - Add input validation in dispatch_with_transaction before envelope deref MINOR fixes: - Align DLQ consumer_group with subscription node_name for traceability - Update idempotency config docstrings: event_id → envelope_id - Update protocol_dispatch_engine tx type-narrowing example to use duck-typing + OnexError per ONEX conventions
Implement idempotency, error classification, and offset commit policies for contract-driven Kafka consumers per ARCH-002.
Changes:
Test coverage:
Related: OMN-1740
Summary by CodeRabbit
New Features
Tests
Chores
✏️ Tip: You can customize this high-level summary in your review settings.