Repository navigation
feat(OMN-1742): Add RequestResponseWiring for Kafka RPC patterns - #231
Conversation
Implements infrastructure-level support for request-response Kafka
communication, replacing bespoke event clients in downstream repos.
Key features:
- Per-instance long-lived consumers with correlation tracking
- Consumer group format: {environment}.rr.{instance_name}.{boot_nonce}
- Automatic correlation_id injection when missing
- InfraTimeoutError on timeout (not InfraUnavailableError)
- MixinAsyncCircuitBreaker integration for publish failures
- Boot nonce generated once per process (uuid4().hex[:8])
Includes 27 unit tests and 12 integration tests covering:
- Correlation ID handling and injection
- Timeout behavior with proper error types
- Consumer group format validation
- Pending map lifecycle and cleanup
- Circuit breaker integration
- Concurrent request isolation
Also bumps INFRA_MAX_UNIONS from 117 to 120 to accommodate
UUID | str correlation_id unions.
📝 WalkthroughWalkthroughAdds a Kafka-backed RequestResponseWiring to the runtime public API implementing per-instance correlation-based request–response flows with timeouts, circuit-breaker integration, boot-nonce consumer-group isolation, and robust cleanup. Also adds comprehensive unit and integration tests and a clarifying comment related to INFRA_MAX_UNIONS (117). Changes
Sequence DiagramsequenceDiagram
participant Client as Client
participant RRW as RequestResponseWiring
participant EB as EventBus
participant KC as KafkaConsumer
participant Remote as RemoteService
Client->>RRW: send_request(instance, payload, timeout)
activate RRW
RRW->>RRW: ensure/inject correlation_id
RRW->>RRW: create pending future
RRW->>EB: publish request (key=correlation_id)
deactivate RRW
EB->>Remote: deliver request
Remote->>EB: publish response on completed/failed topic (includes correlation_id)
EB->>KC: response message
activate KC
KC->>RRW: _consume_responses -> _handle_response_message
deactivate KC
activate RRW
RRW->>RRW: extract correlation_id, match pending future
RRW->>RRW: resolve future or set exception
deactivate RRW
Client->>RRW: await result (future completes)
Estimated Code Review Effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly Related PRs
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing touches
🧪 Generate unit tests (beta)
Comment |
- Update test imports: _RequestResponseInstance → RequestResponseInstanceState - Replace deprecated asyncio.get_event_loop() with get_running_loop() - Fix import sorting in test file
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)
tests/unit/validation/test_validator_defaults.py (1)
35-102:⚠️ Potential issue | 🟡 MinorAdd a pytest unit marker for this module.
This file has no classification marker; add a module‑level
pytestmark(or per‑class markers) to satisfy test taxonomy.
As per coding guidelines: tests/**/*.py: Use pytest markers:@pytest.mark.unit,@pytest.mark.integration,@pytest.mark.slow,@pytest.mark.chaos,@pytest.mark.serialfor test classification.➕ Suggested change
import inspect from collections.abc import Callable from pathlib import Path from unittest.mock import MagicMock, patch +import pytest + +# Mark all tests in this module as unit tests +pytestmark = pytest.mark.unit
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/runtime/request_response_wiring.py`:
- Around line 522-647: send_request currently drops non-UUID correlation IDs
when calling _check_circuit_breaker and _record_circuit_failure and generates a
new UUID in the timeout context, breaking traceability; fix by converting string
correlation_id to a UUID only for the circuit-breaker calls (e.g., build a
corr_for_breaker: UUID | None from correlation_id) while retaining the original
correlation_id (string or UUID) for correlation_key, the
ModelTimeoutErrorContext (timeout_context), and all error messages, and pass
corr_for_breaker into _check_circuit_breaker and _record_circuit_failure so type
expectations are satisfied without losing the original ID.
In `@tests/unit/runtime/test_request_response_wiring.py`:
- Around line 219-339: The tests import and reference a non-existent class
_RequestResponseInstance; update all occurrences to use the real class name
RequestResponseInstanceState (e.g., change from
omnibase_infra.runtime.request_response_wiring import _RequestResponseInstance
to import RequestResponseInstanceState, and replace any variables/constructors
that call _RequestResponseInstance(...) with RequestResponseInstanceState(...))
so the test imports and instantiation match the implementation.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Fix all issues with AI agents
In `@tests/unit/runtime/test_request_response_wiring.py`:
- Around line 325-327: Replace deprecated
asyncio.get_event_loop().create_future() calls with
asyncio.get_running_loop().create_future() in the test file; locate occurrences
that assign futures (the lines using "future: asyncio.Future[dict[str, object]]
= (asyncio.get_event_loop().create_future())" and similar statements around the
symbols "future" in tests (appearing at the noted occurrences) and update them
to call asyncio.get_running_loop().create_future() so they use the running event
loop within the `@pytest.mark.asyncio` async context.
🧹 Nitpick comments (3)
src/omnibase_infra/runtime/request_response_wiring.py (3)
648-680: Inconsistent return type behavior in_ensure_correlation_id.The method returns
UUID | str, but the behavior differs: when an existing value is found, it returnsstr(existing), while when injecting a new ID, it returns theUUIDobject directly. Consider returningstr(new_id)at lines 670 and 680 for consistent behavior, which would simplify callers' type handling.♻️ Suggested fix for consistent return type
# Inject new correlation_id new_id = uuid4() payload[config.field] = str(new_id) - return new_id + return str(new_id) # For header location, we'd need to handle differently # For now, treat as body existing = payload.get(config.field) if existing is not None: return str(existing) new_id = uuid4() payload[config.field] = str(new_id) - return new_id + return str(new_id)
510-518: Header-based correlation location not implemented.The code falls back to body-based extraction with a warning when
config.location == "headers". This is a reasonable approach for initial implementation, but consider adding a TODO comment or tracking this as technical debt if header-based correlation is a planned feature.
398-403: Consumer task exits on unexpected error without restart.If an unexpected exception occurs in
_consume_responses, the task terminates and won't process further responses. Pending requests will hang until timeout. This may be intentional (fail-fast), but consider whether a retry mechanism or at least failing all pending futures would improve resilience.💡 Alternative approach: fail pending futures on consumer crash
except Exception as e: self._logger.exception( "Unexpected error in consumer task for instance %s: %s", instance.name, e, ) # Fail all pending futures so callers don't hang consumer_error = RuntimeError( f"Consumer task crashed for instance '{instance.name}': {e}" ) for future in instance.pending.values(): if not future.done(): future.set_exception(consumer_error) instance.pending.clear()
…g unions - Change correlation_id from UUID | str to UUID throughout - Parse incoming correlation IDs to UUID at input boundaries - Convert to string only at serialization (JSON encoding) - Remove unnecessary isinstance(correlation_id, UUID) checks - Union count back to 117 (UUID | None is optional, doesn't count)
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/runtime/request_response_wiring.py`:
- Around line 665-674: The code parses an existing correlation ID into a UUID
(using UUID(...)) but does not normalize the original payload value back to a
string, causing json.dumps to fail if callers passed a UUID object; update the
branch in request_response_wiring where existing is not None so that after
computing correlation_id = UUID(str(existing)) you assign payload[config.field]
= str(correlation_id) (same normalization done in the uuid4() branch) so payload
always contains a JSON-serializable string; ensure this change references
payload, config.field, correlation_id, UUID and uuid4.
- Around line 565-575: When raising ProtocolConfigurationError for an unwired
instance, propagate a correlation_id into the ModelInfraErrorContext: extract
correlation_id from the incoming payload (e.g., payload.get("correlation_id") or
request.get("correlation_id")), and if absent generate one with uuid.uuid4();
then call ModelInfraErrorContext.with_correlation(..., correlation_id=the_id)
and include that context in the ProtocolConfigurationError constructor
(alongside existing transport_type EnumInfraTransportType.KAFKA and
operation="send_request"), and ensure uuid is imported.
| # Get instance | ||
| instance = self._instances.get(instance_name) | ||
| if instance is None: | ||
| raise ProtocolConfigurationError( | ||
| f"Request-response instance '{instance_name}' not wired", | ||
| context=ModelInfraErrorContext.with_correlation( | ||
| transport_type=EnumInfraTransportType.KAFKA, | ||
| operation="send_request", | ||
| ), | ||
| instance_name=instance_name, | ||
| ) |
There was a problem hiding this comment.
Propagate payload correlation_id into the unwired-instance error context.
Right now the ProtocolConfigurationError context is created without attempting to reuse a correlation_id from the incoming payload. This breaks traceability for failed requests.
🔧 Suggested fix
- if instance is None:
+ if instance is None:
+ # Best-effort correlation propagation for error context
+ corr_for_context = None
+ try:
+ corr_for_context = self._extract_correlation_id(
+ payload, ModelCorrelationConfig()
+ )
+ except Exception:
+ corr_for_context = None
raise ProtocolConfigurationError(
f"Request-response instance '{instance_name}' not wired",
- context=ModelInfraErrorContext.with_correlation(
+ context=ModelInfraErrorContext.with_correlation(
transport_type=EnumInfraTransportType.KAFKA,
operation="send_request",
+ correlation_id=corr_for_context,
),
instance_name=instance_name,
)As per coding guidelines: Propagate correlation_id from incoming requests; auto-generate with uuid4() if missing; include in all error context.
🤖 Prompt for AI Agents
In `@src/omnibase_infra/runtime/request_response_wiring.py` around lines 565 -
575, When raising ProtocolConfigurationError for an unwired instance, propagate
a correlation_id into the ModelInfraErrorContext: extract correlation_id from
the incoming payload (e.g., payload.get("correlation_id") or
request.get("correlation_id")), and if absent generate one with uuid.uuid4();
then call ModelInfraErrorContext.with_correlation(..., correlation_id=the_id)
and include that context in the ProtocolConfigurationError constructor
(alongside existing transport_type EnumInfraTransportType.KAFKA and
operation="send_request"), and ensure uuid is imported.
| existing = payload.get(config.field) | ||
|
|
||
| if existing is not None: | ||
| # Parse existing to UUID - correlation IDs are always UUIDs | ||
| correlation_id = UUID(str(existing)) | ||
| else: | ||
| # Generate new UUID | ||
| correlation_id = uuid4() | ||
| payload[config.field] = str(correlation_id) | ||
|
|
There was a problem hiding this comment.
Normalize existing UUIDs to strings before JSON serialization.
If callers provide a UUID object in payload, json.dumps will raise TypeError because the payload isn’t normalized. You only string-encode correlation_id when it’s missing; existing UUIDs remain un-serialized.
🐛 Proposed fix
- if existing is not None:
- # Parse existing to UUID - correlation IDs are always UUIDs
- correlation_id = UUID(str(existing))
- else:
- # Generate new UUID
- correlation_id = uuid4()
- payload[config.field] = str(correlation_id)
+ if existing is not None:
+ # Parse existing to UUID - correlation IDs are always UUIDs
+ correlation_id = UUID(str(existing))
+ # Normalize payload for JSON serialization
+ payload[config.field] = str(correlation_id)
+ else:
+ # Generate new UUID
+ correlation_id = uuid4()
+ payload[config.field] = str(correlation_id)🤖 Prompt for AI Agents
In `@src/omnibase_infra/runtime/request_response_wiring.py` around lines 665 -
674, The code parses an existing correlation ID into a UUID (using UUID(...))
but does not normalize the original payload value back to a string, causing
json.dumps to fail if callers passed a UUID object; update the branch in
request_response_wiring where existing is not None so that after computing
correlation_id = UUID(str(existing)) you assign payload[config.field] =
str(correlation_id) (same normalization done in the uuid4() branch) so payload
always contains a JSON-serializable string; ensure this change references
payload, config.field, correlation_id, UUID and uuid4.
- Remove duplicate `# type: ignore[arg-type]` comments in integration tests (4 occurrences) - Replace deprecated `asyncio.get_event_loop()` with `get_running_loop()` in unit tests (5 occurrences)
Summary
Implements infrastructure-level support for request-response Kafka communication, replacing bespoke event clients in downstream repos (omniclaude's
routing_event_client.py,intelligence_event_client.py).Key Features:
{environment}.rr.{instance_name}.{boot_nonce}correlation_idinjection when missing (prevents pending map memory leaks)InfraTimeoutErroron timeout (notInfraUnavailableError)MixinAsyncCircuitBreakerintegration for publish failuresuuid4().hex[:8]generated once per processAcceptance Criteria
RequestResponseWiringclass inomnibase_infra/runtime/dict[str, asyncio.Future]send_request(instance_name, payload) -> responsemethod{app}.rr.{instance_name}.{boot_nonce}MixinAsyncCircuitBreakerintegration for publish failuresInfraTimeoutErroron timeoutTest Plan
poetry run pytest tests/unit/runtime/test_request_response_wiring.py -v- 27 tests passpoetry run pytest tests/integration/runtime/test_request_response_wiring_integration.py -v- 12 tests passpoetry run mypy src/omnibase_infra/runtime/request_response_wiring.py- passespoetry run ruff check src/omnibase_infra/runtime/request_response_wiring.py- passesRelated
Summary by CodeRabbit
New Features
Validation
Tests