Repository navigation
feat(OMN-1621): Add contract-driven event bus subscription wiring - #200
Conversation
Implement automatic Kafka topic wiring from handler contract event_bus sections, so runtime subscribes to declared topics without manual wiring. New modules: - EventBusSubcontractWiring: Wires subscribe_topics to Kafka consumers - PublisherTopicScoped: Validates publishing against contract publish_topics - load_event_bus_subcontract(): Loads event_bus section from contract YAML Runtime integration: - ServiceRuntimeHostProcess now calls _wire_event_bus_subscriptions() - Cleanup in stop() unsubscribes from all wired topics - Conditional activation when dispatch_engine is present Architecture (ARCH-002): - Runtime owns all Kafka plumbing - Handlers receive ModelEventEnvelope, not raw Kafka messages - Publishers validate against contract-declared topics Tests: 81 new tests (58 unit, 23 integration)
📝 WalkthroughWalkthroughAdds a contract-driven event-bus wiring component and a topic-scoped publisher, introduces a dispatch protocol, integrates wiring into the runtime host lifecycle, and adds extensive unit and integration tests for wiring, publishing, deserialization, and cleanup. Changes
Sequence Diagram(s)sequenceDiagram
participant Contract as Contract File
participant Runtime as RuntimeHostProcess
participant Wiring as EventBusSubcontractWiring
participant EventBus as Event Bus
participant Dispatch as Message Dispatch Engine
participant Handler as Handler
Contract->>Runtime: load_event_bus_subcontract(contract_path)
Runtime->>Wiring: instantiate wiring(event_bus, dispatch_engine, environment)
Runtime->>Wiring: wire_subscriptions(subcontract, node_name)
Wiring->>EventBus: subscribe(topic, group_id, callback)
EventBus->>Wiring: deliver ProtocolEventMessage
Wiring->>Wiring: _deserialize_to_envelope(message)
Wiring->>Dispatch: dispatch(topic, ModelEventEnvelope)
Dispatch->>Handler: route envelope to handler
sequenceDiagram
participant Client as Client Code
participant Publisher as PublisherTopicScoped
participant Validator as Allowed Topics
participant EventBus as Event Bus
Client->>Publisher: publish(event_type,payload,topic,correlation_id)
Publisher->>Validator: check topic in allowed_topics
alt allowed
Publisher->>Publisher: resolve_topic(env,topic)
Publisher->>Publisher: serialize payload -> JSON bytes
Publisher->>EventBus: publish(full_topic, key, value)
EventBus-->>Publisher: ack
Publisher-->>Client: return True
else disallowed
Publisher-->>Client: raise ProtocolConfigurationError
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing touches
Comment |
There was a problem hiding this comment.
Actionable comments posted: 3
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/runtime/event_bus_subcontract_wiring.py`:
- Around line 259-280: The callback function should stop re-raising raw
exceptions and instead wrap them in an OnexError subclass; catch
json.JSONDecodeError and other Exception inside callback, create/raise a
specific OnexError (e.g., MessageDeserializationError or MessageDispatchError)
using "raise YourOnexError(...) from e" so the original cause is preserved, and
keep the existing self._logger.exception calls; update any tests that assert
JSONDecodeError to expect the new OnexError subclass. Ensure references:
callback, _deserialize_to_envelope, and _dispatch_engine.dispatch are the
locations to implement the OnexError wrapping.
In `@src/omnibase_infra/runtime/publisher_topic_scoped.py`:
- Around line 169-219: The publish method currently assumes correlation_id is a
str and calls correlation_id.encode(), which will fail if a UUID is passed;
update the method signature/usage to accept UUID | str for correlation_id and
normalize it before encoding (e.g., use str(correlation_id) when building key),
ensuring key = str(correlation_id).encode("utf-8") if correlation_id else None;
keep all existing checks (topic validation using self._allowed_topics and
resolve_topic) unchanged.
- Around line 205-212: Replace the two ValueError raises in PublisherTopicScoped
(the "topic is required for PublisherTopicScoped" and the "Topic '{topic}' not
in contract's publish_topics..." checks against self._allowed_topics) with an
OnexError subclass such as ProtocolConfigurationError, preserving the existing
error messages and f-strings; add the necessary import for
ProtocolConfigurationError from your project's OnexError module and update any
tests that assert ValueError to expect ProtocolConfigurationError instead.
🧹 Nitpick comments (5)
src/omnibase_infra/runtime/publisher_topic_scoped.py (1)
118-145: Freeze allowed_topics to prevent post-init mutation
allowed_topicsis stored as the caller’s mutable set, so external mutation can silently expand permissions. Consider freezing internally and returning it directly.♻️ Proposed change
- self._allowed_topics = allowed_topics + self._allowed_topics = frozenset(allowed_topics) @@ - return frozenset(self._allowed_topics) + return self._allowed_topicsAlso applies to: 237-248
src/omnibase_infra/runtime/event_bus_subcontract_wiring.py (1)
284-315: Prefer JsonType for envelope payload typing
ModelEventEnvelope[object]loses JSON typing; usingJsonTypematches the payload contract. As per coding guidelines, preferJsonTypefor JSON-compatible values.♻️ Proposed change
+from omnibase_core.types import JsonType @@ - ) -> ModelEventEnvelope[object]: + ) -> ModelEventEnvelope[JsonType]: @@ - return ModelEventEnvelope[object].model_validate(data) + return ModelEventEnvelope[JsonType].model_validate(data)tests/unit/runtime/test_event_bus_subcontract_wiring.py (3)
35-36: Unused import:Generatoris imported but never used.The
Generatortype imported in theTYPE_CHECKINGblock is not used anywhere in this test file.🧹 Remove unused import
-if TYPE_CHECKING: - from collections.abc import Generator
282-296: Consider strengthening the envelope assertion.The test only verifies
envelope is not None, which is a weak assertion. Consider verifying the envelope type and that the deserialized payload matches expectations.💡 Suggested enhancement
call_args = mock_dispatch_engine.dispatch.call_args envelope = call_args[0][1] - # Envelope should be deserialized from message - assert envelope is not None + # Envelope should be deserialized from message with correct payload + assert envelope is not None + assert envelope.payload == {"key": "value"} + assert envelope.event_type == "test.event"
521-532: Remove unused variablecontract_file.The variable
contract_fileis defined on line 523 but never used. The test intentionally uses a hardcoded nonexistent path to trigger the warning.🧹 Remove unused variable
def test_uses_provided_logger(self, tmp_path: Path) -> None: """Test function uses provided logger.""" - contract_file = tmp_path / "contract.yaml" - # Non-existent file to trigger warning - mock_logger = MagicMock() load_event_bus_subcontract( Path("/nonexistent/contract.yaml"), logger=mock_logger, ) mock_logger.warning.assert_called()Also, consider removing
tmp_pathfrom the fixture parameters since it's not needed after this fix.
… type: ignore Replace duck-typed object parameters with existing protocol types: - event_bus_subcontract_wiring.py: Use ProtocolEventBusSubscriber - publisher_topic_scoped.py: Use ProtocolEventBusLike - Remove redundant deferred import in service_runtime_host_process.py This improves type safety and removes # type: ignore[attr-defined] comments.
…tWiring Update EventBusSubcontractWiring to use the ProtocolDispatchEngine protocol instead of `object` with `type: ignore`. This provides proper type safety while maintaining duck typing per ONEX conventions. Changes: - Import ProtocolDispatchEngine from omnibase_infra.protocols - Update dispatch_engine parameter type annotation - Remove type: ignore comment from dispatch call - Export ProtocolDispatchEngine from protocols __init__.py
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/runtime/service_runtime_host_process.py`:
- Around line 2475-2479: ProtocolDispatchEngine.dispatch currently declares a
None return but MessageDispatchEngine.dispatch returns ModelDispatchResult,
causing a protocol mismatch when MessageDispatchEngine is passed into
EventBusSubcontractWiring; change the ProtocolDispatchEngine.dispatch signature
to return ModelDispatchResult (i.e., update the protocol definition for
ProtocolDispatchEngine.dispatch to -> ModelDispatchResult) so implementations
like MessageDispatchEngine.dispatch remain compliant and
EventBusSubcontractWiring can rely on the routing result.
- Add environment parameter validation to EventBusSubcontractWiring and PublisherTopicScoped to reject empty/whitespace-only values - Narrow exception handling in load_event_bus_subcontract from broad Exception to specific ValidationError from pydantic - Improve runtime wiring log messages to include actual topic names for better operational visibility and debugging - Add explicit @pytest.mark.asyncio decorators to 15 async test methods for clarity (tests pass due to asyncio mode=AUTO but explicit is better)
There was a problem hiding this comment.
Actionable comments posted: 5
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/protocols/protocol_dispatch_engine.py`:
- Around line 58-68: Replace the TypeError raise in the protocol compliance
example with a ProtocolConfigurationError constructed via
ModelInfraErrorContext.with_correlation so examples follow the OnexError
subclass policy; locate the conditional that checks hasattr(engine, 'dispatch')
and callable(engine.dispatch) (related to engine/dispatch and
EventBusSubcontractWiring) and change the else branch to raise
ProtocolConfigurationError using ModelInfraErrorContext.with_correlation,
ensuring the error communicates the object does not implement
ProtocolDispatchEngine.
In `@src/omnibase_infra/runtime/event_bus_subcontract_wiring.py`:
- Around line 412-431: The code assumes the YAML root is a mapping and calls
contract_data.get(...), which will raise AttributeError for non-mapping YAML;
update the contract loading block (around contract_path, contract_data, and the
return of ModelEventBusSubcontract.model_validate) to first check
isinstance(contract_data, dict) (or similar mapping check), log a warning via
_logger that the contract file has an unexpected structure and the path, and
return None when the root is not a mapping so invalid contracts degrade
gracefully instead of raising.
- Around line 125-156: In the __init__ of EventBusSubcontractWiring (the
constructor shown), replace the ValueError check with normalizing the
environment (call environment = environment.strip()) and if the result is empty
raise ProtocolConfigurationError constructed via
ModelInfraErrorContext.with_correlation (include a clear message like
"environment must be a non-empty string"); assign the stripped value to
self._environment; ensure you import or reference ProtocolConfigurationError and
ModelInfraErrorContext.with_correlation so the raised error conforms to the
OnexError-only policy.
In `@src/omnibase_infra/runtime/service_runtime_host_process.py`:
- Around line 2463-2471: The current dict branch always converts the environment
to a string (str(event_bus_config.get("environment", "local"))) which turns None
into "None" and bypasses later validation; change it so you first retrieve the
raw value (e.g., val = event_bus_config.get("environment", None)), treat None as
missing by leaving environment as None, and only call str(...) when val is not
None (so empty/whitespace values still propagate for EventBusSubcontractWiring
to validate). Update the branch that sets environment (referencing
event_bus_config and environment) to implement this conditional conversion.
- Around line 691-699: The _dispatch_engine field is never set so
_wire_event_bus_subscriptions() returns early; fix by wiring a
MessageDispatchEngine during initialization or via constructor injection: add an
optional parameter (e.g., dispatch_engine: MessageDispatchEngine | None) to
__init__ and assign it to self._dispatch_engine, or resolve/instantiate
MessageDispatchEngine inside __init__ before EventBusSubcontractWiring is used;
ensure start() can rely on self._dispatch_engine being non-None when
_wire_event_bus_subscriptions() runs and update any container resolution logic
to supply the engine if you choose constructor injection (refer to symbols:
__init__, self._dispatch_engine, _wire_event_bus_subscriptions, start,
MessageDispatchEngine, EventBusSubcontractWiring).
| .. code-block:: python | ||
|
|
||
| # Verify required method exists and is callable | ||
| if hasattr(engine, 'dispatch') and callable(engine.dispatch): | ||
| wiring = EventBusSubcontractWiring( | ||
| event_bus=event_bus, | ||
| dispatch_engine=engine, | ||
| environment="dev", | ||
| ) | ||
| else: | ||
| raise TypeError("Object does not implement ProtocolDispatchEngine") |
There was a problem hiding this comment.
Avoid TypeError in the protocol compliance example
Line 68 raises TypeError, but the project mandates OnexError subclasses for raised errors. Please update the example to use ProtocolConfigurationError (with ModelInfraErrorContext.with_correlation) so the docs match policy. As per coding guidelines, only OnexError subclasses should be raised.
♻️ Docstring tweak
- raise TypeError("Object does not implement ProtocolDispatchEngine")
+ raise ProtocolConfigurationError(
+ "Object does not implement ProtocolDispatchEngine",
+ context=ModelInfraErrorContext.with_correlation(
+ transport_type=EnumInfraTransportType.RUNTIME,
+ operation="dispatch_engine_validation",
+ target_name=type(engine).__name__,
+ ),
+ )🤖 Prompt for AI Agents
In `@src/omnibase_infra/protocols/protocol_dispatch_engine.py` around lines 58 -
68, Replace the TypeError raise in the protocol compliance example with a
ProtocolConfigurationError constructed via
ModelInfraErrorContext.with_correlation so examples follow the OnexError
subclass policy; locate the conditional that checks hasattr(engine, 'dispatch')
and callable(engine.dispatch) (related to engine/dispatch and
EventBusSubcontractWiring) and change the else branch to raise
ProtocolConfigurationError using ModelInfraErrorContext.with_correlation,
ensuring the error communicates the object does not implement
ProtocolDispatchEngine.
| def __init__( | ||
| self, | ||
| event_bus: ProtocolEventBusSubscriber, | ||
| dispatch_engine: ProtocolDispatchEngine, | ||
| environment: str, | ||
| ) -> None: | ||
| """Initialize event bus wiring. | ||
|
|
||
| Args: | ||
| event_bus: The event bus implementation (EventBusKafka or EventBusInmemory). | ||
| Must implement subscribe(topic, group_id, on_message) -> unsubscribe callable. | ||
| Duck typed per ONEX patterns. | ||
| dispatch_engine: Engine to dispatch received messages to handlers. | ||
| Must implement ProtocolDispatchEngine interface. | ||
| Must be frozen (registrations complete) before wiring subscriptions. | ||
| environment: Environment prefix for topics (e.g., 'dev', 'prod'). | ||
| Used to resolve topic suffixes to full topic names. | ||
|
|
||
| Note: | ||
| The dispatch_engine should be frozen before wiring subscriptions. | ||
| Attempting to dispatch to an unfrozen engine will raise an error. | ||
|
|
||
| Raises: | ||
| ValueError: If environment is empty or whitespace-only. | ||
| """ | ||
| if not environment or not environment.strip(): | ||
| raise ValueError("environment must be a non-empty string") | ||
|
|
||
| self._event_bus = event_bus | ||
| self._dispatch_engine = dispatch_engine | ||
| self._environment = environment | ||
| self._unsubscribe_callables: list[Callable[[], Awaitable[None]]] = [] |
There was a problem hiding this comment.
Raise ProtocolConfigurationError and normalize environment
Line 150 raises ValueError, which violates the OnexError-only policy, and the value isn’t normalized. Strip the environment and raise ProtocolConfigurationError with ModelInfraErrorContext.with_correlation. As per coding guidelines, only OnexError subclasses should be raised.
🐛 Proposed fix
-from typing import TYPE_CHECKING
+from typing import TYPE_CHECKING
+
+from omnibase_infra.enums import EnumInfraTransportType
+from omnibase_infra.errors import ModelInfraErrorContext, ProtocolConfigurationError
@@
- if not environment or not environment.strip():
- raise ValueError("environment must be a non-empty string")
+ environment = environment.strip()
+ if not environment:
+ raise ProtocolConfigurationError(
+ "environment must be a non-empty string",
+ context=ModelInfraErrorContext.with_correlation(
+ transport_type=EnumInfraTransportType.RUNTIME,
+ operation="event_bus_subcontract_wiring_init",
+ target_name="event_bus_subcontract_wiring",
+ ),
+ )
@@
- self._environment = environment
+ self._environment = environment🤖 Prompt for AI Agents
In `@src/omnibase_infra/runtime/event_bus_subcontract_wiring.py` around lines 125
- 156, In the __init__ of EventBusSubcontractWiring (the constructor shown),
replace the ValueError check with normalizing the environment (call environment
= environment.strip()) and if the result is empty raise
ProtocolConfigurationError constructed via
ModelInfraErrorContext.with_correlation (include a clear message like
"environment must be a non-empty string"); assign the stripped value to
self._environment; ensure you import or reference ProtocolConfigurationError and
ModelInfraErrorContext.with_correlation so the raised error conforms to the
OnexError-only policy.
| try: | ||
| with contract_path.open() as f: | ||
| contract_data = yaml.safe_load(f) | ||
|
|
||
| if contract_data is None: | ||
| _logger.warning( | ||
| "Empty contract file: %s", | ||
| contract_path, | ||
| ) | ||
| return None | ||
|
|
||
| event_bus_data = contract_data.get("event_bus") | ||
| if not event_bus_data: | ||
| _logger.debug( | ||
| "No event_bus section in contract: %s", | ||
| contract_path, | ||
| ) | ||
| return None | ||
|
|
||
| return ModelEventBusSubcontract.model_validate(event_bus_data) |
There was a problem hiding this comment.
Guard against non-mapping contract YAML
If the YAML root isn’t a mapping, contract_data.get(...) will raise AttributeError and break startup. Add a type guard to log and return None so invalid contracts degrade gracefully.
🐛 Proposed fix
- event_bus_data = contract_data.get("event_bus")
+ if not isinstance(contract_data, dict):
+ _logger.warning(
+ "Contract root must be a mapping, got %s: %s",
+ type(contract_data).__name__,
+ contract_path,
+ )
+ return None
+
+ event_bus_data = contract_data.get("event_bus")📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| try: | |
| with contract_path.open() as f: | |
| contract_data = yaml.safe_load(f) | |
| if contract_data is None: | |
| _logger.warning( | |
| "Empty contract file: %s", | |
| contract_path, | |
| ) | |
| return None | |
| event_bus_data = contract_data.get("event_bus") | |
| if not event_bus_data: | |
| _logger.debug( | |
| "No event_bus section in contract: %s", | |
| contract_path, | |
| ) | |
| return None | |
| return ModelEventBusSubcontract.model_validate(event_bus_data) | |
| try: | |
| with contract_path.open() as f: | |
| contract_data = yaml.safe_load(f) | |
| if contract_data is None: | |
| _logger.warning( | |
| "Empty contract file: %s", | |
| contract_path, | |
| ) | |
| return None | |
| try: | |
| event_bus_data = contract_data.get("event_bus") | |
| except AttributeError: | |
| _logger.warning( | |
| "Contract root must be a mapping: %s", | |
| contract_path, | |
| ) | |
| return None | |
| if not event_bus_data: | |
| _logger.debug( | |
| "No event_bus section in contract: %s", | |
| contract_path, | |
| ) | |
| return None | |
| return ModelEventBusSubcontract.model_validate(event_bus_data) |
🤖 Prompt for AI Agents
In `@src/omnibase_infra/runtime/event_bus_subcontract_wiring.py` around lines 412
- 431, The code assumes the YAML root is a mapping and calls
contract_data.get(...), which will raise AttributeError for non-mapping YAML;
update the contract loading block (around contract_path, contract_data, and the
return of ModelEventBusSubcontract.model_validate) to first check
isinstance(contract_data, dict) (or similar mapping check), log a warning via
_logger that the contract file has an unexpected structure and the path, and
return None when the root is not a mapping so invalid contracts degrade
gracefully instead of raising.
| # Event bus subcontract wiring for contract-driven subscriptions (OMN-1621) | ||
| # Bridges contract-declared topics to Kafka subscriptions. | ||
| # None until wired during start() when dispatch_engine is available. | ||
| self._event_bus_wiring: EventBusSubcontractWiring | None = None | ||
|
|
||
| # Message dispatch engine for routing received messages (OMN-1621) | ||
| # Used by event_bus_wiring to dispatch messages to handlers. | ||
| # None = not configured, wiring will be skipped | ||
| self._dispatch_engine: MessageDispatchEngine | None = None |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Find all assignments to _dispatch_engine
rg -n "_dispatch_engine\s*=" src/omnibase_infra/runtime/service_runtime_host_process.py
echo "---"
# Find all references to _dispatch_engine
rg -n "_dispatch_engine" src/omnibase_infra/runtime/service_runtime_host_process.py | head -20
echo "---"
# Check MessageDispatchEngine usage and imports
rg -n "MessageDispatchEngine" src/omnibase_infra/runtime/service_runtime_host_process.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 511
🏁 Script executed:
#!/bin/bash
# Get the implementation of _wire_event_bus_subscriptions() around lines 2451-2455
sed -n '2445,2465p' src/omnibase_infra/runtime/service_runtime_host_process.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 794
🏁 Script executed:
#!/bin/bash
# Get the start() method around lines 1003-1006
sed -n '995,1015p' src/omnibase_infra/runtime/service_runtime_host_process.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 1046
🏁 Script executed:
#!/bin/bash
# Check if there's any initialization logic for _dispatch_engine elsewhere in the file
rg -n "dispatch_engine" src/omnibase_infra/runtime/service_runtime_host_process.py | head -30Repository: OmniNode-ai/omnibase_infra
Length of output: 709
🏁 Script executed:
#!/bin/bash
# Get the __init__ method to see if _dispatch_engine is passed as parameter
sed -n '200,350p' src/omnibase_infra/runtime/service_runtime_host_process.py | grep -A 200 "def __init__"Repository: OmniNode-ai/omnibase_infra
Length of output: 7006
🏁 Script executed:
#!/bin/bash
# Search for any assignment patterns with dispatch_engine (including indirect)
rg -n "dispatch_engine\s*=" src/omnibase_infra/runtime/service_runtime_host_process.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 125
🏁 Script executed:
#!/bin/bash
# Check the class definition and __init__ signature more carefully
head -400 src/omnibase_infra/runtime/service_runtime_host_process.py | tail -250Repository: OmniNode-ai/omnibase_infra
Length of output: 11257
🏁 Script executed:
#!/bin/bash
# Look for any method that might set _dispatch_engine
rg -n "def.*dispatch" src/omnibase_infra/runtime/service_runtime_host_process.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 52
🏁 Script executed:
#!/bin/bash
# Get lines with __init__ definition and next 100 lines to see the full initialization
rg -n "def __init__" src/omnibase_infra/runtime/service_runtime_host_process.py -A 100 | head -150Repository: OmniNode-ai/omnibase_infra
Length of output: 8276
🏁 Script executed:
#!/bin/bash
# Also search for any method that might initialize dispatch_engine
rg -n "self._dispatch_engine" src/omnibase_infra/runtime/service_runtime_host_process.py -B 2 -A 2Repository: OmniNode-ai/omnibase_infra
Length of output: 1013
🏁 Script executed:
#!/antml/bash
# Check the full file to understand the context better - get line count first
wc -l src/omnibase_infra/runtime/service_runtime_host_process.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 133
🏁 Script executed:
#!/bin/bash
# Get the full __init__ method - it starts at line 366
sed -n '366,750p' src/omnibase_infra/runtime/service_runtime_host_process.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 17936
🏁 Script executed:
#!/bin/bash
# Search the entire codebase for any assignment to _dispatch_engine
rg "self\._dispatch_engine\s*=" --type pyRepository: OmniNode-ai/omnibase_infra
Length of output: 176
🏁 Script executed:
#!/bin/bash
# Get the start() method implementation (should be around line 1003-1006 range based on comments)
sed -n '995,1050p' src/omnibase_infra/runtime/service_runtime_host_process.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 2402
🏁 Script executed:
#!/bin/bash
# Look for any property setter for _dispatch_engine or MessageDispatchEngine instantiation
rg "MessageDispatchEngine" src/omnibase_infra/runtime/service_runtime_host_process.py -B 2 -A 2Repository: OmniNode-ai/omnibase_infra
Length of output: 404
🏁 Script executed:
#!/an/bash
# Check the exact implementation of _wire_event_bus_subscriptions to understand the flow
sed -n '2420,2490p' src/omnibase_infra/runtime/service_runtime_host_process.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 3152
Contract wiring never activates: _dispatch_engine has no initialization path
self._dispatch_engine is initialized to None (line 699) with no parameter in __init__ and no assignment anywhere in the codebase. When start() calls _wire_event_bus_subscriptions() (line 1006), the method short-circuits at line 2451 (if not self._dispatch_engine: return) and never wires contract-declared topics. The precondition comment at line 2429 states "_dispatch_engine must be set" but provides no mechanism to do so. Either instantiate MessageDispatchEngine in __init__ or accept it as a constructor parameter with container resolution.
🤖 Prompt for AI Agents
In `@src/omnibase_infra/runtime/service_runtime_host_process.py` around lines 691
- 699, The _dispatch_engine field is never set so
_wire_event_bus_subscriptions() returns early; fix by wiring a
MessageDispatchEngine during initialization or via constructor injection: add an
optional parameter (e.g., dispatch_engine: MessageDispatchEngine | None) to
__init__ and assign it to self._dispatch_engine, or resolve/instantiate
MessageDispatchEngine inside __init__ before EventBusSubcontractWiring is used;
ensure start() can rely on self._dispatch_engine being non-None when
_wire_event_bus_subscriptions() runs and update any container resolution logic
to supply the engine if you choose constructor injection (refer to symbols:
__init__, self._dispatch_engine, _wire_event_bus_subscriptions, start,
MessageDispatchEngine, EventBusSubcontractWiring).
| # Get environment from config, defaulting to "local" | ||
| config = self._config or {} | ||
| event_bus_config = config.get("event_bus", {}) | ||
| if isinstance(event_bus_config, dict): | ||
| environment = str(event_bus_config.get("environment", "local")) | ||
| else: | ||
| # Could be a ModelEventBusConfig or similar | ||
| environment = getattr(event_bus_config, "environment", "local") | ||
|
|
There was a problem hiding this comment.
Avoid str(None) for event_bus.environment
In the dict branch, str(None) becomes "None", which bypasses the environment validation and prefixes topics with "None.". Treat None as missing and let empty/whitespace values fall through to validation in EventBusSubcontractWiring.
🐛 Suggested fix
- if isinstance(event_bus_config, dict):
- environment = str(event_bus_config.get("environment", "local"))
- else:
- # Could be a ModelEventBusConfig or similar
- environment = getattr(event_bus_config, "environment", "local")
+ if isinstance(event_bus_config, dict):
+ env_raw = event_bus_config.get("environment", "local")
+ else:
+ # Could be a ModelEventBusConfig or similar
+ env_raw = getattr(event_bus_config, "environment", "local")
+ environment = "local" if env_raw is None else str(env_raw)🤖 Prompt for AI Agents
In `@src/omnibase_infra/runtime/service_runtime_host_process.py` around lines 2463
- 2471, The current dict branch always converts the environment to a string
(str(event_bus_config.get("environment", "local"))) which turns None into "None"
and bypasses later validation; change it so you first retrieve the raw value
(e.g., val = event_bus_config.get("environment", None)), treat None as missing
by leaving environment as None, and only call str(...) when val is not None (so
empty/whitespace values still propagate for EventBusSubcontractWiring to
validate). Update the branch that sets environment (referencing event_bus_config
and environment) to implement this conditional conversion.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/runtime/publisher_topic_scoped.py`:
- Around line 148-149: Replace the raised ValueError with
ProtocolConfigurationError in the environment validation branch: where the code
currently checks "if not environment or not environment.strip(): raise
ValueError(...)", change the exception to ProtocolConfigurationError (keeping or
slightly refining the same error message) so the module raises an OnexError
subclass; ensure the imported ProtocolConfigurationError is used and no other
behavior changes.
| if not environment or not environment.strip(): | ||
| raise ValueError("environment must be a non-empty string") |
There was a problem hiding this comment.
Use ProtocolConfigurationError instead of ValueError for OnexError policy compliance.
Per coding guidelines, only OnexError subclasses should be raised. The environment validation should use ProtocolConfigurationError which is already imported.
🐛 Proposed fix
if not environment or not environment.strip():
- raise ValueError("environment must be a non-empty string")
+ raise ProtocolConfigurationError(
+ "environment must be a non-empty string"
+ )📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| if not environment or not environment.strip(): | |
| raise ValueError("environment must be a non-empty string") | |
| if not environment or not environment.strip(): | |
| raise ProtocolConfigurationError( | |
| "environment must be a non-empty string" | |
| ) |
🤖 Prompt for AI Agents
In `@src/omnibase_infra/runtime/publisher_topic_scoped.py` around lines 148 - 149,
Replace the raised ValueError with ProtocolConfigurationError in the environment
validation branch: where the code currently checks "if not environment or not
environment.strip(): raise ValueError(...)", change the exception to
ProtocolConfigurationError (keeping or slightly refining the same error message)
so the module raises an OnexError subclass; ensure the imported
ProtocolConfigurationError is used and no other behavior changes.
…UUID handling - Wrap raw exceptions in OnexError subclasses per CLAUDE.md policy: - JSONDecodeError -> RuntimeHostError in event_bus_subcontract_wiring.py - ValueError -> ProtocolConfigurationError in publisher_topic_scoped.py - Fix UUID correlation_id bug: accept str | UUID | None, normalize with str() - Add _normalize_correlation_id() helper for explicit intent - Freeze allowed_topics internally to prevent mutation - Remove unused Generator import from test file - Update tests to expect new exception types - Reduce union count by replacing list|set and list|tuple with Collection/Sequence Code review findings from CodeRabbit addressed.
ea45f4d to
2583a79
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
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/handlers/models/model_filesystem_config.py (1)
62-78: Reject string inputs forallowed_paths.
Sequence[str]includesstr, sotuple(v)will split a path string into individual characters and undermine the whitelist. Add an explicit guard forstr/bytes:🔒 Proposed fix
`@field_validator`("allowed_paths", mode="before") `@classmethod` def coerce_to_tuple(cls, v: Sequence[str]) -> tuple[str, ...]: - # Always return tuple to satisfy return type - return tuple(v) + if isinstance(v, (str, bytes)): + raise ValueError("allowed_paths must be a sequence of paths, not a string") + return tuple(v)
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/runtime/event_bus_subcontract_wiring.py`:
- Around line 251-308: The callback in _create_dispatch_callback currently
treats ValidationError from deserialization as a dispatch error; update the
exception handling so that ValidationError (raised by
ModelEventEnvelope.model_validate inside _deserialize_to_envelope) is caught
alongside json.JSONDecodeError: add an except ValidationError (or a tuple except
(json.JSONDecodeError, ValidationError)) branch that logs with the same message
and raises RuntimeHostError using
ModelInfraErrorContext.with_correlation(transport_type=EnumInfraTransportType.KAFKA,
operation="event_bus_deserialize") so deserialization failures are correctly
classified, leaving the existing generic except Exception for true dispatch
errors from await self._dispatch_engine.dispatch.
| def _create_dispatch_callback( | ||
| self, | ||
| topic: str, | ||
| ) -> Callable[[ProtocolEventMessage], Awaitable[None]]: | ||
| """Create callback that bridges Kafka consumer to dispatch engine. | ||
|
|
||
| Creates an async callback function that: | ||
| 1. Receives ProtocolEventMessage from the Kafka consumer | ||
| 2. Deserializes the message value to ModelEventEnvelope | ||
| 3. Dispatches the envelope to the MessageDispatchEngine | ||
|
|
||
| Error Handling: | ||
| - Deserialization errors are logged and the message is skipped | ||
| - Dispatch errors are propagated (handled by the event bus DLQ logic) | ||
|
|
||
| Args: | ||
| topic: The full topic name for routing context in logs. | ||
|
|
||
| Returns: | ||
| Async callback function compatible with event bus subscribe(). | ||
| """ | ||
|
|
||
| async def callback(message: ProtocolEventMessage) -> None: | ||
| """Process incoming Kafka message and dispatch to engine.""" | ||
| try: | ||
| envelope = self._deserialize_to_envelope(message) | ||
| # Dispatch via ProtocolDispatchEngine interface | ||
| await self._dispatch_engine.dispatch(topic, envelope) | ||
| except json.JSONDecodeError as e: | ||
| self._logger.exception( | ||
| "Failed to deserialize message from topic '%s': %s", | ||
| topic, | ||
| e, | ||
| ) | ||
| # Wrap in OnexError per CLAUDE.md: "OnexError Only" | ||
| raise RuntimeHostError( | ||
| f"Failed to deserialize message from topic '{topic}'", | ||
| context=ModelInfraErrorContext.with_correlation( | ||
| transport_type=EnumInfraTransportType.KAFKA, | ||
| operation="event_bus_deserialize", | ||
| ), | ||
| ) from e | ||
| except Exception as e: | ||
| self._logger.exception( | ||
| "Failed to dispatch message from topic '%s': %s", | ||
| topic, | ||
| e, | ||
| ) | ||
| # Wrap in OnexError per CLAUDE.md: "OnexError Only" | ||
| raise RuntimeHostError( | ||
| f"Failed to dispatch message from topic '{topic}'", | ||
| context=ModelInfraErrorContext.with_correlation( | ||
| transport_type=EnumInfraTransportType.KAFKA, | ||
| operation="event_bus_dispatch", | ||
| ), | ||
| ) from e | ||
|
|
||
| return callback |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
# Check if the file exists and get its structure
wc -l src/omnibase_infra/runtime/event_bus_subcontract_wiring.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 132
🏁 Script executed:
# Read the file to see the _deserialize_to_envelope method and imports
head -50 src/omnibase_infra/runtime/event_bus_subcontract_wiring.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 2107
🏁 Script executed:
# Search for _deserialize_to_envelope method implementation
rg -n "_deserialize_to_envelope" src/omnibase_infra/runtime/event_bus_subcontract_wiring.py -A 15Repository: OmniNode-ai/omnibase_infra
Length of output: 1645
🏁 Script executed:
# Search for ModelEventEnvelope definition and usage
rg -n "ModelEventEnvelope" src/omnibase_infra/runtime/event_bus_subcontract_wiring.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 688
🏁 Script executed:
# Check current imports in the file
rg -n "^import |^from " src/omnibase_infra/runtime/event_bus_subcontract_wiring.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 832
🏁 Script executed:
# Search for ValidationError usage in the file
rg -n "ValidationError" src/omnibase_infra/runtime/event_bus_subcontract_wiring.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 229
🏁 Script executed:
# Get the complete _deserialize_to_envelope method
sed -n '310,345p' src/omnibase_infra/runtime/event_bus_subcontract_wiring.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 1271
🏁 Script executed:
# Check the full callback implementation to understand exception flow
sed -n '268,308p' src/omnibase_infra/runtime/event_bus_subcontract_wiring.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 1847
Catch ValidationError in deserialization exception handler.
ValidationError from ModelEventEnvelope.model_validate() (line 341) is currently caught by the generic except Exception handler and logged with operation="event_bus_dispatch". Since it's a deserialization error, it should be caught alongside json.JSONDecodeError and logged with operation="event_bus_deserialize" for accurate context.
Proposed fix
- except json.JSONDecodeError as e:
+ except (json.JSONDecodeError, ValidationError) as e:
self._logger.exception(
"Failed to deserialize message from topic '%s': %s",
topic,
e,
)
# Wrap in OnexError per CLAUDE.md: "OnexError Only"
raise RuntimeHostError(
f"Failed to deserialize message from topic '{topic}'",
context=ModelInfraErrorContext.with_correlation(
transport_type=EnumInfraTransportType.KAFKA,
operation="event_bus_deserialize",
),
) from e🤖 Prompt for AI Agents
In `@src/omnibase_infra/runtime/event_bus_subcontract_wiring.py` around lines 251
- 308, The callback in _create_dispatch_callback currently treats
ValidationError from deserialization as a dispatch error; update the exception
handling so that ValidationError (raised by ModelEventEnvelope.model_validate
inside _deserialize_to_envelope) is caught alongside json.JSONDecodeError: add
an except ValidationError (or a tuple except (json.JSONDecodeError,
ValidationError)) branch that logs with the same message and raises
RuntimeHostError using
ModelInfraErrorContext.with_correlation(transport_type=EnumInfraTransportType.KAFKA,
operation="event_bus_deserialize") so deserialization failures are correctly
classified, leaving the existing generic except Exception for true dispatch
errors from await self._dispatch_engine.dispatch.
The test_scaling_with_worker_count performance test fails intermittently in CI due to resource contention and context-switch overhead on GitHub Actions runners. The test measures asyncio.Lock serialization behavior which varies significantly with shared CI resources. - Add @pytest.mark.skipif(IS_CI) decorator matching other perf tests - Move is_ci_environment import to module level for skip decorator - Keep threshold logic for potential manual --run-ci-tests override
Code reviewNo issues found. Checked for bugs and CLAUDE.md compliance. |
Summary
Implement automatic Kafka topic wiring from handler contract
event_bussections, so the runtime subscribes to declared topics without manual wiring.Key Changes:
EventBusSubcontractWiring: Wiressubscribe_topicsto Kafka consumers with environment-prefixed topic resolutionPublisherTopicScoped: Contract-validated publisher that enforcespublish_topicsallowlistload_event_bus_subcontract(): Utility to load event_bus section from contract YAMLServiceRuntimeHostProcesswith proper lifecycle managementArchitecture (ARCH-002 Compliance)
The runtime owns all Kafka plumbing. Nodes/handlers:
contract.yamlModelEventEnvelope(not raw Kafka messages)set_kafka_producer()or similarTest Plan
EventBusSubcontractWiring(27 tests)PublisherTopicScoped(31 tests)Related
ModelEventBusSubcontractexists in omnibase_core 0.9.6)Summary by CodeRabbit
New Features
Documentation
Tests
Chores
✏️ Tip: You can customize this high-level summary in your review settings.