Skip to content

feat(kafka): wire session timeout to all consumers [OMN-5445] - #905

Merged
jonahgabriel merged 4 commits into
mainfrom
jonahgabriel/omn-5445-kafka-consumer-session-timeout-fix
Mar 20, 2026
Merged

jonahgabriel merged 4 commits into
mainfrom
jonahgabriel/omn-5445-kafka-consumer-session-timeout-fix

Conversation

@jonahgabriel

@jonahgabriel jonahgabriel commented Mar 19, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

  • Add session_timeout_ms (30s) and heartbeat_interval_ms (10s) to ModelKafkaEventBusConfig with env var overrides (KAFKA_SESSION_TIMEOUT_MS, KAFKA_HEARTBEAT_INTERVAL_MS) and advisory heartbeat/session ratio validator
  • Wire session_timeout_ms, heartbeat_interval_ms, reconnect_backoff_ms, reconnect_backoff_max_ms into EventBusKafka consumer constructor; wire reconnect_backoff_ms/reconnect_backoff_max_ms into both producer constructors (previously configured but never passed)
  • Add session timeout fields to all 6 standalone consumer BaseSettings configs and wire into their AIOKafkaConsumer constructors
  • Hardcode 30000/10000 defaults in 3 ephemeral consumers (no config surface)
  • Add pattern validation exemption for ModelKafkaEventBusConfig method count (pre-existing borderline at 10 methods)

Test plan

  • 19 new tests for ModelKafkaEventBusConfig session timeout (defaults, constraints, env overrides, advisory validator, JSON schema)
  • 7 new tests for EventBusKafka consumer/producer kwargs wiring (session_timeout, heartbeat, reconnect backoff)
  • 24 new tests for standalone consumer config fields (parametrized across all 6 configs)
  • 11 existing reconnect backoff tests pass (no regression)
  • All 38 pre-commit hooks pass
  • ruff check + format clean

Summary by CodeRabbit

  • New Features

    • Configurable Kafka consumer session timeout and heartbeat interval exposed across event bus and observability services (defaults 30000ms / 10000ms, with range constraints).
    • Producer/consumer reconnect/backoff parameters forwarded into Kafka clients.
    • New platform DLQ topic suffix added for skill-lifecycle and used as the default DLQ where applicable.
  • Validation

    • Advisory validator warns when heartbeat is too large relative to session timeout; validation exemptions updated.
  • Tests

    • Added unit tests for defaults, bounds, env overrides, and wiring of timeout/reconnect configs.

… consumers [OMN-5445]

Eliminate implicit reliance on aiokafka's aggressive 10s default session
timeout across all AIOKafkaConsumer instances in the platform. The new
platform defaults (30s session / 10s heartbeat) prevent rebalance storms
caused by UnknownMemberIdError during transient processing delays.

Three-surface fix:
- ModelKafkaEventBusConfig: new fields with env var overrides
  (KAFKA_SESSION_TIMEOUT_MS, KAFKA_HEARTBEAT_INTERVAL_MS) and advisory
  heartbeat/session ratio validator
- EventBusKafka: wire session/heartbeat/reconnect kwargs into consumer
  constructor; wire reconnect backoff kwargs into both producer constructors
  (previously configured but never passed)
- 6 standalone consumers: add fields to BaseSettings configs and wire
  into AIOKafkaConsumer constructors
- 3 ephemeral consumers: hardcode 30000/10000 defaults (no config surface)

Includes 50 new unit tests covering config defaults, env overrides,
field constraints, advisory validator, and constructor kwarg wiring.
@coderabbitai

coderabbitai Bot commented Mar 19, 2026 •

Copy link
Copy Markdown

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Pro

Run ID: c927626b-a80c-449c-9b32-b2d1e4dc6864

📥 Commits

Reviewing files that changed from the base of the PR and between f3db0bd and 49bc385.

📒 Files selected for processing (1)
  • tests/unit/topics/test_platform_topic_suffixes.py
✅ Files skipped from review due to trivial changes (1)
  • tests/unit/topics/test_platform_topic_suffixes.py

📝 Walkthrough

Walkthrough

Adds Kafka session/heartbeat and reconnect tuning: new Pydantic fields and env overrides, propagation of session/heartbeat and reconnect_backoff kwargs into AIOKafka producer/consumer constructors, addition/export of a DLQ topic suffix, validation exemption, and tests for the new wiring.

Changes

