Skip to content

Archive legacy codebase for clean migration - #10

Merged
jonahgabriel merged 1 commit into
mainfrom
feature/archive-separation
Sep 18, 2025
Merged

jonahgabriel merged 1 commit into
mainfrom
feature/archive-separation

Conversation

@jonahgabriel

Copy link
Copy Markdown
Collaborator

Archive Legacy Codebase

This PR archives all legacy ONEX infrastructure components to enable clean domain-specific PR reviews.

What's Being Archived (286 files)

  • Complete legacy source code from src_archived/
  • Infrastructure nodes: postgres, consul, kafka adapters and orchestrators
  • Models: health, infrastructure, observability, security, workflow
  • Tests: integration, unit, load testing, and e2e
  • Scripts: validation, migration, and intelligence hooks
  • Documentation: templates, guides, and architectural decisions
  • Deployment configs: Docker, database migrations, infrastructure

Why This Approach

  • Clean Review Process: Separates code preservation from new implementation
  • Manageable PR Size: Reduces main PR from 411 files to ~125 implementation files
  • Clear Intent: Archive shows what's being replaced vs. what's being added
  • Historical Preservation: All legacy code remains accessible in archive/

Next Steps After Merge

  1. Create domain-specific PRs for new implementation:
    • Core health models (protocol-based architecture)
    • Infrastructure enums and shared components
    • Testing infrastructure and validation
  2. Each domain PR will be focused and reviewable (~20-30 files)

Legacy Components Archived

  • Service Adapters: Consul, Kafka, PostgreSQL, Vault integrations
  • Processing Nodes: Compute, aggregation, wrapper components
  • Architecture: Reducers, orchestrators, workflow coordination
  • Monitoring: Health monitors, observability, metrics collection
  • Security: Audit logging, encryption, rate limiting, TLS
  • Testing: Comprehensive test suites for all components

This archival enables Phase 2 infrastructure migration to contract-driven ONEX architecture with strong typing and zero backwards compatibility.

Archive all legacy ONEX infrastructure components including:
- Complete source code from src_archived/ (286 files)
- Infrastructure tools, models, and configurations
- Legacy tests, scripts, and documentation
- Database migrations and deployment configurations

This archival enables clean domain-specific PR reviews by
separating legacy code preservation from new implementation.

Legacy components archived:
- Infrastructure nodes (postgres, consul, kafka adapters)
- Models (health, infrastructure, observability, security)
- Testing frameworks and integration tests
- Deployment and migration scripts
- Documentation and templates

Archival supports Phase 2 infrastructure migration to
contract-driven ONEX architecture with strong typing.
@github-actions

Copy link
Copy Markdown
Contributor

🔍 Pull Request Review for Archive Separation

I've completed a comprehensive review of PR #10 - Archive legacy codebase for clean migration.

✅ Overall Assessment: APPROVED with Recommendations

This is a well-structured archival PR that properly preserves legacy code while enabling clean migration to the ONEX contract-driven architecture. The approach aligns perfectly with the ZERO BACKWARDS COMPATIBILITY policy outlined in CLAUDE.md.


🎯 Strengths

1. Clean Separation Strategy

  • ✅ Excellent approach to archive 286 legacy files rather than deleting them
  • ✅ Preserves historical context while enabling fresh implementation
  • ✅ Makes domain-specific PRs more manageable (~20-30 files each)

2. Architectural Alignment

  • ✅ Follows ONEX infrastructure migration plan from CLAUDE.md
  • ✅ Supports transition to contract-driven node architecture
  • ✅ Enables strong typing enforcement (no Any types)

3. Code Organization

  • ✅ Proper archival structure with clear categorization
  • ✅ Maintains test suites for reference during migration
  • ✅ Preserves documentation and templates for future use

🔒 Security Observations

Positive Security Practices

  • ✅ No hardcoded credentials found in sampled files
  • ✅ Proper use of environment variables and Docker secrets
  • ✅ Security-first approach with credential manager fallback pattern

Example from postgres_connection_manager.py:79-100:

  • Correctly reads credentials from environment or Docker secrets file
  • Proper error handling with OnexError chaining
  • No default passwords or connection strings

⚠️ Areas Requiring Attention

1. Import Path Inconsistencies

Some files show partial migration to omnibase_core:

  • Line 19 in postgres_connection_manager.py: Correctly uses omnibase_core
  • However, relative imports on line 54 suggest incomplete migration
  • Recommendation: Since this is archive, these inconsistencies are acceptable but should be noted for new implementation

2. Pydantic V2 Migration

Observed use of legacy Pydantic patterns:

  • Config class instead of model_config
  • Field() patterns that may need updating
  • Impact: Not critical for archive, but new implementations must use Pydantic V2

3. Missing Shared Model Pattern

Current structure doesn't fully implement the shared model pattern from CLAUDE.md:

  • Models are node-specific rather than shared across infrastructure
  • Next Steps: Phase 2 implementation should extract shared models to src/omnibase_infra/models/{service}/

📊 Performance Considerations

Connection Pooling Configuration

The ConnectionConfig shows good defaults:

  • Min connections: 5
  • Max connections: 50
  • Inactive connection lifetime: 300s
  • Recommendation: Document these settings for new implementation reference

🚀 Recommendations for Next Steps

1. Phase 1: PostgreSQL Adapter Node (Per CLAUDE.md)

Create message bus bridge pattern implementation:

  • Contract-first approach with contract.yaml
  • Shared models in src/omnibase_infra/models/postgres/
  • Node-specific adapter models only

2. Domain-Specific PR Strategy

Suggested PR breakdown for Phase 2:

  1. Core health models with protocol-based architecture
  2. Infrastructure enums and shared components
  3. Service adapters (Consul, Kafka, Vault)
  4. Processing nodes and orchestrators
  5. Testing infrastructure

3. Migration Validation Checklist

For each new implementation PR:

  • Contract-driven with complete contract.yaml
  • Strong typing (zero Any usage)
  • Proper OnexError chaining with CoreErrorCode
  • Shared model pattern implementation
  • Container injection pattern (ONEXContainer)

💡 Best Practices Observed

  1. Error Handling: Proper use of OnexError with chaining
  2. Documentation: Comprehensive docstrings and type hints
  3. Security: No hardcoded secrets, proper credential management
  4. Testing: Preserved test suites for migration reference

✅ Conclusion

This archival PR successfully preserves legacy infrastructure while setting up for clean contract-driven migration. The approach reduces review complexity and maintains historical context.

Approval Recommendation: Merge this PR to unblock domain-specific implementation work.


Review conducted following ONEX Infrastructure standards and ZERO BACKWARDS COMPATIBILITY policy

@jonahgabriel
jonahgabriel merged commit b1062d2 into main Sep 18, 2025
1 of 2 checks passed
@jonahgabriel
jonahgabriel deleted the feature/archive-separation branch September 18, 2025 14:21
jonahgabriel added a commit that referenced this pull request Apr 17, 2026
…-9034]

Extracts audit logic from inline python3 HEREDOCs in the shell script
into a testable Python lib so Check A / Check B / fix-payload can be
exercised with dependency injection instead of bash-subprocess mocking
that never worked.

Thread-by-thread:
- #1-5 (CodeQL unused locals): removed. The old tests created variables
  like `protection`, `commits_data`, `check_runs_data` and never asserted
  on them. New tests assert on audit_repo() return values directly.
- #6 (cross-repo PAT): workflow now uses
  `secrets.CROSS_REPO_PAT || secrets.GITHUB_TOKEN` (matches env-parity.yml
  pattern) + preflight check step with ::warning:: when absent. Without
  the PAT, 9 sibling repos will [SKIP] — documented in workflow header.
- #7 (pagination per_page=50): lib.PAGE_SIZE = 100 (GitHub API max).
  collect_seen_check_run_names now paginates until empty or short page.
- #8 (mock doesn't intercept bash): audit logic lives in
  scripts/audit_branch_protection_lib.py with a GhCaller injection seam.
  Tests import the lib and pass fake `gh` callables — no subprocesses
  at unit-test time.
- #9 (hardcoded /Volumes in test_rac_violation_detected): entire test
  removed as part of rewrite; no more subprocess.run + cwd=...
- #10 (smoke-test returncode in (0,1)): new tests assert on explicit
  status/rac/orphan_contexts/message fields, not returncodes.

Lib surface:
  parse_required_approving_review_count(protection_json) -> int
  parse_required_contexts(protection_json) -> list[str]
  build_fix_payload(protection_json) -> dict
  collect_seen_check_run_names(owner, repo, commits, gh) -> set[str]
  find_orphan_contexts(required, seen) -> list[str]
  audit_repo(owner, repo, gh, commits_to_scan=5) -> dict

Shell script calls scripts/audit_branch_protection_lib_cli.py for the
audit step and the --fix payload construction; the `gh api PUT` side
effect stays in bash.

Verification:
  uv run pytest tests/ci/test_branch_protection_audit.py -v
  = 19 passed in 0.19s
  shellcheck scripts/audit-branch-protection.sh = clean
  bash -n scripts/audit-branch-protection.sh = syntax ok
  uv run mypy scripts/audit_branch_protection_lib*.py = Success
  CI-matching pytest (split 1/15, -m "not slow and not chaos and not kafka")
  = 1346 passed, 2 env-dependent Postgres failures (no local Postgres)
github-merge-queue Bot pushed a commit that referenced this pull request Apr 17, 2026
* fix(ci): branch-protection-audit gate (OMN-9034)

Adds periodic CI audit of branch protection settings across all OmniNode-ai
repos. Catches two invariants that caused overnight failures: (A) non-zero
required_approving_review_count that blocks the solo-dev merge workflow, and
(B) orphaned required status check contexts that no CI job ever satisfies.

- scripts/audit-branch-protection.sh — shellcheck-clean, MIT SPDX, --dry-run
  default, --fix mode for automated remediation
- tests/ci/test_branch_protection_audit.py — 11 unit tests (pytest.mark.unit)
  covering clean/rac-violation/orphan-context/fix-mutation cases
- .github/workflows/branch-protection-audit.yml — schedule 23 */4 * * * +
  workflow_dispatch; fails workflow on any violation (report-only, no --fix)
- CLAUDE.md: ## Branch protection section documenting dry-run gate rule

* fix(tests): remove hardcoded /Volumes path in test_clean_repo [OMN-9034]

CI Split 1/15 failed with FileNotFoundError on
'/Volumes/PRO-G40/Code/omni_worktrees/OMN-BP-AUDIT/omnibase_infra'
because the prior commit baked the author's local worktree path into
the test's subprocess cwd.

Fix: resolve script + cwd relative to the test file via
Path(__file__).resolve().parents[2], matching the pattern required by
CLAUDE.md Rule 6 (no hardcoded absolute paths).

Verified locally: uv run pytest tests/ci/test_branch_protection_audit.py
= 11 passed in 6.01s.

* fix(ci): resolve 10 CR/CodeQL threads on branch-protection-audit [OMN-9034]

Extracts audit logic from inline python3 HEREDOCs in the shell script
into a testable Python lib so Check A / Check B / fix-payload can be
exercised with dependency injection instead of bash-subprocess mocking
that never worked.

Thread-by-thread:
- #1-5 (CodeQL unused locals): removed. The old tests created variables
  like `protection`, `commits_data`, `check_runs_data` and never asserted
  on them. New tests assert on audit_repo() return values directly.
- #6 (cross-repo PAT): workflow now uses
  `secrets.CROSS_REPO_PAT || secrets.GITHUB_TOKEN` (matches env-parity.yml
  pattern) + preflight check step with ::warning:: when absent. Without
  the PAT, 9 sibling repos will [SKIP] — documented in workflow header.
- #7 (pagination per_page=50): lib.PAGE_SIZE = 100 (GitHub API max).
  collect_seen_check_run_names now paginates until empty or short page.
- #8 (mock doesn't intercept bash): audit logic lives in
  scripts/audit_branch_protection_lib.py with a GhCaller injection seam.
  Tests import the lib and pass fake `gh` callables — no subprocesses
  at unit-test time.
- #9 (hardcoded /Volumes in test_rac_violation_detected): entire test
  removed as part of rewrite; no more subprocess.run + cwd=...
- #10 (smoke-test returncode in (0,1)): new tests assert on explicit
  status/rac/orphan_contexts/message fields, not returncodes.

Lib surface:
  parse_required_approving_review_count(protection_json) -> int
  parse_required_contexts(protection_json) -> list[str]
  build_fix_payload(protection_json) -> dict
  collect_seen_check_run_names(owner, repo, commits, gh) -> set[str]
  find_orphan_contexts(required, seen) -> list[str]
  audit_repo(owner, repo, gh, commits_to_scan=5) -> dict

Shell script calls scripts/audit_branch_protection_lib_cli.py for the
audit step and the --fix payload construction; the `gh api PUT` side
effect stays in bash.

Verification:
  uv run pytest tests/ci/test_branch_protection_audit.py -v
  = 19 passed in 0.19s
  shellcheck scripts/audit-branch-protection.sh = clean
  bash -n scripts/audit-branch-protection.sh = syntax ok
  uv run mypy scripts/audit_branch_protection_lib*.py = Success
  CI-matching pytest (split 1/15, -m "not slow and not chaos and not kafka")
  = 1346 passed, 2 env-dependent Postgres failures (no local Postgres)

---------

Co-authored-by: jonahgabriel <jonahgabriel@users.noreply.github.com>
jonahgabriel added a commit that referenced this pull request Aug 8, 2026
…ateway forwarder

The gateway forwarder's delivery-loop supervision previously raced
delivery.wait() against shutdown_event.wait() with no retry: any
cloud-leg delivery failure (broker drop, transient network fault)
propagated straight out of run_gateway_forwarder and terminated the
process (docs/design/2026-08-08-gateway-node-architecture-lift.md
axis #10, live restart-count evidence in that doc's 1.2).

- runtime/gateway_forwarder.py: new _supervise_gateway_delivery loop
  wraps the existing delivery.wait()/shutdown race. A delivery
  failure now retries with bounded exponential backoff + jitter
  instead of exiting. Reconnect state only resets after the delivery
  loop survives a full heartbeat_interval_seconds "recovery confirm"
  window without failing again (a bare delivery.start() succeeding
  only proves tasks were scheduled, not that the cloud leg is back).
  Process still exits on shutdown or on an unrecoverable
  (non-connectivity) failure -- config/heartbeat-task errors.
- Once a failure window crosses the contract-declared
  degraded_after_seconds, one DEGRADED status event publishes via
  the new ServiceGatewayForwarder.publish_status(), reusing the
  existing gateway-heartbeat.v1 event shape/topic
  (ModelGatewayHeartbeat gains status="degraded" +
  consecutive_failures + detail). Published on the LOCAL bus (not
  cloud) so DEGRADED stays observable while the cloud leg that
  caused it is itself unreachable.
- contract.yaml: liveness block gains
  reconnect_backoff_initial_seconds / _max_seconds / _jitter_seconds
  and degraded_after_seconds (contract_version/node_version patch
  bump 0.1.0 -> 0.1.1); ModelGatewayForwarderConfig gains the
  matching typed fields + a max>=initial validator.
- Data-plane loop (NodeGatewayDelivery, KafkaTransport) is untouched;
  supervision wraps it per scope.

Also vendors 2 omnimarket node migrations
(node_canary_score_reducer/0003, node_projection_registration/0004)
via scripts/sync-node-migrations.sh -- required by the repo-wide
node-migration-sync pre-commit/CI hook on any omnibase_infra branch.
infra#2683 (OMN-15732, split-migration rework) was NOT merged at
push time; these 2 files duplicate that in-flight resolution and
this PR must be rebased onto dev (picking up #2683's shape) once it
lands. No _LEGACY_DEFAULT_SCHEMA_SQL_EXACT_PATHS or exemption/
validator logic touched.

Ticket: OMN-15742 (G2, gateway-lift Phase 0)
jonahgabriel added a commit that referenced this pull request Aug 9, 2026
…ateway forwarder

The gateway forwarder's delivery-loop supervision previously raced
delivery.wait() against shutdown_event.wait() with no retry: any
cloud-leg delivery failure (broker drop, transient network fault)
propagated straight out of run_gateway_forwarder and terminated the
process (docs/design/2026-08-08-gateway-node-architecture-lift.md
axis #10, live restart-count evidence in that doc's 1.2).

- runtime/gateway_forwarder.py: new _supervise_gateway_delivery loop
  wraps the existing delivery.wait()/shutdown race. A delivery
  failure now retries with bounded exponential backoff + jitter
  instead of exiting. Reconnect state only resets after the delivery
  loop survives a full heartbeat_interval_seconds "recovery confirm"
  window without failing again (a bare delivery.start() succeeding
  only proves tasks were scheduled, not that the cloud leg is back).
  Process still exits on shutdown or on an unrecoverable
  (non-connectivity) failure -- config/heartbeat-task errors.
- Once a failure window crosses the contract-declared
  degraded_after_seconds, one DEGRADED status event publishes via
  the new ServiceGatewayForwarder.publish_status(), reusing the
  existing gateway-heartbeat.v1 event shape/topic
  (ModelGatewayHeartbeat gains status="degraded" +
  consecutive_failures + detail). Published on the LOCAL bus (not
  cloud) so DEGRADED stays observable while the cloud leg that
  caused it is itself unreachable.
- contract.yaml: liveness block gains
  reconnect_backoff_initial_seconds / _max_seconds / _jitter_seconds
  and degraded_after_seconds (contract_version/node_version patch
  bump 0.1.0 -> 0.1.1); ModelGatewayForwarderConfig gains the
  matching typed fields + a max>=initial validator.
- Data-plane loop (NodeGatewayDelivery, KafkaTransport) is untouched;
  supervision wraps it per scope.

Also vendors 2 omnimarket node migrations
(node_canary_score_reducer/0003, node_projection_registration/0004)
via scripts/sync-node-migrations.sh -- required by the repo-wide
node-migration-sync pre-commit/CI hook on any omnibase_infra branch.
infra#2683 (OMN-15732, split-migration rework) was NOT merged at
push time; these 2 files duplicate that in-flight resolution and
this PR must be rebased onto dev (picking up #2683's shape) once it
lands. No _LEGACY_DEFAULT_SCHEMA_SQL_EXACT_PATHS or exemption/
validator logic touched.

Ticket: OMN-15742 (G2, gateway-lift Phase 0)
jonahgabriel added a commit that referenced this pull request Aug 9, 2026
…ateway forwarder (#2688)

* fix(OMN-15742): bounded reconnect backoff + DEGRADED bus status for gateway forwarder

The gateway forwarder's delivery-loop supervision previously raced
delivery.wait() against shutdown_event.wait() with no retry: any
cloud-leg delivery failure (broker drop, transient network fault)
propagated straight out of run_gateway_forwarder and terminated the
process (docs/design/2026-08-08-gateway-node-architecture-lift.md
axis #10, live restart-count evidence in that doc's 1.2).

- runtime/gateway_forwarder.py: new _supervise_gateway_delivery loop
  wraps the existing delivery.wait()/shutdown race. A delivery
  failure now retries with bounded exponential backoff + jitter
  instead of exiting. Reconnect state only resets after the delivery
  loop survives a full heartbeat_interval_seconds "recovery confirm"
  window without failing again (a bare delivery.start() succeeding
  only proves tasks were scheduled, not that the cloud leg is back).
  Process still exits on shutdown or on an unrecoverable
  (non-connectivity) failure -- config/heartbeat-task errors.
- Once a failure window crosses the contract-declared
  degraded_after_seconds, one DEGRADED status event publishes via
  the new ServiceGatewayForwarder.publish_status(), reusing the
  existing gateway-heartbeat.v1 event shape/topic
  (ModelGatewayHeartbeat gains status="degraded" +
  consecutive_failures + detail). Published on the LOCAL bus (not
  cloud) so DEGRADED stays observable while the cloud leg that
  caused it is itself unreachable.
- contract.yaml: liveness block gains
  reconnect_backoff_initial_seconds / _max_seconds / _jitter_seconds
  and degraded_after_seconds (contract_version/node_version patch
  bump 0.1.0 -> 0.1.1); ModelGatewayForwarderConfig gains the
  matching typed fields + a max>=initial validator.
- Data-plane loop (NodeGatewayDelivery, KafkaTransport) is untouched;
  supervision wraps it per scope.

Also vendors 2 omnimarket node migrations
(node_canary_score_reducer/0003, node_projection_registration/0004)
via scripts/sync-node-migrations.sh -- required by the repo-wide
node-migration-sync pre-commit/CI hook on any omnibase_infra branch.
infra#2683 (OMN-15732, split-migration rework) was NOT merged at
push time; these 2 files duplicate that in-flight resolution and
this PR must be rebased onto dev (picking up #2683's shape) once it
lands. No _LEGACY_DEFAULT_SCHEMA_SQL_EXACT_PATHS or exemption/
validator logic touched.

Ticket: OMN-15742 (G2, gateway-lift Phase 0)

* fix(OMN-15742): stop outbound consumer loop re-forwarding local-only publishes

Reconciliation finding D1: publish_status (DEGRADED) publishes directly onto
the local bus's canonical outbound topic -- exactly the topic the
forwarder's own outbound consumer (NodeGatewayDelivery polling
local_consumer, the SAME KafkaTransport object local_bus publishes into in
runtime/gateway_forwarder.py:run_gateway_forwarder) is subscribed to. The
untagged envelope does not match _forward_outbound_message's loopback skip
(only "cloud-to-local"), so it falls through to _prepare_outbound and leaks
to the cloud leg -- directly contradicting
test_publish_status_degraded_goes_to_local_bus_not_cloud, which only passes
today because it uses a bare _MockGatewayBus with no consumer loop
(feedback_real_dispatch_path_tests gap).

Fix: local-only publishes are now stamped with
gateway_direction="local-mirror"; _forward_outbound_message and
validate_outbound_message skip both "cloud-to-local" (existing) and
"local-mirror" (new) via a shared _LOCAL_ONLY_DIRECTIONS set. The stamping
helper is a module-level function, not a method, to stay under the
ONEX Pattern Validation method-count threshold on ServiceGatewayForwarder.

New test test_real_outbound_consumer_loop_does_not_reforward_degraded_status
wires the REAL NodeGatewayDelivery consumer loop against a fake transport
that actually connects publish (local_bus) to poll (local_consumer) -- the
same object playing both roles, matching the real runtime wiring -- unlike
every existing test in this file, which uses two disconnected fakes and so
cannot observe this class of bug. Verified RED against the parent commit
(f0cdfe4, pre-fix): DEGRADED payload reached the cloud leg. GREEN after
this fix.

The matching heartbeat-local-mirror half of this bug (#2692, OMN-15570)
gets the same fix applied on that branch separately, since it stacks on
this one and its own diff must be independently correct.

Ticket: OMN-15742
Evidence-Ticket: OMN-15742

* fix(OMN-15742): stop test _Source.poll() from starving the asyncio event loop

Root cause of every anomalously long/hanging local pre-push run on this
branch throughout this session (observed: 10h33m on the original Mac before
this session started per the rolling ledger, then repeated 33min-2h+ hangs
on .200 across multiple push attempts -- all traced to the SAME mechanism
via a --timeout=60 diagnostic run that caught the stack mid-hang inside
NodeGatewayDelivery._run_direction's `while True:` loop).

_Source.poll() (a test fake in test_gateway_delivery_service.py) returned
`[]` with zero internal `await` suspension points. A coroutine with no real
suspension point does not yield control back to the asyncio event loop when
awaited -- so when test_real_outbound_consumer_loop_does_not_reforward_degraded_status
(added by the prior commit on this branch) drives a REAL
NodeGatewayDelivery task loop with _Source as the "idle" cloud consumer,
`_run_direction`'s `while True: await source.poll(...)` busy-spins the
event loop forever. `task.cancel()` only takes effect at the next real
suspension point, so the spinning task is uncancellable: `delivery.stop()`
(called in the test's `finally`) awaits `asyncio.gather(*tasks)` on tasks
that can never be scheduled to observe their own cancellation, deadlocking
the entire test process indefinitely -- explaining every "still running,
100% CPU, no progress" observation this session mistook for slowness rather
than a hang.

Fix: _Source.poll() now honors timeout_ms via asyncio.sleep(timeout_ms /
1000) before returning empty, matching how a real transport (and this same
file's _SharedLocalTransport.poll(), which correctly uses
asyncio.wait_for(..., timeout=...)) actually behaves. Every other _Source
usage in this file calls deliver_message() directly and never invokes
poll() at all, so this costs exactly one test ~50ms and changes nothing
else (verified: full file 5/5 passed in 0.20s, vs. hanging past a 60s
per-test timeout pre-fix).

Ticket: OMN-15742

* fix(OMN-15742): wire required canary field into remaining gateway test fixtures

Two gaps surfaced by the fixed full local suite run (previous commit
removed the busy-spin hang that was masking these):

- test_gateway_forwarder_config.py: the two reconnect_backoff tests
  (added by this branch's own reconnect-supervision commit, authored
  before this branch was rebased onto dev's #2690/OMN-15741 G1 commit)
  never got canary=_canary() added, unlike every other constructor call
  in the same file.
- test_gateway_forwarder_service_handler_seam.py: pre-existing dev-tip
  gap (unrelated to this branch's diff) -- same class of fix already
  applied independently on the OMN-15781 branch.

Ticket: OMN-15742
jonahgabriel added a commit that referenced this pull request Aug 9, 2026
…hdog (#2700)

* fix(OMN-15742): bounded reconnect backoff + DEGRADED bus status for gateway forwarder

The gateway forwarder's delivery-loop supervision previously raced
delivery.wait() against shutdown_event.wait() with no retry: any
cloud-leg delivery failure (broker drop, transient network fault)
propagated straight out of run_gateway_forwarder and terminated the
process (docs/design/2026-08-08-gateway-node-architecture-lift.md
axis #10, live restart-count evidence in that doc's 1.2).

- runtime/gateway_forwarder.py: new _supervise_gateway_delivery loop
  wraps the existing delivery.wait()/shutdown race. A delivery
  failure now retries with bounded exponential backoff + jitter
  instead of exiting. Reconnect state only resets after the delivery
  loop survives a full heartbeat_interval_seconds "recovery confirm"
  window without failing again (a bare delivery.start() succeeding
  only proves tasks were scheduled, not that the cloud leg is back).
  Process still exits on shutdown or on an unrecoverable
  (non-connectivity) failure -- config/heartbeat-task errors.
- Once a failure window crosses the contract-declared
  degraded_after_seconds, one DEGRADED status event publishes via
  the new ServiceGatewayForwarder.publish_status(), reusing the
  existing gateway-heartbeat.v1 event shape/topic
  (ModelGatewayHeartbeat gains status="degraded" +
  consecutive_failures + detail). Published on the LOCAL bus (not
  cloud) so DEGRADED stays observable while the cloud leg that
  caused it is itself unreachable.
- contract.yaml: liveness block gains
  reconnect_backoff_initial_seconds / _max_seconds / _jitter_seconds
  and degraded_after_seconds (contract_version/node_version patch
  bump 0.1.0 -> 0.1.1); ModelGatewayForwarderConfig gains the
  matching typed fields + a max>=initial validator.
- Data-plane loop (NodeGatewayDelivery, KafkaTransport) is untouched;
  supervision wraps it per scope.

Also vendors 2 omnimarket node migrations
(node_canary_score_reducer/0003, node_projection_registration/0004)
via scripts/sync-node-migrations.sh -- required by the repo-wide
node-migration-sync pre-commit/CI hook on any omnibase_infra branch.
infra#2683 (OMN-15732, split-migration rework) was NOT merged at
push time; these 2 files duplicate that in-flight resolution and
this PR must be rebased onto dev (picking up #2683's shape) once it
lands. No _LEGACY_DEFAULT_SCHEMA_SQL_EXACT_PATHS or exemption/
validator logic touched.

Ticket: OMN-15742 (G2, gateway-lift Phase 0)

* fix(OMN-15742): stop outbound consumer loop re-forwarding local-only publishes

Reconciliation finding D1: publish_status (DEGRADED) publishes directly onto
the local bus's canonical outbound topic -- exactly the topic the
forwarder's own outbound consumer (NodeGatewayDelivery polling
local_consumer, the SAME KafkaTransport object local_bus publishes into in
runtime/gateway_forwarder.py:run_gateway_forwarder) is subscribed to. The
untagged envelope does not match _forward_outbound_message's loopback skip
(only "cloud-to-local"), so it falls through to _prepare_outbound and leaks
to the cloud leg -- directly contradicting
test_publish_status_degraded_goes_to_local_bus_not_cloud, which only passes
today because it uses a bare _MockGatewayBus with no consumer loop
(feedback_real_dispatch_path_tests gap).

Fix: local-only publishes are now stamped with
gateway_direction="local-mirror"; _forward_outbound_message and
validate_outbound_message skip both "cloud-to-local" (existing) and
"local-mirror" (new) via a shared _LOCAL_ONLY_DIRECTIONS set. The stamping
helper is a module-level function, not a method, to stay under the
ONEX Pattern Validation method-count threshold on ServiceGatewayForwarder.

New test test_real_outbound_consumer_loop_does_not_reforward_degraded_status
wires the REAL NodeGatewayDelivery consumer loop against a fake transport
that actually connects publish (local_bus) to poll (local_consumer) -- the
same object playing both roles, matching the real runtime wiring -- unlike
every existing test in this file, which uses two disconnected fakes and so
cannot observe this class of bug. Verified RED against the parent commit
(f0cdfe4, pre-fix): DEGRADED payload reached the cloud leg. GREEN after
this fix.

The matching heartbeat-local-mirror half of this bug (#2692, OMN-15570)
gets the same fix applied on that branch separately, since it stacks on
this one and its own diff must be independently correct.

Ticket: OMN-15742
Evidence-Ticket: OMN-15742

* fix(OMN-15742): stop test _Source.poll() from starving the asyncio event loop

Root cause of every anomalously long/hanging local pre-push run on this
branch throughout this session (observed: 10h33m on the original Mac before
this session started per the rolling ledger, then repeated 33min-2h+ hangs
on .200 across multiple push attempts -- all traced to the SAME mechanism
via a --timeout=60 diagnostic run that caught the stack mid-hang inside
NodeGatewayDelivery._run_direction's `while True:` loop).

_Source.poll() (a test fake in test_gateway_delivery_service.py) returned
`[]` with zero internal `await` suspension points. A coroutine with no real
suspension point does not yield control back to the asyncio event loop when
awaited -- so when test_real_outbound_consumer_loop_does_not_reforward_degraded_status
(added by the prior commit on this branch) drives a REAL
NodeGatewayDelivery task loop with _Source as the "idle" cloud consumer,
`_run_direction`'s `while True: await source.poll(...)` busy-spins the
event loop forever. `task.cancel()` only takes effect at the next real
suspension point, so the spinning task is uncancellable: `delivery.stop()`
(called in the test's `finally`) awaits `asyncio.gather(*tasks)` on tasks
that can never be scheduled to observe their own cancellation, deadlocking
the entire test process indefinitely -- explaining every "still running,
100% CPU, no progress" observation this session mistook for slowness rather
than a hang.

Fix: _Source.poll() now honors timeout_ms via asyncio.sleep(timeout_ms /
1000) before returning empty, matching how a real transport (and this same
file's _SharedLocalTransport.poll(), which correctly uses
asyncio.wait_for(..., timeout=...)) actually behaves. Every other _Source
usage in this file calls deliver_message() directly and never invokes
poll() at all, so this costs exactly one test ~50ms and changes nothing
else (verified: full file 5/5 passed in 0.20s, vs. hanging past a 60s
per-test timeout pre-fix).

Ticket: OMN-15742

* fix(OMN-15742): wire required canary field into remaining gateway test fixtures

Two gaps surfaced by the fixed full local suite run (previous commit
removed the busy-spin hang that was masking these):

- test_gateway_forwarder_config.py: the two reconnect_backoff tests
  (added by this branch's own reconnect-supervision commit, authored
  before this branch was rebased onto dev's #2690/OMN-15741 G1 commit)
  never got canary=_canary() added, unlike every other constructor call
  in the same file.
- test_gateway_forwarder_service_handler_seam.py: pre-existing dev-tip
  gap (unrelated to this branch's diff) -- same class of fix already
  applied independently on the OMN-15781 branch.

Ticket: OMN-15742

* fix(OMN-15748): gateway poison-pill quarantine + membership-loss watchdog

decode_message() ran before the try/except that nacks, so any undecodable
record on a mirrored source topic crashed the process (live-observed
2026-08-08T16:50Z). Nacking it (the ordinary failure path) would also be
wrong for a permanently malformed record: nack seeks back to the same
offset and re-crashes forever. Quarantine instead: log, best-effort
dead-letter via the DLQ topic derived from the record's own topic, commit
past the offset, keep the loop alive.

Also adds an independent membership-loss watchdog to NodeGatewayDelivery.
The live 2026-08-09T10:03:07Z incident was a clean single LeaveGroup with
zero exception (aiokafka's client-side max_poll_interval_ms idle-eviction,
fired from a separate asyncio task than the poll loop) -- the existing
exception-triggered reconnect supervision in runtime/gateway_forwarder.py
cannot observe a task that is alive, hung, and never exceptions. The
watchdog tracks per-direction last-progress timestamps (updated on poll()
completion and deliver_message completion) plus a best-effort
consumer.assignment() membership probe, and on staleness force-recreates
just the affected direction's transport (a fresh client forces a new group
join, sidestepping the permanently-stuck idle clock) and publishes a
DEGRADED status. Reuses the previously contract-declared but unwired
max_silence_window_seconds config field as the staleness threshold.

TDD: both regression tests are RED against pre-fix source (verified via
git stash of the source-only diff) and GREEN after.

OMN-15748

* fix(OMN-15748): consumer-scoped watchdog recovery, no CancelledError escape

Fixes two verifier-confirmed defects in the membership-loss watchdog
(mergesweep-0809-poisonpath-verify):

- wait()/_run_direction: a watchdog stale_task.cancel() made the
  supervisor's delivery.wait() (a frozen asyncio.gather snapshot) observe
  a cancelled child and propagate CancelledError -- a BaseException that
  escapes gateway_forwarder.py's except Exception and kills the process on
  every recovery. _run_direction now catches CancelledError and returns
  cleanly when its direction is watchdog-recovering; wait() re-reads the
  live task set after every completion instead of gathering a one-time
  snapshot, so a swapped-in replacement task stays supervised.

- Recovery no longer calls close()+start() (which stops BOTH the consumer
  and producer on a KafkaTransport instance shared across directions --
  one direction's consumer is the same object the other direction's
  publish and status/heartbeat emission use as a producer). Adds
  KafkaTransport.restart_consumer(), a consumer-only recreate that never
  touches the producer; _recover_stalled_direction now calls it
  exclusively.

De-vacuums the watchdog regression test (arms delivery.wait() before the
stall, matching production ordering) and adds direct KafkaTransport-level
coverage proving the producer is untouched across a restart_consumer()
call.

* fix(OMN-15748): address CodeRabbit findings on watchdog wait() + log sanitization

Two live (non-outdated) CodeRabbit threads on infra#2700:

- wait()'s asyncio.wait was unbounded, waking only on a tracked task's
  completion. _recover_stalled_direction splices the replacement task into
  self._tasks strictly AFTER the stale task has already finished and this
  loop has already woken/rebuilt `watched` from the stale-only list -- so an
  unbounded wait then blocks on the remaining old tasks with no further
  wakeup until one of THEM completes. A real exception on the recovered
  direction would go unsupervised (only surfacing as an "exception was never
  retrieved" warning at GC). Fixed by bounding the wait with
  _watchdog_tick_seconds so the loop always revisits self._tasks, and by also
  picking up an already-completed-with-exception replacement on that rescan.
  Added test_watchdog_replacement_task_failure_still_reaches_wait, verified
  RED against the unbounded wait() (reverted locally, re-applied after) and
  GREEN after.

- _quarantine_undecodable_message logged the raw exception via exc_info=error,
  bypassing the sanitize_error_message() the DLQ payload path already uses a
  few lines up. A pydantic ValidationError from decode_message() embeds the
  offending record content, which is untrusted tenant bus traffic that can
  carry PII/credential-bearing values. Now logs the sanitized message,
  matching the DLQ path.

* test(OMN-15748): harden test_monitor_logs.py subprocess timeouts against CI fleet contention

The self-hosted CI fleet runs many concurrent pytest-split shards per job;
interpreter spawn/import can occasionally queue behind that CPU contention.
A hardcoded 10s subprocess.run(..., timeout=10) tripped this twice live on
unrelated PRs' Split 12/15 (infra#2697 and this PR, both times as three
subprocess.TimeoutExpired failures at exactly the 10s bound, unrelated to
either PR's own diff -- this file isn't touched by either PR's product
diff). Bumped to a shared _SUBPROCESS_TIMEOUT_SECONDS = 30 constant across
all three call sites: still tight enough to catch a genuine hang, with real
headroom against fleet contention.
jonahgabriel added a commit that referenced this pull request Aug 25, 2026
…ce now present)

Empty commit to re-trigger the pull_request synchronize event so ci.yml's
occ-preflight job resolves the PR body live (Evidence-Source: OCC#7074,
already present as of the prior push) on its first attempt, proving site
#10 (application-database-domain-enforcement) without waiting on the
unrelated OCC Companion Merged Gate / shadow test-suite completion.
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