Repository navigation
feat(dlq): configure DLQ for permanently failing intents [OMN-949] - #90
jonahgabriel merged 14 commits into
Conversation
Implements comprehensive Dead Letter Queue support for the KafkaEventBus:
## DLQ Topic Configuration
- Add topic_constants.py with ONEX-compliant DLQ topic naming
- Pattern: {env}.dlq.{category}.v1 (e.g., dev.dlq.intents.v1)
- Utility functions: build_dlq_topic(), parse_dlq_topic(), is_dlq_topic()
- Add get_dlq_topic() helper to ModelKafkaEventBusConfig
## Alerting Integration
- Add ModelDlqEvent for strongly-typed callback payloads
- Add ModelDlqMetrics for aggregate metrics tracking
- Add register_dlq_callback() for custom alerting hooks
- Structured logging at WARNING (success) and ERROR (failure) levels
- Metrics: total/successful/failed publishes, latency, per-topic/error counts
## No Silent Drops
- Fix deserialization errors to route to DLQ via _publish_raw_to_dlq()
- Add retry exhaustion check before DLQ routing
- Verify all permanent failure paths route to DLQ
## Documentation
- Add docs/architecture/DLQ_MESSAGE_FORMAT.md with payload schema
- Add docs/operations/DLQ_REPLAY_GUIDE.md with replay procedures
- Add scripts/dlq_replay.py CLI utility skeleton
## Tests
- Add 44 tests for topic_constants.py
- Add 6 tests for DLQ routing behavior
- All 71 event_bus tests pass
## Validation
- Update INFRA_MAX_UNIONS: 580 -> 590 (legitimate X | None patterns)
|
Warning Rate limit exceeded@jonahgabriel has exceeded the limit for the number of commits that can be reviewed per hour. Please wait 14 minutes and 40 seconds before requesting another review. ⌛ How to resolve this issue?After the wait time has elapsed, a review can be triggered using the We recommend that you space out your commits to avoid hitting the rate limit. 🚦 How do rate limits work?CodeRabbit enforces hourly rate limits for each developer per organization. Our paid plans have higher rate limits than the trial, open-source and free plans. In all cases, we re-allow further reviews after a brief timeout. Please see our FAQ for further information. 📒 Files selected for processing (10)
📝 WalkthroughWalkthroughAdds comprehensive Dead Letter Queue (DLQ) support: topic naming utilities, DLQ message schema docs, runtime DLQ publishing and metrics, replay tooling, centralized error sanitization, protocol/type adjustments, and unit/integration tests covering DLQ behavior and utilities. Changes
Sequence DiagramsequenceDiagram
autonumber
participant Handler as Event Handler
participant Bus as KafkaEventBus
participant DLQProd as DLQ Producer
participant Metrics as ModelDlqMetrics
participant Callbacks as DLQ Callbacks
Handler->>Bus: deliver message (includes retry_count, correlation_id)
Bus->>Handler: invoke handler
Handler--xBus: handler raises exception
Bus->>Bus: evaluate retries_exhausted?
alt retries exhausted
Bus->>DLQProd: _publish_raw_to_dlq / _publish_to_dlq (sanitized error, original_message, headers)
DLQProd->>DLQProd: build payload + headers
DLQProd-->>Bus: publish result (success/failure) + timing
Bus->>Metrics: record_dlq_publish(original_topic, error_type, success, duration_ms)
Bus->>Callbacks: invoke registered callbacks(ModelDlqEvent)
Callbacks-->>Bus: callback return (errors logged, isolated)
else retries remain
Bus->>Bus: log transient failure, continue (no DLQ)
end
Bus->>Bus: continue consume loop
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Poem
Comment |
PR Review: DLQ Implementation for Permanently Failing IntentsSummaryThis PR implements comprehensive Dead Letter Queue (DLQ) support for the KafkaEventBus. The implementation is well-architected and production-ready with excellent documentation, strong typing, comprehensive testing, and proper ONEX compliance. Overall Assessment: ✅ APPROVED with minor suggestions for enhancement. 🎯 Strengths1. Excellent Documentation
2. Strong Type Safety ✅
3. Comprehensive Testing ✅
4. ONEX Compliance ✅
5. Robust Error Handling
🔍 Issues Found1.
|
| Criterion | Status | Notes |
|---|---|---|
| DLQ topic configured | ✅ | {env}.dlq.{category}.v1 pattern implemented |
| Alerting integrated | Models exist but callbacks not invoked | |
| No silent drops | ✅ | Deserialization + retry exhaustion both route to DLQ |
| DLQ message format documented | ✅ | Excellent DLQ_MESSAGE_FORMAT.md |
| Replay mechanism defined | Documented but script is stub |
Partial criteria met: Alerting and replay need follow-up work.
🚀 Deployment Readiness
✅ Safe to Merge
This PR is production-ready with proper fallbacks:
- DLQ is optional (
dead_letter_topic: str | None) - DLQ failures don't crash consumer (best-effort pattern)
- All tests pass, pre-commit hooks pass
- Breaking changes: None (backwards compatible)
📝 Post-Merge Recommendations
-
Immediate (Before Production):
- Add DLQ callback invocation
- Wire up metrics tracking
- Add circuit breaker protection
-
Follow-Up (Phase 2):
- Complete DLQ replay script implementation
- Add integration tests with live Kafka
- Set up monitoring dashboards for DLQ metrics
-
Documentation:
- Add operational runbook for DLQ monitoring thresholds
- Document DLQ topic retention policies
- Create troubleshooting guide for DLQ publish failures
🎯 Final Verdict
APPROVED ✅
This is excellent work with strong architectural foundations. The identified issues are:
- Circuit breaker: Consistency with ONEX patterns
- Callbacks/Metrics: Complete the feature implementation
- Replay script: Known stub, documented in ticket
None of these block merging since:
- DLQ is optional (safe default: disabled)
- DLQ failures are logged and don't crash the system
- The foundation is solid for incremental enhancement
Great job on the comprehensive documentation and testing! 🎉
📚 References
- OMN-949: DLQ Configuration ticket
- CLAUDE.md: ONEX infrastructure patterns
docs/patterns/circuit_breaker_implementation.md: Circuit breaker requirementsdocs/patterns/error_handling_patterns.md: Error handling standards
…re-dlq-for-permanently-failing-intents
After merging main branch, the ProtocolEventBusLike protocol was deleted. This protocol is required by TimeoutEmitter and other components. - Add protocol_event_bus_like.py to protocols directory (proper location) - Re-export from protocols/__init__.py - Re-export from mixins/__init__.py for backwards compatibility
There was a problem hiding this comment.
Actionable comments posted: 2
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/protocols/__init__.py (1)
8-11: Update module docstring to include ProtocolEventBusLike.The module docstring lists the exported protocols but doesn't mention the newly added
ProtocolEventBusLike. Please add it to the list for completeness.Suggested documentation update
Protocols: + - ProtocolEventBusLike: Interface for event bus publishing (minimal duck-typed interface) - ProtocolIdempotencyStore: Interface for idempotency checking and deduplication - ProtocolPluginCompute: Interface for deterministic compute plugins - ProtocolSnapshotPublisher: Interface for snapshot publishing services (F2)
🧹 Nitpick comments (5)
src/omnibase_infra/mixins/__init__.py (1)
24-24: Consider removing protocol re-export from mixins module.Per ONEX conventions, protocols should be imported directly from
omnibase_infra.protocols.protocol_event_bus_likerather than through the mixins module. Re-exporting the protocol here creates an alternate import path that may lead to inconsistent usage patterns across the codebase.Recommended approach
Consumers should import the protocol directly:
from omnibase_infra.protocols import ProtocolEventBusLikeRather than:
from omnibase_infra.mixins import ProtocolEventBusLikeConsider removing lines 24 and 35 to enforce a single canonical import path.
Based on learnings: protocols should be imported from dedicated
protocol_*module paths.Also applies to: 35-35
docs/architecture/DLQ_MESSAGE_FORMAT.md (1)
356-362: Consider removing specific line number references.The implementation reference section includes specific line numbers (e.g., "lines 1765-1860") that will become stale as the codebase evolves. Consider linking to method names or using GitHub permalink URLs instead.
- **KafkaEventBus**: `src/omnibase_infra/event_bus/kafka_event_bus.py` - - `_publish_to_dlq()` method (lines 1765-1860) + - `_publish_to_dlq()` methodsrc/omnibase_infra/event_bus/kafka_event_bus.py (1)
2065-2101: Consider adding metrics and callback integration to_publish_raw_to_dlq.Unlike
_publish_to_dlq, this method doesn't update_dlq_metricsor invoke DLQ callbacks. This creates an observability gap for deserialization failures. Consider adding similar metrics/callback handling for consistency.scripts/dlq_replay.py (2)
78-128: Consider using Pydantic model instead of dataclass.Per coding guidelines, "All data structures must be proper Pydantic models." The
DLQMessagedataclass could be converted to a Pydantic model for consistency with the rest of the codebase. However, since this is a CLI script skeleton that may not need full validation, this could be deferred.
278-280: Unusual generator stub pattern.The
returnfollowed byyieldon line 280 is valid but unconventional. Consider using a more explicit pattern:- # Stub: Yield no messages (for skeleton) - return - yield # Make this a generator + # Stub: Yield no messages (for skeleton) + if False: # noqa: SIM223 + yield # type: ignore[misc]Or simply remove the unreachable
yieldsince modern Python type checkers handle empty async generators better.
📜 Review details
Configuration used: defaults
Review profile: CHILL
Plan: Lite
📒 Files selected for processing (17)
docs/architecture/DLQ_MESSAGE_FORMAT.mddocs/operations/README.mdscripts/dlq_replay.pysrc/omnibase_infra/event_bus/__init__.pysrc/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/event_bus/models/__init__.pysrc/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.pysrc/omnibase_infra/event_bus/models/model_dlq_event.pysrc/omnibase_infra/event_bus/models/model_dlq_metrics.pysrc/omnibase_infra/event_bus/topic_constants.pysrc/omnibase_infra/mixins/__init__.pysrc/omnibase_infra/protocols/__init__.pysrc/omnibase_infra/protocols/protocol_event_bus_like.pysrc/omnibase_infra/validation/infra_validators.pytests/unit/event_bus/test_kafka_event_bus.pytests/unit/event_bus/test_topic_constants.pytests/unit/validation/test_validator_defaults.py
🧰 Additional context used
📓 Path-based instructions (4)
**/*.py
📄 CodeRabbit inference engine (CLAUDE.md)
**/*.py: NEVER useAnytype - Always use specific types. All data structures must be proper Pydantic models.
UseEnumMessageCategoryfor message routing, topic parsing, and dispatcher selection (values: EVENT, COMMAND, INTENT). UseEnumNodeOutputTypefor execution shape validation and handler return type validation (values: EVENT, COMMAND, INTENT, PROJECTION).
Use PEP 604 union syntaxX | Nonefor nullable types instead ofOptional[X]. Example:def get_user(id: str) -> User | None:instead ofdef get_user(id: str) -> Optional[User]:
All services MUST use ModelONEXContainer for dependency injection. Bootstrap withcontainer = ModelONEXContainer()and resolve services viacontainer.service_registry.resolve_service(ServiceType)
Always propagatecorrelation_idfrom incoming requests to error context. Auto-generate usinguuid4()if no correlation_id exists. Use UUID format for all new correlation IDs.
NEVER include in error messages or context: passwords, API keys, tokens, secrets, full connection strings with credentials, PII, internal IP addresses, private keys, certificates, session tokens, or cookies. SAFE to include: service names, operation names, correlation IDs, error codes, sanitized hostnames, port numbers, retry counts, timeout values, resource identifiers (non-sensitive).
UseProtocolConfigurationErrorfor invalid config,SecretResolutionErrorfor missing secrets,InfraConnectionErrorfor connection failures,InfraTimeoutErrorfor operation timeouts,InfraAuthenticationErrorfor auth failures,InfraUnavailableErrorfor unavailable resources.
Use graceful degradation forInfraTimeoutError. Pattern: try primary source with timeout, fall back to cache/secondary source on timeout, aggregate results with degradation flag.
Files:
src/omnibase_infra/event_bus/models/__init__.pytests/unit/event_bus/test_topic_constants.pysrc/omnibase_infra/validation/infra_validators.pysrc/omnibase_infra/event_bus/models/model_dlq_event.pysrc/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/mixins/__init__.pyscripts/dlq_replay.pysrc/omnibase_infra/event_bus/__init__.pysrc/omnibase_infra/protocols/protocol_event_bus_like.pysrc/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.pytests/unit/event_bus/test_kafka_event_bus.pysrc/omnibase_infra/protocols/__init__.pytests/unit/validation/test_validator_defaults.pysrc/omnibase_infra/event_bus/topic_constants.pysrc/omnibase_infra/event_bus/models/model_dlq_metrics.py
**/*infra*.py
📄 CodeRabbit inference engine (CLAUDE.md)
**/*infra*.py: IncludeModelInfraErrorContextwith every infrastructure error. Context should include:transport_type(EnumInfraTransportType),operation,target_name, andcorrelation_id. Example:raise InfraConnectionError('Failed to connect', context=context)
UseEnumInfraTransportTypefor transport identification in error context. Values include: HTTP, DATABASE, KAFKA, CONSUL, VAULT, VALKEY, GRPC
Files:
src/omnibase_infra/validation/infra_validators.py
**/model_*.py
📄 CodeRabbit inference engine (CLAUDE.md)
**/model_*.py: One model per file - Each file contains exactly oneModel*class
Usemodel_<name>.pyfile naming pattern withModel<Name>class pattern for Pydantic models. Example:model_kafka_message.py→ModelKafkaMessage
Files:
src/omnibase_infra/event_bus/models/model_dlq_event.pysrc/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.pysrc/omnibase_infra/event_bus/models/model_dlq_metrics.py
**/protocol*.py
📄 CodeRabbit inference engine (CLAUDE.md)
Use
protocol_<name>.pyfile naming pattern for standalone protocols orprotocols.pyfor domain-grouped protocols withProtocol<Name>class pattern. Example:protocol_event_bus.pycontainsProtocolEventBus
Files:
src/omnibase_infra/protocols/protocol_event_bus_like.py
🧠 Learnings (29)
📓 Common learnings
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to **/*.py : Messages exceeding `KAFKA_MAX_REQUEST_SIZE` (default 10MB) are automatically routed to oversized DLQ topic. Use `scripts/process_oversized_dlq.py` to handle these messages. See `docs/guides/DLQ_HANDLING.md` for complete documentation.
Learnt from: CR
Repo: OmniNode-ai/omniintelligence PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-07T17:50:13.678Z
Learning: Implement Kafka event-driven architecture with proper topic naming using prefix dev.archon-intelligence. and proper event flow pattern with Effect nodes consuming events, processing, and publishing results with Dead Letter Queue routing
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to {services/**/*.py,scripts/bulk_ingest_repository.py} : Implement fail-closed configuration for security hardening. All external requests must validate URLs, implement DLQ routing, and handle failures gracefully.
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T17:23:49.777Z
Learning: All ONEX nodes must use the node_kafka_event_bus as a secondary reference only for complex backend and event bus logic and advanced configuration patterns
📚 Learning: 2025-12-06T22:21:32.649Z
Learnt from: CR
Repo: OmniNode-ai/omniagent PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-06T22:21:32.649Z
Learning: Applies to nodes/**/*.py : Use event bus mixins from `omnibase_core` for Kafka publishing instead of direct Kafka clients
Applied to files:
src/omnibase_infra/event_bus/models/__init__.pysrc/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/mixins/__init__.pysrc/omnibase_infra/event_bus/__init__.pysrc/omnibase_infra/protocols/protocol_event_bus_like.pytests/unit/event_bus/test_kafka_event_bus.pysrc/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to **/*.py : Publish intelligence requests to Kafka event bus using topics: dev.archon-intelligence.intelligence.code-analysis-{requested,completed,failed}.v1 for consistency and event-driven architecture
Applied to files:
tests/unit/event_bus/test_topic_constants.pytests/unit/event_bus/test_kafka_event_bus.py
📚 Learning: 2025-11-29T17:13:38.776Z
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to **/*.py : Messages exceeding `KAFKA_MAX_REQUEST_SIZE` (default 10MB) are automatically routed to oversized DLQ topic. Use `scripts/process_oversized_dlq.py` to handle these messages. See `docs/guides/DLQ_HANDLING.md` for complete documentation.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pyscripts/dlq_replay.pydocs/architecture/DLQ_MESSAGE_FORMAT.mdtests/unit/event_bus/test_kafka_event_bus.py
📚 Learning: 2025-12-07T17:50:13.678Z
Learnt from: CR
Repo: OmniNode-ai/omniintelligence PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-07T17:50:13.678Z
Learning: Implement Kafka event-driven architecture with proper topic naming using prefix dev.archon-intelligence. and proper event flow pattern with Effect nodes consuming events, processing, and publishing results with Dead Letter Queue routing
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pydocs/architecture/DLQ_MESSAGE_FORMAT.mdtests/unit/event_bus/test_kafka_event_bus.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: Applies to **/*infra*.py : Include `ModelInfraErrorContext` with every infrastructure error. Context should include: `transport_type` (EnumInfraTransportType), `operation`, `target_name`, and `correlation_id`. Example: `raise InfraConnectionError('Failed to connect', context=context)`
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: Applies to **/*.py : Use `ProtocolConfigurationError` for invalid config, `SecretResolutionError` for missing secrets, `InfraConnectionError` for connection failures, `InfraTimeoutError` for operation timeouts, `InfraAuthenticationError` for auth failures, `InfraUnavailableError` for unavailable resources.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/nodes/**/*.py : Import mixins from omnibase_core.mixins.* and use Mixin* naming pattern (e.g., MixinHealthCheck, MixinMetrics, MixinEventBus) - never use local custom mixins unless experimental and documented
Applied to files:
src/omnibase_infra/mixins/__init__.pysrc/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Applies to **/*.py : Import protocols from `omnibase.protocol.protocol_*` paths
Applied to files:
src/omnibase_infra/mixins/__init__.pysrc/omnibase_infra/protocols/protocol_event_bus_like.pysrc/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-11-24T17:24:41.687Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T17:24:41.687Z
Learning: Applies to **/*.py : Import protocols from `omnibase.protocol.protocol_<name>` module paths
Applied to files:
src/omnibase_infra/mixins/__init__.pysrc/omnibase_infra/protocols/protocol_event_bus_like.pysrc/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-11-24T16:32:55.606Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T16:32:55.606Z
Learning: Node implementations must use mixin-based composition from `omnibase_core.mixins` (e.g., `MixinHealthCheck`, `MixinNodeExecutor`) to add capabilities
Applied to files:
src/omnibase_infra/mixins/__init__.py
📚 Learning: 2025-11-29T17:13:38.776Z
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to {services/**/*.py,scripts/bulk_ingest_repository.py} : Implement fail-closed configuration for security hardening. All external requests must validate URLs, implement DLQ routing, and handle failures gracefully.
Applied to files:
scripts/dlq_replay.pydocs/architecture/DLQ_MESSAGE_FORMAT.mdtests/unit/event_bus/test_kafka_event_bus.py
📚 Learning: 2025-11-24T17:23:49.777Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T17:23:49.777Z
Learning: Applies to **/node_*/v[0-9]*_[0-9]*_[0-9]*/SCHEMA_DECISIONS.md : Each versioned ONEX node implementation directory must include a `SCHEMA_DECISIONS.md` file documenting schema-specific design decisions, implementation notes, and validation strategies
Applied to files:
docs/architecture/DLQ_MESSAGE_FORMAT.md
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/events/**/*.py : Kafka event publishing MUST use OnexEnvelopeV1 format with 13 topics for event streaming at all workflow lifecycle stages
Applied to files:
docs/architecture/DLQ_MESSAGE_FORMAT.mdsrc/omnibase_infra/event_bus/topic_constants.py
📚 Learning: 2025-12-24T17:28:15.618Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.618Z
Learning: Applies to **/protocol*.py : Use `protocol_<name>.py` file naming pattern for standalone protocols or `protocols.py` for domain-grouped protocols with `Protocol<Name>` class pattern. Example: `protocol_event_bus.py` contains `ProtocolEventBus`
Applied to files:
src/omnibase_infra/protocols/protocol_event_bus_like.pysrc/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-12-08T00:48:30.737Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_spi PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-08T00:48:30.737Z
Learning: Applies to src/omnibase_spi/protocols/**/*.py : Protocols must inherit from `typing.Protocol` and use `...` (ellipsis) for method bodies
Applied to files:
src/omnibase_infra/protocols/protocol_event_bus_like.pysrc/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-11-24T17:24:41.687Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T17:24:41.687Z
Learning: Applies to **/protocols/protocol_*.py : Use Protocol for tool interfaces and plugin APIs based on method shape (structural typing), not Pydantic models with inheritance
Applied to files:
src/omnibase_infra/protocols/protocol_event_bus_like.pysrc/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-11-24T17:22:32.195Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/canonical_patterns.mdc:0-0
Timestamp: 2025-11-24T17:22:32.195Z
Learning: Applies to **/protocols/protocol_*.py : Use Protocol from typing module for all interface definitions; never use ABC (Abstract Base Classes) for service interfaces
Applied to files:
src/omnibase_infra/protocols/protocol_event_bus_like.pysrc/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-11-24T17:24:41.687Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T17:24:41.687Z
Learning: Applies to **/protocols/protocol_*.py : Protocol class names must follow the pattern `Protocol<Name>` (e.g., `ProtocolFileGenerator`)
Applied to files:
src/omnibase_infra/protocols/protocol_event_bus_like.py
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to scripts/tests/**/*.sh : Implement comprehensive test suites in scripts/tests/ with separate test files for Kafka, PostgreSQL, Intelligence, and Routing functionality
Applied to files:
tests/unit/event_bus/test_kafka_event_bus.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: KafkaEventBus intentionally violates pattern validator thresholds: 14 methods (threshold: 10) for lifecycle/pub-sub/circuit breaker requirements and 10 __init__ parameters (threshold: 5) for backwards compatibility during config migration. This complexity is acceptable and documented.
Applied to files:
tests/unit/event_bus/test_kafka_event_bus.py
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to agents/**/*.py : Use routing_event_client from agents/lib/routing_event_client.py for agent routing via Kafka with route_via_events() function
Applied to files:
tests/unit/event_bus/test_kafka_event_bus.py
📚 Learning: 2025-12-06T22:21:32.649Z
Learnt from: CR
Repo: OmniNode-ai/omniagent PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-06T22:21:32.649Z
Learning: Applies to shared/protocols/protocol_*.py : Use `protocol_*` prefix for protocol interface files in `shared/protocols/` directory
Applied to files:
src/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Applies to **/protocols/protocol_*.py : Use TYPE_CHECKING guards and forward references for circular import prevention in protocol files
Applied to files:
src/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-12-08T00:48:30.737Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_spi PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-08T00:48:30.737Z
Learning: Applies to src/omnibase_spi/protocols/nodes/*.py : Use Protocol naming convention `Protocol{Type}Node` for node protocols (e.g., `ProtocolComputeNode`, `ProtocolEffectNode`)
Applied to files:
src/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-11-24T16:32:55.606Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T16:32:55.606Z
Learning: Node communication must use event-driven patterns through `ModelEventEnvelope` from `omnibase_core.models.events.model_event_envelope`
Applied to files:
src/omnibase_infra/protocols/__init__.py
📚 Learning: 2025-12-20T04:09:41.832Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_core PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-20T04:09:41.832Z
Learning: Applies to **/*.py : Use PEP 604 union syntax (str | None) instead of Optional or Union types
Applied to files:
tests/unit/validation/test_validator_defaults.py
📚 Learning: 2025-11-28T18:58:53.781Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/canonical_patterns.mdc:0-0
Timestamp: 2025-11-28T18:58:53.781Z
Learning: Applies to **/*.py : Use proper union type definitions and discriminated unions where appropriate
Applied to files:
tests/unit/validation/test_validator_defaults.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: Applies to **/*.py : Use `EnumMessageCategory` for message routing, topic parsing, and dispatcher selection (values: EVENT, COMMAND, INTENT). Use `EnumNodeOutputType` for execution shape validation and handler return type validation (values: EVENT, COMMAND, INTENT, PROJECTION).
Applied to files:
src/omnibase_infra/event_bus/topic_constants.py
🧬 Code graph analysis (10)
src/omnibase_infra/event_bus/models/__init__.py (2)
src/omnibase_infra/event_bus/models/model_dlq_event.py (1)
ModelDlqEvent(50-207)src/omnibase_infra/event_bus/models/model_dlq_metrics.py (1)
ModelDlqMetrics(61-305)
tests/unit/event_bus/test_topic_constants.py (1)
src/omnibase_infra/event_bus/topic_constants.py (4)
build_dlq_topic(131-186)get_dlq_topic_for_original(259-305)is_dlq_topic(241-256)parse_dlq_topic(209-238)
src/omnibase_infra/event_bus/models/model_dlq_event.py (3)
src/omnibase_infra/event_bus/kafka_event_bus.py (2)
default(472-488)environment(513-519)src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py (1)
default(567-592)tests/helpers/deterministic.py (1)
now(136-147)
src/omnibase_infra/event_bus/kafka_event_bus.py (4)
src/omnibase_infra/event_bus/models/model_dlq_event.py (1)
ModelDlqEvent(50-207)src/omnibase_infra/event_bus/models/model_dlq_metrics.py (3)
ModelDlqMetrics(61-305)create_empty(299-305)record_dlq_publish(188-249)src/omnibase_infra/event_bus/models/model_event_headers.py (1)
ModelEventHeaders(16-98)src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py (1)
ModelKafkaEventBusConfig(117-722)
src/omnibase_infra/mixins/__init__.py (1)
src/omnibase_infra/protocols/protocol_event_bus_like.py (1)
ProtocolEventBusLike(36-108)
src/omnibase_infra/event_bus/__init__.py (1)
src/omnibase_infra/event_bus/topic_constants.py (4)
build_dlq_topic(131-186)get_dlq_topic_for_original(259-305)is_dlq_topic(241-256)parse_dlq_topic(209-238)
tests/unit/event_bus/test_kafka_event_bus.py (3)
src/omnibase_infra/event_bus/kafka_event_bus.py (5)
environment(513-519)KafkaEventBus(207-2226)config(495-501)_publish_raw_to_dlq(2065-2226)_publish_to_dlq(1836-2033)src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py (1)
ModelKafkaEventBusConfig(117-722)src/omnibase_infra/event_bus/models/model_event_headers.py (1)
ModelEventHeaders(16-98)
src/omnibase_infra/protocols/__init__.py (1)
src/omnibase_infra/protocols/protocol_event_bus_like.py (1)
ProtocolEventBusLike(36-108)
src/omnibase_infra/event_bus/topic_constants.py (1)
src/omnibase_infra/enums/enum_message_category.py (3)
EnumMessageCategory(34-196)from_topic(121-163)topic_suffix(100-118)
src/omnibase_infra/event_bus/models/model_dlq_metrics.py (4)
src/omnibase_infra/event_bus/kafka_event_bus.py (1)
default(472-488)src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py (1)
default(567-592)tests/helpers/deterministic.py (1)
now(136-147)src/omnibase_infra/projectors/snapshot_publisher_registration.py (1)
topic(224-226)
🔇 Additional comments (30)
src/omnibase_infra/protocols/__init__.py (1)
45-45: LGTM! Protocol import and export follow conventions.The import from the dedicated protocol module and alphabetical placement in
__all__are correct.Also applies to: 55-55
src/omnibase_infra/protocols/protocol_event_bus_like.py (4)
1-33: Excellent module documentation and thread-safety guidance.The module header provides comprehensive thread-safety requirements and recommended patterns for implementers. The references to related documentation and components add valuable context.
35-72: LGTM! Protocol definition follows all conventions.The protocol class adheres to ONEX standards:
- Uses
@runtime_checkabledecorator for structural type checking- Inherits from
typing.Protocol- Follows
Protocol<Name>naming convention- Provides exceptional thread-safety documentation with clear implementer contracts
The comprehensive docstring detailing thread-safety guarantees, locking recommendations, and caller assumptions is particularly valuable for ensuring correct implementations.
74-108: Method signatures are well-designed and follow conventions.Both async methods adhere to protocol standards:
- Use ellipsis (
...) for method bodies per PEP 544- Include thread-safety guarantees in docstrings
- Use PEP 604 union syntax (
bytes | None)- The
envelope: objectparameter is appropriately broad for structural typing, allowing different implementations to use their specific envelope types
110-111: LGTM! Clean export with type annotation.The
__all__declaration correctly exports the protocol with an explicit type annotation.tests/unit/validation/test_validator_defaults.py (1)
43-61: LGTM! Test updates are consistent with the constant change.The test documentation and assertions are correctly updated to reflect the new
INFRA_MAX_UNIONSthreshold of 590. The changes follow the established pattern for documenting baseline history and maintain consistency with the constant update ininfra_validators.py.The assertion message clearly references OMN-949 for traceability.
src/omnibase_infra/validation/infra_validators.py (1)
369-373: LGTM! The constant update is well-documented and properly positioned.The increase from 580 to 590 (+10 unions) aligns perfectly with the OMN-949 DLQ implementation scope. The ~584 baseline is clearly documented with only a 6-union buffer remaining for near-term growth. If substantial development is planned that adds more unions, consider scheduling a proactive threshold bump with the next related ticket.
docs/operations/README.md (1)
9-9: LGTM!The new DLQ Replay Guide entry is properly documented in the runbooks table with a clear description that aligns with the broader DLQ feature introduction.
tests/unit/event_bus/test_topic_constants.py (1)
1-328: LGTM!Comprehensive test coverage for DLQ topic constants and utilities. The tests cover:
- Constant value verification
- Topic building with singular/plural normalization
- Topic parsing and validation
- Error conditions and edge cases
- Integration with
ModelKafkaEventBusConfig.get_dlq_topic()The 44 tests mentioned in the PR summary appear well-structured and follow pytest conventions.
src/omnibase_infra/event_bus/models/__init__.py (1)
17-29: LGTM!Clean public API extension for the DLQ models. The imports and
__all__exports are properly organized, and the module docstring is updated to reflect the newModelDlqEventandModelDlqMetricsexports.docs/architecture/DLQ_MESSAGE_FORMAT.md (1)
1-369: Well-documented DLQ message format specification.The documentation is comprehensive, covering JSON schema, headers, workflow diagrams, monitoring recommendations, and security considerations. This will be valuable for operators and developers working with the DLQ system.
src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py (1)
676-722: LGTM!The
get_dlq_topic()method is well-implemented:
- Explicit
dead_letter_topictakes precedence (documented behavior)- Falls back to ONEX-compliant topic naming via
build_dlq_topic()- Local import to avoid circular dependencies is appropriate
- Comprehensive docstring with examples
The method aligns with the coding guidelines for using specific types and follows the established patterns in this codebase.
tests/unit/event_bus/test_kafka_event_bus.py (3)
1457-1464: Excellent test class documentation.The docstring clearly ties these tests to the OMN-949 acceptance criteria ("No silent drops") and lists the specific scenarios being tested. This makes it easy to understand the purpose and verify coverage.
1606-1686: Comprehensive test for retry behavior.This test correctly verifies that messages with remaining retries are NOT routed to DLQ, which is critical for the retry semantics. The filtering logic at lines 1679-1684 properly isolates DLQ sends from other producer calls.
1743-1806: Critical resilience test for consumer stability.This test ensures that DLQ publish failures don't crash the consumer, which is essential for system stability. The health check verification at the end confirms the bus remains operational.
src/omnibase_infra/event_bus/topic_constants.py (1)
1-322: Well-structured DLQ topic utilities.The module provides a clean API for DLQ topic naming with:
- Proper use of
typing.Finalfor constants- Case-insensitive regex pattern for validation
- Comprehensive input validation in
build_dlq_topic()- Support for both singular and plural category forms
- Integration with
EnumMessageCategoryfor topic inferenceThe constants and utilities follow ONEX naming conventions as specified in the documentation.
src/omnibase_infra/event_bus/models/model_dlq_event.py (1)
50-210: LGTM!Well-designed Pydantic model following ONEX patterns:
- Immutable with
frozen=Truefor thread-safe callback sharing- Comprehensive field documentation
- Proper type hints using PEP 604 syntax (
str | None)to_log_context()for structured logging with proper serializationis_criticalproperty for alerting logicmetric_labelsproperty following Prometheus/OpenTelemetry conventionsThe model provides complete context for DLQ event observability without coupling to specific alerting implementations.
src/omnibase_infra/event_bus/__init__.py (1)
35-67: LGTM! Clean module re-export pattern for DLQ utilities.The new DLQ constants and helper functions are properly imported and exported with clear organization separating Event Bus, Topic Constants, and Topic Functions categories.
src/omnibase_infra/event_bus/models/model_dlq_metrics.py (3)
61-93: LGTM! Well-designed metrics model following copy-on-write pattern.The model correctly implements immutable-style updates with proper Pydantic configuration (
extra="forbid",validate_assignment=True). Thread-safety is well-documented.
155-186: LGTM! Computed properties handle edge cases correctly.Division by zero is properly guarded, and the default values (0.0 for avg_latency/failure_rate, 1.0 for success_rate) are sensible.
188-249: LGTM! Copy-on-write implementation is correct.The method properly creates shallow copies of dictionaries and uses
model_copy(update=...)to return a new instance without mutating the original.src/omnibase_infra/event_bus/kafka_event_bus.py (6)
192-202: LGTM! Clean type alias and import additions.The
DlqCallbackTypealias properly defines the async callback signature for DLQ event hooks.
412-418: LGTM! Proper initialization with separate locks for DLQ state.Using dedicated locks for metrics and callbacks prevents contention with the main event bus lock and avoids potential deadlocks.
542-586: LGTM! Well-designed callback registration with proper cleanup.The async unregister function correctly checks for callback presence before removal, and all operations are protected by the callbacks lock.
1338-1347: LGTM! Deserialization errors correctly routed to DLQ.This ensures no silent message drops when message conversion fails, which addresses the PR objective of preventing silent failures.
1358-1398: LGTM! Retry exhaustion check correctly implements ModelEventHeaders contract.The logic
retry_count >= max_retriesaligns with the documented behavior inModelEventHeaders. The note about republishing being the caller's responsibility is helpful.
2035-2063: LGTM! Callback invocation with proper error isolation.Each callback is invoked sequentially with exceptions logged but not propagated, ensuring one failing callback doesn't prevent others from executing.
scripts/dlq_replay.py (3)
168-174: LGTM! Appropriate CLI error handling for invalid input.Exiting with code 1 and logging a warning for invalid correlation ID format is reasonable CLI behavior.
788-815: LGTM! Clean async CLI entry point.The main function properly handles argument parsing, verbose logging toggle, command dispatch, and exit code propagation.
205-214: LGTM! Good non-retryable error type set.The
NON_RETRYABLE_ERRORSset aligns with ONEX error types and correctly identifies permanent failures that shouldn't be replayed.
PR Review: DLQ Implementation for KafkaEventBus (OMN-949)SummaryThis PR implements comprehensive Dead Letter Queue (DLQ) support for the KafkaEventBus. The implementation is well-structured and follows ONEX conventions with strong typing, proper error handling, and good documentation. Strengths1. Excellent ONEX Compliance
2. Comprehensive Documentation
3. Topic Naming ArchitectureThe topic_constants.py module is well-designed with clear separation of constants and utilities, regex validation, and helper functions. 4. Test Coverage
Issues and Recommendations1. MEDIUM: DLQ Replay Script ImplementationLocation: scripts/dlq_replay.py Issue: The replay script contains only stub implementations - DLQConsumer.consume_messages() and DLQProducer.replay_message() have no actual Kafka integration. Recommendation: Add a clear STUB IMPLEMENTATION - NOT PRODUCTION READY warning in the file header. Consider moving full implementation to a follow-up ticket (OMN-949-replay). The acceptance criteria state replay mechanism defined which is satisfied, but operators should know this is not functional. 2. MINOR: Error Sanitization VerificationLocation: ModelDlqEvent.error_message field Question: Verify that KafkaEventBus._publish_to_dlq() actually implements sanitization per CLAUDE.md guidelines (remove passwords, API keys, tokens, PII, connection strings). 3. MINOR: Missing Environment Name ValidationLocation: topic_constants.py:186 in build_dlq_topic() Current: Only validates environment is non-empty Recommendation: Add environment format validation to match DLQ_TOPIC_PATTERN requirements (alphanumeric, underscores, hyphens only). 4. MINOR: Generator Stub PatternLocation: dlq_replay.py:263 The stub generator pattern with return followed by unreachable yield is confusing. Use if False: yield instead for clearer intent. Security ConsiderationsGood Practices
Verify Before Merge
Test Coverage AssessmentWell-Covered
Needs Verification
Acceptance Criteria Review
Recommendations for MergeRecommended Actions Before Merge:
Follow-Up Tickets:
Overall AssessmentRating: Strong Implementation - Approve with Minor Recommendations This PR demonstrates excellent ONEX architecture understanding with strong typing, immutability patterns, comprehensive documentation, and good test coverage for topic utilities. Recommendation: APPROVE - This is solid infrastructure work that will improve operational resilience. The replay script stub is acceptable for this PR as the mechanism is defined per acceptance criteria. Excellent work on this feature! The DLQ architecture is well-designed and follows ONEX patterns consistently. |
- Fix timing concern: move await future outside producer lock - Add error sanitization to _publish_to_dlq with 47 sensitive patterns - Add STUB IMPLEMENTATION warning to dlq_replay.py header - Add environment name validation to build_dlq_topic() - Complete dlq_replay.py with aiokafka integration - Add 17 DLQ integration tests - Remove protocol re-export from mixins module - Remove fragile line number references from SQL comments - Add metrics/callback integration to _publish_raw_to_dlq - Convert dlq_replay.py dataclasses to Pydantic models - Fix unusual generator stub pattern with clear documentation
…re-dlq-for-permanently-failing-intents Resolved conflicts: - src/omnibase_infra/utils/__init__.py: Combined imports from both branches (util_error_sanitization + util_semver exports) - src/omnibase_infra/validation/infra_validators.py: Combined INFRA_MAX_UNIONS threshold (603 = baseline + DLQ + main branch additions) - tests/unit/validation/test_validator_defaults.py: Updated test assertions to match combined threshold
PR Review: Dead Letter Queue (DLQ) Implementation [OMN-949]SummaryThis PR implements comprehensive Dead Letter Queue (DLQ) support for permanently failing messages in the KafkaEventBus. The implementation is well-structured and follows ONEX conventions with strong typing, proper error handling, and comprehensive testing. Below are detailed findings across code quality, security, performance, and test coverage. ✅ Strengths1. Excellent ONEX Compliance
2. Security Best Practices
3. Comprehensive Documentation
4. Testing Coverage
🔍 Issues FoundCritical IssuesNone found - This is a solid implementation. High Priority Issues1. Potential Race Condition in DLQ Metrics Update (kafka_event_bus.py:2036+)The metrics update uses copy-on-write, but there's a potential race if multiple consumers publish to DLQ concurrently: # Current implementation (line 2036+)
async with self._dlq_lock:
current_metrics = self._dlq_metrics
# ... metrics update logic ...
self._dlq_metrics = updated_metricsIssue: If Recommendation: This is likely acceptable given the lock, but consider documenting the threading model or using atomic counters for high-throughput scenarios. Reference: CLAUDE.md "Circuit Breaker Thread Safety" pattern (similar lock-holding requirements) 2. Missing DLQ Topic Validation at Configuration Time (model_kafka_event_bus_config.py)The dead_letter_topic: str | None = Field(
default=None,
description="Dead letter queue topic name for failed messages",
)Issue: Invalid topic names (e.g., spaces, special characters) will only fail at runtime when publishing to DLQ, not at configuration validation time. Recommendation: Add a Pydantic validator to check topic naming conventions: @field_validator("dead_letter_topic")
@classmethod
def validate_dlq_topic(cls, v: str | None) -> str | None:
if v is not None and not DLQ_TOPIC_PATTERN.match(v):
# Consider allowing both ONEX format and free-form for backwards compatibility
logger.warning(f"DLQ topic does not match ONEX pattern: {v}")
return vReference: 3. Error Sanitization May Be Overly Aggressive (util_error_sanitization.py:42-92)The SENSITIVE_PATTERNS: tuple[str, ...] = (
"auth", # Too broad - matches "authentication", "author", "authorize"
"bearer", # Matches legitimate error context
...
)Issue: This could redact legitimate error messages that don't contain credentials:
Recommendation: Make patterns more specific or use word boundaries: # More specific patterns
"password=",
"token=",
"authorization:", # Header context
"bearer ", # Token prefixTrade-off: Security vs. debuggability. Current approach is conservative (better safe than sorry), but consider documenting the trade-off. Medium Priority Issues4. DLQ Replay Script Hardcodes
|
| Criterion | Status | Evidence |
|---|---|---|
| DLQ topic configured | ✅ PASS | topic_constants.py, build_dlq_topic() |
| Alerting integrated | ✅ PASS | ModelDlqEvent, callback hooks in kafka_event_bus.py |
| No silent drops | ✅ PASS | Deserialization errors fixed, retry exhaustion check added |
| DLQ message format documented | ✅ PASS | DLQ_MESSAGE_FORMAT.md with JSON schema |
| Replay mechanism defined | ✅ PASS | dlq_replay.py CLI utility, filtering, rate limiting |
All acceptance criteria met.
🎯 Recommendations Summary
Must Fix Before Merge
- None - This PR is production-ready as-is.
Should Consider
- Add DLQ topic validation to
ModelKafkaEventBusConfig(High Priority Add Claude Code GitHub Workflow #2) - Review error sanitization patterns for false positives (High Priority feat: RedPanda Event Bus Integration with Fail-Fast Infrastructure #3)
- Add test for concurrent DLQ metrics (Missing Test Archive legacy codebase for clean migration #10)
- Fix offset type inconsistency (Medium Priority feat: Complete Hook Node Protocol Integration and Production Readiness #6)
Nice to Have
- Document race condition potential in metrics update (High Priority feat: PostgreSQL Adapter with Comprehensive Tests and Structured Logging #1)
- Extract replay helper methods for better testability (Low Priority feat: migrate 96 models to domain-organized structure #8)
- Add DLQ publish timeout test (Missing Test Core Domain Models: Protocol-Based Health Architecture #11)
- Add malformed message handling test (Missing Test feat: Infrastructure Domain Models with Service Integration Architecture #12)
✨ Final Verdict
APPROVED with minor suggestions
This is a high-quality implementation that demonstrates:
- Strong adherence to ONEX coding standards
- Comprehensive security considerations
- Thorough documentation and testing
- Production-ready error handling
The issues identified are minor improvements rather than blockers. The PR successfully addresses all acceptance criteria for OMN-949 and provides a solid foundation for DLQ-based message recovery.
Great work! 🚀
📚 References
- CLAUDE.md: ONEX coding standards (strong typing, naming conventions, error patterns)
docs/patterns/error_handling_patterns.md: Error hierarchy and usagedocs/patterns/retry_backoff_compensation_strategy.md: Retry policiesdocs/architecture/DLQ_MESSAGE_FORMAT.md: DLQ schema and security (added in this PR)
There was a problem hiding this comment.
Actionable comments posted: 7
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/event_bus/kafka_event_bus.py (1)
1327-1399: Use message-levelcorrelation_idfor DLQ publishing instead of consumer-task UUIDIn
_consume_loop, DLQ routing for handler failures currently passes the consumer-loopcorrelation_id(created in_start_consumer_for_topic) into_publish_to_dlq, and that value is then propagated into:
- DLQ payload field
correlation_id- DLQ headers
ModelEventHeaders(correlation_id=...)ModelDlqEvent.correlation_idThis loses the original message’s
correlation_idfromevent_message.headersand makes it harder to correlate DLQ entries back to upstream requests, conflicting with the guideline to always propagate the incomingcorrelation_idwhen available. Based on learnings, correlation IDs should reflect the message, not the consumer task.A targeted adjustment keeps the consumer correlation for consumer-loop logs but uses the message-level
correlation_idfor DLQ:Proposed change (core idea)
# In _consume_loop, inside the subscriber exception block - if retries_exhausted: - await self._publish_to_dlq( - original_topic=topic, - failed_message=event_message, - error=e, - correlation_id=correlation_id, - ) + if retries_exhausted: + await self._publish_to_dlq( + original_topic=topic, + failed_message=event_message, + error=e, + correlation_id=event_message.headers.correlation_id, + )You may also want to consider using
event_message.headers.correlation_idin the log context of the handler failure (in addition to the consumer-loop correlation ID) for better traceability.Based on learnings, correlation IDs from
ModelEventHeadersshould flow into all error contexts and DLQ artifacts whenever available.Also applies to: 1837-1862, 1890-1908, 2019-2034
🧹 Nitpick comments (10)
src/omnibase_infra/utils/util_error_sanitization.py (1)
39-92: Error sanitizer implementation aligns with security guidelinesThe pattern set and
sanitize_error_messagebehavior (case-insensitive scanning, full-message redaction on hit, truncation with marker, and type-prefix format) are aligned with the “when in doubt, redact” guidance and avoidAnyin the API. This is a solid central utility for DLQ/log usage.If you later need environment-specific or service-specific patterns, consider allowing injection/extension of
SENSITIVE_PATTERNS(e.g., via a small registration API) while keeping the default tuple immutable.Also applies to: 95-155, 157-160
tests/integration/event_bus/test_dlq_integration.py (1)
316-396: DLQ publish test does not assert that a DLQ message was received
test_dlq_publish_on_handler_failureis documented as verifying that messages are published to DLQ after retry exhaustion, but it only asserts that the failing handler was called at least once and never inspectsdlq_messages_received.To make this test actually validate DLQ behavior (while still being tolerant to timing):
Possible improvement
- # Verify handler was called and failed - assert handler_call_count >= 1, "Handler should have been called at least once" + # Verify handler was called and failed + assert handler_call_count >= 1, "Handler should have been called at least once" + + # Best-effort DLQ assertion (only if we observed a DLQ message) + if dlq_messages_received: + assert len(dlq_messages_received) >= 1Or, if flakiness is a concern, consider renaming the test to reflect that it only asserts handler failure unless a DLQ message is explicitly checked.
src/omnibase_infra/event_bus/kafka_event_bus.py (1)
413-420: DLQ metrics and callback infrastructure is well-structuredInitializing DLQ metrics with
ModelDlqMetrics.create_empty(), guarding updates with_dlq_metrics_lock, and exposing a copy via thedlq_metricsproperty gives a clean, mutation-safe surface. Theregister_dlq_callbackAPI with an async unregister function and_invoke_dlq_callbackserror isolation pattern is also solid and straightforward to consume.If you later see heavy concurrent access to
dlq_metrics, you might also guard the getter with_dlq_metrics_lockto guarantee a snapshot aligned exactly with the last write, although the current copy-on-write pattern is already safe enough for most use cases.Also applies to: 531-587
src/omnibase_infra/event_bus/topic_constants.py (1)
65-101: DLQ topic utilities are consistent with ONEX naming and enums
build_dlq_topic,parse_dlq_topic,is_dlq_topic, andget_dlq_topic_for_originalcleanly encode the<env>.dlq.<category>.<version>convention and integrate properly withEnumMessageCategory.from_topicviacategory.topic_suffix. Environment validation viaENV_PATTERNmatches the documented constraints, and the mapping supports both singular and plural category forms.You may later want to validate the optional
versionargument inbuild_dlq_topicagainst thev\d+pattern (or reuseDLQ_TOPIC_PATTERN) to prevent accidental non-version strings from slipping into DLQ topics.Also applies to: 112-143, 145-215, 237-266, 269-285, 287-351
scripts/dlq_replay.py (6)
113-146: Consider stricter validation instead of silent type coercion.The
from_kafka_messagemethod uses defensive type coercion (str(),int(str())) which could hide data quality issues in the DLQ payload. If the DLQ message schema is malformed, this will silently convert incorrect types rather than failing fast.Consider:
- Log a warning when type coercion occurs (e.g., when
correlation_idis invalid and a new UUID is generated).- Validate the payload schema more strictly and raise descriptive errors for malformed messages.
- Add a counter metric for malformed DLQ messages to detect payload issues.
Example: Add logging for type coercion
correlation_id_str = payload.get("correlation_id", "") try: correlation_id = UUID(str(correlation_id_str)) except (ValueError, AttributeError): + logger.warning( + f"Invalid correlation_id in DLQ message, generating new UUID", + extra={"invalid_value": str(correlation_id_str)[:50]} + ) correlation_id = uuid4()
392-399: Extract hardcoded max_request_size to configuration.The 10MB
max_request_sizeis hardcoded but should be configurable or reference a shared constant. According to the learnings, messages exceedingKAFKA_MAX_REQUEST_SIZE(default 10MB) are routed to an oversized DLQ topic. The replay script should respect this same limit.🔎 Proposed fix
Add to
ModelReplayConfig:class ModelReplayConfig(BaseModel): """Configuration for DLQ replay operation.""" bootstrap_servers: str = "localhost:9092" dlq_topic: str = "dlq-events" max_replay_count: int = 5 rate_limit_per_second: float = 100.0 + max_message_size_bytes: int = 10485760 # 10MB default dry_run: bool = FalseThen use in producer:
self._producer = AIOKafkaProducer( bootstrap_servers=self.config.bootstrap_servers, acks="all", enable_idempotence=True, - max_request_size=10485760, # 10MB max message size + max_request_size=self.config.max_message_size_bytes, request_timeout_ms=30000, # 30 second timeout )Based on learnings, oversized messages are automatically routed to a separate DLQ topic.
442-446: Consider optimizing rate limiting to avoid delaying the first message.The rate limiting logic sleeps before publishing, which means the first message is unnecessarily delayed. Standard rate limiting should allow the first message to go through immediately and only throttle subsequent messages.
🔎 Proposed fix
# Rate limiting elapsed = (datetime.now(UTC) - self._last_publish).total_seconds() - if elapsed < self._interval: + if elapsed < self._interval and self._last_publish != datetime.min.replace(tzinfo=UTC): await asyncio.sleep(self._interval - elapsed)Or initialize
_last_publishto a past timestamp:- self._last_publish = datetime.now(UTC) + self._last_publish = datetime.min.replace(tzinfo=UTC)
673-677: Use safe text truncation for multi-byte characters.Line 676 uses
[:80]to truncatefailure_reason, which can split multi-byte UTF-8 characters and cause display issues or crashes. Use a safer truncation method.🔎 Proposed fix
Add a helper function:
def safe_truncate(text: str, max_length: int) -> str: """Safely truncate text without breaking multi-byte characters.""" if len(text) <= max_length: return text # Encode and decode to ensure we don't break multi-byte sequences return text[:max_length].encode('utf-8', errors='ignore').decode('utf-8', errors='ignore') + '...'Then use it:
- print(f" Reason: {message.failure_reason[:80]}...") + print(f" Reason: {safe_truncate(message.failure_reason, 77)}")
945-950: Consider more specific type annotation for signal handler frame parameter.The
frameparameter is typed asobject, but the standard library usestypes.FrameType | Nonefor signal handlers. Whileobjectis correct and works, a more specific type improves IDE support.🔎 Proposed fix
+import types + def signal_handler(signum: int, frame: object) -> None: - """Handle shutdown signals (SIGINT/SIGTERM).""" +def signal_handler(signum: int, frame: types.FrameType | None) -> None: + """Handle shutdown signals (SIGINT/SIGTERM)."""
148-169: Consider integrating with ONEX DLQ topic naming conventions.The script accepts
--dlq-topicas a plain string, but according to the PR summary, DLQ topics follow the ONEX naming pattern{env}.dlq.{category}.v1as defined intopic_constants.py. Consider adding:
- Topic validation to ensure the provided DLQ topic matches the expected pattern.
- Helper to construct DLQ topic names from environment and category.
- Auto-discovery of available DLQ topics from Kafka metadata.
This would prevent users from accidentally pointing the replay script at non-DLQ topics or using incorrect topic names.
Example: Topic validation
def validate_dlq_topic(topic: str) -> bool: """Validate that topic matches ONEX DLQ naming pattern.""" import re pattern = r'^(dev|staging|prod)\.dlq\.[a-z]+\.v\d+$' return re.match(pattern, topic) is not None # In ModelReplayConfig.from_args(): if not validate_dlq_topic(args.dlq_topic): logger.warning( f"DLQ topic '{args.dlq_topic}' doesn't match expected pattern: " "{{env}}.dlq.{{category}}.v{{version}}" )
📜 Review details
Configuration used: defaults
Review profile: CHILL
Plan: Lite
📒 Files selected for processing (14)
docs/architecture/DLQ_MESSAGE_FORMAT.mdscripts/dlq_replay.pysrc/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/event_bus/topic_constants.pysrc/omnibase_infra/protocols/__init__.pysrc/omnibase_infra/runtime/message_dispatch_engine.pysrc/omnibase_infra/schemas/schema_registration_projection.sqlsrc/omnibase_infra/services/timeout_emitter.pysrc/omnibase_infra/utils/__init__.pysrc/omnibase_infra/utils/util_error_sanitization.pysrc/omnibase_infra/validation/infra_validators.pytests/integration/event_bus/test_dlq_integration.pytests/unit/utils/test_util_error_sanitization.pytests/unit/validation/test_validator_defaults.py
✅ Files skipped from review due to trivial changes (2)
- src/omnibase_infra/schemas/schema_registration_projection.sql
- docs/architecture/DLQ_MESSAGE_FORMAT.md
🚧 Files skipped from review as they are similar to previous changes (2)
- tests/unit/validation/test_validator_defaults.py
- src/omnibase_infra/protocols/init.py
🧰 Additional context used
📓 Path-based instructions (3)
**/*.py
📄 CodeRabbit inference engine (CLAUDE.md)
**/*.py: NEVER useAnytype - Always use specific types. All data structures must be proper Pydantic models.
UseEnumMessageCategoryfor message routing, topic parsing, and dispatcher selection (values: EVENT, COMMAND, INTENT). UseEnumNodeOutputTypefor execution shape validation and handler return type validation (values: EVENT, COMMAND, INTENT, PROJECTION).
Use PEP 604 union syntaxX | Nonefor nullable types instead ofOptional[X]. Example:def get_user(id: str) -> User | None:instead ofdef get_user(id: str) -> Optional[User]:
All services MUST use ModelONEXContainer for dependency injection. Bootstrap withcontainer = ModelONEXContainer()and resolve services viacontainer.service_registry.resolve_service(ServiceType)
Always propagatecorrelation_idfrom incoming requests to error context. Auto-generate usinguuid4()if no correlation_id exists. Use UUID format for all new correlation IDs.
NEVER include in error messages or context: passwords, API keys, tokens, secrets, full connection strings with credentials, PII, internal IP addresses, private keys, certificates, session tokens, or cookies. SAFE to include: service names, operation names, correlation IDs, error codes, sanitized hostnames, port numbers, retry counts, timeout values, resource identifiers (non-sensitive).
UseProtocolConfigurationErrorfor invalid config,SecretResolutionErrorfor missing secrets,InfraConnectionErrorfor connection failures,InfraTimeoutErrorfor operation timeouts,InfraAuthenticationErrorfor auth failures,InfraUnavailableErrorfor unavailable resources.
Use graceful degradation forInfraTimeoutError. Pattern: try primary source with timeout, fall back to cache/secondary source on timeout, aggregate results with degradation flag.
Files:
src/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/validation/infra_validators.pysrc/omnibase_infra/utils/util_error_sanitization.pysrc/omnibase_infra/services/timeout_emitter.pytests/integration/event_bus/test_dlq_integration.pysrc/omnibase_infra/event_bus/topic_constants.pytests/unit/utils/test_util_error_sanitization.pyscripts/dlq_replay.pysrc/omnibase_infra/runtime/message_dispatch_engine.pysrc/omnibase_infra/utils/__init__.py
**/*infra*.py
📄 CodeRabbit inference engine (CLAUDE.md)
**/*infra*.py: IncludeModelInfraErrorContextwith every infrastructure error. Context should include:transport_type(EnumInfraTransportType),operation,target_name, andcorrelation_id. Example:raise InfraConnectionError('Failed to connect', context=context)
UseEnumInfraTransportTypefor transport identification in error context. Values include: HTTP, DATABASE, KAFKA, CONSUL, VAULT, VALKEY, GRPC
Files:
src/omnibase_infra/validation/infra_validators.py
**/util_*.py
📄 CodeRabbit inference engine (CLAUDE.md)
Use
util_<name>.pyfile naming pattern for utility functions. Example:util_retry.py→retry_with_backoff()
Files:
src/omnibase_infra/utils/util_error_sanitization.py
🧠 Learnings (26)
📓 Common learnings
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to **/*.py : Messages exceeding `KAFKA_MAX_REQUEST_SIZE` (default 10MB) are automatically routed to oversized DLQ topic. Use `scripts/process_oversized_dlq.py` to handle these messages. See `docs/guides/DLQ_HANDLING.md` for complete documentation.
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to {services/**/*.py,scripts/bulk_ingest_repository.py} : Implement fail-closed configuration for security hardening. All external requests must validate URLs, implement DLQ routing, and handle failures gracefully.
Learnt from: CR
Repo: OmniNode-ai/omniintelligence PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-07T17:50:13.678Z
Learning: Implement Kafka event-driven architecture with proper topic naming using prefix dev.archon-intelligence. and proper event flow pattern with Effect nodes consuming events, processing, and publishing results with Dead Letter Queue routing
📚 Learning: 2025-11-29T17:13:38.776Z
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to **/*.py : Messages exceeding `KAFKA_MAX_REQUEST_SIZE` (default 10MB) are automatically routed to oversized DLQ topic. Use `scripts/process_oversized_dlq.py` to handle these messages. See `docs/guides/DLQ_HANDLING.md` for complete documentation.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pytests/integration/event_bus/test_dlq_integration.pyscripts/dlq_replay.py
📚 Learning: 2025-12-06T22:21:32.649Z
Learnt from: CR
Repo: OmniNode-ai/omniagent PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-06T22:21:32.649Z
Learning: Applies to nodes/**/*.py : Use event bus mixins from `omnibase_core` for Kafka publishing instead of direct Kafka clients
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-12-07T17:50:13.678Z
Learnt from: CR
Repo: OmniNode-ai/omniintelligence PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-07T17:50:13.678Z
Learning: Implement Kafka event-driven architecture with proper topic naming using prefix dev.archon-intelligence. and proper event flow pattern with Effect nodes consuming events, processing, and publishing results with Dead Letter Queue routing
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pytests/integration/event_bus/test_dlq_integration.py
📚 Learning: 2025-11-29T17:13:38.776Z
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to {services/**/*.py,scripts/bulk_ingest_repository.py} : Implement fail-closed configuration for security hardening. All external requests must validate URLs, implement DLQ routing, and handle failures gracefully.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to **/*.py : Publish intelligence requests to Kafka event bus using topics: dev.archon-intelligence.intelligence.code-analysis-{requested,completed,failed}.v1 for consistency and event-driven architecture
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pytests/integration/event_bus/test_dlq_integration.pysrc/omnibase_infra/event_bus/topic_constants.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: Applies to **/*infra*.py : Include `ModelInfraErrorContext` with every infrastructure error. Context should include: `transport_type` (EnumInfraTransportType), `operation`, `target_name`, and `correlation_id`. Example: `raise InfraConnectionError('Failed to connect', context=context)`
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: Applies to **/*.py : Use `ProtocolConfigurationError` for invalid config, `SecretResolutionError` for missing secrets, `InfraConnectionError` for connection failures, `InfraTimeoutError` for operation timeouts, `InfraAuthenticationError` for auth failures, `InfraUnavailableError` for unavailable resources.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/events/**/*.py : Kafka event publishing MUST use OnexEnvelopeV1 format with 13 topics for event streaming at all workflow lifecycle stages
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/event_bus/topic_constants.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: KafkaEventBus intentionally violates pattern validator thresholds: 14 methods (threshold: 10) for lifecycle/pub-sub/circuit breaker requirements and 10 __init__ parameters (threshold: 5) for backwards compatibility during config migration. This complexity is acceptable and documented.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: Applies to **/*.py : NEVER include in error messages or context: passwords, API keys, tokens, secrets, full connection strings with credentials, PII, internal IP addresses, private keys, certificates, session tokens, or cookies. SAFE to include: service names, operation names, correlation IDs, error codes, sanitized hostnames, port numbers, retry counts, timeout values, resource identifiers (non-sensitive).
Applied to files:
src/omnibase_infra/utils/util_error_sanitization.pytests/unit/utils/test_util_error_sanitization.pysrc/omnibase_infra/utils/__init__.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Applies to **/*.py : Import protocols from `omnibase.protocol.protocol_*` paths
Applied to files:
src/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-11-24T17:24:41.687Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T17:24:41.687Z
Learning: Applies to **/*.py : Import protocols from `omnibase.protocol.protocol_<name>` module paths
Applied to files:
src/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-12-08T00:48:30.737Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_spi PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-08T00:48:30.737Z
Learning: Applies to src/omnibase_spi/**/*.py : SPI modules may import from `omnibase_core` for type hints and model runtime usage (allowed and required)
Applied to files:
src/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Applies to **/protocols/protocol_*.py : Use TYPE_CHECKING guards and forward references for circular import prevention in protocol files
Applied to files:
src/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Use Protocol for interface definitions when implementations may live outside core codebase; use Pydantic models only for base classes with shared logic
Applied to files:
src/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-11-24T17:22:32.195Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/canonical_patterns.mdc:0-0
Timestamp: 2025-11-24T17:22:32.195Z
Learning: Applies to **/protocols/protocol_*.py : Use Protocol from typing module for all interface definitions; never use ABC (Abstract Base Classes) for service interfaces
Applied to files:
src/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-12-08T00:48:30.737Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_spi PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-08T00:48:30.737Z
Learning: Applies to src/omnibase_spi/protocols/**/*.py : Protocols must inherit from `typing.Protocol` and use `...` (ellipsis) for method bodies
Applied to files:
src/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-12-08T00:48:30.737Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_spi PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-08T00:48:30.737Z
Learning: Import `omnibase_core` models and types only for type hints and runtime usage - follow the SPI → Core dependency direction
Applied to files:
src/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-11-24T17:24:41.687Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T17:24:41.687Z
Learning: Applies to **/protocols/protocol_*.py : Use Protocol for tool interfaces and plugin APIs based on method shape (structural typing), not Pydantic models with inheritance
Applied to files:
src/omnibase_infra/services/timeout_emitter.py
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to scripts/tests/**/*.sh : Implement comprehensive test suites in scripts/tests/ with separate test files for Kafka, PostgreSQL, Intelligence, and Routing functionality
Applied to files:
tests/integration/event_bus/test_dlq_integration.py
📚 Learning: 2025-11-29T22:07:25.230Z
Learnt from: CR
Repo: OmniNode-ai/omniintelligence PR: 0
File: migration_sources/omniarchon/CLAUDE.md:0-0
Timestamp: 2025-11-29T22:07:25.230Z
Learning: Applies to migration_sources/omniarchon/**/tests/**/*.py : All integration tests must verify correct Kafka port usage for context (9092 for Docker, 29092 for host). Test both local (qdrant, memgraph) and remote (PostgreSQL, Redpanda) database connectivity. Never assume test environment configuration.
Applied to files:
tests/integration/event_bus/test_dlq_integration.py
📚 Learning: 2025-11-24T17:24:41.687Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T17:24:41.687Z
Learning: Applies to src/omnibase/constants/*_constants.py : Constants files must follow the naming pattern `<domain>_constants.py` and be located in `src/omnibase/constants/`
Applied to files:
src/omnibase_infra/event_bus/topic_constants.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Applies to src/omnibase/constants/*_constants.py : Constants files must follow the naming pattern `<domain>_constants.py` and be located in `src/omnibase/constants/` directory
Applied to files:
src/omnibase_infra/event_bus/topic_constants.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: Applies to **/*.py : Use `EnumMessageCategory` for message routing, topic parsing, and dispatcher selection (values: EVENT, COMMAND, INTENT). Use `EnumNodeOutputType` for execution shape validation and handler return type validation (values: EVENT, COMMAND, INTENT, PROJECTION).
Applied to files:
src/omnibase_infra/event_bus/topic_constants.pysrc/omnibase_infra/runtime/message_dispatch_engine.py
📚 Learning: 2025-12-24T17:28:15.619Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-24T17:28:15.619Z
Learning: Applies to **/*.py : Always propagate `correlation_id` from incoming requests to error context. Auto-generate using `uuid4()` if no correlation_id exists. Use UUID format for all new correlation IDs.
Applied to files:
src/omnibase_infra/utils/__init__.py
🧬 Code graph analysis (6)
src/omnibase_infra/event_bus/kafka_event_bus.py (8)
src/omnibase_infra/runtime/runtime_host_process.py (1)
event_bus(348-354)src/omnibase_infra/event_bus/models/model_dlq_event.py (1)
ModelDlqEvent(50-207)src/omnibase_infra/event_bus/models/model_dlq_metrics.py (3)
ModelDlqMetrics(61-305)create_empty(299-305)record_dlq_publish(188-249)src/omnibase_infra/event_bus/models/model_event_headers.py (1)
ModelEventHeaders(16-98)src/omnibase_infra/utils/util_error_sanitization.py (1)
sanitize_error_message(95-154)src/omnibase_infra/projectors/snapshot_publisher_registration.py (1)
topic(224-226)src/omnibase_infra/services/timeout_emitter.py (1)
environment(286-288)src/omnibase_infra/event_bus/inmemory_event_bus.py (1)
environment(178-184)
src/omnibase_infra/services/timeout_emitter.py (1)
src/omnibase_infra/protocols/protocol_event_bus_like.py (1)
ProtocolEventBusLike(36-108)
src/omnibase_infra/event_bus/topic_constants.py (1)
src/omnibase_infra/enums/enum_message_category.py (3)
EnumMessageCategory(34-196)from_topic(121-163)topic_suffix(100-118)
tests/unit/utils/test_util_error_sanitization.py (1)
src/omnibase_infra/utils/util_error_sanitization.py (1)
sanitize_error_message(95-154)
src/omnibase_infra/runtime/message_dispatch_engine.py (2)
src/omnibase_infra/enums/enum_message_category.py (1)
EnumMessageCategory(34-196)src/omnibase_infra/utils/util_error_sanitization.py (1)
sanitize_error_message(95-154)
src/omnibase_infra/utils/__init__.py (1)
src/omnibase_infra/utils/util_error_sanitization.py (1)
sanitize_error_message(95-154)
🔇 Additional comments (10)
src/omnibase_infra/services/timeout_emitter.py (1)
42-42: LGTM! Import follows protocol patterns.The import of
ProtocolEventBusLikeis correctly sourced fromomnibase_infra.protocolsand aligns with established protocol import patterns. The addition supports the broader DLQ protocol surface introduced in this PR.Based on learnings, protocols should be imported from
omnibase_infra.protocolsmodule paths.src/omnibase_infra/validation/infra_validators.py (1)
371-375: LGTM! Threshold update is well-documented.The INFRA_MAX_UNIONS increase from 586 to 603 is appropriate for the DLQ configuration changes in this PR. The documentation clearly tracks the history, provides rationale, and maintains a small buffer above the current baseline (~597). The validation tests have already passed per the PR summary, confirming the accuracy of this threshold.
src/omnibase_infra/utils/__init__.py (1)
7-9: Centralized export of error sanitization utilities looks goodRe-exporting
SENSITIVE_PATTERNSandsanitize_error_messagefromomnibase_infra.utilskeeps callers decoupled from the concrete module path and aligns with how other utilities are surfaced. No issues noticed.Also applies to: 18-21, 28-39
src/omnibase_infra/runtime/message_dispatch_engine.py (1)
149-151: Consistent use of shared sanitizer in dispatch errorsReplacing local sanitization with
sanitize_error_messagefromomnibase_infra.utilsin dispatcher failures, output validation, and result-construction errors centralizes security logic and avoids duplicate pattern lists. The way you pass the sanitized strings into both logs and metrics matches the error-sanitization guidelines.Also applies to: 1148-1152, 1306-1313, 1364-1387
tests/unit/utils/test_util_error_sanitization.py (1)
20-267: Sanitizer unit test coverage is comprehensiveThe tests exercise all critical behaviors of
sanitize_error_messageandSENSITIVE_PATTERNS(credential patterns, connection strings, PEM headers, truncation, case-insensitivity, and DLQ-style error cases). This gives good confidence that new DLQ paths won’t leak obvious secrets.tests/integration/event_bus/test_dlq_integration.py (1)
381-389: No changes needed. In Python 3.11+,asyncio.wait_forraises the built-inTimeoutError, not a distinctasyncio.TimeoutError. Theasyncio.TimeoutErroris now a deprecated alias of the built-in exception. Catchingexcept TimeoutError:is the correct approach and will properly handle timeouts.Likely an incorrect or invalid review comment.
scripts/dlq_replay.py (4)
237-347: LGTM - DLQConsumer implementation is robust.The consumer implementation correctly handles:
- Connection errors with proper exception types
- Graceful shutdown with
asyncio.CancelledError- Null message values and JSON decode errors
- Independent consumer groups using PID
493-530: LGTM - Filtering logic is clear and correct.The
should_replayfunction correctly implements:
- Max replay count enforcement
- Non-retryable error type filtering
- Topic, error type, and correlation ID filters
537-643: LGTM - Executor orchestration is well-structured.The replay executor correctly:
- Manages consumer and producer lifecycle
- Implements dry-run mode
- Tracks replay results with proper correlation IDs
- Handles exceptions gracefully during replay
- Enforces limit on messages processed
963-1018: LGTM - Main entry point has robust error handling.The main function correctly:
- Handles all expected exception types with appropriate exit codes
- Provides user-friendly error messages for common issues
- Sets up graceful shutdown signal handlers
- Dispatches to command handlers properly
CI Fixes: - Fix circular import in protocol_plugin_compute.py by importing directly from model submodules instead of protocol re-exports - Fix ONEX validators contract import paths by using src/ as source path instead of src/omnibase_infra/ for proper module resolution DLQ Replay Script (scripts/dlq_replay.py): - Remove undefined error types from NON_RETRYABLE_ERRORS - Replace sys.exit(1) with ValueError in factory method - Add UTF-8 encoding error handling with errors='replace' - Remove --dry-run/--rate-limit from list and stats commands - Add STUB IMPLEMENTATION warning to header - Extract max_request_size and request_timeout_ms to config - Add explicit retry_count validation (no silent type coercion) - Optimize rate limiting to not delay first message - Add safe_truncate() for multi-byte character handling - Fix signal handler type annotation (FrameType | None) Kafka Event Bus (kafka_event_bus.py): - Sanitize error messages in _publish_raw_to_dlq to prevent secret leakage - Sanitize callback error in _invoke_dlq_callbacks - Document producer lock timing behavior (intentional design) Tests: - Add assertion in DLQ integration test to verify message received
Code Review: DLQ Implementation (OMN-949)SummaryThis PR implements comprehensive Dead Letter Queue support for the KafkaEventBus. The implementation is high quality with excellent documentation, strong typing, and proper error handling. Strengths
Critical Issues (BLOCKERS)1. File Naming Convention Violation
2. Error Sanitization Over-Redaction
3. DLQ Disabled Silently
High Priority Issues4. Protocol Naming
5. Misleading Stub Warning
Medium Priority6. Performance: dlq_metrics Copy Overhead
7. Rate Limiter Uses Wrong Clock
Documentation Improvements8. Missing Idempotency Warning
Security Hardening9. DLQ Replay Script Authentication
Overall AssessmentRating: 8.5/10 - Excellent work with fixable issues Production-ready after addressing the 4 blockers. The architecture is sound, code is well-tested, and documentation is comprehensive. Reviewed by: Claude Sonnet 4.5 |
- Fix NON_RETRYABLE_ERRORS comments to properly attribute error sources - Fix offset type inconsistency with explicit str conversion - Fix duplicate phrase in get_dlq_topic_for_original docstring - Remove ModelIntrospectionConfig re-export from mixins module
PR Review: DLQ Implementation (OMN-949)This PR implements comprehensive Dead Letter Queue (DLQ) support for the KafkaEventBus with proper error handling, metrics, and observability. The implementation follows ONEX patterns well and addresses all acceptance criteria. StrengthsExcellent ONEX Compliance:
Security Best Practices:
Comprehensive Documentation:
Solid Test Coverage:
Metrics and Observability:
Issues and RecommendationsCRITICAL - Issue 1: message_offset type annotation Location: src/omnibase_infra/event_bus/models/model_dlq_event.py:117-120 Problem: message_offset is typed as str | None but Kafka offsets are integers Fix: Change to message_offset: int | None Impact: Type safety violation, should be fixed before merge CRITICAL - Issue 2: Verify Awaitable import Location: src/omnibase_infra/event_bus/kafka_event_bus.py:543-544 Check that Callable and Awaitable are properly imported (ONEX prefers collections.abc) MEDIUM - Issue 3: Consider DLQ mixin extraction KafkaEventBus has 14 methods and 10 init parameters. Consider extracting DLQ functionality into MixinDlqPublisher to reduce complexity and improve testability. Not blocking, but recommended for future refactor. LOW - Issue 4: DLQ replay script status The scripts/dlq_replay.py (1139 lines) is described as skeleton. Consider adding TODO comment or creating follow-up ticket. Acceptance Criteria
RecommendationAPPROVE with required changes:
Excellent work on a complex feature! The error sanitization, metrics, and documentation are particularly strong. ONEX Compliance Score: 9.5/10 Review based on ONEX infrastructure patterns from CLAUDE.md |
…iance [OMN-949] Address PR #90 nitpick feedback to use Pydantic models instead of dataclasses: - Convert IntrospectionPerformanceMetrics to ModelIntrospectionPerformanceMetrics - Convert ValidationResult to ModelValidationResult - Convert HandlerInfo to ModelDetectedNodeInfo (renamed to avoid anti-pattern) Add backwards compatibility aliases to avoid breaking existing code.
…re-dlq-for-permanently-failing-intents Resolves merge conflicts from concurrent main branch updates: Conflicts resolved: - src/omnibase_infra/mixins/__init__.py: Keep both exports - src/omnibase_infra/mixins/mixin_node_introspection.py: Use imported ModelIntrospectionPerformanceMetrics (removed inline dupe), keep PerformanceMetricsCacheDict TypedDict - src/omnibase_infra/validation/infra_validators.py: Combine union histories (586 + 17 DLQ + 3 main = 606) - tests/unit/event_bus/test_kafka_event_bus.py: Keep both test classes (DLQ routing + topic validation) - tests/unit/validation/test_validator_defaults.py: Update assertion to 606 Combined history from both branches: - OMN-949: DLQ configuration and routing (~17 unions) - OMN-881: PR review fixes and _EventBusType alias (+1 union) - OMN-811: RegistryCompute merge (+2 unions) [OMN-949]
PR Review: Dead Letter Queue (DLQ) Implementation [OMN-949]Overall AssessmentVerdict: ✅ APPROVE with minor suggestions This is an excellent implementation of comprehensive DLQ support for the KafkaEventBus. The code follows ONEX conventions, demonstrates strong attention to security, and includes thorough testing. The implementation addresses all acceptance criteria and provides robust infrastructure for handling permanently failing messages. ✅ Strengths1. Excellent Model Design
2. Security-First Error SanitizationThe
3. Topic Naming Conventions
4. Comprehensive Documentation
5. Testing Excellence
📋 Suggestions (Non-Blocking)1. DLQ Replay Script EnhancementThe 2. Add DLQ Metrics to Health CheckExpose 3. DLQ Callback Error HandlingConsider wrapping callback invocations in try/except to prevent callback failures from blocking DLQ publishing. 🔒 Security Review ✅
🧪 Test Coverage: Excellent ✅
📦 ONEX Compliance: All Guidelines Followed ✅
🎯 Acceptance Criteria: All Met ✅
🏁 Final Recommendation✅ APPROVE This PR demonstrates excellent engineering practices, strong security focus, comprehensive testing, and full ONEX compliance. The implementation is production-ready. Great work! 🎉 Reviewed with: Claude Sonnet 4.5 | Standards: ONEX Infrastructure Guidelines (CLAUDE.md) |
The DLQ replay script is fully functional for Kafka operations. PostgreSQL tracking is tracked separately in OMN-1032. - Remove misleading "STUB IMPLEMENTATION" warning - Remove implementation status section (no longer needed) - Clean up environment variables (remove unused PostgreSQL vars) - Add reference to OMN-1032 for planned PostgreSQL tracking [OMN-949]
PR Review: Dead Letter Queue Implementation [OMN-949]SummaryThis is a well-architected and comprehensive implementation of Dead Letter Queue support for the KafkaEventBus. The PR demonstrates excellent adherence to ONEX principles with strong typing, thorough documentation, and comprehensive test coverage. ✅ Strengths1. Excellent Code Quality
2. Security-First DesignThe error sanitization implementation (
3. Comprehensive Documentation
4. Test Coverage
5. Operational Excellence
🔍 Areas for Consideration1. DLQ Replay Script - PostgreSQL TrackingThe replay script includes PostgreSQL connection parameters but the tracking is not yet implemented: File: Recommendation:
2. DLQ Topic Versioning StrategyFile: The DLQ schema is hardcoded to DLQ_TOPIC_VERSION: Final[str] = "v1"Question: What happens when the DLQ schema evolves?
Recommendation: Add a comment or doc note about schema evolution strategy (e.g., "Schema changes require new topic version; replay script must handle both during migration") 3. Error Sanitization Pattern MaintenanceFile: The Recommendation:
4. DLQ Message Size LimitsFile: Large failed messages could cause DLQ publishing to fail if they exceed Kafka's Observation: The replay script has configurable Recommendation: Consider adding a size check or truncation for very large messages before DLQ publishing (with a log warning) 5. Callback Error HandlingFile: Question: What happens if a registered callback raises an exception? Recommendation: Ensure callbacks are wrapped in try/except to prevent one failing callback from blocking others (verify this is already implemented) 🎯 ONEX Compliance
Note on Agent-Driven Development: CLAUDE.md mandates "ALL CODING TASKS MUST USE SUB-AGENTS - NO EXCEPTIONS". The commit messages don't reference agent invocations. This may be a policy violation unless agents were used but not documented. 🧪 Testing NotesIntegration Test QualityThe integration tests are well-designed:
Suggested Additional Tests
📊 Metrics
🚀 RecommendationAPPROVE with minor suggestions This PR represents high-quality infrastructure engineering. The implementation is:
The suggestions above are minor refinements, not blockers. The code quality exceeds typical standards. Before Merge
Post-Merge Recommendations
Excellent work on this implementation! 🎉 The attention to detail, especially around security (error sanitization) and operational concerns (replay script, metrics, callbacks), demonstrates a mature understanding of production system requirements. |
There was a problem hiding this comment.
Actionable comments posted: 2
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
src/omnibase_infra/mixins/mixin_node_introspection.py (1)
1075-1124: Critical:discover_capabilities_msmetric is never populated.The
discover_capabilities_msvariable is initialized to0.0at line 1077 but never updated. Inget_capabilities()at line 884,discover_elapsed_msis calculated but not returned or made accessible toget_introspection_data().Impact:
- The
discover_capabilities_msfield in performance metrics will always be0.0- Observability is degraded - cannot monitor the performance of the class-level method signature cache discovery
Root cause:
The conversion to frozen Pydantic models is correct, but the refactoring failed to capture thediscover_elapsed_msvalue fromget_capabilities().🔎 Proposed fix
Option 1: Return timing from get_capabilities()
Modify
get_capabilities()to return timing information:- async def get_capabilities(self) -> CapabilitiesDict: + async def get_capabilities(self) -> tuple[CapabilitiesDict, float]: """Extract node capabilities via reflection.""" self._ensure_initialized() start_time = time.perf_counter() discover_start = time.perf_counter() cached_signatures = self._get_class_method_signatures() discover_elapsed_ms = (time.perf_counter() - discover_start) * 1000 # ... rest of method ... elapsed_ms = (time.perf_counter() - start_time) * 1000 - return capabilities + return capabilities, discover_elapsed_msThen update callers:
# Build fresh introspection data with timing for each component cap_start = time.perf_counter() - capabilities = await self.get_capabilities() + capabilities, discover_capabilities_ms = await self.get_capabilities() get_capabilities_ms = (time.perf_counter() - cap_start) * 1000 - - # Extract method count from capabilities - method_sigs = capabilities.get("method_signatures", {}) - method_count = len(method_sigs) if isinstance(method_sigs, dict) else 0Option 2: Store as instance variable
Add
_last_discover_capabilities_msinstance variable and set it inget_capabilities(), then read inget_introspection_data(). This is less clean but requires fewer changes.src/omnibase_infra/event_bus/kafka_event_bus.py (1)
1326-1388: Use the message's correlation_id in DLQ events instead of the consumer task's correlation_idThe current implementation loses the original
ModelEventHeaders.correlation_idfrom incoming messages when publishing to the dead letter queue. Instead, both_publish_to_dlqand_publish_raw_to_dlquse the consumer task'scorrelation_idparameter—making it impossible to trace failures end-to-end through downstream systems and violating the guideline to propagate correlation_id from incoming requests to error context.For the normal handler-failure path: Pass
event_message.headers.correlation_idto_publish_to_dlq(the message's original correlation_id) instead of the consumer's task correlation_id.For the deserialization-failure path: In
_publish_raw_to_dlq, extract the correlation_id fromraw_msg.headersusing the existing_kafka_headers_to_modelextraction logic, falling back to the passed-in task correlation_id only if headers are missing or invalid.This ensures DLQ payloads,
ModelDlqEventheaders, and logs all carry the same correlation_id that original requesters and downstream handlers expect, while preserving the consumer task id for internal consumer-loop diagnostics if needed.
🧹 Nitpick comments (4)
src/omnibase_infra/nodes/reducers/registration_reducer.py (1)
462-463: Consider exportingModelValidationResultin__all__for forward compatibility.The alias provides good backwards compatibility. However, the new canonical name
ModelValidationResultis not exported in__all__(line 1136), only the aliasValidationResultis. For consistency with theModel*naming convention used elsewhere in ONEX, consider also exportingModelValidationResultto allow callers to migrate to the new name.🔎 Proposed addition to __all__
__all__ = [ "RegistrationReducer", # Validation types (for tests and custom validators) + "ModelValidationResult", "ValidationResult", "ValidationErrorCode", # Performance threshold constants (for tests and monitoring) "PERF_THRESHOLD_REDUCE_MS", "PERF_THRESHOLD_INTENT_BUILD_MS", "PERF_THRESHOLD_IDEMPOTENCY_CHECK_MS", ]tests/integration/event_bus/test_dlq_integration.py (1)
616-678: Strengthen DLQ metrics increment assertion to catch regressions
test_dlq_metrics_increment_on_publishonly checksfinal_metrics.total_publishes >= initial_total, so it will still pass if DLQ metrics never update. Given you already wait for processing, consider asserting thatfinal_metrics.total_publishes > initial_total(or at least logging when it does not change) so a broken metrics path can be detected instead of silently passing.src/omnibase_infra/event_bus/topic_constants.py (1)
65-207: DLQ topic helpers align well with ONEX naming; consider normalizing parse outputThe DLQ constants and helpers (
DLQ_TOPIC_PATTERN,build_dlq_topic,parse_dlq_topic,is_dlq_topic,get_dlq_topic_for_original) correctly enforce the<env>.dlq.<category>.<version>pattern, validate environments, and reuseEnumMessageCategory/topic_suffixfor category handling. This is a solid foundation for DLQ routing.One optional improvement:
DLQ_TOPIC_PATTERNis case‑insensitive, butparse_dlq_topiccurrently returnsenvironmentandcategoryin whatever case was matched. If callers expect normalized lowercase values, you may want tolower()those fields in the returned dict (and possibly forENV_PATTERNvalidation too) to avoid subtle casing differences between built and parsed topics.Also applies to: 237-285, 287-334
scripts/dlq_replay.py (1)
998-1054: GracefulShutdown wiring is unused and may interfere with normal SIGINT behavior
setup_signal_handlers()installs custom handlers forSIGINT/SIGTERMthat just log and set an internal event;GracefulShutdownitself is never consulted by the replay commands (nowait_for_shutdown()oris_shutdown_requestedchecks in the main loops). This adds complexity without affecting behavior, and custom signal handlers that don’t raise can also change how Ctrl‑C behaves underasyncio.run.Two options to simplify and make behavior clearer:
- Either wire the shutdown handler into your long‑running operations (e.g., have
DLQReplayExecutor.executeor theconsume_messagesloop periodically checkshutdown_handler.is_shutdown_requestedand break early), and rely on that for graceful termination, or- Drop the custom signal handling entirely and let
KeyboardInterrupt/asyncio.runhandle Ctrl‑C in the usual way.Also, in the Kafka connection error branch, consider using the resolved
config.bootstrap_serversrather thanargs.bootstrap_serversso the message reflects any environment overrides actually used to connect.Also applies to: 1062-1093
📜 Review details
Configuration used: defaults
Review profile: CHILL
Plan: Lite
📒 Files selected for processing (16)
scripts/dlq_replay.pyscripts/validate.pysrc/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/event_bus/topic_constants.pysrc/omnibase_infra/mixins/__init__.pysrc/omnibase_infra/mixins/mixin_node_introspection.pysrc/omnibase_infra/nodes/reducers/registration_reducer.pysrc/omnibase_infra/plugins/plugin_compute_base.pysrc/omnibase_infra/protocols/protocol_plugin_compute.pysrc/omnibase_infra/services/timeout_emitter.pysrc/omnibase_infra/validation/execution_shape_validator.pysrc/omnibase_infra/validation/infra_validators.pytests/integration/event_bus/test_dlq_integration.pytests/unit/event_bus/test_kafka_event_bus.pytests/unit/mixins/test_mixin_node_introspection.pytests/unit/validation/test_validator_defaults.py
💤 Files with no reviewable changes (1)
- src/omnibase_infra/mixins/init.py
🚧 Files skipped from review as they are similar to previous changes (3)
- src/omnibase_infra/services/timeout_emitter.py
- tests/unit/event_bus/test_kafka_event_bus.py
- tests/unit/validation/test_validator_defaults.py
🧰 Additional context used
📓 Path-based instructions (5)
**/*.py
📄 CodeRabbit inference engine (CLAUDE.md)
**/*.py: Never useAnytypes in Python code. Always use specific types or useobjectfor generic dispatchers that must accept any payload.
Use PEP 604 union syntaxX | Nonefor nullable types instead ofOptional[X]in Python code.
UseModelEventEnvelope[object]instead ofAnyfor generic dispatchers that must accept any payload type. UseModelEventEnvelope[SpecificType]for dispatchers that know the exact payload type.
All services must use ModelONEXContainer for dependency injection. Bootstrap using container.service_registry.resolve_service(ServiceType).
Never include passwords, API keys, tokens, secrets, full connection strings with credentials, PII, internal IP addresses, private keys, certificates, or session tokens in error messages or context. Only include service names, operation names, correlation IDs, error codes, sanitized hostnames, port numbers, retry counts, and resource identifiers.
Always propagate correlation_id from incoming requests to error context. If no correlation_id exists, generate one using uuid4(). Include correlation_id in all error context for distributed tracing.
Use EnumMessageCategory for message routing decisions and topic parsing. Use EnumNodeOutputType for validating node output shapes and handler return types. PROJECTION enum value only exists in EnumNodeOutputType and is only valid for REDUCER nodes.
Use EnumInfraTransportType for transport identification in error context: HTTP, DATABASE, KAFKA, CONSUL, VAULT, VALKEY, GRPC.
Files:
tests/unit/mixins/test_mixin_node_introspection.pysrc/omnibase_infra/nodes/reducers/registration_reducer.pysrc/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/protocols/protocol_plugin_compute.pysrc/omnibase_infra/validation/execution_shape_validator.pysrc/omnibase_infra/validation/infra_validators.pysrc/omnibase_infra/mixins/mixin_node_introspection.pyscripts/dlq_replay.pytests/integration/event_bus/test_dlq_integration.pysrc/omnibase_infra/event_bus/topic_constants.pysrc/omnibase_infra/plugins/plugin_compute_base.pyscripts/validate.py
**/nodes/**/*.py
📄 CodeRabbit inference engine (CLAUDE.md)
**/nodes/**/*.py: Node base classes and I/O models must be imported from omnibase_core.nodes (NodeEffect, NodeCompute, NodeReducer, NodeOrchestrator, and their respective Input/Output/Transaction models). Never define new node archetypes in omnibase_infra.
Node Introspection via MixinNodeIntrospection automatically discovers node capabilities using Python reflection. Prefix internal/sensitive methods with _ to exclude from introspection. Use generic operation names and parameter names to avoid exposing implementation details.
MixinNodeIntrospection is designed for single-threaded asyncio usage. Cache operations (_introspection_cache, invalidate_introspection_cache) are synchronous and NOT thread-safe without external synchronization. If multi-threading is required, use external locks.
In multi-threaded environments, synchronize access to MixinNodeIntrospection using threading.Lock for instance-level cache (_introspection_cache). For asyncio applications, prefer single event loop with cooperative multitasking over multi-threading.
Files:
src/omnibase_infra/nodes/reducers/registration_reducer.py
**/kafka_event_bus.py
📄 CodeRabbit inference engine (CLAUDE.md)
KafkaEventBus intentionally violates pattern validator thresholds (14 methods, 10 init parameters) due to event bus pattern requirements and backwards compatibility during config migration. This complexity exception is documented.
Files:
src/omnibase_infra/event_bus/kafka_event_bus.py
**/protocol_*.py
📄 CodeRabbit inference engine (CLAUDE.md)
Use protocol files with pattern
protocol_<name>.pyfor standalone protocols orprotocols.pyfor domain-grouped protocols within modules, with class namingProtocol<Name>.
Files:
src/omnibase_infra/protocols/protocol_plugin_compute.py
**/mixin_*.py
📄 CodeRabbit inference engine (CLAUDE.md)
Use mixin files with pattern
mixin_<name>.pyand class namingMixin<Name>for all mixin classes.
Files:
src/omnibase_infra/mixins/mixin_node_introspection.py
🧠 Learnings (35)
📓 Common learnings
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to {services/**/*.py,scripts/bulk_ingest_repository.py} : Implement fail-closed configuration for security hardening. All external requests must validate URLs, implement DLQ routing, and handle failures gracefully.
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to **/*.py : Messages exceeding `KAFKA_MAX_REQUEST_SIZE` (default 10MB) are automatically routed to oversized DLQ topic. Use `scripts/process_oversized_dlq.py` to handle these messages. See `docs/guides/DLQ_HANDLING.md` for complete documentation.
Learnt from: CR
Repo: OmniNode-ai/omniintelligence PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-07T17:50:13.678Z
Learning: Implement Kafka event-driven architecture with proper topic naming using prefix dev.archon-intelligence. and proper event flow pattern with Effect nodes consuming events, processing, and publishing results with Dead Letter Queue routing
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/dag_based_tool_bootstrapping.mdc:0-0
Timestamp: 2025-11-24T17:22:55.666Z
Learning: Ensure DAG system includes event-driven coordination, comprehensive validation (100% success rate), retry policies, error recovery, health monitoring, rollback capabilities, and Tool-as-a-Service ready architecture
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/nodes/**/*.py : Node Introspection via MixinNodeIntrospection automatically discovers node capabilities using Python reflection. Prefix internal/sensitive methods with _ to exclude from introspection. Use generic operation names and parameter names to avoid exposing implementation details.
Applied to files:
tests/unit/mixins/test_mixin_node_introspection.pysrc/omnibase_infra/mixins/mixin_node_introspection.py
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/nodes/**/*.py : Import mixins from omnibase_core.mixins.* and use Mixin* naming pattern (e.g., MixinHealthCheck, MixinMetrics, MixinEventBus) - never use local custom mixins unless experimental and documented
Applied to files:
tests/unit/mixins/test_mixin_node_introspection.pysrc/omnibase_infra/mixins/mixin_node_introspection.py
📚 Learning: 2025-12-20T04:09:41.832Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_core PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-20T04:09:41.832Z
Learning: Applies to src/omnibase_core/models/**/*.py : Add from_attributes=True to ConfigDict for immutable value objects nested in Pydantic models used with pytest-xdist parallel execution
Applied to files:
src/omnibase_infra/nodes/reducers/registration_reducer.py
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/nodes/**/model_*.py : All node implementations must define strongly-typed input and output models following convention Model<NodeName>Input and Model<NodeName>Output. Import base models from omnibase_core.nodes.
Applied to files:
src/omnibase_infra/nodes/reducers/registration_reducer.py
📚 Learning: 2025-11-29T17:13:38.776Z
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to **/*.py : Messages exceeding `KAFKA_MAX_REQUEST_SIZE` (default 10MB) are automatically routed to oversized DLQ topic. Use `scripts/process_oversized_dlq.py` to handle these messages. See `docs/guides/DLQ_HANDLING.md` for complete documentation.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pyscripts/dlq_replay.pytests/integration/event_bus/test_dlq_integration.py
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to **/*.py : Publish intelligence requests to Kafka event bus using topics: dev.archon-intelligence.intelligence.code-analysis-{requested,completed,failed}.v1 for consistency and event-driven architecture
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pytests/integration/event_bus/test_dlq_integration.pysrc/omnibase_infra/event_bus/topic_constants.py
📚 Learning: 2025-12-06T22:21:32.649Z
Learnt from: CR
Repo: OmniNode-ai/omniagent PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-06T22:21:32.649Z
Learning: Applies to nodes/**/*.py : Use event bus mixins from `omnibase_core` for Kafka publishing instead of direct Kafka clients
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-11-29T17:13:38.776Z
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Use event-driven architecture with Kafka topics for asynchronous processing: enrichment, code analysis, manifest processing, and entity embedding pipelines.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-12-07T17:50:13.678Z
Learnt from: CR
Repo: OmniNode-ai/omniintelligence PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-07T17:50:13.678Z
Learning: Implement Kafka event-driven architecture with proper topic naming using prefix dev.archon-intelligence. and proper event flow pattern with Effect nodes consuming events, processing, and publishing results with Dead Letter Queue routing
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pytests/integration/event_bus/test_dlq_integration.py
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/kafka_event_bus.py : KafkaEventBus intentionally violates pattern validator thresholds (14 methods, 10 __init__ parameters) due to event bus pattern requirements and backwards compatibility during config migration. This complexity exception is documented.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-11-29T17:13:38.776Z
Learnt from: CR
Repo: OmniNode-ai/omniarchon PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-29T17:13:38.776Z
Learning: Applies to {services/**/*.py,scripts/bulk_ingest_repository.py} : Implement fail-closed configuration for security hardening. All external requests must validate URLs, implement DLQ routing, and handle failures gracefully.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pyscripts/dlq_replay.py
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/events/**/*.py : Kafka event publishing MUST use OnexEnvelopeV1 format with 13 topics for event streaming at all workflow lifecycle stages
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.pysrc/omnibase_infra/event_bus/topic_constants.py
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/*.py : Never include passwords, API keys, tokens, secrets, full connection strings with credentials, PII, internal IP addresses, private keys, certificates, or session tokens in error messages or context. Only include service names, operation names, correlation IDs, error codes, sanitized hostnames, port numbers, retry counts, and resource identifiers.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/errors/**/*.py : Use ModelInfraErrorContext when raising infrastructure errors, including transport_type, operation, target_name, and correlation_id fields.
Applied to files:
src/omnibase_infra/event_bus/kafka_event_bus.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Applies to **/*.py : Import protocols from `omnibase.protocol.protocol_*` paths
Applied to files:
src/omnibase_infra/protocols/protocol_plugin_compute.py
📚 Learning: 2025-11-24T17:24:41.687Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T17:24:41.687Z
Learning: Applies to **/*.py : Import protocols from `omnibase.protocol.protocol_<name>` module paths
Applied to files:
src/omnibase_infra/protocols/protocol_plugin_compute.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Applies to **/protocols/protocol_*.py : Use TYPE_CHECKING guards and forward references for circular import prevention in protocol files
Applied to files:
src/omnibase_infra/protocols/protocol_plugin_compute.pysrc/omnibase_infra/plugins/plugin_compute_base.py
📚 Learning: 2025-11-24T17:24:41.687Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T17:24:41.687Z
Learning: Applies to **/protocols/protocol_*.py : Use Protocol for tool interfaces and plugin APIs based on method shape (structural typing), not Pydantic models with inheritance
Applied to files:
src/omnibase_infra/protocols/protocol_plugin_compute.pysrc/omnibase_infra/plugins/plugin_compute_base.py
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/nodes/**/*.py : Node base classes and I/O models must be imported from omnibase_core.nodes (NodeEffect, NodeCompute, NodeReducer, NodeOrchestrator, and their respective Input/Output/Transaction models). Never define new node archetypes in omnibase_infra.
Applied to files:
src/omnibase_infra/protocols/protocol_plugin_compute.py
📚 Learning: 2025-12-08T00:48:30.737Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_spi PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-08T00:48:30.737Z
Learning: Applies to src/omnibase_spi/protocols/**/*.py : All public protocols must be decorated with `runtime_checkable`
Applied to files:
src/omnibase_infra/protocols/protocol_plugin_compute.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Use Protocol for interface definitions when implementations may live outside core codebase; use Pydantic models only for base classes with shared logic
Applied to files:
src/omnibase_infra/protocols/protocol_plugin_compute.pysrc/omnibase_infra/plugins/plugin_compute_base.py
📚 Learning: 2025-11-28T18:58:53.781Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/canonical_patterns.mdc:0-0
Timestamp: 2025-11-28T18:58:53.781Z
Learning: Applies to **/*.py : Use direct field access in validators instead of `info.data` fallbacks in Pydantic models
Applied to files:
src/omnibase_infra/validation/execution_shape_validator.py
📚 Learning: 2025-12-06T22:21:32.649Z
Learnt from: CR
Repo: OmniNode-ai/omniagent PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-06T22:21:32.649Z
Learning: Applies to nodes/**/*.py : Use `omnibase_infra` handlers for OmniIntelligence queries via HttpRestAdapter envelope pattern
Applied to files:
src/omnibase_infra/validation/execution_shape_validator.py
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/nodes/**/*.py : MixinNodeIntrospection is designed for single-threaded asyncio usage. Cache operations (_introspection_cache, invalidate_introspection_cache) are synchronous and NOT thread-safe without external synchronization. If multi-threading is required, use external locks.
Applied to files:
src/omnibase_infra/mixins/mixin_node_introspection.py
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/nodes/**/*.py : In multi-threaded environments, synchronize access to MixinNodeIntrospection using threading.Lock for instance-level cache (_introspection_cache). For asyncio applications, prefer single event loop with cooperative multitasking over multi-threading.
Applied to files:
src/omnibase_infra/mixins/mixin_node_introspection.py
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to agents/**/*.py : Use correlation_id UUID for end-to-end traceability across all agent routing, manifest injection, and execution events
Applied to files:
scripts/dlq_replay.py
📚 Learning: 2025-11-24T17:23:49.777Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T17:23:49.777Z
Learning: Applies to **/node_*/v[0-9]*_[0-9]*_[0-9]*/models/error_codes.py : All ONEX node error handling must use auto-generated error codes defined in `models/error_codes.py` from contract definitions
Applied to files:
scripts/dlq_replay.py
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/errors/**/*.py : Select error classes based on scenario: Use ProtocolConfigurationError for config validation, SecretResolutionError for missing secrets, InfraConnectionError for connection failures, InfraTimeoutError for timeouts, InfraAuthenticationError for auth failures, InfraUnavailableError for unavailable resources.
Applied to files:
scripts/dlq_replay.py
📚 Learning: 2025-12-25T18:32:08.629Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T18:32:08.629Z
Learning: Applies to **/*.py : Use EnumInfraTransportType for transport identification in error context: HTTP, DATABASE, KAFKA, CONSUL, VAULT, VALKEY, GRPC.
Applied to files:
scripts/dlq_replay.py
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to scripts/tests/**/*.sh : Implement comprehensive test suites in scripts/tests/ with separate test files for Kafka, PostgreSQL, Intelligence, and Routing functionality
Applied to files:
tests/integration/event_bus/test_dlq_integration.py
📚 Learning: 2025-11-24T17:24:41.687Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T17:24:41.687Z
Learning: Applies to src/omnibase/constants/*_constants.py : Constants files must follow the naming pattern `<domain>_constants.py` and be located in `src/omnibase/constants/`
Applied to files:
src/omnibase_infra/event_bus/topic_constants.py
📚 Learning: 2025-11-24T16:33:32.747Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/standards.mdc:0-0
Timestamp: 2025-11-24T16:33:32.747Z
Learning: Applies to src/omnibase/constants/*_constants.py : Constants files must follow the naming pattern `<domain>_constants.py` and be located in `src/omnibase/constants/` directory
Applied to files:
src/omnibase_infra/event_bus/topic_constants.py
📚 Learning: 2025-12-08T00:48:30.737Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_spi PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-08T00:48:30.737Z
Learning: Import `omnibase_core` models and types only for type hints and runtime usage - follow the SPI → Core dependency direction
Applied to files:
src/omnibase_infra/plugins/plugin_compute_base.py
📚 Learning: 2025-12-08T00:48:30.737Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_spi PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-08T00:48:30.737Z
Learning: Applies to src/omnibase_spi/**/*.py : SPI modules may import from `omnibase_core` for type hints and model runtime usage (allowed and required)
Applied to files:
src/omnibase_infra/plugins/plugin_compute_base.py
🧬 Code graph analysis (7)
src/omnibase_infra/event_bus/kafka_event_bus.py (4)
src/omnibase_infra/event_bus/models/model_dlq_event.py (1)
ModelDlqEvent(50-207)src/omnibase_infra/event_bus/models/model_dlq_metrics.py (3)
ModelDlqMetrics(61-305)create_empty(299-305)record_dlq_publish(188-249)src/omnibase_infra/event_bus/models/model_event_headers.py (1)
ModelEventHeaders(16-98)src/omnibase_infra/utils/util_error_sanitization.py (1)
sanitize_error_message(95-154)
src/omnibase_infra/protocols/protocol_plugin_compute.py (3)
src/omnibase_infra/plugins/models/model_plugin_context.py (1)
ModelPluginContext(20-99)src/omnibase_infra/plugins/models/model_plugin_input_data.py (1)
ModelPluginInputData(20-82)src/omnibase_infra/plugins/models/model_plugin_output_data.py (1)
ModelPluginOutputData(20-86)
src/omnibase_infra/validation/execution_shape_validator.py (4)
src/omnibase_infra/enums/enum_handler_type.py (1)
EnumHandlerType(37-69)src/omnibase_infra/handlers/handler_db.py (1)
handler_type(86-88)src/omnibase_infra/handlers/handler_consul.py (1)
handler_type(194-204)src/omnibase_infra/handlers/handler_http.py (1)
handler_type(92-94)
src/omnibase_infra/mixins/mixin_node_introspection.py (1)
src/omnibase_infra/models/discovery/model_introspection_performance_metrics.py (1)
ModelIntrospectionPerformanceMetrics(26-166)
tests/integration/event_bus/test_dlq_integration.py (4)
src/omnibase_infra/event_bus/models/model_dlq_event.py (1)
ModelDlqEvent(50-207)src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py (1)
ModelKafkaEventBusConfig(117-722)src/omnibase_infra/event_bus/topic_constants.py (3)
build_dlq_topic(145-214)parse_dlq_topic(237-266)is_dlq_topic(269-284)src/omnibase_infra/event_bus/models/model_event_headers.py (1)
ModelEventHeaders(16-98)
src/omnibase_infra/event_bus/topic_constants.py (1)
src/omnibase_infra/enums/enum_message_category.py (3)
EnumMessageCategory(34-196)from_topic(121-163)topic_suffix(100-118)
src/omnibase_infra/plugins/plugin_compute_base.py (3)
src/omnibase_infra/plugins/models/model_plugin_context.py (1)
ModelPluginContext(20-99)src/omnibase_infra/plugins/models/model_plugin_input_data.py (1)
ModelPluginInputData(20-82)src/omnibase_infra/plugins/models/model_plugin_output_data.py (1)
ModelPluginOutputData(20-86)
🔇 Additional comments (15)
src/omnibase_infra/nodes/reducers/registration_reducer.py (2)
356-356: LGTM!The Pydantic import is appropriate for the new
ModelValidationResultmodel definition.
421-459: LGTM!The conversion from dataclass to Pydantic model is well-implemented:
- Uses
ConfigDict(frozen=True)for immutability, aligning with ONEX requirements for immutable value objects- Properly uses PEP 604 union syntax (
ValidationErrorCode | None) per coding guidelines- Factory methods (
success(),failure()) maintain clean API semanticssrc/omnibase_infra/validation/infra_validators.py (1)
373-377: LGTM! Threshold update is well-documented and justified.The
INFRA_MAX_UNIONSupdate from 589 to 606 is clearly documented with the OMN-949 ticket reference and explanation of the DLQ-related union additions. The threshold history follows the established pattern and aligns with the PR's scope of adding comprehensive DLQ support.src/omnibase_infra/validation/execution_shape_validator.py (3)
81-81: LGTM: Pydantic import for model conversion.The import addition supports the dataclass-to-Pydantic conversion below.
257-286: Excellent backwards-compatible refactoring to Pydantic.The conversion from dataclass to Pydantic BaseModel is well-executed:
ConfigDict(arbitrary_types_allowed=True)correctly handles AST node types- Field renamed
handler_type→node_typefor ONEX terminology alignment- Read-only
handler_typeproperty maintains full backwards compatibility- Public alias
HandlerInfo = ModelDetectedNodeInfopreserves the API surface- All existing code using
handler.handler_typecontinues to work transparentlyThis allows gradual migration to the new field name while ensuring zero breaking changes.
456-456: Constructor calls correctly updated to new field name.Both construction sites properly use
node_type=handler_typewherehandler_typeis the local variable from detection logic. Existing usages throughout the file continue to access via the backwards-compatiblehandler.handler_typeproperty (lines 630, 691, 777, 790, etc.).Also applies to: 468-468
scripts/validate.py (1)
182-186: Code change is correct—no other packages undersrc/exist to validate.The change from
Path("src/omnibase_infra/")toPath("src/")does not expand validation to other packages as initially assumed. The repository contains only theomnibase_infrapackage undersrc/. This change simply qualifies module names at the project level (e.g., "omnibase_infra.clients" instead of "clients"), which is more accurate for Python's import system. The implementation is correct and the comments clearly explain the technical rationale.tests/unit/mixins/test_mixin_node_introspection.py (2)
1974-1974: LGTM! Correct Pydantic v2 API usage.The change from
to_dict()tomodel_dump()correctly reflects the Pydantic v2 migration. Themodel_dump()method is the canonical serialization API in Pydantic v2.Also applies to: 2737-2737
2738-2752: LGTM! Test expectations correctly updated.The addition of
"captured_at"toexpected_keyscorrectly reflects the field added toModelIntrospectionPerformanceMetrics. The error message also correctly referencesmodel_dump().src/omnibase_infra/mixins/mixin_node_introspection.py (3)
210-211: LGTM! Good backward-compatibility pattern.The alias preserves the public API surface while migrating to the centralized
ModelIntrospectionPerformanceMetricstype. This maintains compatibility for external consumers.
1157-1181: Correct frozen model construction pattern.The pattern of collecting metric values in local variables before constructing the frozen
IntrospectionPerformanceMetricsobject is correct. The threshold detection logic properly identifies slow operations.Note: This segment is implemented correctly, but it cannot populate
discover_capabilities_msdue to the issue flagged in the previous comment.
1186-1197: LGTM! Simplified event construction.The removal of the conversion comment and direct assignment of
performance_metrics=metricsis correct. Themetricsvariable is already aModelIntrospectionPerformanceMetrics(via the alias), so no conversion is needed.src/omnibase_infra/protocols/protocol_plugin_compute.py (1)
109-113: Direct model imports correctly break the plugins circular dependencySwitching to direct imports from
plugins.models.*while keeping the public aliases and__all__intact cleanly resolves the cycle without changing the external protocol surface. No further changes needed here.src/omnibase_infra/plugins/plugin_compute_base.py (1)
162-171: Backwards-compatible aliases to shared plugin models look goodImporting the concrete plugin models directly and exposing local
PluginContext/PluginInputData/PluginOutputDataaliases keeps this base class in sync with the protocol and avoids the circular dependency. This is a clean, non-breaking refactor.src/omnibase_infra/event_bus/kafka_event_bus.py (1)
413-420: DLQ metrics and callback plumbing are well-structuredInitializing DLQ metrics with a copy-on-write pattern and exposing a read‑only
dlq_metricsview, plus guarded registration/unregistration of async DLQ callbacks, gives a clean extension surface for observability/alerting without leaking internal state. The lock usage around_dlq_metricsand_dlq_callbacksis appropriate.Also applies to: 531-588
…949] Merge Conflicts Resolved: - docs/operations/README.md: Keep both DLQ Replay Guide + Thread Pool entries - infra_validators.py: Include OMN-816 + OMN-949 in threshold history - test_validator_defaults.py: Update baseline to ~603 unions Hardcoded Values Fix (User Priority): - Remove localhost:9092 default from dlq_replay.py - Make KAFKA_BOOTSTRAP_SERVERS env var required (no fallback) - Update docstring and argparse help to reflect requirement Security Fixes: - Sanitize bootstrap_servers in all log messages - Factory method raises ValueError instead of sys.exit() Robustness Fixes: - Add rate_limit_per_second validation (prevent division by zero) - Add UTF-8 decoding error handling with replacement chars Code Quality: - Remove unnecessary TypedDict re-exports from mixins module - Add CapabilitiesTypedDict to models/discovery exports
PR Review: Dead Letter Queue (DLQ) Implementation [OMN-949]🎯 Overall AssessmentSTRONG APPROVE ✅ - This is an exemplary implementation that demonstrates deep understanding of ONEX principles, infrastructure resilience patterns, and production-grade error handling. The PR successfully addresses all acceptance criteria with comprehensive testing, documentation, and security considerations. ✅ Strengths1. Architecture & Design Excellence
2. Security & Error Sanitization
3. Production Readiness
4. Documentation Quality
5. Testing Coverage
6. Code Quality
🔍 Technical Deep DiveDLQ Routing LogicThe implementation correctly handles all failure paths:
Error Sanitization Flow
Thread SafetyCircuit breaker integration follows proper lock pattern with I/O operations outside locks 📋 Minor Observations (Not Blocking)1. NON_RETRYABLE_ERRORS MappingQuestion: Should Recommendation: Consider creating 2. DLQ Replay Script - Environment VariableExcellent security decision to remove 3. Validation Threshold UpdatesWell-documented rationale for 🎓 ONEX Compliance Checklist
🚀 Performance ConsiderationsPositive
Potential Optimization (Future)
🔒 Security ReviewStrengths
Documentation Note
📊 Test Coverage AssessmentUnit Tests (Strong)
Integration Tests (Excellent)
🏆 Highlights for Other Developers
✅ Final RecommendationMERGE - This PR exceeds quality standards and demonstrates:
Minor observations above are not blocking and can be addressed in follow-up work. Congratulations on an exemplary implementation! 🎉 📚 References for Future Work
Reviewed by: Claude Sonnet 4.5 (Code Quality Analysis) |
…N-949] Create EnumNonRetryableErrorCategory to prevent divergence between dlq_replay.py and kafka_event_bus.py on retry eligibility logic. New enum in omnibase_infra.enums with: - AUTHENTICATION_ERROR: Invalid credentials won't become valid - CONFIGURATION_ERROR: Schema/config errors require code changes - SECRET_RESOLUTION_ERROR: Missing secrets won't appear on retry - VALIDATION_ERROR: Invalid input won't become valid Helper methods: - is_non_retryable(error_type) -> bool - get_all_values() -> frozenset[str] - get_description(error_type) -> str Updates dlq_replay.py to use centralized enum instead of inline frozenset.
Comprehensive PR Review: Dead Letter Queue (DLQ) ImplementationOverviewThis PR implements comprehensive DLQ support for the KafkaEventBus (OMN-949). The implementation is well-designed and production-ready, with excellent attention to security, observability, and operational concerns. Below is detailed feedback organized by category. ✅ Strengths1. Security & Error Sanitization ⭐The introduction of
Suggestion: Consider adding sanitization tests for edge cases:
2. Non-Retryable Error Categorization ⭐The new
Minor: The docstring mentions "thread-safe" but enums are inherently thread-safe due to immutability. This is correct but could note that it's due to Python's enum immutability guarantees. 3. DLQ Replay Tool ⭐The
Excellent pattern: The 4. Documentation Quality ⭐
5. Test Coverage
🔍 Issues & Concerns1. Type Safety -
|
| Criterion | Status | Notes |
|---|---|---|
| Error handling | Address retry_count parsing, DLQ publish failure observability | |
| Security | Add payload sanitization option, document ACL requirements | |
| Observability | ✅ | Good logging, metrics definitions in docs |
| Documentation | ✅ | Comprehensive docs, runbooks, examples |
| Test coverage | Good unit tests, add edge case coverage | |
| Performance | ✅ | Rate limiting, async patterns correct |
| ONEX compliance | Fix original_offset type (str → int) |
🎯 Recommended Actions Before Merge
High Priority (Blocking):
- Fix
original_offsettype - Change fromstr | Nonetoint | None - Add DLQ publish failure metrics - Implement
dlq_publish_failures_totalcounter - Handle
retry_countparsing errors - Add try-except with default fallback
Medium Priority (Strongly Recommended):
- Add payload sanitization config -
sanitize_dlq_payloadsoption in config - Add tests for sanitization edge cases - Multi-byte UTF-8, boundary conditions
- Document Kafka ACL requirements - Add to DLQ_MESSAGE_FORMAT.md security section
Low Priority (Nice to Have):
- Extract DLQ logic to submodule - Reduce
kafka_event_bus.pycomplexity - Add Prometheus query examples - In DLQ_MESSAGE_FORMAT.md monitoring section
🌟 Overall Assessment
Verdict: Approve with minor changes
This is a well-engineered, production-ready implementation of DLQ functionality. The code demonstrates:
- Strong security awareness (sanitization utilities)
- Excellent observability design (metrics, callbacks, logs)
- Thoughtful operational tooling (replay script with safety features)
- Comprehensive documentation
The identified issues are minor and easily addressable. The high-priority fixes (offset type, error handling) are small changes that significantly improve type safety and resilience.
Recommended merge path:
- Address high-priority issues (1-3)
- Merge to staging for integration testing
- Address medium-priority issues in follow-up PR
- Production rollout with monitoring dashboard setup
Excellent work on this complex feature! The attention to security (sanitization), operational concerns (replay tool), and observability (metrics, callbacks) demonstrates mature engineering practices. 🎉
📚 Related Work
Per ONEX Agent-Driven Development policy, future enhancements should delegate to:
agent-security-audit: Review DLQ security (payload sanitization, ACLs)agent-testing: Expand edge case coverage for sanitization and replayagent-performance: Profile sanitization overhead if DLQ volume increases
OMN-949: ✅ All acceptance criteria met
OMN-1032: Referenced for PostgreSQL tracking integration (future work)
…en tests [OMN-949] Changes: - Remove unused GracefulShutdown class and setup_signal_handlers() from dlq_replay.py - asyncio's default SIGINT handling via CancelledError is sufficient and was already being used - Export ModelValidationResult in registration_reducer.py __all__ for forward compatibility (ONEX naming convention) - Strengthen DLQ metrics test assertions to verify increment counts, per-topic metrics, and per-error-type metrics - Add test for DLQ metrics on full consumer flow All 3844 unit tests pass.
PR Review: Dead Letter Queue (DLQ) Implementation [OMN-949]SummaryThis PR implements comprehensive DLQ support for permanently failing messages in the KafkaEventBus. The implementation is well-architected and follows ONEX patterns closely. Overall assessment: Approve with minor recommendations. ✅ Strengths1. Excellent Architecture & Design
2. ONEX Compliance✅ No 3. Test Coverage
4. Observability & Operations
🔍 Code Quality IssuesCritical IssuesNone identified - Code is production-ready. Medium Priority Recommendations1. DLQ Replay Script: Bootstrap Servers ValidationFile: The current implementation requires # Current: No validation after retrieval
bootstrap_servers = os.environ.get("KAFKA_BOOTSTRAP_SERVERS")
if bootstrap_servers is None:
# raises ValueErrorRecommendation: Add basic format validation to catch typos early: if not bootstrap_servers or not bootstrap_servers.strip():
raise ValueError("KAFKA_BOOTSTRAP_SERVERS cannot be empty")
if ',' in bootstrap_servers:
# Validate each server has host:port format
for server in bootstrap_servers.split(','):
if ':' not in server.strip():
raise ValueError(f"Invalid server format: {server}")2. Error Sanitization: URL Scheme DetectionFile: Current SENSITIVE_PATTERNS = (
"mongodb://", "postgres://", "mysql://", "redis://", ...
)Recommendation: Consider a regex pattern for generic # Add to SENSITIVE_PATTERNS or create separate validation
r"[a-z]+://[^@]+:[^@]+@" # Matches scheme://user:pass@This would catch custom database schemes like 3. Topic Constants: Environment ValidationFile: The def build_dlq_topic(environment: str, category: EnumMessageCategory) -> str:
return f"{environment}.dlq.{category.value}.v1"Recommendation: Add validation to prevent invalid topic names: def build_dlq_topic(environment: str, category: EnumMessageCategory) -> str:
if not environment or '.' in environment:
raise ValueError(f"Invalid environment: {environment}")
return f"{environment}.dlq.{category.value}.v1"Low Priority Suggestions4. ModelDlqMessage: Retry Count Validation Error MessageFile: The # Current
raise ValueError(f"Invalid retry_count value: '{value}' is not a valid integer")
# Suggestion: Include remediation hint
raise ValueError(
f"Invalid retry_count value: '{value}' is not a valid integer. "
"Expected integer >= 0. DLQ message may be corrupted."
)5. KafkaEventBus: DLQ Metrics CollectionFile: The dlq_event = ModelDlqEvent(
# ... existing fields ...
duration_ms=duration_ms, # Add this field to ModelDlqEvent
)This would help identify DLQ publish performance issues. 🔒 Security Review✅ Excellent Security Practices
|
| Area | Test File | Coverage |
|---|---|---|
| DLQ Routing | test_kafka_event_bus.py |
6 tests (retry exhaustion, publish success/failure) |
| Topic Naming | test_topic_constants.py |
44 tests (all categories, edge cases) |
| Error Sanitization | test_util_error_sanitization.py |
267 lines (patterns, edge cases) |
| Integration | test_dlq_integration.py |
760 lines (end-to-end flow) |
Missing Test Cases (Optional)
-
Circuit Breaker + DLQ Interaction
- Verify DLQ publish when circuit breaker is OPEN
- Expected: DLQ publish should succeed even if circuit is open
-
Concurrent DLQ Publishes
- Multiple handlers failing simultaneously
- Verify producer lock prevents race conditions
-
DLQ Replay Error Handling
- Replay with malformed DLQ messages
- Replay with missing
original_topicfield
🎯 ONEX Pattern Compliance
✅ Fully Compliant
| Pattern | Status | Evidence |
|---|---|---|
No Any types |
✅ | All models use specific types |
| Pydantic models | ✅ | ModelDlqEvent, ModelDlqMetrics, ModelReplayConfig |
| Naming conventions | ✅ | model_dlq_event.py → ModelDlqEvent |
| Error hierarchy | ✅ | Uses OnexError subclasses (InfraConnectionError, etc.) |
| Protocol resolution | ✅ | Circuit breaker mixin pattern |
| Strong typing | ✅ | Type hints throughout, mypy passing |
Exception Noted
EnumNonRetryableErrorCategory Location:
- Currently in
omnibase_infra.enums - Per ticket OMN-1032, should eventually move to
omnibase_core.enums
Rationale: This is a known migration plan and doesn't block the PR.
📝 Documentation Quality
✅ Excellent
-
DLQ_MESSAGE_FORMAT.md(399 lines)- Complete JSON schema with examples
- Security considerations section
- Reprocessing guide with code examples
- Monitoring and alerting recommendations
-
Inline Documentation
- Comprehensive docstrings in
_publish_to_dlq() - Clear design notes (e.g., producer lock pattern explanation)
- Type hints aid IDE autocomplete
- Comprehensive docstrings in
-
Operations Guide
DLQ_REPLAY_GUIDE.mdprovides clear procedures- CLI examples with dry-run workflow
Minor Documentation Gap
File: scripts/dlq_replay.py
The script header mentions:
# See Also:
# OMN-1032 - PostgreSQL tracking integration (planned)Recommendation: Add a TODO comment in the relevant function:
async def replay_message(self, message: ModelDlqMessage) -> ModelReplayResult:
# TODO(OMN-1032): Integrate with PostgreSQL tracking for replay state persistence
# Currently stores state in memory only🚀 Final Verdict
APPROVE ✅
This PR represents high-quality production-ready code with:
- Excellent architecture and ONEX compliance
- Strong security posture with credential sanitization
- Comprehensive test coverage (71 event_bus + 44 topic_constants + integration tests)
- Well-documented DLQ format and replay procedures
- No critical issues identified
Recommended Actions Before Merge
-
Address medium-priority recommendations (optional but valuable):
- Add bootstrap servers validation in replay script
- Consider URL scheme regex pattern for sanitization
- Add environment validation in
build_dlq_topic()
-
Update operations documentation with security warning about
original_message.valuepreservation -
Verify integration test with live Kafka (mentioned in test plan checklist)
Post-Merge Follow-Up
- OMN-1032: Migrate
EnumNonRetryableErrorCategorytoomnibase_core - Performance monitoring: Track DLQ publish latency in production
- Alert tuning: Calibrate thresholds based on production traffic
📦 Approval Summary
Code Quality: ⭐⭐⭐⭐⭐ (5/5)
Test Coverage: ⭐⭐⭐⭐⭐ (5/5)
Documentation: ⭐⭐⭐⭐⭐ (5/5)
Security: ⭐⭐⭐⭐⭐ (5/5)
ONEX Compliance: ⭐⭐⭐⭐⭐ (5/5)
Overall: ⭐⭐⭐⭐⭐ LGTM - Excellent work!
Review conducted according to ONEX coding standards and CLAUDE.md guidelines.
Add TODO comment in replay_message() referencing OMN-1032 for future PostgreSQL state persistence integration.
Code Review: DLQ Implementation (OMN-949)This is a well-architected DLQ implementation with strong ONEX compliance. Excellent documentation and test coverage. Strengths
Issues Requiring AttentionCRITICAL: Type Annotation Consistency
HIGH: DLQ Replay Script Status
MEDIUM: EnumNonRetryableErrorCategory Completeness
MEDIUM: DLQ Metrics Thread Safety
Security Review
Pre-Merge ChecklistMust Fix:
Should Fix:
SummaryExcellent work with strong ONEX architecture. Fix type annotations and clarify replay script status before merge. Recommendation: Approve with required changes Reviewed by: Claude Code (Sonnet 4.5) |
Summary
Implements comprehensive Dead Letter Queue (DLQ) support for permanently failing intents in the KafkaEventBus, addressing all acceptance criteria from OMN-949.
Changes
topic_constants.pywith ONEX-compliant naming pattern{env}.dlq.{category}.v1ModelDlqEventandModelDlqMetricswith callback hooks for custom alertingDLQ_MESSAGE_FORMAT.mdpayload schema andDLQ_REPLAY_GUIDE.mdoperations guidescripts/dlq_replay.pyCLI utility skeletonFiles Changed
model_dlq_event.py,model_dlq_metrics.pytopic_constants.pyDLQ_MESSAGE_FORMAT.md,DLQ_REPLAY_GUIDE.mdscripts/dlq_replay.pykafka_event_bus.py,model_kafka_event_bus_config.pytest_topic_constants.py(44 tests), 6 DLQ routing testsAcceptance Criteria
{env}.dlq.{category}.v1)Test plan
Linear Issue
Closes OMN-949
Summary by CodeRabbit
New Features
Documentation
Tests
✏️ Tip: You can customize this high-level summary in your review settings.