Cohort / File(s) Summary
Core Kafka Config
src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py
Added session_timeout_ms and heartbeat_interval_ms fields (defaults/constraints), post-validator that warns if heartbeat > session/3, and env-var overrides (KAFKA_SESSION_TIMEOUT_MS, KAFKA_HEARTBEAT_INTERVAL_MS).
Event Bus Kafka
src/omnibase_infra/event_bus/event_bus_kafka.py
Producer creation and lazy-producer recreation now forward reconnect_backoff_ms and reconnect_backoff_max_ms; consumer creation forwards session_timeout_ms, heartbeat_interval_ms, and reconnect backoff kwargs.
Service Consumer Configs
src/omnibase_infra/services/.../config.py
agent_actions, context_audit, injection_effectiveness, llm_cost_aggregation, skill_lifecycle, session
Added session_timeout_ms and heartbeat_interval_ms fields (defaults/constraints) to multiple consumer config classes; some DLQ defaults switched to centralized suffix constants.
Service Consumer Implementations
src/omnibase_infra/services/.../consumer.py
agent_actions, context_audit, injection_effectiveness, llm_cost_aggregation, skill_lifecycle, session
Consumer start() implementations now pass session_timeout_ms and heartbeat_interval_ms into AIOKafkaConsumer(...).
Standalone Consumers & Wiring
src/omnibase_infra/projectors/snapshot_publisher_registration.py, src/omnibase_infra/runtime/request_response_wiring.py, src/omnibase_infra/tui/consumers/consumer_status.py
Standalone consumer constructions updated to include session_timeout_ms=30000 and heartbeat_interval_ms=10000.
Topic Suffixes & Exports
src/omnibase_infra/topics/platform_topic_suffixes.py, src/omnibase_infra/topics/__init__.py
Added/exported SUFFIX_OMNICLAUDE_SKILL_LIFECYCLE_DLQ and included it in observability DLQ suffixes used for platform topic specs.
Scripts & Baselines
scripts/check_contract_topic_parity.py, scripts/validation/topic_literal_baseline.txt
Added legacy allowlist entry for the new DLQ suffix and adjusted topic-literal baseline line entries.
Validation Exemptions
src/omnibase_infra/validation/validation_exemptions.yaml
Added exemption for ModelKafkaEventBusConfig concerning method-count validator due to Pydantic cross-field validators and computed properties (OMN-5445).
Tests
tests/unit/event_bus/test_kafka_config_session_timeout.py, tests/unit/event_bus/test_kafka_timeout_kwargs.py, tests/unit/services/test_consumer_session_timeout.py, tests/unit/topics/test_platform_topic_suffixes.py
Added tests asserting defaults, bounds, validator warning behavior, env overrides, JSON schema presence, that timeout/reconnect kwargs are forwarded to AIOKafka constructors, and the new DLQ suffix partition spec.

Sequence Diagram(s)

(Skipped — changes are configuration wiring and small constructor argument additions; no new multi-component control flow requiring a sequence diagram.)

Estimated code review effort

🎯 3 (Moderate) | ⏱️ ~25 minutes

Possibly related PRs

Poem

🐰
I tuned the heartbeats, set the time just right,
Producers reconnect, consumers hum through night.
Config carrots lined in tidy rows,
Little rabbit hops where the heartbeat flows. 🥕

