Repository navigation
fix(OMN-15395): close the SECOND CreateTopics path — contract-driven RF on the operator CLI, capped drift, memoized probe, loud RF refusal - #2552
Conversation
…moized probe, loud RF refusal Residual round against PR #2550's own adversarial review. Four defects, each proven RED by reverting its own fix hunk (9 mutations, all caught). D2 (HIGH) — scripts/create_kafka_topics.py was a SECOND live CreateTopics path with a flat `--replication-factor` default of 1 that discarded every contract's declared topic_config.replication_factor and never consulted the fail-closed policy. It is the path docs/operations/README.md tells operators to run and the one compare_environments.py names in its topic-parity fix_hint, so the documented runbook reproduced AWS_KAFKA_HIGH_RISK_CONFIG_RF_EQUALS_ONE by hand. It now resolves per-topic through the same ModelTopicProvisioningPolicy, with the same measured capacity ceiling (read off the list_topics() response the diff already needs — no extra round trip) and the same fail-closed batch check before the first create. The flag is gone; the lane partition cap applies too, so the two paths cannot disagree about a topic's shape. Batch resolution is now the shared module-level `resolve_specs_for_creation`, used by both paths. New static guard `_raw_create_topics_offenders` scans src/ AND scripts/ for any NewTopic construction that bypasses the policy or hardcodes a literal replication factor; the swallow guard now scans scripts/ as well. Executed against the pre-fix file, the guard reports `scripts/create_kafka_topics.py:329: module resolves no replication policy` — the exact line the review cited. D3 (HIGH) — _report_spec_drift compared broker partitions against the UNCAPPED spec.partitions while creation applies the ONEX_TOPIC_PROVISIONER_MAX_PARTITIONS cap (live at 1 on dev/stability/judge), so the provisioner reported drift against topics it had itself just created correctly. Drift now compares against _creation_spec — the resolved, capped effective spec. Divergence the cap explains (a topic created before the cap was lowered; Kafka cannot reduce partitions) is labelled partition_cap_suppressed and kept out of the feed the operator-gated reassignment lane consumes, never emitted as partition_mismatch. D4 (LOW) — the capacity memo keyed on `policy.broker_count is None`, which cannot bound an unmeasurable cluster: the field stays None, so the probe re-ran on every entrypoint (measured 3 calls / 3 entrypoints; now 1). The ATTEMPT is memoized. D5 (LOW) — "fails loudly at CreateTopics" was false: the broker's INVALID_REPLICATION_FACTOR landed in the per-topic `except Exception`, became a warning plus a name in `failed`, and the pass returned status="partial" with the topic silently absent. It is now classified and raised as a typed TopicReplicationPolicyError naming the topic, the refused value and whether capacity was measurable — on both runtime paths, and as a classified ERROR on the operator CLI and the managed-staging checker. A transient create failure still degrades to best-effort (negative control). Ticket: OMN-15395
📝 WalkthroughWalkthroughKafka topic creation now derives replication factors from contract entries and runtime provisioning policy, binds policy decisions to measured broker capacity, applies partition caps, and fails closed on unhostable replication factors. Provisioner drift detection and operator workflows share these rules, with expanded regression coverage. ChangesKafka topic durability
Estimated code review effort: 4 (Complex) | ~60 minutes Sequence Diagram(s)sequenceDiagram
participant OperatorCLI
participant ProvisioningPolicy
participant KafkaAdmin
OperatorCLI->>KafkaAdmin: list_topics metadata
KafkaAdmin-->>OperatorCLI: existing topics and broker count
OperatorCLI->>ProvisioningPolicy: resolve missing contract specs
ProvisioningPolicy-->>OperatorCLI: resolved replication factors
OperatorCLI->>KafkaAdmin: create_topics with capped partitions
KafkaAdmin-->>OperatorCLI: creation result or replication refusal
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✨ Finishing Touches 💡 1📝 Generate docstrings 💡
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Actionable comments posted: 4
🧹 Nitpick comments (5)
tests/unit/scripts/test_create_kafka_topics.py (1)
374-384: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueAssert against the managed-floor constant instead of the literal
2. The sibling provisioner test importsMANAGED_MINIMUM_REPLICATION_FACTOR; using it here keeps this test tied to the policy rather than to a copy of its current value.♻️ Suggested tweak
+ from omnibase_infra.topics.model_topic_provisioning_policy import ( + MANAGED_MINIMUM_REPLICATION_FACTOR, + ) + assert _run_live([_entry(replication_factor=None)], recorder) == 0 - assert [t.replication_factor for t in recorder.requested] == [2] + assert [t.replication_factor for t in recorder.requested] == [ + MANAGED_MINIMUM_REPLICATION_FACTOR + ]🤖 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/scripts/test_create_kafka_topics.py` around lines 374 - 384, Update test_undeclared_replication_resolves_to_the_managed_floor_not_one to assert against the imported MANAGED_MINIMUM_REPLICATION_FACTOR constant instead of the literal 2, ensuring the test follows the managed replication policy.tests/unit/event_bus/test_topic_provisioner_rf_policy_omn15395.py (3)
1604-1631: 🎯 Functional Correctness | 🔵 Trivial | ⚡ Quick winAssert the
partition_cap_suppressedlabel, not just the absence of drift. The docstring's stated property is "not silently dropped", but the test only checks the drift feed is empty — it would still pass if_report_spec_driftdropped the divergence entirely. Capturing the info log would pin the distinct label.♻️ Add the log assertion
async def test_partitions_above_the_cap_are_not_a_reassignment_target( self, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, ) -> None: @@ - with _patched_admin(recorder): - result = await provisioner.ensure_provisioned_topics_exist() + with caplog.at_level(logging.INFO): + with _patched_admin(recorder): + result = await provisioner.ensure_provisioned_topics_exist() assert [entry for entry in result["drift"] if TOPIC in entry] == [] + assert "partition_cap_suppressed" in caplog.text🤖 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/event_bus/test_topic_provisioner_rf_policy_omn15395.py` around lines 1604 - 1631, The test test_partitions_above_the_cap_are_not_a_reassignment_target currently verifies only that drift is empty; capture the info log emitted during provisioner.ensure_provisioned_topics_exist() and assert it contains the partition_cap_suppressed label for TOPIC. Keep the existing empty-drift assertion to verify the divergence is excluded from the repair feed.
805-838: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value
_FailingAdminduplicates the whole_FakeAdminboundary. Acreate_topics_error: BaseException | Nonefield on_AdminRecorder(raised at the top of_FakeAdmin.create_topics) would let this test reuse_patched_adminand keep one fake broker to maintain.🤖 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/event_bus/test_topic_provisioner_rf_policy_omn15395.py` around lines 805 - 838, Add a create_topics_error: BaseException | None field to _AdminRecorder, update _FakeAdmin.create_topics to raise it before normal behavior when set, and remove the duplicated _FailingAdmin class. Configure the existing _patched_admin fixture/context with _TransientBrokerError so the test reuses the shared fake broker boundary.
439-461: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueOffender labels drop the directory path.
f"{root.name}/{path.name}"collapsessrc/omnibase_infra/a/b/service_x.pytoomnibase_infra/service_x.py, so two same-named modules in different packages are indistinguishable in the failure message — the sibling guard usespath.relative_to(src_root). Also note the fixed 12-line window can attribute a literalreplication_factor=from adjacent code to the wrongNewTopic(site.♻️ Path-qualified label
- label = f"{root.name}/{path.name}:{index + 1}" + label = f"{root.name}/{path.relative_to(root)}:{index + 1}"Note:
test_create_topics_guard_sees_a_planted_third_pathexpectsscripts/planted_creator.py:2, which this keeps satisfied.🤖 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/event_bus/test_topic_provisioner_rf_policy_omn15395.py` around lines 439 - 461, Update the offender-label construction in the scan loop to use each file’s path relative to its corresponding source root, preserving the existing root-relative format such as scripts/planted_creator.py. Replace the fixed 12-line window used by _LITERAL_RF_RE with logic scoped to the current NewTopic construction so replication_factor literals from adjacent code are not attributed to the wrong site.src/omnibase_infra/topics/managed_staging_topic_checker.py (1)
266-267: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueNon-durability failures log no cause at all.
The warning names only the topic, so a transient create failure here is undiagnosable from the log.
♻️ Suggested tweak
- else: - logger.warning("Failed to create catalog topic %s", name) + else: + logger.warning( + "Failed to create catalog topic %s: %s", + name, + type(exc).__name__, + )🤖 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 `@src/omnibase_infra/topics/managed_staging_topic_checker.py` around lines 266 - 267, Update the warning in the topic-creation fallback branch to include the underlying exception or failure cause along with the topic name. Preserve the existing warning behavior while making failures handled by the else branch of the managed staging topic checker diagnostically actionable.
🤖 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 `@scripts/create_kafka_topics.py`:
- Around line 410-419: Load and apply ~/.omnibase/.env before calling
ModelTopicProvisioningPolicy.from_env() in the topic provisioning flow, ensuring
KAFKA_SASL_MECHANISM is available when resolving the policy profile. Preserve
the existing bind_policy_to_broker_count and logging behavior after the
environment is loaded.
In `@src/omnibase_infra/event_bus/service_topic_manager.py`:
- Around line 243-250: Serialize capacity probing in the policy measurement flow
using an asyncio.Lock initialized alongside _capacity_probed in __init__. Update
_measured_policy and the bind_policy_to_broker_capacity call so concurrent
callers wait for the in-flight probe, then reuse its measured or UNMEASURABLE
result instead of returning the policy before probing completes; preserve the
single-attempt memoization and coroutine-safety behavior.
- Around line 151-158: The TopicReplicationPolicyError construction in the topic
provisioning failure path must sanitize the broker exception before exposing it
in the message. Replace the raw cause representation in this return with
sanitize_error_message applied to the exception, matching the module’s existing
error-boundary usage while preserving the surrounding context.
In `@src/omnibase_infra/topics/model_topic_provisioning_policy.py`:
- Around line 534-538: Update the aggregate TopicReplicationPolicyError raise in
the batch provisioning path to attach a ModelInfraErrorContext created via
ModelInfraErrorContext.with_correlation(...), preserving the correlation id,
transport, and profile extras already available from the per-spec _violation
context. Keep the existing violation message and fail-closed behavior unchanged.
---
Nitpick comments:
In `@src/omnibase_infra/topics/managed_staging_topic_checker.py`:
- Around line 266-267: Update the warning in the topic-creation fallback branch
to include the underlying exception or failure cause along with the topic name.
Preserve the existing warning behavior while making failures handled by the else
branch of the managed staging topic checker diagnostically actionable.
In `@tests/unit/event_bus/test_topic_provisioner_rf_policy_omn15395.py`:
- Around line 1604-1631: The test
test_partitions_above_the_cap_are_not_a_reassignment_target currently verifies
only that drift is empty; capture the info log emitted during
provisioner.ensure_provisioned_topics_exist() and assert it contains the
partition_cap_suppressed label for TOPIC. Keep the existing empty-drift
assertion to verify the divergence is excluded from the repair feed.
- Around line 805-838: Add a create_topics_error: BaseException | None field to
_AdminRecorder, update _FakeAdmin.create_topics to raise it before normal
behavior when set, and remove the duplicated _FailingAdmin class. Configure the
existing _patched_admin fixture/context with _TransientBrokerError so the test
reuses the shared fake broker boundary.
- Around line 439-461: Update the offender-label construction in the scan loop
to use each file’s path relative to its corresponding source root, preserving
the existing root-relative format such as scripts/planted_creator.py. Replace
the fixed 12-line window used by _LITERAL_RF_RE with logic scoped to the current
NewTopic construction so replication_factor literals from adjacent code are not
attributed to the wrong site.
In `@tests/unit/scripts/test_create_kafka_topics.py`:
- Around line 374-384: Update
test_undeclared_replication_resolves_to_the_managed_floor_not_one to assert
against the imported MANAGED_MINIMUM_REPLICATION_FACTOR constant instead of the
literal 2, ensuring the test follows the managed replication policy.
🪄 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: dc214c76-0480-4944-b317-8f76b1a6908e
📒 Files selected for processing (7)
scripts/create_kafka_topics.pysrc/omnibase_infra/event_bus/service_topic_manager.pysrc/omnibase_infra/topics/broker_capacity_probe.pysrc/omnibase_infra/topics/managed_staging_topic_checker.pysrc/omnibase_infra/topics/model_topic_provisioning_policy.pytests/unit/event_bus/test_topic_provisioner_rf_policy_omn15395.pytests/unit/scripts/test_create_kafka_topics.py
| policy = bind_policy_to_broker_count( | ||
| ModelTopicProvisioningPolicy.from_env(), | ||
| broker_count_from_cluster_metadata(cluster_metadata), | ||
| ) | ||
| print( | ||
| f"Replication policy: profile={policy.profile.value} " | ||
| f"floor={policy.minimum_replication_factor} " | ||
| f"measured_brokers={policy.broker_count} " | ||
| f"ceiling={policy.capacity_replication_factor}" | ||
| ) |
There was a problem hiding this comment.
🔒 Security & Privacy | 🟠 Major | ⚡ Quick win
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Does this script (or a shared bootstrap it imports) load ~/.omnibase/.env?
sed -n '1,145p' scripts/create_kafka_topics.py
rg -n --type=py -C2 'omnibase/\.env|load_dotenv|dotenv' scripts/ src/omnibase_infra/event_bus/models/config/ | head -50
# How does the config derive sasl_mechanism / bootstrap servers from env?
ast-grep outline src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py --items allRepository: OmniNode-ai/omnibase_infra
Length of output: 11089
🏁 Script executed:
#!/bin/bash
set -euo pipefail
echo "== create_kafka_topics.py: env/bootstrap handling =="
rg -n -C 3 'load_dotenv|dotenv|omnibase/.env|~/.omnibase/.env|sasl_mechanism|bootstrap_servers|from_env\(' scripts/create_kafka_topics.py src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.py src/omnibase_infra/event_bus/models/config -g '!**/__pycache__/**'
echo
echo "== create_kafka_topics.py: CLI entry and policy usage =="
sed -n '1,220p' scripts/create_kafka_topics.py
echo
echo "== model_kafka_event_bus_config.py: env resolution details =="
sed -n '176,420p' src/omnibase_infra/event_bus/models/config/model_kafka_event_bus_config.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 50384
🏁 Script executed:
#!/bin/bash
set -euo pipefail
echo "== policy model: from_env and profile selection =="
fd -a 'model_topic_provisioning_policy.py' src scripts
echo
sed -n '1,260p' src/omnibase_infra/topics/model_topic_provisioning_policy.py
echo
echo "== script prologue for env loading =="
sed -n '1,120p' scripts/create_kafka_topics.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 18605
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Narrow search for env loading or explicit policy validation in the topic-creation script.
rg -n -C 2 'load_dotenv|dotenv|~/.omnibase/.env|from_env\(|sasl_mechanism|profile' scripts/create_kafka_topics.py src/omnibase_infra/topics -g '!**/__pycache__/**' | head -200Repository: OmniNode-ai/omnibase_infra
Length of output: 17059
🏁 Script executed:
#!/bin/bash
set -euo pipefail
echo "== locate policy file =="
fd -a 'model_topic_provisioning_policy.py' src
echo
echo "== inspect policy definition and from_env =="
sed -n '1,260p' src/omnibase_infra/topics/model_topic_provisioning_policy.py
echo
echo "== inspect the part of create_kafka_topics.py that binds the policy =="
sed -n '395,430p' scripts/create_kafka_topics.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 14889
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Read the policy model directly; keep output narrow.
fd -a 'model_topic_provisioning_policy.py' src
sed -n '1,240p' src/omnibase_infra/topics/model_topic_provisioning_policy.py
# Show the exact call site in the script.
sed -n '406,418p' scripts/create_kafka_topics.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 12960
Load ~/.omnibase/.env before deriving the replication policy
ModelTopicProvisioningPolicy.from_env() keys off KAFKA_SASL_MECHANISM, but this CLI never sources the operator env. Running it from a clean shell can classify an MSK broker as self_hosted, bypassing the RF1 fail-closed gate. Source the env here or assert the resolved profile matches the target before creating topics.
🤖 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 `@scripts/create_kafka_topics.py` around lines 410 - 419, Load and apply
~/.omnibase/.env before calling ModelTopicProvisioningPolicy.from_env() in the
topic provisioning flow, ensuring KAFKA_SASL_MECHANISM is available when
resolving the policy profile. Preserve the existing bind_policy_to_broker_count
and logging behavior after the environment is loaded.
Source: Coding guidelines
| return TopicReplicationPolicyError( | ||
| f"Broker refused to create topic {topic!r} with " | ||
| f"replication_factor={requested_replication_factor}: " | ||
| f"INVALID_REPLICATION_FACTOR. {measured}. This is a hard provisioning " | ||
| "failure, not a best-effort miss — the topic does NOT exist. Fix the " | ||
| f"owning contract's topic_config.replication_factor, or run against a " | ||
| "cluster whose describe_cluster is reachable so the measured capacity " | ||
| f"ceiling can reduce it (OMN-15395). Broker error: {cause!r}", |
There was a problem hiding this comment.
🔒 Security & Privacy | 🟡 Minor | ⚡ Quick win
Route the broker exception through sanitize_error_message instead of embedding {cause!r}.
Every other boundary in this module logs broker exceptions via sanitize_error_message(e); this message interpolates the raw repr into a user/operator-visible error string.
🛡️ Proposed fix
- f"ceiling can reduce it (OMN-15395). Broker error: {cause!r}",
+ f"ceiling can reduce it (OMN-15395). Broker error: "
+ f"{type(cause).__name__}: {sanitize_error_message(cause)}",As per coding guidelines: "Error messages must never expose passwords, API keys, PII, or credential-bearing connection strings; use the repository sanitization utilities."
📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| return TopicReplicationPolicyError( | |
| f"Broker refused to create topic {topic!r} with " | |
| f"replication_factor={requested_replication_factor}: " | |
| f"INVALID_REPLICATION_FACTOR. {measured}. This is a hard provisioning " | |
| "failure, not a best-effort miss — the topic does NOT exist. Fix the " | |
| f"owning contract's topic_config.replication_factor, or run against a " | |
| "cluster whose describe_cluster is reachable so the measured capacity " | |
| f"ceiling can reduce it (OMN-15395). Broker error: {cause!r}", | |
| return TopicReplicationPolicyError( | |
| f"Broker refused to create topic {topic!r} with " | |
| f"replication_factor={requested_replication_factor}: " | |
| f"INVALID_REPLICATION_FACTOR. {measured}. This is a hard provisioning " | |
| "failure, not a best-effort miss — the topic does NOT exist. Fix the " | |
| f"owning contract's topic_config.replication_factor, or run against a " | |
| f"cluster whose describe_cluster is reachable so the measured capacity " | |
| f"ceiling can reduce it (OMN-15395). Broker error: " | |
| f"{type(cause).__name__}: {sanitize_error_message(cause)}", |
🤖 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 `@src/omnibase_infra/event_bus/service_topic_manager.py` around lines 151 -
158, The TopicReplicationPolicyError construction in the topic provisioning
failure path must sanitize the broker exception before exposing it in the
message. Replace the raw cause representation in this return with
sanitize_error_message applied to the exception, matching the module’s existing
error-boundary usage while preserving the surrounding context.
Source: Coding guidelines
| # OMN-15395 (D4): memoize the capacity probe ATTEMPT, not merely a | ||
| # successful one. Keying the "already probed?" test on | ||
| # ``policy.broker_count is None`` re-probed an UNMEASURABLE cluster on | ||
| # every entrypoint — three entrypoints, three describe_cluster round | ||
| # trips, none of which could ever succeed — which is the per-call | ||
| # fan-out (d) exists to eliminate, reintroduced on the failure path. A | ||
| # policy supplied already-measured counts as probed. | ||
| self._capacity_probed = self._policy.broker_count is not None |
There was a problem hiding this comment.
🩺 Stability & Availability | 🟠 Major | ⚡ Quick win
Sentinel set before the await makes a concurrent second caller skip measurement permanently.
_capacity_probed = True is assigned before bind_policy_to_broker_capacity completes, so a coroutine entering _measured_policy while the first probe is in flight returns the still-unmeasured policy and proceeds to CreateTopics with no capacity ceiling. The class docstring advertises coroutine safety, and with D5 in place that unreduced RF now aborts fail-closed rather than degrading. The old broker_count is None test re-probed in that window; the sentinel removes the retry.
🔒️ Proposed fix: serialize the single attempt
- if self._capacity_probed:
- return self._policy
- self._capacity_probed = True
- self._policy = await bind_policy_to_broker_capacity(admin, self._policy)
- return self._policy
+ async with self._capacity_probe_lock:
+ if self._capacity_probed:
+ return self._policy
+ self._policy = await bind_policy_to_broker_capacity(admin, self._policy)
+ self._capacity_probed = True
+ return self._policywith the lock created in __init__ alongside the sentinel:
self._capacity_probe_lock = asyncio.Lock()Also applies to: 298-302
🤖 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 `@src/omnibase_infra/event_bus/service_topic_manager.py` around lines 243 -
250, Serialize capacity probing in the policy measurement flow using an
asyncio.Lock initialized alongside _capacity_probed in __init__. Update
_measured_policy and the bind_policy_to_broker_capacity call so concurrent
callers wait for the in-flight probe, then reuse its measured or UNMEASURABLE
result instead of returning the policy before probing completes; preserve the
single-attempt memoization and coroutine-safety behavior.
| raise TopicReplicationPolicyError( | ||
| f"Refusing to provision {len(violations)} topic(s) under the " | ||
| f"{policy.profile.value} replication policy; no CreateTopics was " | ||
| "issued. Violations: " + " | ".join(shown) + suffix | ||
| ) |
There was a problem hiding this comment.
📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win
Attach ModelInfraErrorContext.with_correlation(...) to the aggregate error.
The per-spec path (_violation, Line 470) builds a correlation-bearing context, but this batch-level raise drops it, so the fail-closed error that actually reaches operators/CLI carries no correlation id, transport, or profile extras.
🛡️ Proposed fix
raise TopicReplicationPolicyError(
f"Refusing to provision {len(violations)} topic(s) under the "
f"{policy.profile.value} replication policy; no CreateTopics was "
- "issued. Violations: " + " | ".join(shown) + suffix
+ "issued. Violations: " + " | ".join(shown) + suffix,
+ context=ModelInfraErrorContext.with_correlation(
+ transport_type=EnumInfraTransportType.KAFKA,
+ operation="resolve_topic_replication_factor",
+ ),
+ profile=policy.profile.value,
+ minimum_replication_factor=policy.minimum_replication_factor,
)Based on learnings and coding guidelines: "Use the narrowest matching infrastructure error class and create error context with ModelInfraErrorContext.with_correlation(...)" — and this is not a pre-request-context startup guard, so the guideline applies.
📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| raise TopicReplicationPolicyError( | |
| f"Refusing to provision {len(violations)} topic(s) under the " | |
| f"{policy.profile.value} replication policy; no CreateTopics was " | |
| "issued. Violations: " + " | ".join(shown) + suffix | |
| ) | |
| raise TopicReplicationPolicyError( | |
| f"Refusing to provision {len(violations)} topic(s) under the " | |
| f"{policy.profile.value} replication policy; no CreateTopics was " | |
| "issued. Violations: " + " | ".join(shown) + suffix, | |
| context=ModelInfraErrorContext.with_correlation( | |
| transport_type=EnumInfraTransportType.KAFKA, | |
| operation="resolve_topic_replication_factor", | |
| ), | |
| profile=policy.profile.value, | |
| minimum_replication_factor=policy.minimum_replication_factor, | |
| ) |
🤖 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 `@src/omnibase_infra/topics/model_topic_provisioning_policy.py` around lines
534 - 538, Update the aggregate TopicReplicationPolicyError raise in the batch
provisioning path to attach a ModelInfraErrorContext created via
ModelInfraErrorContext.with_correlation(...), preserving the correlation id,
transport, and profile extras already available from the per-spec _violation
context. Keep the existing violation message and fail-closed behavior unchanged.
Source: Coding guidelines
#5534) * evidence: OCC companion pass 1 for OmniNode-ai/omnibase_infra#2552 * evidence: OCC companion self-bind for #5534 --------- Co-authored-by: node-occ-companion-effect <occ-companion-effect@omninode.ai>
…ot the module (#2554) The static guard shipped in #2552 computed admissibility once per FILE: resolves = ("ModelTopicProvisioningPolicy" in text and _POLICY_RESOLVER_RE.search(text)) so every NewTopic(...) in any module that mentions the policy anywhere got a blanket pass unless its RF was an integer literal. Executed against the reconstructed b2ca4fa (#2543) tree it reported only the operator CLI and returned NOTHING for service_topic_manager.py:758's `replication_factor=config.replication_factor` -- that lineage's own defect -- because the module mentions the policy five times elsewhere. That is the same "the guard certified a property it could not see" failure the guard was introduced to remediate, reintroduced in its replacement. Admissibility is now computed from the NewTopic call's own argument expression by AST provenance. A replication factor is admissible only when it traces, via provenance-preserving operations (attribute, subscript, iteration, method call on a resolved receiver, container literal/comprehension, single-arg builtin repackaging), back to a resolve_spec / resolve_specs_for_creation / resolve_replication_factor call. A module-local wrapper qualifies by its BODY (every return resolved), never by its name, so a stub `_resolve_spec` that returns its argument confers nothing. Name lookup is line-ordered per lexical scope, which is what keeps managed_staging_topic_checker admissible: it binds `spec` twice, once from a raw walrus and once from the resolved mapping, and only the nearest preceding binding counts. Also now refused rather than waved through: a literal RF in positional slot 3 (NewTopic('t', 6, 1)), an omitted RF, and a **kwargs splat. Analysis is deliberately conservative: provenance is not tracked through a mutated accumulator. That shape is reported, and the documented remedy is the batch helper resolve_specs_for_creation -- not a loosening. Pinned by a test so it stays a decision rather than a surprise. Test-only change; no src/ or scripts/ behaviour is modified. Ticket: OMN-15395 Evidence-Ticket: OMN-15395 Co-authored-by: t <t@t.t>
OMN-15395 — remediation round 4 (residuals of #2550)
Adversarial review of PR #2550 found four defects. #2550 merged mid-round (squash
cbc42e44, and the merge deleted its branch), so these land on a fresh branch offdev(e2050a20). Every fix below is against code that is already ondev.Headline: #2550 hardened one CreateTopics path. There were three, and the one an operator is actually told to run —
scripts/create_kafka_topics.py— still created every topic at a flat--replication-factordefault of 1.Defects fixed (each proven RED by reverting its own fix hunk)
scripts/create_kafka_topics.py:174-179,329-336hardcoded--replication-factordefault=1 and built a flatNewTopic(t, …, replication_factor=rf), discarding every contract's declaredtopic_config.replication_factor. NoModelTopicProvisioningPolicy, no measured capacity ceiling, no fail-closed RF1 rejection. This is the commanddocs/operations/README.md:41gives operators and the onecompare_environments.py:741names in itskafka_topic_parityfix_hint— so the documented runbook reproducedAWS_KAFKA_HIGH_RISK_CONFIG_RF_EQUALS_ONEby hand, with #2550's gate never consulted. Both existing static guards were blind: the swallow guard only watches calls into the provisioner, and it never scannedscripts/at all.ModelTopicProvisioningPolicy.from_env()→ measured ceiling → fail-closed batch resolution before the firstCreateTopics. Broker count is read off thelist_topics()response the diff already makes (zero extra round trips).--replication-factoris deleted; the lane capONEX_TOPIC_PROVISIONER_MAX_PARTITIONSnow applies here too, so the two paths cannot disagree about a topic's shape. Batch resolution extracted to the sharedresolve_specs_for_creation. New sibling static guard_raw_create_topics_offendersscanssrc/andscripts/for anyNewTopicconstruction that bypasses the policy or hardcodes a literal RF; the swallow guard now scansscripts/too.test_contract_declared_replication_factor_reaches_create_topics,test_managed_staging_rf1_is_rejected_before_any_create,test_replication_factor_is_no_longer_a_cli_flag,test_lane_partition_cap_applies_on_the_operator_path_too,test_every_create_topics_site_resolves_through_the_policy+ 2 planted-shape controls + 1 negative controlservice_topic_manager.py:387—_report_spec_driftresolved the RF but compared broker partitions against the UNCAPPEDspec.partitions, while creation applies_creation_partitions(the env cap, live at 1 on dev / stability-test / judge). The provisioner reported drift against topics it had itself just created, correctly — 159 boguspartition_mismatchentries per pass into the operator-gated WS-M reassignment feed._creation_spec— the resolved and capped effective spec, i.e. exactly what creation would produce. Divergence the cap explains (a topic created before the cap was lowered; Kafka cannot reduce partitions) is labelledpartition_cap_suppressedand kept out of the drift feed rather than handed to the repair lane as an impossible instruction.test_capped_lane_reports_no_partition_drift_for_a_6_partition_contract,test_partitions_above_the_cap_are_not_a_reassignment_target, +test_genuine_partition_drift_is_still_reported(negative control on an uncapped lane)service_topic_manager.py:224— the capacity memo keyed onpolicy.broker_count is None, which cannot bound an unmeasurable cluster: the field staysNoneforever, so the probe re-ran on every entrypoint. Measured 3describe_clustercalls across 3 entrypoints, every one guaranteed to fail — the per-call fan-out (d) exists to eliminate, reappearing on the error path._capacity_probedsentinel, pre-set when a caller supplies an already-measured policy). One probe per instance, success or failure.test_capacity_probe_is_attempted_once_even_when_unmeasurable(RED:assert 3 == 1)broker_capacity_probe.py:19-23and #2550's body both claimed an unhostable RF "fails loudly atCreateTopics". It did not. The broker'sINVALID_REPLICATION_FACTORlanded in the per-topicexcept Exception, became alogger.warning+ a name infailed, and the pass returnedstatus="partial"— indistinguishable from a transient blip, with the topic silently absent. A probe failure therefore left topics silently uncreated.is_invalid_replication_factor_error, matched on the wire error identity so it survives driver wrapping) and raised as a typedTopicReplicationPolicyErrornaming the topic, the refused value, and whether capacity was measurable — on both runtime paths. The operator CLI and the managed-staging checker emit it as a classified ERROR distinct from the best-effort bucket. The claim is now true.test_unmeasurable_cluster_does_not_guess_a_ceiling(fake admin that actually raisesINVALID_REPLICATION_FACTOR),test_unhostable_replication_aborts_the_batch_pass_too,test_unhostable_replication_factor_is_an_error_not_a_warning, +test_a_transient_create_failure_is_still_best_effort(negative control: startup stays best-effort for everything that is not a durability violation)Acceptance-criteria mapping — corrected
#2550's body asserted (a) and (b) as satisfied. At module scope that was true; at repository scope it was false, and the review proved it. Restated honestly:
DEFAULT_EVENT_TOPIC_REPLICATION_FACTORremains deleted"service_topic_manager, butscripts/create_kafka_topics.pycarried an equivalent flat RF1 default inargparseand applied it to every topic. Deleting the constant from one module while an argparse default did the same job in another is not "no flat default".replication_factor=at aNewTopicsite acrosssrc/+scripts/.test_replication_factor_is_no_longer_a_cli_flag;test_contract_declared_replication_factor_reaches_create_topics(RED under M-D2a); guard executed against the pre-fix file belowCreateTopicsTestManagedStagingRejectsRf1plus the newconfig=case"CreateTopicson MSK with zero policy resolution, so the documented reconcile command bypassed the gate entirely.resolve_specs_for_creationand aborts with zero creates issued.test_managed_staging_rf1_is_rejected_before_any_create— assertscreate_calls == 0(RED under M-D2b)_creation_spec, the same value creation applies.test_capped_lane_reports_no_partition_drift_for_a_6_partition_contractCreateTopics, one probe per instancetest_capacity_probe_is_attempted_once_even_when_unmeasurableONEX_BOOT_UNIVERSE_PROVISIONuntouchedRED-before evidence (mutation proof, executed)
Each fix hunk reverted in place, guarding test re-run:
M-D5c initially passed under mutation — the assertion
"ERROR" in stderrwas not discriminating, because the pre-existing "still missing after creation" block prints ERROR either way and the exception's own text already contains the stringINVALID_REPLICATION_FACTOR. The test was tightened to assert the classification (REFUSED by the broker) and the absence of the genericSome topics failed to createbucket, after which the mutation is caught. Recorded because a guard that passes under its own mutation is not a guard.The new static guard, executed against the actual pre-fix file (
git show <pre-fix>:scripts/create_kafka_topics.pyinto a scratch tree):Line 329 is exactly the
NewTopic(construction the review cited. The guard is not a claim that it would have caught the defect — it was run against the defect.D4 measured before/after:
describe_cluster_callsacross 3 entrypoints on an unmeasurable cluster —3before,1after.Seam definition — what changed at the boundary
topic_config)replication_factorNone= undeclared)resolve_specs_for_creation→ floor, then measured ceilingpolicy.resolve_specpartitionsmin(spec, ONEX_TOPIC_PROVISIONER_MAX_PARTITIONS)spec.partitionsdescribe_cluster(async), memoized per instancelist_topics().brokers(sync) — new, same binderbind_policy_to_broker_countdescribe_clusterTopicReplicationPolicyErrorpartitionsGates — run on
.200(stickybeatz-studio), patch-transfer verifiedLocal commit →
git format-patch→git amon the.200worktree. Tree hash identical on both sides (633077ab8dddc0c2e549c7e39efc328e8ab40604) and per-filesha256identical for all 7 changed files, so the gates ran on the same bytes that were pushed.ruff format --check src/ tests/ scripts/— 4790 files already formattedruff check src/ tests/ scripts/— All checks passedmypy src/omnibase_infra/— Success, no issues in 2632 source filesdetect_test_paths.py) →is_full_suite: true(shared_module), so the full suite ran with no-knarrowingpre-commit run --all-filestests/unit/event_bus tests/unit/topics tests/unit/scripts): 1754 passed23 failed, 27020 passed, 493 skipped, 1 xfailed, 42 errorsin 10m34sHonest state of the two non-green gates — neither is from this diff:
Full-suite failures/errors are environmental on
.200, and the count improved against fix(OMN-15395): measure broker capacity, never assume it — RF ceiling from describe_cluster #2550's own baseline on the same host:24 failed / 27017 passed / 492 skipped / 1 xfailed / 42 errorshere vs28 failed / 26994 passed / 42 errorson fix(OMN-15395): measure broker capacity, never assume it — RF ceiling from describe_cluster #2550. Every failure is integration (Postgres/Kafka/ledger — no broker on.200), performance (LLM endpoint SLO), or one of the four env-mutation/caplogunit files sensitive to xdist ordering (test_dependency_materializer,test_kafka_bootstrap_no_localhost_fallback,test_model_kafka_producer_config_max_request_size,test_mcp_server_lifecycle) — 99/99 passing in isolation on this branch, and none import anything this PR touches.pre-commit run --all-filesfails on three pre-existing offenders outside this diff, each verifiedUNCHANGED vs origin/devwithgit diff --quiet origin/dev...HEAD -- <path>:2026headers ontests/ci/test_runner_routing_audit.pyandtests/scripts/test_deploy_runtime_core_contracts_resolution.py;check-required-env-varswantsGITHUB_TOKENin.200's~/.omnibase/.env(machine env gap, not a repo defect);check-url-authorityreports localhost literals inhandler_service_validate.py,adapter_llm_provider_openai.py,seed-infisical.py,model_registry_request.py,registry_api/main.py,handler_event_forward.py— all six untouched by this branch.Staged-scope pre-commit passed clean at commit time (every hook
Passed). Not fixed here — different files, unrelated cause; flagged rather than swept into this diff.A first
--all-filesrun also reportedprod-promotion-lineage-guardandcheck-ai-slopas "files were modified by this hook" — that run overlapped the full pytest run on the same worktree, and the lineage guard's own output in the same log reads27 passed in 5.43s. Re-run uncontended it is the three above only, withgit status --shortempty and the tree hash still633077ab…. Recorded rather than quietly dropped, perfeedback_precommit_modified_by_hook_misattribution.No skip tokens, no
--no-verify, no-knarrowing.Ticket: OMN-15395
Evidence-Ticket: OMN-15395
Evidence-Source: OCC#5534
Summary by CodeRabbit