Repository navigation
feat: Runtime Health Event Pipeline - Waves 0-2 [OMN-5529] - #911
Conversation
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
📝 WalkthroughWalkthroughAdds consumer-health and runtime-error telemetry: DB migrations, Kafka topics and topic specs, Pydantic models/enums, emitters (consumer-health + runtime-log bridge), triage effect nodes/handlers (consumer health and runtime error) with escalation logic, consumer/service integrations, and accompanying unit/integration tests. Changes
Sequence Diagram(s)sequenceDiagram
participant Consumer as Kafka Consumer
participant EventBus as EventBusKafka
participant HealthEmitter as ConsumerHealthEmitter
participant Kafka as Kafka Broker
participant Triage as Consumer Health Triage
participant DB as PostgreSQL
participant Slack as Slack API
Consumer->>EventBus: start consumer
EventBus->>HealthEmitter: emit_event(CONSUMER_STARTED)
HealthEmitter->>Kafka: send(ModelConsumerHealthEvent)
Kafka->>Triage: deliver ModelConsumerHealthEvent
Triage->>DB: UPSERT consumer_health_triage (fingerprint)
DB-->>Triage: return occurrence_count
alt occurrence_count == 1
Triage->>Slack: send WARNING
else occurrence_count >= 3 and auto-restart enabled
Triage->>Kafka: emit restart command
Triage->>DB: update incident_state=restart_pending
else
Triage->>Slack: send REPEATED
end
sequenceDiagram
participant Logger as App Logger
participant Bridge as RuntimeLogEventBridge
participant Queue as Async Queue
participant Kafka as Kafka Broker
participant RuntimeTriage as Runtime Error Triage
participant DB as PostgreSQL
participant Linear as Linear API
Logger->>Bridge: WARNING/ERROR log record
Bridge->>Bridge: templatize & categorize
Bridge->>Queue: enqueue ModelRuntimeErrorEvent (rate-limited)
Queue->>Kafka: send(ModelRuntimeErrorEvent)
Kafka->>RuntimeTriage: deliver ModelRuntimeErrorEvent
RuntimeTriage->>DB: UPSERT runtime_error_triage (fingerprint)
DB-->>RuntimeTriage: occurrence_count, correlated_consumer_fingerprint?
alt matched rule == suppress
RuntimeTriage->>DB: update incident_state=suppressed
else matched rule == ticket
RuntimeTriage->>Linear: create issue
RuntimeTriage->>DB: update incident_state=ticketed
else
RuntimeTriage->>Slack: send alert
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~75 minutes Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 2 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (2 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
📝 Coding Plan
Comment |
| from __future__ import annotations | ||
|
|
||
| import logging | ||
| import os |
Check notice
Code scanning / CodeQL
Unused import Note
Show autofix suggestion
Hide autofix suggestion
Copilot Autofix
AI 7 months ago
To fix an unused import, remove the import statement for the unused module. This reduces unnecessary dependencies and slightly improves readability and import time.
Specifically, in src/omnibase_infra/event_bus/mixin_consumer_health.py, delete the import os line (line 26) and leave the remaining imports unchanged. No additional methods, imports, or definitions are needed; we are only simplifying the import list. This does not change any existing functionality because os is not referenced anywhere in the shown code.
| @@ -23,7 +23,6 @@ | ||
| from __future__ import annotations | ||
|
|
||
| import logging | ||
| import os | ||
| from typing import TYPE_CHECKING | ||
|
|
||
| from omnibase_infra.event_bus.consumer_health_emitter import ConsumerHealthEmitter |
| self._drain_task.cancel() | ||
| try: | ||
| await self._drain_task | ||
| except asyncio.CancelledError: |
Check notice
Code scanning / CodeQL
Empty except Note
Show autofix suggestion
Hide autofix suggestion
Copilot Autofix
AI 7 months ago
Generally, empty except blocks should either be removed (if the exception should propagate) or replaced with concrete handling such as logging, updating state, or at least a comment explaining why it is intentionally ignored.
Here, we want to preserve the behavior that stop() doesn’t raise if _drain_task is cancelled, because cancellation is expected after self._drain_task.cancel(). The best minimal fix is to add a small piece of handling inside the except asyncio.CancelledError: block, such as a debug log indicating that the drain task was cancelled during shutdown. This keeps the method’s outward behavior (no raised exception) but improves observability and satisfies the static analysis rule.
The change is confined to src/omnibase_infra/observability/runtime_log_event_bridge.py, within the stop method: replace the pass in the except asyncio.CancelledError: block with a debug log using the already-defined _bridge_logger. No new imports or helper methods are needed.
| @@ -278,7 +278,7 @@ | ||
| try: | ||
| await self._drain_task | ||
| except asyncio.CancelledError: | ||
| pass | ||
| _bridge_logger.debug("Drain task cancelled during stop()", exc_info=True) | ||
|
|
||
| async def _drain_loop(self) -> None: | ||
| """Background task: drain queue and emit to Kafka.""" |
|
|
||
| from __future__ import annotations | ||
|
|
||
| import asyncio |
Check notice
Code scanning / CodeQL
Unused import Note test
Show autofix suggestion
Hide autofix suggestion
Copilot Autofix
AI 7 months ago
To fix an unused import, the general approach is to remove the import statement for the unused module, leaving all other functionality intact. This eliminates an unnecessary dependency and slightly simplifies the module.
In this case, the best fix is to delete the import asyncio line from tests/unit/event_bus/test_consumer_health_emitter.py. No other code changes are needed because the rest of the file relies on pytest, AsyncMock, and project-specific imports, none of which depend on asyncio being imported here.
Concretely:
- In
tests/unit/event_bus/test_consumer_health_emitter.py, remove line 10:import asyncio. - Keep all other imports and code unchanged.
- No additional methods, imports, or definitions are required.
| @@ -7,7 +7,6 @@ | ||
|
|
||
| from __future__ import annotations | ||
|
|
||
| import asyncio | ||
| from unittest.mock import AsyncMock, patch | ||
|
|
||
| import pytest |
|
|
||
| import logging | ||
| from collections.abc import Awaitable, Callable | ||
| from datetime import UTC, datetime |
Check notice
Code scanning / CodeQL
Unused import Note
Show autofix suggestion
Hide autofix suggestion
Copilot Autofix
AI 7 months ago
To fix unused imports, remove the specific names that are not referenced anywhere in the file. This eliminates unnecessary dependencies without changing runtime behavior.
In this case, both UTC and datetime are reported unused, so the best fix is to delete the entire from datetime import UTC, datetime line at line 23 in src/omnibase_infra/nodes/node_runtime_error_triage_effect/handlers/handler_runtime_error_triage.py. No other code changes are needed, as nothing in the shown snippet depends on these names. We are not changing any existing functionality; we are only cleaning up an unnecessary import.
| @@ -20,7 +20,6 @@ | ||
|
|
||
| import logging | ||
| from collections.abc import Awaitable, Callable | ||
| from datetime import UTC, datetime | ||
| from typing import TYPE_CHECKING, cast | ||
|
|
||
| from pydantic import BaseModel, ConfigDict, Field |
There was a problem hiding this comment.
Actionable comments posted: 18
🧹 Nitpick comments (3)
docker/migrations/forward/054_create_consumer_health_triage.sql (1)
79-81: Redundant index on UNIQUE constraint columns.The index
idx_consumer_restart_state_consumeron(consumer_id, consumer_group, topic)duplicates coverage already provided by the UNIQUE constraint at line 56. PostgreSQL automatically creates an index for UNIQUE constraints.This is harmless but adds minor write overhead. Consider removing if cleanup is desired, or keep for explicit documentation.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@docker/migrations/forward/054_create_consumer_health_triage.sql` around lines 79 - 81, The CREATE INDEX statement for idx_consumer_restart_state_consumer on consumer_restart_state (consumer_id, consumer_group, topic) is redundant because the UNIQUE constraint on consumer_restart_state (consumer_id, consumer_group, topic) already creates the same index; remove the CREATE INDEX IF NOT EXISTS idx_consumer_restart_state_consumer ... statement from the migration (or, if you prefer to keep it for documentation, replace it with a commented note explaining the duplication) so you avoid unnecessary write overhead.src/omnibase_infra/event_bus/topic_constants.py (1)
580-617: Consider adding new health/runtime topics to wiring-health monitoring tuple.These topics are registered and exported, but they are not currently part of
WIRING_HEALTH_MONITORED_TOPICS, which can leave emission/consumption drift on the new pipeline untracked.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/omnibase_infra/event_bus/topic_constants.py` around lines 580 - 617, The new topics TOPIC_CONSUMER_HEALTH, TOPIC_CONSUMER_RESTART_CMD, and TOPIC_RUNTIME_ERROR are defined but not included in WIRING_HEALTH_MONITORED_TOPICS; update the wiring-health monitoring tuple/list named WIRING_HEALTH_MONITORED_TOPICS to include these constants so emissions and subscriptions are monitored. Locate the WIRING_HEALTH_MONITORED_TOPICS definition and append (or insert) TOPIC_CONSUMER_HEALTH, TOPIC_CONSUMER_RESTART_CMD, and TOPIC_RUNTIME_ERROR ensuring import/visibility is correct and the tuple/list remains immutable if intended (convert to tuple if needed). Ensure no duplicate entries are added and run tests/lint to validate.tests/unit/observability/test_runtime_log_event_bridge.py (1)
159-185: Use a deterministic wait mechanism instead of fixedsleep(0.1).The test will be collected and run because
asyncio_mode = "auto"is configured inpyproject.toml, so unmarked async tests are supported. However, relying on a fixed 0.1-second sleep before the assertion is fragile on slower CI workers. Consider using a polling or retry-based approach to deterministically wait for the event to be emitted rather than a time-based sleep.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/unit/observability/test_runtime_log_event_bridge.py` around lines 159 - 185, Replace the fixed asyncio.sleep(0.1) in test_drain_loop_emits_to_kafka with a deterministic wait that polls for the expected condition (e.g., bridge.events_emitted == 1 or producer.send.called) with a short sleep inside a loop and an overall timeout (or use asyncio.wait_for) to fail fast; locate this in the test where RuntimeLogEventBridge is started with bridge.start() and stopped with bridge.stop(), and wait until the bridge has emitted the event instead of sleeping a fixed duration.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@src/omnibase_infra/event_bus/event_bus_kafka.py`:
- Around line 673-678: The ConsumerHealthEmitter instance captures
self._producer in its constructor but is only created once, so when
_ensure_producer() replaces self._producer after a publish timeout the emitter
keeps a stale producer reference and health events fail; also _health_emitter is
never cleared on close(). Fix by re-binding/updating _health_emitter whenever
_ensure_producer() replaces self._producer (create a new
ConsumerHealthEmitter(self._producer) or update its producer) and ensure
_health_emitter is set to None/cleared in close() so no stale reference remains;
update the code paths that replace the producer and the close() method
accordingly (refer to ConsumerHealthEmitter, _ensure_producer, self._producer,
_health_emitter, and close()).
- Around line 1609-1624: The awaited best-effort health emit in the consumer
startup path (the await self._health_emitter.emit_event call inside the
consumer-started block) must not block while holding self._lock or during
teardown; change it to a non-blocking background call or enforce a short
timeout. Locate the emit in _start_consumer_for_topic (called from subscribe())
and the similar emit in the loop-error handler, and replace the direct await
with one of: schedule it with asyncio.create_task to fire-and-forget (optionally
logging exceptions via task.add_done_callback), or wrap the emit in
asyncio.wait_for with a small timeout and swallow TimeoutError so slow publishes
don't stall subscriber registration or shutdown. Ensure any background task
captures and logs exceptions but does not re-raise to the caller.
- Around line 1612-1619: The consumer identity currently uses only the topic
(consumer_identity=f"eventbus.{topic}") which collapses multiple consumers into
one; change it so the identity is unique per logical consumer by including the
effective_group_id and an instance-specific identifier (for example include
effective_group_id and a process/instance/host id such as self._instance_id or
hostname) when calling self._health_emitter.emit_event (where
EnumConsumerHealthEventType.CONSUMER_STARTED/CONSUMER_STOPPED are used), and
apply the same change to the other emit_event call around the 1920–1924 region
so consumer_identity distinguishes group+instance+topic.
- Around line 1926-1929: The health-event code currently uses str(e) directly
for error_message; replace that with a sanitized, truncated version by passing
the exception (or its text) through the existing sanitize_error_message function
before slicing to 500 chars. Update the construction that sets error_message (in
event_bus_kafka.py, around the health event / Kafka publish code that sets
correlation_id/error_type/hostname) to call sanitize_error_message(e) and then
apply [:500], preserving error_type and hostname as-is.
In `@src/omnibase_infra/event_bus/mixin_consumer_health.py`:
- Around line 5-16: Update the module docstring to reflect the current API:
remove the misleading statement that the mixin "handles restart commands" and
instead state it only provides health-event emission helpers; fix the usage
example to call _init_health_emitter() without await (since _init_health_emitter
is synchronous and returns None), and remove or clarify the example of
_handle_restart so it doesn't imply the mixin requires implementing restart
handling; reference MixinConsumerHealth and the _init_health_emitter() method in
the doc edits.
In `@src/omnibase_infra/nodes/node_consumer_health_triage_effect/contract.yaml`:
- Around line 67-91: The retry policy currently lists retryable errors in
retry_policy.retry_on but omits "LINEAR_API_ERROR" even though error_types
contains a recoverable LINEAR_API_ERROR with exponential_backoff; update the
retry_policy.retry_on array to include "LINEAR_API_ERROR" so the declared
recoverable Linear failures will follow the configured exponential_backoff retry
strategy (ensure the string exactly matches the name in error_types).
In
`@src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py`:
- Around line 86-105: The handler currently publishes Kafka messages itself
(owns _producer and calls _emit_restart_command()), which violates the handlers
rule against direct event-bus access; remove the AIOKafkaProducer dependency
from HandlerConsumerHealthTriage (drop the producer param and _producer
attribute), change _emit_restart_command() to build and return a restart command
object (or a simple dict) instead of publishing, and update all call sites
inside the class that previously invoked _emit_restart_command() to
return/propagate that command to the caller; move actual publication logic into
an external adapter/service used by the orchestration layer so the handler only
produces the command payload and does not perform any publish.
- Around line 142-168: The code uses occurrence_count (and _RESTART_THRESHOLD)
as a lifetime counter so a third event at any later time will trigger a restart;
change the logic in the handler around occurrence_count, _RESTART_THRESHOLD and
_RESTART_WINDOW_MINUTES so the threshold is evaluated against events within the
sliding time window instead of the lifetime count — e.g., compute recent_count
by querying/filtering sightings for event.fingerprint within
_RESTART_WINDOW_MINUTES (or reset occurrence_count based on last_seen timestamp)
and use recent_count for the conditional that calls _check_restart_rate_limit,
_emit_restart_command, _create_linear_ticket and for building the returned
ModelTriageResult/incident_state. Ensure the same windowed-count change is
applied in the other block referenced (lines ~188-250) that also uses
occurrence_count for restart vs ticket decisions.
- Around line 147-156: The handler currently calls await
self._emit_restart_command(event) and immediately marks the incident
RESTART_PENDING and returns a restart_command result even when no produce
happened; change _emit_restart_command to return a tuple (success: bool, reason:
Optional[str/Enum]) that distinguishes rate-limited vs publish failures (and
returns success=False when self._producer is None), then update handle() to
branch on that return: only call _update_incident_state(event.fingerprint,
EnumConsumerIncidentState.RESTART_PENDING) and return
ModelTriageResult(action="restart_command", incident_state=RESTART_PENDING, ...)
when success is True; when success is False, set an appropriate
incident_state/action (e.g., no_restart or rate_limited), include the reason in
the ModelTriageResult, and avoid marking RESTART_PENDING so callers see that no
restart was issued; apply the same pattern to the other restart-emitting blocks
in the 300-355 range.
In `@src/omnibase_infra/nodes/node_consumer_health_triage_effect/node.py`:
- Around line 29-30: Replace the direct archetype import for NodeEffect with the
sanctioned archetype path: in node.py change the import "from
omnibase_core.nodes.node_effect import NodeEffect" to "from omnibase_core.nodes
import NodeEffect" (keep the existing ModelONEXContainer import unchanged) so
node archetypes are imported via omnibase_core.nodes; ensure any references to
NodeEffect in the file remain valid.
In `@src/omnibase_infra/nodes/node_runtime_error_triage_effect/contract.yaml`:
- Around line 67-115: The contract's error_handling only models DB pool failures
while the declared services slack_handler and linear_handler are untyped; add
explicit error_types for these channels (e.g., "SLACK_API_FAILURE" and
"LINEAR_API_FAILURE") with appropriate description, recoverable flags and
retry_strategy (for example recoverable:true with "exponential_backoff" if
retries are desired), and include those names in retry_policy.retry_on so the
retry rules apply; update the error_handling.retry_policy and
error_handling.error_types entries (referencing the existing symbols
error_handling, retry_policy, retry_on, DB_POOL_UNAVAILABLE, DB_POOL_TIMEOUT,
slack_handler, linear_handler) to ensure alert/ticket failures are typed and
have retry behavior.
- Around line 59-65: The contract lists two different output field names for the
same concept—handler_routing.output_fields uses matched_rule while
io_operations.output_fields uses rule_name—so pick one canonical field name
(e.g., matched_rule) and update all occurrences in this contract (both the
handler_routing and io_operations sections) and any consumers/validators to the
chosen name; ensure the schema, any references to
correlated_consumer_fingerprint/occurrence_count, and downstream extraction
logic expect the same field identifier to avoid validation/extraction failures.
In
`@src/omnibase_infra/nodes/node_runtime_error_triage_effect/handlers/handler_runtime_error_triage.py`:
- Around line 199-219: The current _correlate_with_layer1 method returns the
most-recent consumer_health_triage fingerprint without using the provided event,
causing incorrect correlations; update _correlate_with_layer1 to only attempt
layer‑1 correlation when the incoming ModelRuntimeErrorEvent contains an
explicit consumer identity/correlation key (e.g., a consumer_id, group_id, or
topic/partition info from model_runtime_error_event.py), otherwise return None;
if present, modify the SQL to filter consumer_health_triage by that identity
(e.g., WHERE consumer_id = $1 OR group_id = $1 / appropriate column) and pass
the event value as a parameter when acquiring the fingerprint.
- Around line 135-158: The current early return when incident_state ==
"suppressed" ignores the rule's suppress_duration_minutes so suppression never
expires; update the logic in handler_runtime_error_triage.py to, when
incident_state == "suppressed", fetch the timestamp when the incident entered
suppressed state (or the incident record) and compare now against
matched_rule.suppress_duration_minutes, and only short-circuit if the
suppression window has not yet elapsed; if it has elapsed, treat it like the
"First time or suppression expired" path (call _update_incident_state and
_send_slack_notification and return the suppressed result). Use
event.fingerprint, matched_rule.suppress_duration_minutes, and the existing
_update_incident_state/_send_slack_notification helpers to implement this check.
In
`@src/omnibase_infra/nodes/node_runtime_error_triage_effect/models/model_triage_rule.py`:
- Around line 25-32: The ModelTriageRule Pydantic model is missing the required
from_attributes flag in its model_config; update the model_config for class
ModelTriageRule to include from_attributes=True (i.e., use
ConfigDict(frozen=True, extra="forbid", from_attributes=True)) so the model
supports attribute-based parsing while keeping existing frozen and extra
settings.
In `@src/omnibase_infra/observability/runtime_log_event_bridge.py`:
- Around line 273-299: The stop() flow can return before backlog is flushed
because _drain_loop only runs while _running; change the shutdown so the drain
loop always drains the queue: either push a sentinel into _queue (e.g., None)
and have _drain_loop treat that as shutdown after emitting remaining items, or
modify _drain_loop's loop condition to `while self._running or not
self._queue.empty()` and avoid cancelling the _drain_task immediately; update
stop() to set _running=False, enqueue the sentinel (if using sentinel approach)
or await the _drain_task completion after signaling, and only cancel if it fails
to finish, ensuring _emit_to_kafka is called for all queued events.
- Around line 164-166: The emit path is enqueuing into an asyncio.Queue from
arbitrary logging threads which is unsafe; capture the target event loop in
start() (e.g., store self._loop = asyncio.get_running_loop() or the provided
loop) and in emit() use self._loop.call_soon_threadsafe(self._queue.put_nowait,
event) to schedule thread-safe enqueueing; alternatively replace the
thread-facing part with logging.handlers.QueueHandler + queue.Queue and have the
async consumer move items into the asyncio.Queue on the event loop. Ensure
references: update start() to set self._loop and change emit() to call
call_soon_threadsafe with self._queue.put_nowait and the ModelRuntimeErrorEvent.
- Around line 205-206: The current fingerprinting uses self.format(record) which
includes formatter-added data; replace that with record.getMessage() when
building raw_message for _templatize_message to ensure stable fingerprints
(i.e., change the raw_message assignment so it uses record.getMessage() instead
of self.format(record) and only append stack traces from
traceback.format_exception(*record.exc_info) when record.exc_info is present),
and avoid relying on record.exc_text; update references around raw_message and
_templatize_message in runtime_log_event_bridge.py to use record.getMessage()
and explicit traceback.format_exception handling.
---
Nitpick comments:
In `@docker/migrations/forward/054_create_consumer_health_triage.sql`:
- Around line 79-81: The CREATE INDEX statement for
idx_consumer_restart_state_consumer on consumer_restart_state (consumer_id,
consumer_group, topic) is redundant because the UNIQUE constraint on
consumer_restart_state (consumer_id, consumer_group, topic) already creates the
same index; remove the CREATE INDEX IF NOT EXISTS
idx_consumer_restart_state_consumer ... statement from the migration (or, if you
prefer to keep it for documentation, replace it with a commented note explaining
the duplication) so you avoid unnecessary write overhead.
In `@src/omnibase_infra/event_bus/topic_constants.py`:
- Around line 580-617: The new topics TOPIC_CONSUMER_HEALTH,
TOPIC_CONSUMER_RESTART_CMD, and TOPIC_RUNTIME_ERROR are defined but not included
in WIRING_HEALTH_MONITORED_TOPICS; update the wiring-health monitoring
tuple/list named WIRING_HEALTH_MONITORED_TOPICS to include these constants so
emissions and subscriptions are monitored. Locate the
WIRING_HEALTH_MONITORED_TOPICS definition and append (or insert)
TOPIC_CONSUMER_HEALTH, TOPIC_CONSUMER_RESTART_CMD, and TOPIC_RUNTIME_ERROR
ensuring import/visibility is correct and the tuple/list remains immutable if
intended (convert to tuple if needed). Ensure no duplicate entries are added and
run tests/lint to validate.
In `@tests/unit/observability/test_runtime_log_event_bridge.py`:
- Around line 159-185: Replace the fixed asyncio.sleep(0.1) in
test_drain_loop_emits_to_kafka with a deterministic wait that polls for the
expected condition (e.g., bridge.events_emitted == 1 or producer.send.called)
with a short sleep inside a loop and an overall timeout (or use
asyncio.wait_for) to fail fast; locate this in the test where
RuntimeLogEventBridge is started with bridge.start() and stopped with
bridge.stop(), and wait until the bridge has emitted the event instead of
sleeping a fixed duration.
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: 75c6e941-128d-4f14-9ca1-1fdeeb628a38
⛔ Files ignored due to path filters (1)
src/omnibase_infra/enums/generated/enum_omnibase_infra_topic.pyis excluded by!**/generated/**
📒 Files selected for processing (37)
docker/migrations/forward/054_create_consumer_health_triage.sqldocker/migrations/forward/055_create_runtime_error_triage.sqlscripts/check_contract_topic_parity.pysrc/omnibase_infra/event_bus/consumer_health_emitter.pysrc/omnibase_infra/event_bus/event_bus_kafka.pysrc/omnibase_infra/event_bus/mixin_consumer_health.pysrc/omnibase_infra/event_bus/topic_constants.pysrc/omnibase_infra/models/health/__init__.pysrc/omnibase_infra/models/health/enum_consumer_health_event_type.pysrc/omnibase_infra/models/health/enum_consumer_health_severity.pysrc/omnibase_infra/models/health/enum_consumer_incident_state.pysrc/omnibase_infra/models/health/enum_runtime_error_category.pysrc/omnibase_infra/models/health/enum_runtime_error_severity.pysrc/omnibase_infra/models/health/model_consumer_health_event.pysrc/omnibase_infra/models/health/model_consumer_restart_command.pysrc/omnibase_infra/models/health/model_runtime_error_event.pysrc/omnibase_infra/nodes/node_consumer_health_triage_effect/__init__.pysrc/omnibase_infra/nodes/node_consumer_health_triage_effect/contract.yamlsrc/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/__init__.pysrc/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.pysrc/omnibase_infra/nodes/node_consumer_health_triage_effect/node.pysrc/omnibase_infra/nodes/node_runtime_error_triage_effect/__init__.pysrc/omnibase_infra/nodes/node_runtime_error_triage_effect/contract.yamlsrc/omnibase_infra/nodes/node_runtime_error_triage_effect/handlers/__init__.pysrc/omnibase_infra/nodes/node_runtime_error_triage_effect/handlers/handler_runtime_error_triage.pysrc/omnibase_infra/nodes/node_runtime_error_triage_effect/models/__init__.pysrc/omnibase_infra/nodes/node_runtime_error_triage_effect/models/model_triage_rule.pysrc/omnibase_infra/nodes/node_runtime_error_triage_effect/node.pysrc/omnibase_infra/observability/runtime_log_event_bridge.pysrc/omnibase_infra/topics/platform_topic_suffixes.pytests/unit/event_bus/test_consumer_health_emitter.pytests/unit/models/health/__init__.pytests/unit/models/health/test_model_consumer_health_event.pytests/unit/models/health/test_model_runtime_error_event.pytests/unit/nodes/test_handler_consumer_health_triage.pytests/unit/nodes/test_handler_runtime_error_triage.pytests/unit/observability/test_runtime_log_event_bridge.py
There was a problem hiding this comment.
Actionable comments posted: 8
🧹 Nitpick comments (1)
src/omnibase_infra/services/observability/injection_effectiveness/consumer.py (1)
639-649: Inconsistent error logging in cleanup block.The health producer cleanup omits error details and
correlation_id, unlike the adjacent Kafka consumer and PostgreSQL pool cleanup blocks (lines 655-681) which captureException as eand includeerror: str(e)andcorrelation_idin the log extras. For debugging consistency, consider aligning with the existing pattern in this file.♻️ Suggested fix for consistency
# Stop health producer (OMN-5523) if self._health_producer is not None: try: await self._health_producer.stop() - except Exception: # noqa: BLE001 — boundary: logs warning and degrades + except Exception as e: # noqa: BLE001 — boundary: logs warning and degrades logger.warning( "Error stopping health producer", - extra={"consumer_id": self._consumer_id}, + extra={ + "consumer_id": self._consumer_id, + "correlation_id": str(correlation_id), + "error": str(e), + }, ) finally: self._health_producer = None🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/omnibase_infra/services/observability/injection_effectiveness/consumer.py` around lines 639 - 649, The cleanup for self._health_producer currently swallows the exception and logs only a generic message; change the except block to capture the exception (except Exception as e) and include the error text and correlation_id in the logger.warning extras (e.g., extra={"consumer_id": self._consumer_id, "error": str(e), "correlation_id": correlation_id}) to match the logging pattern used for the Kafka consumer and PostgreSQL pool cleanup; keep finally setting self._health_producer = None and maintain the same warning severity and message text.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@src/omnibase_infra/runtime/service_kernel.py`:
- Around line 972-977: The Kafka bridge producer is started before environment
variables from ~/.omnibase/.env are loaded, so modify bootstrap() to source that
file near its top (before any Kafka/DB/Infisical clients are created) and ensure
any code that constructs AIOKafkaProducer (e.g., the _BridgeProducer
instantiation and await _bridge_producer.start() lines) runs only after the env
has been loaded; update bootstrap() to load the file once and rely on os.environ
for subsequent client configuration.
- Around line 966-969: The variable runtime_log_bridge (type
RuntimeLogEventBridge | None) must be initialized before entering the main try
block to avoid UnboundLocalError in the finally cleanup; move the
runtime_log_bridge = None initialization so it sits alongside the other pre-try
cleanup-guarded resource initializations (the same area where other resources
are declared) rather than inside the try, ensuring the finally block that
references runtime_log_bridge can always access it.
- Around line 974-1008: The try/except leaks partially-initialized resources:
_BridgeProducer started before creating RuntimeLogEventBridge and
attach_to_loggers may be called before start(), leaving producers running and
loggers attached on failure; fix by declaring tracking vars (e.g.
_bridge_producer = None, runtime_log_bridge = None, allowlist = []) before the
try, create and start the RuntimeLogEventBridge (call
runtime_log_bridge.start()) before calling
runtime_log_bridge.attach_to_loggers(allowlist), and in the except block if
runtime_log_bridge exists call a detach method (or reverse attach) and if
_bridge_producer exists call its stop/close method to ensure the producer is
stopped; ensure you still set runtime_log_bridge = None after cleanup.
In `@src/omnibase_infra/services/observability/agent_actions/consumer.py`:
- Around line 738-755: The startup currently allows exceptions when
creating/starting the optional _health_producer to bubble up and tear down the
consumer, which breaks retries because _running was set earlier; wrap the
dedicated health producer creation/start and _init_health_emitter call
(references: _health_producer, _dlq_producer, AIOKafkaProducer,
_init_health_emitter, ConsumerHealthEmitter.is_enabled()) in a try/except: on
exception log the error and skip health emitter setup (leave _health_producer as
None) so the main consumer startup proceeds; do NOT let exceptions propagate or
call the teardown path that leaves _running inconsistent—alternatively ensure
any cleanup path clears _running (reference: _cleanup_resources and start()) if
you must teardown, but preferred fix is to swallow errors for this best-effort
component and continue.
In `@src/omnibase_infra/services/observability/context_audit/consumer.py`:
- Around line 342-350: The ConsumerHealthEmitter check currently only runs when
self._producer exists (created only if DLQ is enabled), so
ENABLE_CONSUMER_HEALTH_EMITTER is ignored when DLQ is off; change the logic in
consumer.py to call self._init_health_emitter whenever
ConsumerHealthEmitter.is_enabled() is True regardless of self._producer, passing
None or a no-op/placeholder for the producer when DLQ is disabled, and ensure
_init_health_emitter can accept an optional producer (or create the emitter
without producer) so health emission is enabled for ContextAuditConsumer even
when DLQ is disabled.
In `@src/omnibase_infra/services/observability/llm_cost_aggregation/consumer.py`:
- Around line 490-503: The auxiliary health producer startup must not block or
raise during the consumer's critical startup: wrap the AIOKafkaProducer
creation, await self._health_producer.start(), and self._init_health_emitter
calls inside a try/except so any exception is caught and logged (do not
re-raise), and ensure self._health_producer is left as None on failure;
alternatively spawn a background task to initialize the health producer so the
main start() can return successfully. Also update _cleanup_resources to
defensively handle a missing or partially-initialized self._health_producer and
make sure it clears the _running flag so the instance can be retried. Reference
symbols: ConsumerHealthEmitter.is_enabled(), AIOKafkaProducer, await
self._health_producer.start(), _init_health_emitter, _cleanup_resources, and
_running.
In `@src/omnibase_infra/services/observability/skill_lifecycle/consumer.py`:
- Around line 367-375: Health emitter initialization is incorrectly gated on
self._producer (which is only created when dlq_enabled is true), so
ConsumerHealthEmitter.is_enabled() becomes a no-op if DLQ is off; to fix, ensure
the emitter is initialized regardless of DLQ by creating a dedicated health
producer attribute (add self._health_producer = None in __init__), call
self._init_health_emitter(...) using that producer when
ConsumerHealthEmitter.is_enabled() (use
consumer_identity="skill-lifecycle-consumer",
consumer_group=self.config.kafka_group_id, topic=",".join(self.config.topics"),
service_label="SkillLifecycleConsumer"), and ensure you clean up the health
producer in stop() (stop/close self._health_producer) so emitter lifecycle is
independent of dlq_enabled and self._producer.
In `@src/omnibase_infra/services/session/consumer.py`:
- Around line 380-393: The health-producer initialization can raise after the
consumer has already been started (start() and _running set), leaving a leaked
consumer; wrap the ConsumerHealthEmitter block that creates
self._health_producer, awaits self._health_producer.start(), and calls
_init_health_emitter in a try/except/finally so that on any exception you: stop
and await the already-started consumer (the consumer started earlier in
start()), clear self._running, and ensure any partially-started
self._health_producer is stopped/closed before re-raising the exception;
reference ConsumerHealthEmitter, AIOKafkaProducer, self._health_producer,
_init_health_emitter, start(), and _running when making the changes.
---
Nitpick comments:
In
`@src/omnibase_infra/services/observability/injection_effectiveness/consumer.py`:
- Around line 639-649: The cleanup for self._health_producer currently swallows
the exception and logs only a generic message; change the except block to
capture the exception (except Exception as e) and include the error text and
correlation_id in the logger.warning extras (e.g., extra={"consumer_id":
self._consumer_id, "error": str(e), "correlation_id": correlation_id}) to match
the logging pattern used for the Kafka consumer and PostgreSQL pool cleanup;
keep finally setting self._health_producer = None and maintain the same warning
severity and message text.
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: 44ea4625-a1da-41e9-81b4-5fd31989462a
📒 Files selected for processing (7)
src/omnibase_infra/runtime/service_kernel.pysrc/omnibase_infra/services/observability/agent_actions/consumer.pysrc/omnibase_infra/services/observability/context_audit/consumer.pysrc/omnibase_infra/services/observability/injection_effectiveness/consumer.pysrc/omnibase_infra/services/observability/llm_cost_aggregation/consumer.pysrc/omnibase_infra/services/observability/skill_lifecycle/consumer.pysrc/omnibase_infra/services/session/consumer.py
There was a problem hiding this comment.
Actionable comments posted: 3
🧹 Nitpick comments (3)
tests/integration/test_runtime_log_bridge_pipeline.py (2)
212-236: Consider addingbridge.stop()for consistency, though not strictly required.Unlike other tests, this test doesn't call
bridge.start()orbridge.stop(). While this is technically fine since the bridge is never started (so there's no drain task to cancel), addingbridge.stop()would maintain consistency with other tests and be more defensive if the test is later modified.Proposed fix for consistency
# No events should be emitted or queued assert bridge.events_emitted == 0 bridge.detach_from_loggers([test_logger_name]) + await bridge.stop()🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/integration/test_runtime_log_bridge_pipeline.py` around lines 212 - 236, The test test_bridge_disabled_by_default creates a RuntimeLogEventBridge but never calls bridge.stop(), which is inconsistent with other tests and could leave resources if the test is modified later; update the test to call bridge.stop() at the end (after bridge.detach_from_loggers([...])) to mirror the start/stop pattern used elsewhere and ensure any internal tasks are safely cleaned up by invoking RuntimeLogEventBridge.stop().
74-82: Remove unused_EnableBridgeCtxcontext manager.This context manager class is defined but never used in any of the tests. All tests manually set/unset
ENABLE_RUNTIME_LOG_BRIDGEusing try/finally blocks. Either use this context manager in the tests or remove it.Option 1: Remove unused class
-class _EnableBridgeCtx: - """Context manager to enable the runtime log bridge feature flag.""" - - def __enter__(self) -> _EnableBridgeCtx: - os.environ["ENABLE_RUNTIME_LOG_BRIDGE"] = "true" - return self - - def __exit__(self, *args: object) -> None: - os.environ.pop("ENABLE_RUNTIME_LOG_BRIDGE", None) - - class TestRuntimeLogBridgeIntegration:Option 2: Use it in tests (e.g., test_log_error_produces_kafka_event)
async def test_log_error_produces_kafka_event( self, kafka_producer: AIOKafkaProducer, kafka_consumer: AIOKafkaConsumer, ) -> None: """Verify an ERROR log record flows through the bridge to Kafka.""" - os.environ["ENABLE_RUNTIME_LOG_BRIDGE"] = "true" - try: + with _EnableBridgeCtx(): bridge = RuntimeLogEventBridge( # ... rest of test ... - finally: - os.environ.pop("ENABLE_RUNTIME_LOG_BRIDGE", None)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/integration/test_runtime_log_bridge_pipeline.py` around lines 74 - 82, The _EnableBridgeCtx context manager is defined but unused; either remove the class or switch tests that manually toggle ENABLE_RUNTIME_LOG_BRIDGE (for example test_log_error_produces_kafka_event) to use it. To fix, delete the _EnableBridgeCtx class definition if you choose Option 1, or update tests that set os.environ["ENABLE_RUNTIME_LOG_BRIDGE"] in try/finally blocks to use "with _EnableBridgeCtx():" and remove the manual env pop in those tests so the context manager handles cleanup. Ensure references to _EnableBridgeCtx (class name) or the updated tests (e.g., test_log_error_produces_kafka_event) are updated accordingly.tests/integration/test_consumer_health_pipeline.py (1)
3-17: Module docstring currently overstates test coverage.The header mentions triage handler and restart/ticket flows that are not implemented in this file. Aligning the docstring with actual assertions will reduce maintenance confusion.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/integration/test_consumer_health_pipeline.py` around lines 3 - 17, The module docstring at the top of the file overstates coverage (it references triage handlers, restart/ticket flows and other behaviors not implemented in this test module); update the module docstring text to accurately list only the actual assertions and behaviors exercised by the tests in this file (e.g., emit→consume flow, rate limiting, emitter self-metrics accuracy) and remove or reword references to triage/restart/ticket flows and any unrelated items (leave identifiers like OMN-5524 if needed for tracking but ensure the descriptive bullets match current tests).
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@tests/integration/test_consumer_health_pipeline.py`:
- Around line 111-114: The assertions in the test check uppercase literal
strings but the schema emits lowercase enum values; update the checks on
payload["event_type"] and payload["severity"] to compare against the enum .value
members instead of hardcoded uppercase strings (e.g., use
EventType.HEARTBEAT_FAILURE.value and Severity.ERROR.value or the actual enum
classes used in the codebase) while keeping the payload["consumer_identity"]
assertion unchanged.
- Around line 50-78: Replace the hardcoded BOOTSTRAP_SERVERS and immediate-use
of Kafka in the fixtures by sourcing the required env first and reading the
bootstrap servers from an environment variable; specifically remove or stop
using the BOOTSTRAP_SERVERS constant and update kafka_producer and
kafka_consumer to (1) ensure the env file (e.g., ~/.omnibase/.env) is loaded
before any Kafka client is created, (2) read bootstrap servers from a variable
like KAFKA_BOOTSTRAP_SERVERS via os.environ.get and raise a clear exception if
missing, and (3) keep TOPIC_CONSUMER_HEALTH usage but rely on the env-provided
bootstrap value when constructing AIOKafkaProducer/AIOKafkaConsumer.
In `@tests/integration/test_runtime_log_bridge_pipeline.py`:
- Line 43: The BOOTSTRAP_SERVERS constant is hardcoded; change it to read from
an environment variable (e.g. os.environ.get("BOOTSTRAP_SERVERS",
"localhost:19092")) so tests use the configured broker address with a sane
default for local dev; update the top of
tests/integration/test_runtime_log_bridge_pipeline.py to import os if missing
and replace the BOOTSTRAP_SERVERS assignment accordingly.
---
Nitpick comments:
In `@tests/integration/test_consumer_health_pipeline.py`:
- Around line 3-17: The module docstring at the top of the file overstates
coverage (it references triage handlers, restart/ticket flows and other
behaviors not implemented in this test module); update the module docstring text
to accurately list only the actual assertions and behaviors exercised by the
tests in this file (e.g., emit→consume flow, rate limiting, emitter self-metrics
accuracy) and remove or reword references to triage/restart/ticket flows and any
unrelated items (leave identifiers like OMN-5524 if needed for tracking but
ensure the descriptive bullets match current tests).
In `@tests/integration/test_runtime_log_bridge_pipeline.py`:
- Around line 212-236: The test test_bridge_disabled_by_default creates a
RuntimeLogEventBridge but never calls bridge.stop(), which is inconsistent with
other tests and could leave resources if the test is modified later; update the
test to call bridge.stop() at the end (after bridge.detach_from_loggers([...]))
to mirror the start/stop pattern used elsewhere and ensure any internal tasks
are safely cleaned up by invoking RuntimeLogEventBridge.stop().
- Around line 74-82: The _EnableBridgeCtx context manager is defined but unused;
either remove the class or switch tests that manually toggle
ENABLE_RUNTIME_LOG_BRIDGE (for example test_log_error_produces_kafka_event) to
use it. To fix, delete the _EnableBridgeCtx class definition if you choose
Option 1, or update tests that set os.environ["ENABLE_RUNTIME_LOG_BRIDGE"] in
try/finally blocks to use "with _EnableBridgeCtx():" and remove the manual env
pop in those tests so the context manager handles cleanup. Ensure references to
_EnableBridgeCtx (class name) or the updated tests (e.g.,
test_log_error_produces_kafka_event) are updated accordingly.
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: 2ddfa925-7f8b-4de4-b48b-689ec0dafe4b
📒 Files selected for processing (2)
tests/integration/test_consumer_health_pipeline.pytests/integration/test_runtime_log_bridge_pipeline.py
| from omnibase_infra.models.health.model_consumer_health_event import ( | ||
| ModelConsumerHealthEvent, | ||
| ) |
Check notice
Code scanning / CodeQL
Unused import Note test
Show autofix suggestion
Hide autofix suggestion
Copilot Autofix
AI 7 months ago
To fix an unused-import issue, you remove the import statement (or the unused name within a multi-name import) so that the file only depends on modules it actually uses. This reduces noise, avoids misleading future readers about dependencies, and satisfies static analysis tools.
In this case, ModelConsumerHealthEvent is imported on lines 41–43 from omnibase_infra.models.health.model_consumer_health_event, but there is no usage in the shown file. The best fix that does not change existing functionality is simply to delete this import block. No other code changes are needed, no replacements with other symbols, and no additional imports or definitions are required. Concretely, in tests/integration/test_consumer_health_pipeline.py, remove the three-line import block starting at line 41 so that the remaining imports stay intact and in the same order.
| @@ -38,9 +38,6 @@ | ||
| from omnibase_infra.models.health.enum_consumer_health_severity import ( | ||
| EnumConsumerHealthSeverity, | ||
| ) | ||
| from omnibase_infra.models.health.model_consumer_health_event import ( | ||
| ModelConsumerHealthEvent, | ||
| ) | ||
|
|
||
| pytestmark = [ | ||
| pytest.mark.integration, |
There was a problem hiding this comment.
Actionable comments posted: 2
♻️ Duplicate comments (3)
src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py (3)
87-105:⚠️ Potential issue | 🟠 MajorHandler still performs direct Kafka publication inside the handler boundary.
HandlerConsumerHealthTriageowns_producer(Line 104) and publishes in_emit_restart_command()(Line 328). This keeps event-bus side effects inside a handler class.Suggested boundary-safe refactor
- def __init__( - self, - db_pool: Pool, - producer: AIOKafkaProducer | None = None, - *, - slack_handler: Callable[[str], Awaitable[None]] | None = None, - linear_handler: Callable[..., Awaitable[None]] | None = None, - ) -> None: + def __init__( + self, + db_pool: Pool, + *, + slack_handler: Callable[[str], Awaitable[None]] | None = None, + linear_handler: Callable[..., Awaitable[None]] | None = None, + ) -> None: self._db_pool = db_pool - self._producer = producer self._slack_handler = slack_handler self._linear_handler = linear_handler - async def _emit_restart_command(self, event: ModelConsumerHealthEvent) -> None: - ... - await self._producer.send(TOPIC_CONSUMER_RESTART_CMD, value=payload) + def _build_restart_command( + self, event: ModelConsumerHealthEvent + ) -> ModelConsumerRestartCommand: + return ModelConsumerRestartCommand( + consumer_identity=event.consumer_identity, + consumer_group=event.consumer_group, + topic=event.topic, + reason=f"Graduated triage: {event.event_type.value} occurred {_RESTART_THRESHOLD}+ times", + fingerprint=event.fingerprint, + correlation_id=event.correlation_id, + )As per coding guidelines,
**/handlers/handler_*.py: "Handler classes MUST NOT have direct event bus access - no _bus, _event_bus, _publisher attributes or publish methods."Also applies to: 311-329
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py` around lines 87 - 105, HandlerConsumerHealthTriage currently holds a concrete AIOKafkaProducer in the _producer attribute and performs direct publishes in _emit_restart_command, violating the handler boundary rule; replace the concrete Kafka dependency by injecting a boundary-safe abstraction (e.g., a callable like restart_command_emitter: Callable[[dict], Awaitable[None]] or a simple interface) via the constructor instead of AIOKafkaProducer, remove or stop using _producer, and modify _emit_restart_command to call the injected emitter; update any call sites/tests to provide the emitter and keep the handler free of direct event-bus logic (use dependency injection so the actual Kafka publish happens outside the handler).
154-167:⚠️ Potential issue | 🟠 Major
RESTART_PENDINGis set even when restart emission fails or is skipped.At Line 158,
_emit_restart_command()is awaited, but it returns no success status.handle()then always setsRESTART_PENDING(Line 159) and returnsaction="restart_command"(Line 164), even when producer is missing (Line 313) or send fails (Line 361).Suggested success-aware flow
- async def _emit_restart_command(self, event: ModelConsumerHealthEvent) -> None: + async def _emit_restart_command(self, event: ModelConsumerHealthEvent) -> bool: if self._producer is None: logger.warning("Cannot emit restart command: no producer available") - return + return False ... try: ... await self._producer.send(TOPIC_CONSUMER_RESTART_CMD, value=payload) ... + return True except Exception: logger.warning(..., exc_info=True) + return False - await self._emit_restart_command(event) + emitted = await self._emit_restart_command(event) + if not emitted: + return ModelTriageResult( + fingerprint=event.fingerprint, + action="no_restart", + incident_state=EnumConsumerIncidentState.OPEN, + occurrence_count=occurrence_count, + ) await self._update_incident_state( event.fingerprint, EnumConsumerIncidentState.RESTART_PENDING )Also applies to: 311-367
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py` around lines 154 - 167, The handler sets EnumConsumerIncidentState.RESTART_PENDING and returns action "restart_command" unconditionally after awaiting _emit_restart_command; change the flow so handle() (the method containing this block) verifies the actual success of the restart emission before updating state or returning restart_command: make _emit_restart_command return an explicit success/failure (bool) or raise on fatal failure, then in handle() only call _update_incident_state(event.fingerprint, EnumConsumerIncidentState.RESTART_PENDING) and return ModelTriageResult with action="restart_command" and incident_state=RESTART_PENDING when the emit returned success; if emit failed or was skipped (e.g., missing producer), log/handle the failure and return an appropriate ModelTriageResult (or no state change) instead of marking RESTART_PENDING. Ensure the same change is applied to the duplicate block around lines 311-367.
52-54:⚠️ Potential issue | 🟠 Major“3rd in 30 min” policy is not enforced; threshold uses lifetime count.
Line 154 gates restart on cumulative
occurrence_count, while_upsert_incident()only increments (Line 217) and never window-resets by_RESTART_WINDOW_MINUTES. A third event days later can still trigger restart.Suggested windowed counter update
DO UPDATE SET - occurrence_count = consumer_health_triage.occurrence_count + 1, + occurrence_count = CASE + WHEN consumer_health_triage.last_seen_at < NOW() - INTERVAL '30 minutes' + THEN 1 + ELSE consumer_health_triage.occurrence_count + 1 + END, last_seen_at = NOW(), severity = EXCLUDED.severity, error_message = EXCLUDED.error_message, correlation_id = EXCLUDED.correlation_idAlso applies to: 149-155, 206-223
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py` around lines 52 - 54, The restart policy uses a lifetime occurrence_count so old events can trigger a "3rd in 30 min" restart; update the incident storage and restart check to be windowed: modify _upsert_incident to record timestamps of occurrences (or maintain a last_window_start and window_count) and ensure it prunes or resets counts older than _RESTART_WINDOW_MINUTES, then change the restart gating logic in handler_consumer_health_triage (the block referencing _RESTART_THRESHOLD and occurrence_count around where restart is decided) to use the computed occurrences_in_window (or window_count) instead of the cumulative occurrence_count; keep names _RESTART_THRESHOLD, _RESTART_WINDOW_MINUTES, _upsert_incident, and occurrence_count in your changes so the lookup is obvious.
🧹 Nitpick comments (3)
src/omnibase_infra/runtime/service_kernel.py (2)
2258-2268: Minor: Finally block doesn't detach from loggers unlike normal shutdown.The normal shutdown path (lines 2088) calls
detach_from_loggers(), but the finally cleanup doesn't. At process termination this is acceptable since loggers won't outlive the process, but for completeness and consistency, consider adding the detach call here as well.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/omnibase_infra/runtime/service_kernel.py` around lines 2258 - 2268, The finally cleanup for RuntimeLogEventBridge omits detaching from loggers; update the cleanup block that awaits runtime_log_bridge.stop() (where runtime_log_bridge is checked) to also call runtime_log_bridge.detach_from_loggers() (or the module-level detach_from_loggers() helper used in the normal shutdown path) after stopping the bridge and its _producer, ensuring the same detach logic as the normal shutdown path is executed even during the final cleanup and wrap the detach call in the existing try/except best-effort block to avoid raising on cleanup failures.
2078-2103: Duplicated allowlist parsing and direct access to private_producerattribute.
Duplicated allowlist parsing (lines 2081-2087): This duplicates the parsing logic from startup (lines 986-992). If the env var changes between startup and shutdown, the allowlist could be inconsistent. Consider storing the parsed allowlist as an instance variable when attaching.
Private attribute access (lines 2091-2092):
runtime_log_bridge._produceraccesses an internal attribute. Consider either:
- Exposing a
shutdown()method onRuntimeLogEventBridgethat handles producer cleanup internally, or- Storing the producer reference separately when creating the bridge.
♻️ Suggested approach: Store allowlist and producer reference at startup
Track these at initialization time (around line 994):
# After successful bridge start, store for shutdown _bridge_allowlist = allowlist _bridge_producer_ref = _bridge_producerThen use them at shutdown instead of re-parsing/accessing private attributes.
Alternatively, extend
RuntimeLogEventBridgewith ashutdown()method that handles both detach and producer cleanup.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/omnibase_infra/runtime/service_kernel.py` around lines 2078 - 2103, The shutdown block re-parses RUNTIME_LOG_BRIDGE_ALLOWLIST and accesses the private runtime_log_bridge._producer; instead persist the parsed allowlist and producer reference at bridge startup (e.g., store _bridge_allowlist and _bridge_producer on the kernel or context when creating runtime_log_bridge) and use those stored values in the shutdown logic, or better add a public RuntimeLogEventBridge.shutdown() method that calls detach_from_loggers(allowlist) and stops both the bridge and its producer internally so the shutdown code only calls runtime_log_bridge.shutdown() / await runtime_log_bridge.stop() without touching _producer or re-parsing the env var.src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py (1)
346-353: Use_RESTART_WINDOW_MINUTESconsistently instead of hardcoded30 minutes.Line 346 and Line 351 hardcode the interval while
_RESTART_WINDOW_MINUTESalready exists (Line 53). This can drift silently.Suggested constant-driven SQL parameterization
- restart_count_30min = CASE - WHEN consumer_restart_state.restart_window_start < NOW() - INTERVAL '30 minutes' + restart_count_30min = CASE + WHEN consumer_restart_state.restart_window_start < NOW() - ($4::int * INTERVAL '1 minute') THEN 1 ELSE consumer_restart_state.restart_count_30min + 1 END, restart_window_start = CASE - WHEN consumer_restart_state.restart_window_start < NOW() - INTERVAL '30 minutes' + WHEN consumer_restart_state.restart_window_start < NOW() - ($4::int * INTERVAL '1 minute') THEN NOW() ELSE consumer_restart_state.restart_window_start END,event.consumer_identity, event.consumer_group, event.topic, + _RESTART_WINDOW_MINUTES,🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py` around lines 346 - 353, The SQL uses a hardcoded "30 minutes" interval in the CASE that updates consumer_restart_state (checks restart_window_start < NOW() - INTERVAL '30 minutes' and sets restart_window_start using the same interval); replace those hardcoded literals with the existing constant _RESTART_WINDOW_MINUTES (from this module) so the interval is driven by that single constant (e.g., build the SQL interval using _RESTART_WINDOW_MINUTES) and update both occurrences affecting restart_count_30min and restart_window_start to avoid drift.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In
`@src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py`:
- Around line 156-159: The restart rate-limit check (_check_restart_rate_limit)
and the counter mutation are separate, allowing a race; replace these with an
atomic acquisition method (e.g. implement _try_acquire_restart_slot) that
performs the check-and-increment inside a single DB operation or transaction
(single INSERT ... ON CONFLICT/UPDATE with WHERE restart_count_30min <
_MAX_RESTARTS_PER_WINDOW RETURNING or SELECT ... FOR UPDATE then update) and
return a boolean; call _try_acquire_restart_slot from the handler instead of
_check_restart_rate_limit and only call
_emit_restart_command/_update_incident_state when it returns true so the restart
decision and counter update cannot be bypassed by concurrent handlers.
In
`@src/omnibase_infra/nodes/node_runtime_error_triage_effect/handlers/handler_runtime_error_triage.py`:
- Around line 242-274: The ON CONFLICT upsert on runtime_error_triage using "ON
CONFLICT (fingerprint) WHERE incident_state IN ('open', 'suppressed')" will fail
because there is no matching UNIQUE constraint/index; update the DB migration
that created indexes for runtime_error_triage (migration 055) to add a partial
unique index for the fingerprint when incident_state is open or suppressed
(i.e., create UNIQUE INDEX IF NOT EXISTS
idx_runtime_error_triage_fingerprint_open_suppressed ON
runtime_error_triage(fingerprint) WHERE incident_state IN ('open','suppressed'))
so the upsert in handler_runtime_error_triage.py (the INSERT ... ON CONFLICT
(fingerprint) WHERE incident_state IN ('open', 'suppressed') statement) has a
corresponding unique constraint.
---
Duplicate comments:
In
`@src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py`:
- Around line 87-105: HandlerConsumerHealthTriage currently holds a concrete
AIOKafkaProducer in the _producer attribute and performs direct publishes in
_emit_restart_command, violating the handler boundary rule; replace the concrete
Kafka dependency by injecting a boundary-safe abstraction (e.g., a callable like
restart_command_emitter: Callable[[dict], Awaitable[None]] or a simple
interface) via the constructor instead of AIOKafkaProducer, remove or stop using
_producer, and modify _emit_restart_command to call the injected emitter; update
any call sites/tests to provide the emitter and keep the handler free of direct
event-bus logic (use dependency injection so the actual Kafka publish happens
outside the handler).
- Around line 154-167: The handler sets
EnumConsumerIncidentState.RESTART_PENDING and returns action "restart_command"
unconditionally after awaiting _emit_restart_command; change the flow so
handle() (the method containing this block) verifies the actual success of the
restart emission before updating state or returning restart_command: make
_emit_restart_command return an explicit success/failure (bool) or raise on
fatal failure, then in handle() only call
_update_incident_state(event.fingerprint,
EnumConsumerIncidentState.RESTART_PENDING) and return ModelTriageResult with
action="restart_command" and incident_state=RESTART_PENDING when the emit
returned success; if emit failed or was skipped (e.g., missing producer),
log/handle the failure and return an appropriate ModelTriageResult (or no state
change) instead of marking RESTART_PENDING. Ensure the same change is applied to
the duplicate block around lines 311-367.
- Around line 52-54: The restart policy uses a lifetime occurrence_count so old
events can trigger a "3rd in 30 min" restart; update the incident storage and
restart check to be windowed: modify _upsert_incident to record timestamps of
occurrences (or maintain a last_window_start and window_count) and ensure it
prunes or resets counts older than _RESTART_WINDOW_MINUTES, then change the
restart gating logic in handler_consumer_health_triage (the block referencing
_RESTART_THRESHOLD and occurrence_count around where restart is decided) to use
the computed occurrences_in_window (or window_count) instead of the cumulative
occurrence_count; keep names _RESTART_THRESHOLD, _RESTART_WINDOW_MINUTES,
_upsert_incident, and occurrence_count in your changes so the lookup is obvious.
---
Nitpick comments:
In
`@src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.py`:
- Around line 346-353: The SQL uses a hardcoded "30 minutes" interval in the
CASE that updates consumer_restart_state (checks restart_window_start < NOW() -
INTERVAL '30 minutes' and sets restart_window_start using the same interval);
replace those hardcoded literals with the existing constant
_RESTART_WINDOW_MINUTES (from this module) so the interval is driven by that
single constant (e.g., build the SQL interval using _RESTART_WINDOW_MINUTES) and
update both occurrences affecting restart_count_30min and restart_window_start
to avoid drift.
In `@src/omnibase_infra/runtime/service_kernel.py`:
- Around line 2258-2268: The finally cleanup for RuntimeLogEventBridge omits
detaching from loggers; update the cleanup block that awaits
runtime_log_bridge.stop() (where runtime_log_bridge is checked) to also call
runtime_log_bridge.detach_from_loggers() (or the module-level
detach_from_loggers() helper used in the normal shutdown path) after stopping
the bridge and its _producer, ensuring the same detach logic as the normal
shutdown path is executed even during the final cleanup and wrap the detach call
in the existing try/except best-effort block to avoid raising on cleanup
failures.
- Around line 2078-2103: The shutdown block re-parses
RUNTIME_LOG_BRIDGE_ALLOWLIST and accesses the private
runtime_log_bridge._producer; instead persist the parsed allowlist and producer
reference at bridge startup (e.g., store _bridge_allowlist and _bridge_producer
on the kernel or context when creating runtime_log_bridge) and use those stored
values in the shutdown logic, or better add a public
RuntimeLogEventBridge.shutdown() method that calls
detach_from_loggers(allowlist) and stops both the bridge and its producer
internally so the shutdown code only calls runtime_log_bridge.shutdown() / await
runtime_log_bridge.stop() without touching _producer or re-parsing the env var.
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: b4afe523-fde3-4ff1-869f-f551adf87c39
📒 Files selected for processing (4)
src/omnibase_infra/nodes/node_consumer_health_triage_effect/handlers/handler_consumer_health_triage.pysrc/omnibase_infra/nodes/node_runtime_error_triage_effect/handlers/handler_runtime_error_triage.pysrc/omnibase_infra/runtime/service_kernel.pysrc/omnibase_infra/topics/__init__.py
✅ Files skipped from review due to trivial changes (1)
- src/omnibase_infra/topics/init.py
a3958ec to
b7e8b54
Compare
4e44005 to
5d749e5
Compare
…ns [OMN-5511, OMN-5512, OMN-5513, OMN-5514] Wave 0 of OMN-5529 Runtime Health Event Pipeline.
…15, OMN-5516, OMN-5517]
…N-5521] Wave 2 (partial) of OMN-5529 Runtime Health Event Pipeline.
… effect nodes [OMN-5518, OMN-5520, OMN-5522] OMN-5518: Wire ConsumerHealthEmitter into EventBusKafka - Initialize emitter after producer starts (gated by ENABLE_CONSUMER_HEALTH_EMITTER) - Emit CONSUMER_STARTED on successful consumer creation - Emit CONNECTION_LOST/SESSION_TIMEOUT on consume loop errors - Add health_emitter property for external access OMN-5520: Create NodeConsumerHealthTriageEffect - Graduated response: 1st->Slack WARNING, 2nd->Slack REPEATED, 3rd->restart, fail->Linear - HandlerConsumerHealthTriage with PostgreSQL state tracking - Restart rate limiting via consumer_restart_state table - Gated by ENABLE_CONSUMER_HEALTH_TRIAGE + ENABLE_CONSUMER_AUTO_RESTART OMN-5522: Create NodeRuntimeErrorTriageEffect - First-match-wins triage rule engine with ModelTriageRule - Default rules for aiokafka, asyncpg, aiohttp error patterns - Cross-layer correlation with Layer 1 consumer health incidents - Actions: suppress, alert (Slack), ticket (Linear)
…N-5520, OMN-5522] - Cast dict values to int/str for proper type narrowing from dict[str, object] - Use str truncation instead of sanitize_error_message (expects Exception, not str)
…N-5520, OMN-5522]
…eLogEventBridge into service_kernel [OMN-5523, OMN-5525] OMN-5523: Add MixinConsumerHealth to all 6 standalone consumers (session, agent_actions, injection_effectiveness, llm_cost_aggregation, context_audit, skill_lifecycle). Consumers without an existing producer get a dedicated health producer, gated by ENABLE_CONSUMER_HEALTH_EMITTER. OMN-5525: Wire RuntimeLogEventBridge into service_kernel bootstrap. When ENABLE_RUNTIME_LOG_BRIDGE is enabled, creates a dedicated producer, attaches the bridge to allowlisted loggers (RUNTIME_LOG_BRIDGE_ALLOWLIST), and starts/stops with kernel lifecycle. Best-effort -- never blocks startup.
…e pipelines [OMN-5524, OMN-5526] OMN-5524: Integration tests for Layer 1 consumer health pipeline — emitter-to-Kafka flow, rate limiting, feature flag gating, self-metrics, and MixinConsumerHealth wiring. OMN-5526: Integration tests for Layer 2 runtime log event bridge — log-to-Kafka flow, circular logging prevention, rate limiting, feature flag gating, and bridge metrics accuracy. Both require running Kafka/Redpanda broker, marked with @pytest.mark.kafka.
…ffixes [OMN-5529]
Conflict resolution during rebase dropped the closing triple-quote for TOPIC_VALIDATOR_CATCH docstring, causing a syntax error that cascaded into all CI checks failing.
5d749e5 to
de16a60
Compare
Summary
Foundation for the Runtime Health Event Pipeline (OMN-5529). Implements Waves 0-2 of the 5-wave epic:
New Components
ModelConsumerHealthEventModelRuntimeErrorEventConsumerHealthEmitterENABLE_CONSUMER_HEALTH_EMITTERMixinConsumerHealthENABLE_CONSUMER_HEALTH_EMITTERRuntimeLogEventBridgeENABLE_RUNTIME_LOG_BRIDGENew Topics
onex.evt.omnibase-infra.consumer-health.v1- Consumer health eventsonex.cmd.omnibase-infra.consumer-restart.v1- Consumer restart commandsonex.evt.omnibase-infra.runtime-error.v1- Runtime error eventsMigrations
054_create_consumer_health_triage.sql- Incident tracking + restart rate limiting055_create_runtime_error_triage.sql- Runtime error incident tracking with cross-layer correlationTest plan
Summary by CodeRabbit
New Features
Integration
Tests
Chores