🚥 Pre-merge checks | ✅ 3
✅ Passed checks (3 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title accurately describes the main objective of the changeset: adding Kafka session timeout configuration to all consumers. The ticket reference (OMN-5445) provides additional context.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch jonahgabriel/omn-5445-kafka-consumer-session-timeout-fix
📝 Coding Plan
  • Generate coding plan for human review comments

Comment @coderabbitai help to get the list of available commands and usage tips.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick comments (2)
tests/unit/event_bus/test_kafka_timeout_kwargs.py (1)

37-151: Consider adding cleanup to prevent resource leaks in tests.

The tests create EventBusKafka instances and call bus.start() but don't call bus.close() afterward. While the mocks prevent actual Kafka connections, the EventBusKafka class may create asyncio tasks or other resources that should be cleaned up.

♻️ Suggested: Add cleanup with try/finally or fixture

Option 1 — Add cleanup in each test:

 async def test_consumer_receives_session_timeout_ms(
     self, bus: EventBusKafka, kafka_config: ModelKafkaEventBusConfig
 ) -> None:
     """Consumer constructor must receive session_timeout_ms from config."""
     mock_consumer = MagicMock()
     mock_consumer.start = AsyncMock()

-    with patch(
-        "omnibase_infra.event_bus.event_bus_kafka.AIOKafkaProducer"
-    ) as mock_producer_cls:
-        mock_producer = MagicMock()
-        mock_producer.start = AsyncMock()
-        mock_producer.stop = AsyncMock()
-        mock_producer_cls.return_value = mock_producer
-        await bus.start()
+    try:
+        with patch(
+            "omnibase_infra.event_bus.event_bus_kafka.AIOKafkaProducer"
+        ) as mock_producer_cls:
+            mock_producer = MagicMock()
+            mock_producer.start = AsyncMock()
+            mock_producer.stop = AsyncMock()
+            mock_producer_cls.return_value = mock_producer
+            await bus.start()
+
+        with patch(
+            "omnibase_infra.event_bus.event_bus_kafka.AIOKafkaConsumer",
+            return_value=mock_consumer,
+        ) as mock_consumer_cls:
+            await bus.subscribe(
+                "test-topic", on_message=AsyncMock(), group_id="test-group"
+            )
+            call_kwargs = mock_consumer_cls.call_args
+            assert call_kwargs.kwargs["session_timeout_ms"] == 45000
+    finally:
+        await bus.close()

Option 2 — Update the fixture to yield and cleanup:

`@pytest.fixture`
async def bus(kafka_config: ModelKafkaEventBusConfig) -> AsyncIterator[EventBusKafka]:
    """Create EventBusKafka instance with test config."""
    bus = EventBusKafka(config=kafka_config)
    yield bus
    await bus.close()
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@tests/unit/event_bus/test_kafka_timeout_kwargs.py` around lines 37 - 151,
Tests start EventBusKafka instances via bus.start() but never call bus.close(),
risking leaked asyncio tasks; ensure each test cleans up by calling await
bus.close() in a finally block or convert the bus fixture to a yield-style async
fixture that yields EventBusKafka and calls await bus.close() after yield
(reference: EventBusKafka, bus.start, bus.close, and the bus fixture).
src/omnibase_infra/services/session/config_consumer.py (1)

62-81: Consider adding advisory validation for heartbeat/session ratio.

The defaults are correctly configured (30000ms session timeout with 10000ms heartbeat = exactly the recommended 1:3 ratio). However, unlike ModelKafkaEventBusConfig which has an advisory validator, this config allows users to set incompatible values without warning.

For consistency with the event bus config, consider adding a @model_validator(mode="after") that logs a warning when heartbeat_interval_ms > session_timeout_ms / 3.

♻️ Optional: Add advisory validator for ratio
     `@model_validator`(mode="after")
     def validate_timing_relationships(self) -> Self:
         """Validate timing relationships between configuration values.
         ...
         """
         batch_timeout_seconds = self.batch_timeout_ms / 1000
         min_recommended_circuit_timeout = batch_timeout_seconds * 2

         if self.circuit_breaker_timeout_seconds < min_recommended_circuit_timeout:
             logger.warning(
                 "Circuit breaker timeout (%ds) is less than 2x batch timeout (%.1fs). "
                 ...
             )
+
+        # Kafka recommends heartbeat_interval_ms <= session_timeout_ms / 3
+        if self.heartbeat_interval_ms > self.session_timeout_ms / 3:
+            logger.warning(
+                "heartbeat_interval_ms (%d) exceeds Kafka's recommended maximum "
+                "of session_timeout_ms / 3 (%d). This may cause premature rebalances.",
+                self.heartbeat_interval_ms,
+                self.session_timeout_ms // 3,
+            )
         return self
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/omnibase_infra/services/session/config_consumer.py` around lines 62 - 81,
Add an advisory post-validation to the Pydantic model that contains
session_timeout_ms and heartbeat_interval_ms: implement a
`@model_validator`(mode="after") (same pattern as ModelKafkaEventBusConfig) on
that config class which checks if heartbeat_interval_ms > session_timeout_ms / 3
and logs a warning (use the module/class logger) rather than raising an
exception; reference the fields session_timeout_ms and heartbeat_interval_ms in
the check and include their values in the warning message so users receive a
clear advisory when their ratio is incompatible.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Nitpick comments:
In `@src/omnibase_infra/services/session/config_consumer.py`:
- Around line 62-81: Add an advisory post-validation to the Pydantic model that
contains session_timeout_ms and heartbeat_interval_ms: implement a
`@model_validator`(mode="after") (same pattern as ModelKafkaEventBusConfig) on
that config class which checks if heartbeat_interval_ms > session_timeout_ms / 3
and logs a warning (use the module/class logger) rather than raising an
exception; reference the fields session_timeout_ms and heartbeat_interval_ms in
the check and include their values in the warning message so users receive a
clear advisory when their ratio is incompatible.

In `@tests/unit/event_bus/test_kafka_timeout_kwargs.py`:
- Around line 37-151: Tests start EventBusKafka instances via bus.start() but
never call bus.close(), risking leaked asyncio tasks; ensure each test cleans up
by calling await bus.close() in a finally block or convert the bus fixture to a
yield-style async fixture that yields EventBusKafka and calls await bus.close()
after yield (reference: EventBusKafka, bus.start, bus.close, and the bus
fixture).

ℹ️ Review info
⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Pro

Run ID: be289d71-a7a0-4d8d-84d0-ad604761c0c0

📥 Commits

Reviewing files that changed from the base of the PR and between 8f344fd and 83e1bd0.

📒 Files selected for processing (21)
  • src/omnibase_infra/event_bus/event_bus_kafka.py
  • src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py
  • src/omnibase_infra/projectors/snapshot_publisher_registration.py
  • src/omnibase_infra/runtime/request_response_wiring.py
  • src/omnibase_infra/services/observability/agent_actions/config.py
  • src/omnibase_infra/services/observability/agent_actions/consumer.py
  • src/omnibase_infra/services/observability/context_audit/config.py
  • src/omnibase_infra/services/observability/context_audit/consumer.py
  • src/omnibase_infra/services/observability/injection_effectiveness/config.py
  • src/omnibase_infra/services/observability/injection_effectiveness/consumer.py
  • src/omnibase_infra/services/observability/llm_cost_aggregation/config.py
  • src/omnibase_infra/services/observability/llm_cost_aggregation/consumer.py
  • src/omnibase_infra/services/observability/skill_lifecycle/config.py
  • src/omnibase_infra/services/observability/skill_lifecycle/consumer.py
  • src/omnibase_infra/services/session/config_consumer.py
  • src/omnibase_infra/services/session/consumer.py
  • src/omnibase_infra/tui/consumers/consumer_status.py
  • src/omnibase_infra/validation/validation_exemptions.yaml
  • tests/unit/event_bus/test_kafka_config_session_timeout.py
  • tests/unit/event_bus/test_kafka_timeout_kwargs.py
  • tests/unit/services/test_consumer_session_timeout.py

…MN-5445]

