Repository navigation
fix(OMN-11318): avoid grouped Pattern B terminal waits - #1680
Conversation
📝 WalkthroughWalkthroughThis PR fixes a local ingress hang in Pattern B terminal waiter by switching from a grouped subscription-based Kafka consumer to an ungrouped consumer with explicit partition discovery and assignment. The implementation adds a helper that discovers partitions and assigns them with retry timeout before publishing the worker command. ChangesPattern B Kafka consumer ungrouped partition assignment
Sequence Diagram(s)sequenceDiagram
participant RuntimePatternBBroker
participant AIOKafkaConsumer
participant KafkaCluster
RuntimePatternBBroker->>AIOKafkaConsumer: create(bootstrap_servers, group_id=None)
RuntimePatternBBroker->>AIOKafkaConsumer: partitions_for_topic(terminal_topic)
AIOKafkaConsumer->>KafkaCluster: metadata request for topic partitions
KafkaCluster-->>AIOKafkaConsumer: partition list
RuntimePatternBBroker->>AIOKafkaConsumer: assign(TopicPartition set)
AIOKafkaConsumer->>KafkaCluster: assign partitions
RuntimePatternBBroker->>AIOKafkaConsumer: assignment() with retry until timeout
AIOKafkaConsumer-->>RuntimePatternBBroker: assigned partitions
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~20 minutes Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Actionable comments posted: 1
🧹 Nitpick comments (1)
tests/unit/runtime/test_service_pattern_b_broker.py (1)
114-121: ⚡ Quick winMake partition discovery configurable per topic in the fake consumer.
Current fake behavior always resolves every topic once started, so it can’t catch regressions where assignment happens before all terminal topics are discoverable.
Suggested refactor sketch
class _FakeAIOKafkaConsumer: @@ - def __init__(self, *topics: object, **kwargs: object) -> None: + def __init__(self, *topics: object, **kwargs: object) -> None: self.topics = topics self.kwargs = kwargs self.messages: asyncio.Queue[SimpleNamespace] = asyncio.Queue() self.assigned_partitions: set[object] = set() + self.topic_partitions: dict[str, set[int]] = {} @@ def partitions_for_topic(self, topic: str) -> set[int]: - return {0} if self.started and topic else set() + if not self.started or not topic: + return set() + return set(self.topic_partitions.get(topic, {0}))🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@tests/unit/runtime/test_service_pattern_b_broker.py` around lines 114 - 121, The fake consumer currently resolves every topic once started via partitions_for_topic, causing tests to miss cases where some topics aren't yet discoverable; modify the fake consumer to accept a configurable mapping or callback (e.g., a constructor arg like discoverable_partitions: dict[str, set[int]] or partitions_resolver: Callable[[str], set[int]]) and change partitions_for_topic(self, topic: str) to return discoverable_partitions.get(topic, set()) or call the resolver, leaving assign(self, partitions: set[object]) and assignment(self) untouched so tests can control which topics are discoverable and simulate assignment-before-discovery regressions.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@src/omnibase_infra/runtime/service_pattern_b_broker.py`:
- Around line 309-320: The loop currently returns as soon as any partition is
found, which lets multi-terminal routes assign only a subset of topics; change
the condition so it waits until every topic in topics has at least one partition
before calling consumer.assign(partitions) and returning. Use
consumer.partitions_for_topic(topic) to detect per-topic partitions, build the
TopicPartition set as shown, then verify that all topics are represented (e.g.,
every topic has non-empty topic_partitions) before assigning; keep the same
timeout check using asyncio.get_running_loop().time() >= deadline and raise
TimeoutError if the deadline is exceeded.
---
Nitpick comments:
In `@tests/unit/runtime/test_service_pattern_b_broker.py`:
- Around line 114-121: The fake consumer currently resolves every topic once
started via partitions_for_topic, causing tests to miss cases where some topics
aren't yet discoverable; modify the fake consumer to accept a configurable
mapping or callback (e.g., a constructor arg like discoverable_partitions:
dict[str, set[int]] or partitions_resolver: Callable[[str], set[int]]) and
change partitions_for_topic(self, topic: str) to return
discoverable_partitions.get(topic, set()) or call the resolver, leaving
assign(self, partitions: set[object]) and assignment(self) untouched so tests
can control which topics are discoverable and simulate
assignment-before-discovery regressions.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: 7631286d-09f8-427b-8457-c1580c9b09c5
📒 Files selected for processing (3)
contracts/OMN-11318.yamlsrc/omnibase_infra/runtime/service_pattern_b_broker.pytests/unit/runtime/test_service_pattern_b_broker.py
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@tests/integration/test_pattern_b_broker_terminal_waiter.py`:
- Line 79: Replace the hardcoded broker endpoint assigned to _bootstrap_servers
("redpanda:9092") with a value read from an environment variable (e.g.,
KAFKA_BOOTSTRAP_SERVERS), falling back to the current literal as a default;
update the assignment where _bootstrap_servers is defined so tests pick the
broker from os.environ (or equivalent env API) rather than a hardcoded string.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: fcb9fd9d-e351-4da0-aeff-692d0f784bed
📒 Files selected for processing (2)
contracts/OMN-11318.yamltests/integration/test_pattern_b_broker_terminal_waiter.py
Evidence-Ticket: OMN-11318
Summary
Verification
uv run validate-yaml contracts/OMN-11318.yamluv run pytest tests/unit/runtime/test_service_pattern_b_broker.py -quv run ruff format --check src/omnibase_infra/runtime/service_pattern_b_broker.py tests/unit/runtime/test_service_pattern_b_broker.pyuv run ruff check src/omnibase_infra/runtime/service_pattern_b_broker.py tests/unit/runtime/test_service_pattern_b_broker.pyRuntime Evidence
.201runtime rebuilt from merged infra363b05b2was healthy, and delegate-skill command consumer group was stable with 1 member and 0 lag./skillsmoke correlation86622b02-46b5-4339-ab13-2758b2851255emittedonex.evt.omnimarket.delegate-skill-completed.v1, but the HTTP request did not return, proving this terminal waiter boundary.Post-merge deploy proof still required: rebuild/restart the
.201stability runtime from this merged change, then rerun the/skillsmoke and verify it returns after the terminal event.Summary by CodeRabbit
Bug Fixes
Tests
Documentation
Evidence-Source: OCC#1242