Skip to content

fix(OMN-15395): measure broker capacity, never assume it — RF ceiling from describe_cluster - #2550

Merged
jonahgabriel merged 1 commit into
devfrom
jonah/omn-15395-rf-durability-remediation
Jul 30, 2026
Merged

jonahgabriel merged 1 commit into
devfrom
jonah/omn-15395-rf-durability-remediation

Conversation

@jonahgabriel

@jonahgabriel jonahgabriel commented Jul 29, 2026 •

Copy link
Copy Markdown
Collaborator

OMN-15395 — remediation round 3

Fixes the six defects adversarial review raised against PR #2546 (merged). #2546's branch was already merged, so this lands on a fresh branch off dev (59bbd485).

Headline: the capacity ceiling that reduces a contract-declared replication factor was an assumption, not a measurement. It is now read from the live cluster.

Defects fixed (each proven RED by reverting its own fix hunk)

# Sev Defect Fix RED-before guard
1 HIGH self_hosted() set capacity_replication_factor = 1 unconditionally for every cluster whose sasl_mechanism != "AWS_MSK_IAM", and resolve_replication_factor silently reduced any declared RF to it. Broker count was probed nowhere in src/. ModelKafkaEventBusConfig also accepts PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512 / OAUTHBEARER, so any multi-broker cluster not on MSK IAM — including an MSK cluster fronted by SCRAM — had its contract-declared RF2/RF3 clamped to RF1: the exact AWS_KAFKA_HIGH_RISK_CONFIG_RF_EQUALS_ONE condition this ticket exists to eliminate, reintroduced by the mechanism meant to prevent it. New topics/broker_capacity_probe.py reads a live describe_cluster node count; ModelTopicProvisioningPolicy.with_broker_capacity() is now the only way a ceiling is installed. Config-derived policies are unmeasured (capacity=None, broker_count=None) and reduce nothing. Bound once per provisioner instance, before any spec resolves, on all three creation paths. TestCapacityCeilingIsMeasuredNotAssumed, TestSaslClusterIsNotAssumedSingleNode
2 MED Static swallow guard's regex await\s+_?\w*provisioner\w*\. could not span the . in self._provisioner., nor match any receiver not literally named "provisioner" — and that shape is already in the tree at handler_topic_migration_executor.py:127. The guard certified a property it could not see. Regex anchored on the method, receiver-agnostic; extracted to _provisioning_swallow_offenders(). test_guard_sees_every_receiver_shape (6 planted receiver shapes) + a re-raising negative control
3 MED The config= branch of ensure_topic_exists (finding-4 fix from #2546) shipped with zero coverage — reverting both halves left the impacted selection at 1252 passed. Tests added. TestSnapshotConfigCreationPath — kills both mutations: resolver output (not raw config.replication_factor) reaches NewTopic, and the recorded created_spec gates the readiness poll
4 MED _report_spec_drift built its expectation from the UNRESOLVED self._spec_by_name, so every single-node lane (local Redpanda, CI, .201 stability/prod/judge) reported the 11 RF2 topics as replication_mismatch on every pass — seeding the operator-gated WS-M reassignment queue with targets a one-node cluster cannot host. Drift resolves through the same policy the creation site uses. A spec the policy would refuse is reported, not raised (the fail-closed abort is scoped to creates). TestDriftIsReportedAgainstTheResolvedSpec + a genuine-under-replication negative control
5 LOW Capacity reduction logged at logger.info while the field docstring promised a warning. A durability downgrade below WARNING is invisible under normal log filtering. logger.warning, message now names the measured broker count. test_reduction_is_logged_at_warning_not_info
6 LOW node_topic_migration_executor_effect/contract.yaml still said managed "refuses" an undeclared RF; the managed profile has had a default since #2546. Text corrected to describe resolve-to-floor + refuse-below-floor + measured-capacity reduction. n/a (doc)

Seam definition — topic-spec fields, end to end

Field Contract (topic_config) ModelTopicSpec Policy NewTopic Readiness
replication_factor optional int ≥1 int | None (None = undeclared, never 1) resolve_replication_factor → floor check, then measured ceiling replication_factor=<resolved> _created_specs[topic].replication_factor
partitions optional int int (default 6) carried through untouched num_partitions=<spec ∧ env cap> …partitions
kafka_config optional map Mapping | None carried through untouched topic_configs required-key visibility
broker node count — — broker_count / capacity_replication_factor, from describe_cluster only bounds replication_factor —

Invariants: the ceiling may only ever reduce, never raise; it is never installed below the profile's durability floor (a measurement under the floor leaves it unset so resolution refuses rather than clamping); an unmeasurable cluster gets no ceiling, so a declared RF the broker cannot host fails loudly at CreateTopics instead of being quietly downgraded.

Acceptance-criteria mapping (unchanged from #2543/#2546, re-verified)

  • (a) explicit contract-driven RF, no module-level RF1 default — DEFAULT_EVENT_TOPIC_REPLICATION_FACTOR remains deleted; the only defaults are profile-scoped and ≥ the floor.
  • (b) RF1 in managed staging rejected fail-closed before any CreateTopics — TestManagedStagingRejectsRf1, plus the new config= case.
  • (c) spec passes through every creation and readiness path — now including the drift-report path, which was the remaining site not using the resolver's output.
  • (d) list/diff before CreateTopics — unchanged; the capacity probe adds one describe_cluster metadata request per provisioner instance (test_capacity_is_measured_once_per_provisioner pins it), not a per-topic call. It is a plain Metadata API, not the DescribeTopicDynamicConfiguration MSK IAM denies.
  • (e) ONEX_BOOT_UNIVERSE_PROVISION untouched.
  • (f) RED-before/GREEN-after against the real provisioner — see below.

RED-before evidence (mutation proof, executed)

Each fix hunk reverted in place, guarding test re-run:

[RED (correct)] M1-HIGH:  self_hosted() hardcodes the RF1 capacity ceiling again      2 failed, 6 passed
[RED (correct)] M2-MED:   static guard regex reverts to receiver-anchored pattern     4 failed, 2 passed
[RED (correct)] M3-MED:   config= path hands the RAW config RF to NewTopic            1 failed
[RED (correct)] M3b-MED:  config= path records no created spec for readiness          1 failed
[RED (correct)] M4-MED:   drift compares against the UNRESOLVED contract spec         1 failed
[RED (correct)] M5-LOW:   capacity reduction logged at INFO again                     1 failed
ALL MUTATIONS CAUGHT

Tests drive the real TopicProvisioner from a real contract YAML through the real ContractTopicExtractor to the real aiokafka.admin.NewTopic; the only substitution is the admin client (the network boundary), which now also serves describe_cluster.

Gates — run on .200 (stickybeatz-studio), patch-transfer verified

Local commit → git format-patch → git am on the .200 worktree; tree hash identical on both sides (d6479871018b71b90612c1028fec42920e39c072), so the gates ran on the same bytes that were pushed.

  • ruff format --check src/ tests/ — 4574 files already formatted
  • ruff check src/ tests/ — All checks passed
  • mypy src/omnibase_infra/ — Success, no issues in 2632 source files
  • governed selector (detect_test_paths.py) → is_full_suite: true (shared_module), so the full suite ran, no -k narrowing: 28 failed, 26994 passed, 492 skipped, 1 xfailed, 42 errors
  • pre-commit run --all-files

Honest state of the two non-green gates — neither is from this diff:

  1. Full-suite failures/errors are environmental on .200. Every one is integration (Postgres/Kafka/ledger — Unable connect to "localhost:19092"), performance (LLM endpoint SLO), or one of five unit files that are env-mutation/caplog tests sensitive to xdist ordering. Those five (test_dependency_materializer, test_kafka_bootstrap_no_localhost_fallback, test_model_kafka_producer_config_max_request_size, test_mcp_server_lifecycle, test_util_db_transaction) pass 127/127 in isolation on this branch, and none import anything this PR touches. A full-suite baseline on untouched origin/dev was launched on .200 for a like-for-like comparison.
  2. pre-commit run --all-files fails on two pre-existing offenders outside this diff: tests/ci/test_runner_routing_audit.py and tests/scripts/test_deploy_runtime_core_contracts_resolution.py carry SPDX-FileCopyrightText: 2026 on origin/dev itself (verified with git show origin/dev:<path>), which the SPDX hook rejects; and check-required-env-vars wants GITHUB_TOKEN in .200's ~/.omnibase/.env (a machine env gap, not a repo defect). Staged-scope pre-commit passed clean at commit time. Not fixed here — different files, unrelated cause; flagged rather than silently swept into this diff.

No skip tokens, no --no-verify, no -k narrowing.

Ticket: OMN-15395

Evidence-Ticket: OMN-15395
Evidence-Source: OCC#5528

… from describe_cluster

Remediation round on the RF-policy work. Six adversarial-review defects, each
proven RED by reverting its fix hunk and re-running the guarding test.

HIGH — the capacity ceiling was an assumption, not a measurement.
`self_hosted()` set `capacity_replication_factor = 1` unconditionally for every
cluster whose `sasl_mechanism != "AWS_MSK_IAM"`, and `resolve_replication_factor`
silently reduced any declared RF down to it. Broker count was probed nowhere in
`src/`. `ModelKafkaEventBusConfig` also accepts PLAIN / SCRAM-SHA-256 /
SCRAM-SHA-512 / OAUTHBEARER, so any multi-broker cluster not reached over MSK IAM
had its contract-declared RF2/RF3 clamped to RF1 — the exact
AWS_KAFKA_HIGH_RISK_CONFIG_RF_EQUALS_ONE condition this ticket exists to
eliminate, reintroduced by the mechanism meant to prevent it. The ceiling is now
installed only by `with_broker_capacity()` from a live `describe_cluster` node
count (new `topics/broker_capacity_probe.py`), bound once per provisioner
instance before any spec is resolved, on all three creation paths
(`ensure_provisioned_topics_exist`, `ensure_topic_exists`,
`create_missing_catalog_topics`). Unmeasurable cluster => no ceiling => a
declared RF fails loudly at CreateTopics rather than being quietly downgraded.
A measurement below the durability floor installs no ceiling at all, so
resolution refuses instead of clamping.

MEDIUM — the static swallow guard was blind to `self._provisioner.`. Its regex
`await\s+_?\w*provisioner\w*\.` could not span the dot, nor match any receiver
not literally named "provisioner" — and that shape is already in the tree at
`handler_topic_migration_executor.py:127`. Now anchored on the METHOD, with
parametrised meta-tests that plant each receiver shape and assert the guard
flags it (plus a re-raising negative control).

MEDIUM — the `config=` creation branch of `ensure_topic_exists` shipped with zero
coverage; reverting both halves of its fix left the impacted selection fully
green. Added tests that kill both mutations: resolver output (not the raw
`config.replication_factor`) reaches NewTopic, and the recorded `created_spec`
gates the readiness poll.

MEDIUM — `_report_spec_drift` built its expectation from the UNRESOLVED spec, so
every single-node lane reported the eleven RF2 topics as replication drift on
every pass and seeded the operator-gated WS-M reassignment queue with targets a
one-node cluster cannot host. Drift now resolves through the same policy the
creation site uses; a spec the policy would refuse is reported, not raised.

LOW — capacity reduction logged at INFO while the field docstring promised a
warning: a silent durability downgrade below WARNING is invisible under normal
log filtering. Now `logger.warning`, with a test that pins the level.

LOW — `node_topic_migration_executor_effect/contract.yaml` still said managed
"refuses" an undeclared replication_factor; the managed profile has had a
default since the previous round. Text corrected.

Ticket: OMN-15395
@coderabbitai

coderabbitai Bot commented Jul 29, 2026

Copy link
Copy Markdown

Warning

Review limit reached

You’ve reached a temporary PR review limit under our Fair Usage Limits Policy.

Your recent review volume is higher than typical usage, so adaptive limits are currently applied.

Next review available in: 38 minutes

Enable usage-based reviews in Billing to review now. Otherwise, wait until the next included review is available.
You're only billed for reviews past your plan's rate limits ($0.25/file).

How can I continue?

After more reviews become available, a review can be triggered using the @coderabbitai review command as a PR comment. Alternatively, push new commits to this PR.

To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews.

How do review limits work?

CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability.

For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window.

Please refer docs for additional details.

Review details
⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Pro

Run ID: e8a7751c-6031-42e3-9aec-06545f44d72f

📥 Commits

Reviewing files that changed from the base of the PR and between 59bbd48 and 6efc6b6.

📒 Files selected for processing (10)
  • src/omnibase_infra/event_bus/service_topic_manager.py
  • src/omnibase_infra/nodes/node_topic_migration_executor_effect/contract.yaml
  • src/omnibase_infra/topics/__init__.py
  • src/omnibase_infra/topics/broker_capacity_probe.py
  • src/omnibase_infra/topics/managed_staging_topic_checker.py
  • src/omnibase_infra/topics/model_topic_provisioning_policy.py
  • tests/unit/event_bus/test_topic_provisioner_rf_policy_omn15395.py
  • tests/unit/topics/test_broker_capacity_probe.py
  • tests/unit/topics/test_no_contract_declares_rf1.py
  • tests/unit/topics/test_topic_provisioning_policy.py

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

jonahgabriel pushed a commit to OmniNode-ai/onex_change_control that referenced this pull request Jul 30, 2026
jonahgabriel added a commit to OmniNode-ai/onex_change_control that referenced this pull request Jul 30, 2026
#5528)

* evidence: OCC companion pass 1 for OmniNode-ai/omnibase_infra#2550

* evidence: OCC companion self-bind for #5528

---------

Co-authored-by: node-occ-companion-effect <occ-companion-effect@omninode.ai>
@github-actions

Copy link
Copy Markdown
Contributor

⚠️ Hostile Reviewer — DEGRADED (informational)

Blocking findings (critical): 0
Total findings: 0
Models succeeded: none

Note: All reviewer models failed or were unavailable. Degraded results are informational during the pilot phase (OMN-8468/OMN-8524) and do not block merge. Error: all review endpoints [192.168.86.201:8000 192.168.86.201:8001 ] unreachable — preflight short-circuit (no models available)


Gate semantics (pilot phase)

Verdict Meaning Blocks merge?
passed No critical findings No
blocked CRITICAL findings found Yes
degraded All models unavailable (infra) No (pilot)

Powered by omniintelligence.review_pairing.cli_review — node-based adversarial review via HandlerLlmCliSubprocess (OMN-8468/OMN-8524)

@jonahgabriel
jonahgabriel merged commit cbc42e4 into dev Jul 30, 2026
152 of 159 checks passed
@jonahgabriel
jonahgabriel deleted the jonah/omn-15395-rf-durability-remediation branch July 30, 2026 00:48
jonahgabriel added a commit that referenced this pull request Jul 30, 2026
…moized probe, loud RF refusal (#2552)

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

Co-authored-by: t <t@t.t>
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.

2 participants