Moves two raw topic literal strings in observability consumer config files
to the canonical constants in platform_topic_suffixes.py, fixing the
Arch Invariants (OMN-3343) CI check that blocked this PR.

- agent_actions/config.py: import and use SUFFIX_OMNICLAUDE_AGENT_ACTIONS_DLQ
- skill_lifecycle/config.py: import and use new SUFFIX_OMNICLAUDE_SKILL_LIFECYCLE_DLQ
- platform_topic_suffixes.py: add SUFFIX_OMNICLAUDE_SKILL_LIFECYCLE_DLQ constant
  and include it in _OMNICLAUDE_OBSERVABILITY_DLQ_TOPIC_SUFFIXES for provisioning
- topic_literal_baseline.txt: update grandfathered line numbers after import insertion;
  remove now-fixed DLQ entries (lines 131/134 agent_actions, 124/126 skill_lifecycle)
- check_contract_topic_parity.py: add skill-lifecycle-dlq to legacy allowlist
  (mirrors pattern for agent-actions-dlq and agent-observability-dlq)
@jonahgabriel
jonahgabriel enabled auto-merge March 19, 2026 14:16

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

🤖 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/topics/platform_topic_suffixes.py`:
- Around line 648-650: Update the docstring on
ConfigSkillLifecycleConsumer.dlq_topic to remove the phrase "hardcoded default"
and instead state that the default is supplied via the centralized shared suffix
constant from platform_topic_suffixes (refer to the shared suffix constant, e.g.
SKILL_LIFECYCLE_DLQ_SUFFIX), so the wording reflects centralized configuration
rather than a hardcoded value.

ℹ️ Review info
⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Pro

Run ID: f44b57e2-81d1-4409-bc1f-53c58222f7c8

📥 Commits

Reviewing files that changed from the base of the PR and between 83e1bd0 and 950283c.

📒 Files selected for processing (5)
  • scripts/check_contract_topic_parity.py
  • scripts/validation/topic_literal_baseline.txt
  • src/omnibase_infra/services/observability/agent_actions/config.py
  • src/omnibase_infra/services/observability/skill_lifecycle/config.py
  • src/omnibase_infra/topics/platform_topic_suffixes.py
✅ Files skipped from review due to trivial changes (2)
  • scripts/check_contract_topic_parity.py
  • scripts/validation/topic_literal_baseline.txt
🚧 Files skipped from review as they are similar to previous changes (2)
  • src/omnibase_infra/services/observability/agent_actions/config.py
  • src/omnibase_infra/services/observability/skill_lifecycle/config.py

Comment thread src/omnibase_infra/topics/platform_topic_suffixes.py
@jonahgabriel
jonahgabriel added this pull request to the merge queue Mar 19, 2026
@github-merge-queue
github-merge-queue Bot removed this pull request from the merge queue due to no response for status checks Mar 19, 2026
@jonahgabriel
jonahgabriel added this pull request to the merge queue Mar 19, 2026
Merged via the queue into main with commit 3e2f698 Mar 20, 2026
45 checks passed
@jonahgabriel
jonahgabriel deleted the jonahgabriel/omn-5445-kafka-consumer-session-timeout-fix branch March 20, 2026 00:50
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant