From e3feab86a9e10cb40cfa8152f3eb1fde28b4a12f Mon Sep 17 00:00:00 2001 From: jonahgabriel Date: Mon, 22 Jun 2026 09:30:57 -0400 Subject: [PATCH 1/5] =?UTF-8?q?feat(OMN-13472):=20ARCH-004=20imperative-or?= =?UTF-8?q?chestrator=20ratchet=20=E2=80=94=20cross-file=20rule=20+=20base?= =?UTF-8?q?line=20+=20gate?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds ARCH-004 'Contract-Declared Orchestrator Workflow Must Be Bound To An Executor' to node_architecture_validator. Cross-file node-directory rule that joins contract.yaml (fsm/workflow_coordination), handler_routing.routing_strategy, handler source, and node type/name — catching the delegation-shaped anti-pattern ARCH-003 structurally misses (handler-owned _transition( in a non-*Orchestrator class; declared-but-unbound fsm). - RuleContractDeclaredOrchestratorWorkflow registered in validators/__init__.py - scanner_imperative_orchestrator_ratchet: --check-all/--report, --check-changed/--ratchet, --strict; baseline can only shrink - architecture-handshakes/imperative-orchestrator-baseline.yaml (9 current hard-fails, owner OMN-13471; delegation = sole P0 risk 10) - Wired via OMN-12550 path: scripts/validate.py imperative_orchestrators subcommand + pre-commit hook (changed-node ratchet, blocking) + CI full report (non-blocking initially). Cites OMN-12550 + OMN-13325 in configs. - Tests: ARCH-003-passes / ARCH-004-fails delegation-shape proof + 14 more. Refs OMN-13472 (epic OMN-13471), OMN-12550, OMN-13325. --- .github/workflows/ci.yml | 16 + .pre-commit-config.yaml | 22 +- .../imperative-orchestrator-baseline.yaml | 99 +++ pyproject.toml | 2 + scripts/validate.py | 91 +- .../node_architecture_validator/contract.yaml | 20 + .../validators/__init__.py | 18 + ...scanner_imperative_orchestrator_ratchet.py | 477 +++++++++++ ...contract_declared_orchestrator_workflow.py | 794 ++++++++++++++++++ ...contract_declared_orchestrator_workflow.py | 634 ++++++++++++++ 10 files changed, 2171 insertions(+), 2 deletions(-) create mode 100644 architecture-handshakes/imperative-orchestrator-baseline.yaml create mode 100644 src/omnibase_infra/nodes/node_architecture_validator/validators/scanner_imperative_orchestrator_ratchet.py create mode 100644 src/omnibase_infra/nodes/node_architecture_validator/validators/validator_contract_declared_orchestrator_workflow.py create mode 100644 tests/unit/nodes/node_architecture_validator/test_validator_contract_declared_orchestrator_workflow.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f91d939021..48debb13a8 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -196,6 +196,22 @@ jobs: echo "================================================================" uv run pytest tests/audit/test_io_violations.py::TestCIGateIOPurity -v --tb=short + # ARCH-004 imperative-orchestrator ratchet (OMN-13472). + # Full report over the repo: surfaces every contract-declared-but-unbound + # orchestrator FSM driven imperatively by a handler (the delegation-shaped + # anti-pattern ARCH-003 misses). Non-blocking INITIALLY (continue-on-error) + # while the baseline is above threshold — the changed-node ratchet + # (pre-commit, blocking) stops regressions. Promote to a required gate once + # the baseline is below the agreed threshold; this rides OMN-12550 + # (validator gating) / OMN-13325 (ratchet enforcement), not a fresh hook. + - name: Run imperative-orchestrator ratchet report (ARCH-004) + continue-on-error: true + run: | + echo "================================================================" + echo "Imperative-Orchestrator Ratchet Report (ARCH-004, OMN-13472)" + echo "================================================================" + uv run python scripts/validate.py imperative_orchestrators --verbose + - name: Run markdown link validation run: | echo "================================================================" diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml index a2bdb24cdc..a5b845f226 100644 --- a/.pre-commit-config.yaml +++ b/.pre-commit-config.yaml @@ -256,6 +256,26 @@ repos: types: [python] pass_filenames: false stages: [pre-commit] + # ARCH-004 imperative-orchestrator ratchet (OMN-13472). + # + # Changed-node ratchet (BLOCKING): for any node directory touched by this + # commit, fail if it introduces a NEW / worsened / untracked + # imperative-orchestrator hard-fail (contract-declared-but-unbound FSM + # driven imperatively by a handler — the delegation-shaped anti-pattern + # ARCH-003 structurally misses). Accepted debt is frozen in + # architecture-handshakes/imperative-orchestrator-baseline.yaml (it can + # only shrink). + # + # This rides the validator-gating work of OMN-12550 (wire ARCH-001/002/003 + # as blocking gates) and the ratchet-enforcement epic OMN-13325 — it is NOT + # a parallel/competing hook. Owner epic for the baselined nodes: OMN-13471. + - id: onex-imperative-orchestrator-ratchet + name: ONEX Imperative-Orchestrator Ratchet (ARCH-004) + entry: uv run --frozen python scripts/validate.py imperative_orchestrators + language: system + files: '(contract\.yaml|handlers/handler_.*\.py)$' + pass_filenames: true + stages: [pre-commit] # Architecture layer validation - core/infra separation # Verifies omnibase_core has no infrastructure dependencies (kafka, httpx, etc.) # @@ -709,4 +729,4 @@ ci: # Note: onex-validate-clean-root is NOT skipped (standalone script, no omnibase_core dependency) # Also skip migration freeze (checks staged files via git diff --cached, which is empty in CI; # CI should run: uv run python scripts/validation/validate_migration_freeze.py --check-committed) - skip: [no-env-file, onex-validate-architecture, onex-validate-architecture-layers, onex-validate-contracts, onex-validate-patterns, onex-validate-unions, onex-validate-imports, onex-validate-any-types, onex-validate-migration-freeze, normalization-symmetry] + skip: [no-env-file, onex-validate-architecture, onex-imperative-orchestrator-ratchet, onex-validate-architecture-layers, onex-validate-contracts, onex-validate-patterns, onex-validate-unions, onex-validate-imports, onex-validate-any-types, onex-validate-migration-freeze, normalization-symmetry] diff --git a/architecture-handshakes/imperative-orchestrator-baseline.yaml b/architecture-handshakes/imperative-orchestrator-baseline.yaml new file mode 100644 index 0000000000..bd4a9f7950 --- /dev/null +++ b/architecture-handshakes/imperative-orchestrator-baseline.yaml @@ -0,0 +1,99 @@ +# SPDX-FileCopyrightText: 2026 OmniNode.ai Inc. +# SPDX-License-Identifier: MIT +# +# ARCH-004 imperative-orchestrator ratchet baseline (OMN-13472). +# Generated by: +# uv run python -m omnibase_infra.nodes.node_architecture_validator.validators.scanner_imperative_orchestrator_ratchet \ +# --check-all --report --write-baseline +# Owner epic: OMN-13471. Wiring: OMN-12550 / OMN-13325. +schema_version: + major: 1 + minor: 0 + patch: 0 +repo: omnibase_infra +rule: ARCH-004 +description: 'Accepted imperative-orchestrator hard-fails for the ARCH-004 ratchet (OMN-13472). One entry per current hard-fail. This list can only SHRINK: a new/worsened/untracked finding on a touched node fails the changed-node ratchet. Each entry carries an owner ticket for its decomposition (delegation -> OMN-13471).' +entries: + - repo: omnimarket + node: autopilot_orchestrator + max_handler_path: src/omnimarket/nodes/node_autopilot_orchestrator/handlers/handler_autopilot_orchestrator.py + line_count: 634 + risk_score: 6 + finding_codes: + - H1 + - W3 + - W4 + owner_ticket: OMN-13471 + - repo: omnibase_infra + node: node_chain_orchestrator + max_handler_path: src/omnibase_infra/nodes/node_chain_orchestrator/handlers/handler_chain_replay_complete.py + line_count: 70 + risk_score: 2 + finding_codes: + - H2 + owner_ticket: OMN-13471 + - repo: omnimarket + node: node_delegation_orchestrator + max_handler_path: src/omnimarket/nodes/node_delegation_orchestrator/handlers/handler_delegation_workflow.py + line_count: 1542 + risk_score: 10 + finding_codes: + - H1 + - H2 + - H3 + - W1 + - W2 + - W3 + - W4 + owner_ticket: OMN-13471 + - repo: omnibase_infra + node: node_merge_sweep_workflow_orchestrator + max_handler_path: src/omnibase_infra/nodes/node_merge_sweep_workflow_orchestrator/handlers/handler_auto_merge_complete.py + line_count: 90 + risk_score: 2 + finding_codes: + - H2 + owner_ticket: OMN-13471 + - repo: omnibase_infra + node: node_registration_orchestrator + max_handler_path: src/omnibase_infra/nodes/node_registration_orchestrator/handlers/handler_node_registration_acked.py + line_count: 435 + risk_score: 2 + finding_codes: + - H2 + owner_ticket: OMN-13471 + - repo: omnibase_infra + node: node_routing_orchestrator + max_handler_path: src/omnibase_infra/nodes/node_routing_orchestrator/handlers/handler_health_complete.py + line_count: 70 + risk_score: 2 + finding_codes: + - H2 + owner_ticket: OMN-13471 + - repo: omnibase_infra + node: node_rsd_orchestrator + max_handler_path: src/omnibase_infra/nodes/node_rsd_orchestrator/handlers/handler_rsd_data_fetch_complete.py + line_count: 90 + risk_score: 2 + finding_codes: + - H2 + owner_ticket: OMN-13471 + - repo: omnibase_infra + node: node_scope_workflow_orchestrator + max_handler_path: src/omnibase_infra/nodes/node_scope_workflow_orchestrator/handlers/handler_scope_file_read_complete.py + line_count: 73 + risk_score: 2 + finding_codes: + - H2 + owner_ticket: OMN-13471 + - repo: omnimarket + node: pr_lifecycle_orchestrator + max_handler_path: src/omnimarket/nodes/node_pr_lifecycle_orchestrator/handlers/handler_pr_lifecycle_orchestrator.py + line_count: 1815 + risk_score: 7 + finding_codes: + - H1 + - W1 + - W2 + - W4 + owner_ticket: OMN-13471 diff --git a/pyproject.toml b/pyproject.toml index 536022cf6c..cd97edee90 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -319,6 +319,8 @@ ignore = [ # Intentional stdout sinks / fallback output "src/omnibase_infra/event_bus/service_topic_manager.py" = ["T201"] "src/omnibase_infra/observability/sinks/sink_logging_structured.py" = ["T201"] +# ARCH-004 imperative-orchestrator ratchet CLI (OMN-13472): intentional CLI/report stdout. +"src/omnibase_infra/nodes/node_architecture_validator/validators/scanner_imperative_orchestrator_ratchet.py" = ["T201"] [tool.ruff.lint.isort] # Ensure consistent import sorting between local and CI diff --git a/scripts/validate.py b/scripts/validate.py index 0e7cfc3e08..84ec529fa6 100755 --- a/scripts/validate.py +++ b/scripts/validate.py @@ -535,6 +535,85 @@ def run_declarative_nodes( return True +def run_imperative_orchestrators( + verbose: bool = False, files: list[str] | None = None +) -> bool: + """Run the ARCH-004 imperative-orchestrator ratchet (OMN-13472). + + Detects contract-declared-but-unbound orchestrator FSMs driven imperatively + by a handler (the delegation-shaped anti-pattern ARCH-003 misses). Wired + through OMN-12550 (validator gating) / OMN-13325 (ratchet enforcement). + + Modes: + - With ``files`` (pre-commit): changed-node ratchet — only node + directories containing a changed file are scanned, and the gate fails + on a NEW / worsened / untracked finding (``--check-changed --ratchet``). + - Without ``files`` (CI/manual full report): full report over the repo + (``--check-all --report``), non-blocking on its own. + + The baseline lives at + ``architecture-handshakes/imperative-orchestrator-baseline.yaml`` and can + only shrink. + + Args: + verbose: Enable verbose output. + files: Optional list of changed files (from pre-commit). If provided, + only the node directories containing those files are ratcheted. + """ + try: + from omnibase_infra.nodes.node_architecture_validator.validators.scanner_imperative_orchestrator_ratchet import ( + load_baseline, + node_dirs_for_changed_files, + ratchet_violations, + scan_node_dirs, + ) + except ImportError as e: + print(f"Skipping imperative-orchestrator ratchet: {e}") + return True + + from pathlib import Path as _Path + + repo_root = _Path.cwd() + baseline = load_baseline( + repo_root / "architecture-handshakes" / "imperative-orchestrator-baseline.yaml" + ) + + if files: + node_dirs = node_dirs_for_changed_files(repo_root, files) + if not node_dirs: + if verbose: + print("Imperative Orchestrators: SKIP (no node dirs in changeset)") + return True + result = scan_node_dirs("omnibase_infra", node_dirs, repo_root=repo_root) + failures = ratchet_violations(result.hard_fails, baseline) + passed = not failures + if verbose or not passed: + print(f"Imperative Orchestrators: {'PASS' if passed else 'FAIL'}") + print( + f" Changed node dirs: {len(node_dirs)}, " + f"hard-fails: {len(result.hard_fails)}" + ) + for f in failures: + print(f" - {f}") + return passed + + # Full report mode (non-blocking on its own). + from omnibase_infra.nodes.node_architecture_validator.validators.scanner_imperative_orchestrator_ratchet import ( + discover_node_dirs, + ) + + node_dirs = discover_node_dirs(repo_root) + result = scan_node_dirs("omnibase_infra", node_dirs, repo_root=repo_root) + print( + f"Imperative Orchestrators (full report): " + f"{len(result.hard_fails)} hard-fail node(s) across {len(node_dirs)} dirs." + ) + if verbose: + for line in result.report_lines: + print(line) + return True + + def run_io_audit(verbose: bool = False) -> bool: """Run I/O purity audit for REDUCER and COMPUTE nodes. @@ -1074,6 +1153,7 @@ def main() -> int: "any_types", "localhandler", "declarative_nodes", + "imperative_orchestrators", "io_audit", "imports", "markdown_links", @@ -1084,7 +1164,10 @@ def main() -> int: parser.add_argument( "files", nargs="*", - help="Optional list of files to validate (for declarative_nodes or markdown_links)", + help=( + "Optional list of files to validate (for declarative_nodes, " + "imperative_orchestrators, or markdown_links)" + ), ) parser.add_argument("--verbose", "-v", action="store_true", help="Verbose output") parser.add_argument( @@ -1105,6 +1188,7 @@ def main() -> int: "any_types": run_any_types, "localhandler": run_localhandler, "declarative_nodes": run_declarative_nodes, + "imperative_orchestrators": run_imperative_orchestrators, "io_audit": run_io_audit, "imports": run_imports, "markdown_links": run_markdown_links, @@ -1121,6 +1205,11 @@ def main() -> int: # Pass files to declarative_nodes validator if provided files = args.files if args.files else None success = run_declarative_nodes(args.verbose, files=files) + elif args.validator == "imperative_orchestrators": + # Pass changed files for the changed-node ratchet (pre-commit); without + # files this runs the full non-blocking report. + files = args.files if args.files else None + success = run_imperative_orchestrators(args.verbose, files=files) else: success = validator_map[args.validator](args.verbose) diff --git a/src/omnibase_infra/nodes/node_architecture_validator/contract.yaml b/src/omnibase_infra/nodes/node_architecture_validator/contract.yaml index e4e2c80cf3..e3e424a633 100644 --- a/src/omnibase_infra/nodes/node_architecture_validator/contract.yaml +++ b/src/omnibase_infra/nodes/node_architecture_validator/contract.yaml @@ -133,6 +133,26 @@ validation_rules: suggested_fix: > Remove state transition logic from orchestrator. State transitions belong exclusively to reducer nodes. Orchestrators should call reducers for FSM decisions. + - rule_id: "ARCH-004" + name: "Contract-Declared Orchestrator Workflow Must Be Bound To An Executor" + description: > + Cross-file, node-directory rule (OMN-13472). An orchestrator-like node that declares an fsm:/workflow-state set MUST bind it to a runtime executor. It must NOT leave the contract transition table decorative (ModelContractOrchestrator has no typed fsm/state_machine field, no traverser consumes orchestrator fsm.transitions) while a handler drives the transitions itself. Catches the delegation-shaped anti-pattern (_transition(...) calls in a non-"*Orchestrator" handler, payload_type_match catchall funneling 3+ payload types into one handler, and one handler that both selects the next state and constructs terminal/compat events) that ARCH-003's single-file, class-name-gated AST approach structurally misses. Reducers are exempt. + + # ERROR severity: a declared-but-unbound orchestrator FSM driven imperatively + # by a handler is the OMN-13408-class footgun (terminal/projection corruption). + severity: "ERROR" + detection_strategy: + type: "node_directory_join" + checks: + - decorative_fsm_with_handler_driven_transitions + - payload_type_match_fanin_3plus + - handler_selects_state_and_constructs_events + target_files: + - "contract.yaml" + - "handlers/handler_*.py" + suggested_fix: > + Bind the contract fsm to an executor (migrate to a typed, executor-bound workflow/DAG field per OMN-12835) or decompose the monolithic handler into per-step thin handlers driven by a reducer. See OMN-13471 (delegation decomposition epic). New/worsened cases on touched nodes fail the changed-node ratchet; accepted debt is recorded in architecture-handshakes/imperative-orchestrator-baseline.yaml with an owner ticket. + # Dependencies (protocols this node requires) dependencies: - name: "protocol_ast_analyzer" diff --git a/src/omnibase_infra/nodes/node_architecture_validator/validators/__init__.py b/src/omnibase_infra/nodes/node_architecture_validator/validators/__init__.py index 1a1cea194c..62e6895e66 100644 --- a/src/omnibase_infra/nodes/node_architecture_validator/validators/__init__.py +++ b/src/omnibase_infra/nodes/node_architecture_validator/validators/__init__.py @@ -26,6 +26,15 @@ Reducers own state machines; orchestrators are "reaction planners" that coordinate work based on reducer outputs. + - ARCH-004: Contract-Declared Orchestrator Workflow Must Be Bound To An Executor + Cross-file, node-directory rule (OMN-13472). An orchestrator-like node + that declares an fsm:/workflow-state set must bind it to a runtime + executor; it must not leave the contract table decorative while a + handler drives the transitions itself. Catches the delegation-shaped + anti-pattern (_transition(...) in a non-"*Orchestrator" handler) that + ARCH-003's single-file, class-name-gated AST approach structurally + misses. Reducers are exempt. + Two Interfaces: **1. Function-based validators** - Direct file validation, returns detailed results. @@ -79,6 +88,11 @@ from __future__ import annotations +from omnibase_infra.nodes.node_architecture_validator.validators.validator_contract_declared_orchestrator_workflow import ( + RuleContractDeclaredOrchestratorWorkflow, + analyze_node_directory, + validate_contract_declared_orchestrator_workflow, +) from omnibase_infra.nodes.node_architecture_validator.validators.validator_no_direct_dispatch import ( RuleNoDirectDispatch, validate_no_direct_dispatch, @@ -97,8 +111,12 @@ "validate_no_direct_dispatch", "validate_no_handler_publishing", "validate_no_orchestrator_fsm", + # Functions (node-directory cross-file validators) + "validate_contract_declared_orchestrator_workflow", + "analyze_node_directory", # Classes (protocol-compliant rules) "RuleNoDirectDispatch", "RuleNoHandlerPublishing", "RuleNoOrchestratorFSM", + "RuleContractDeclaredOrchestratorWorkflow", ] diff --git a/src/omnibase_infra/nodes/node_architecture_validator/validators/scanner_imperative_orchestrator_ratchet.py b/src/omnibase_infra/nodes/node_architecture_validator/validators/scanner_imperative_orchestrator_ratchet.py new file mode 100644 index 0000000000..63d6972e99 --- /dev/null +++ b/src/omnibase_infra/nodes/node_architecture_validator/validators/scanner_imperative_orchestrator_ratchet.py @@ -0,0 +1,477 @@ +# SPDX-FileCopyrightText: 2025 OmniNode.ai Inc. +# SPDX-License-Identifier: MIT +"""Repository scanner + ratchet for ARCH-004 (imperative-orchestrator debt). + +This module turns the per-node ARCH-004 analysis +(``validator_contract_declared_orchestrator_workflow.analyze_node_directory``) +into a repository-wide gate with three modes: + + * ``--check-all --report`` — scan every node directory, emit one report row + per hard-fail. Never hides debt; intended for CI reporting and baseline + generation. Exit 0 (reporting) unless ``--ratchet`` is also requested. + * ``--check-changed --ratchet`` — scan only the node directories touched by + a changeset. Fail (exit 1) if a touched node introduces a NEW hard-fail, + WORSENS its risk score / grows its max handler beyond the baseline, or is + an untracked hard-fail with no baseline entry. + * ``--strict`` — after a remediation wave, treat any baselined node that the + live scan no longer reproduces as a recurrence failure if it reappears; + with ``--strict`` a baselined node that still hard-fails is itself a + failure (the baseline is being dropped). Use to ratchet the baseline down. + +Baseline format: ``architecture-handshakes/imperative-orchestrator-baseline.yaml``. +This mirrors the repo's existing handshake-baseline pattern +(``validator-requirements-baseline.yaml``): the baseline can only SHRINK; a new +or worsened finding fails the gate. + +Related: + - OMN-13472 (ARCH-004 ratchet — Workstream B) + - OMN-13471 (delegation decomposition epic — baseline owner ticket) + - OMN-12550 / OMN-13325 (validator wiring / ratchet-enforcement epic) +""" + +from __future__ import annotations + +import argparse +import sys +from dataclasses import dataclass, field +from pathlib import Path + +import yaml + +from omnibase_infra.nodes.node_architecture_validator.validators.validator_contract_declared_orchestrator_workflow import ( + analyze_node_directory, +) + +#: Canonical baseline location relative to a repo root. +BASELINE_RELATIVE_PATH = "architecture-handshakes/imperative-orchestrator-baseline.yaml" + +#: Owner ticket recorded against the delegation hard-fail in the baseline +#: (the dedicated decomposition epic). +DELEGATION_OWNER_TICKET = "OMN-13471" + +#: Owner ticket for every other current hard-fail (same anti-pattern class, +#: tracked under the imperative-orchestrator decomposition epic). The audit +#: (docs/audits/2026-06-22-imperative-orchestrator-audit.md) treats these as +#: baseline debt under the same epic, with delegation as the sole P0. +FLEET_OWNER_TICKET = "OMN-13471" + + +@dataclass(frozen=True) +class BaselineEntry: + """One accepted hard-fail in the ratchet baseline.""" + + repo: str + node: str + max_handler_path: str + line_count: int + risk_score: int + finding_codes: tuple[str, ...] + owner_ticket: str + + def to_dict(self) -> dict[str, object]: + return { + "repo": self.repo, + "node": self.node, + "max_handler_path": self.max_handler_path, + "line_count": self.line_count, + "risk_score": self.risk_score, + "finding_codes": list(self.finding_codes), + "owner_ticket": self.owner_ticket, + } + + @classmethod + def from_dict(cls, data: dict[str, object]) -> BaselineEntry: + def _as_int(value: object) -> int: + if isinstance(value, int): + return value + if isinstance(value, str) and value.isdigit(): + return int(value) + return 0 + + codes_raw = data.get("finding_codes", []) + codes: tuple[str, ...] = ( + tuple(str(c) for c in codes_raw) if isinstance(codes_raw, list) else () + ) + return cls( + repo=str(data.get("repo", "")), + node=str(data["node"]), + max_handler_path=str(data.get("max_handler_path", "")), + line_count=_as_int(data.get("line_count", 0)), + risk_score=_as_int(data.get("risk_score", 0)), + finding_codes=codes, + owner_ticket=str(data.get("owner_ticket", "")), + ) + + +@dataclass +class ScanResult: + """Aggregated scan result across a set of node directories.""" + + repo: str + hard_fails: list[BaselineEntry] = field(default_factory=list) + #: Per-node finding detail for reporting (node -> messages). + report_lines: list[str] = field(default_factory=list) + + +# --------------------------------------------------------------------------- # +# Discovery +# --------------------------------------------------------------------------- # +def discover_node_dirs(repo_root: Path) -> list[Path]: + """Return every ``nodes//`` directory under ``repo_root`` with a contract.""" + src = repo_root / "src" + search_root = src if src.is_dir() else repo_root + node_dirs: set[Path] = set() + for contract in search_root.rglob("contract.yaml"): + parent = contract.parent + # A node directory's parent is conventionally ``nodes/``. + if parent.parent.name == "nodes" or "nodes" in parent.parts: + node_dirs.add(parent) + return sorted(node_dirs) + + +def node_dirs_for_changed_files( + repo_root: Path, changed_files: list[str] +) -> list[Path]: + """Map changed file paths to the node directories that contain them.""" + node_dirs: set[Path] = set() + for raw in changed_files: + if not raw: + continue + candidate = (repo_root / raw).resolve() + # Walk up until we find a directory holding a contract.yaml. + cur = candidate if candidate.is_dir() else candidate.parent + while True: + if (cur / "contract.yaml").is_file() and "nodes" in cur.parts: + node_dirs.add(cur) + break + if cur == repo_root or cur.parent == cur: + break + cur = cur.parent + return sorted(node_dirs) + + +# --------------------------------------------------------------------------- # +# Scanning +# --------------------------------------------------------------------------- # +def _relative_handler_path(handler_path: Path | None, repo_root: Path | None) -> str: + """Render a handler path repo-relative (never a machine-absolute path). + + Absolute ``/Users/...`` / ``/Volumes/...`` paths in committed artifacts are a + portability bug (Operating Rule #6). When ``repo_root`` is known and the + handler lives under it, the stored path is relative to the repo root. + """ + if handler_path is None: + return "" + if repo_root is not None: + try: + return str(handler_path.resolve().relative_to(repo_root.resolve())) + except ValueError: + pass + # Fall back to the path from the ``src/`` segment onward, never absolute. + parts = handler_path.parts + if "src" in parts: + idx = parts.index("src") + return str(Path(*parts[idx:])) + return handler_path.name + + +def scan_node_dirs( + repo: str, node_dirs: list[Path], repo_root: Path | None = None +) -> ScanResult: + """Run ARCH-004 over ``node_dirs`` and collect hard-fails + report lines. + + Args: + repo: Repo name recorded on each baseline entry. + node_dirs: Node directories to analyze. + repo_root: Repo root used to render ``max_handler_path`` repo-relative + (avoids committing machine-absolute paths). + """ + result = ScanResult(repo=repo) + for node_dir in node_dirs: + analysis = analyze_node_directory(node_dir) + if analysis is None: + continue + if not analysis.has_hard_fail: + continue + entry = BaselineEntry( + repo=repo, + node=analysis.node_name, + max_handler_path=_relative_handler_path( + analysis.max_handler_path, repo_root + ), + line_count=analysis.max_handler_lines, + risk_score=analysis.risk_score, + finding_codes=tuple(analysis.finding_codes), + owner_ticket=( + DELEGATION_OWNER_TICKET + if "delegation" in analysis.node_name + else FLEET_OWNER_TICKET + ), + ) + result.hard_fails.append(entry) + for v in analysis.violations(): + result.report_lines.append( + f" [{analysis.node_name}] {v.severity.value.upper()}: {v.message}" + ) + return result + + +# --------------------------------------------------------------------------- # +# Baseline I/O +# --------------------------------------------------------------------------- # +def _baseline_key(repo: str, node: str) -> str: + """Composite baseline key (``repo::node``) so cross-repo nodes never collide.""" + return f"{repo}::{node}" + + +def load_baseline(baseline_path: Path) -> dict[str, BaselineEntry]: + """Load the baseline keyed by ``repo::node``. Missing file => empty baseline.""" + if not baseline_path.is_file(): + return {} + raw = yaml.safe_load(baseline_path.read_text(encoding="utf-8")) + if not isinstance(raw, dict): + return {} + entries = raw.get("entries", []) + out: dict[str, BaselineEntry] = {} + if isinstance(entries, list): + for item in entries: + if isinstance(item, dict) and "node" in item: + entry = BaselineEntry.from_dict(item) + out[_baseline_key(entry.repo, entry.node)] = entry + return out + + +def render_baseline_yaml(repo: str, entries: list[BaselineEntry]) -> str: + """Render a baseline YAML document from scan entries.""" + doc = { + "schema_version": {"major": 1, "minor": 0, "patch": 0}, + "repo": repo, + "rule": "ARCH-004", + "description": ( + "Accepted imperative-orchestrator hard-fails for the ARCH-004 " + "ratchet (OMN-13472). One entry per current hard-fail. This list can " + "only SHRINK: a new/worsened/untracked finding on a touched node " + "fails the changed-node ratchet. Each entry carries an owner ticket " + "for its decomposition (delegation -> OMN-13471)." + ), + "entries": [e.to_dict() for e in sorted(entries, key=lambda e: e.node)], + } + header = ( + "# SPDX-FileCopyrightText: 2026 OmniNode.ai Inc.\n" + "# SPDX-License-Identifier: MIT\n" + "#\n" + "# ARCH-004 imperative-orchestrator ratchet baseline (OMN-13472).\n" + "# Generated by:\n" + "# uv run python -m omnibase_infra.nodes.node_architecture_validator." + "validators.scanner_imperative_orchestrator_ratchet \\\n" + "# --check-all --report --write-baseline\n" + "# Owner epic: OMN-13471. Wiring: OMN-12550 / OMN-13325.\n" + ) + return header + yaml.safe_dump(doc, sort_keys=False, default_flow_style=False) + + +def write_baseline( + baseline_path: Path, repo: str, entries: list[BaselineEntry] +) -> None: + """Write the baseline YAML to ``baseline_path`` (creating the dir if needed).""" + baseline_path.parent.mkdir(parents=True, exist_ok=True) + baseline_path.write_text(render_baseline_yaml(repo, entries), encoding="utf-8") + + +# --------------------------------------------------------------------------- # +# Ratchet comparison +# --------------------------------------------------------------------------- # +def ratchet_violations( + scanned: list[BaselineEntry], baseline: dict[str, BaselineEntry] +) -> list[str]: + """Return ratchet failures for the scanned (changed) hard-fails. + + A scanned hard-fail fails the ratchet when it: + * is NOT in the baseline (new / untracked hard-fail), OR + * worsens its risk score vs the baseline, OR + * materially grows its max handler line count vs the baseline, OR + * adds a finding code not present in the baseline entry, OR + * has a baseline entry with no owner ticket. + """ + failures: list[str] = [] + for entry in scanned: + base = baseline.get(_baseline_key(entry.repo, entry.node)) + if base is None: + failures.append( + f"{entry.node}: NEW imperative-orchestrator hard-fail " + f"(codes={list(entry.finding_codes)}, risk={entry.risk_score}, " + f"lines={entry.line_count}) — not in baseline. Decompose it or, " + f"if intentional debt, add a baseline entry with an owner ticket." + ) + continue + if entry.risk_score > base.risk_score: + failures.append( + f"{entry.node}: risk score WORSENED " + f"{base.risk_score} -> {entry.risk_score}." + ) + if entry.line_count > base.line_count: + failures.append( + f"{entry.node}: max handler GREW " + f"{base.line_count} -> {entry.line_count} lines." + ) + new_codes = set(entry.finding_codes) - set(base.finding_codes) + if new_codes: + failures.append( + f"{entry.node}: NEW finding codes {sorted(new_codes)} " + f"(baseline had {list(base.finding_codes)})." + ) + if not base.owner_ticket: + failures.append( + f"{entry.node}: baseline entry has no owner_ticket; " + f"every accepted hard-fail must cite a decomposition ticket." + ) + return failures + + +def strict_violations( + scanned: list[BaselineEntry], baseline: dict[str, BaselineEntry] +) -> list[str]: + """Return strict-mode failures: a baselined node that still hard-fails. + + In ``--strict`` mode the baseline is being dropped: any node still on the + baseline that the live scan reproduces is a failure (it must be remediated + before the baseline entry is removed). + """ + scanned_keys = {_baseline_key(e.repo, e.node) for e in scanned} + return [ + f"{entry.node}: still hard-fails ARCH-004 but is baselined; " + f"--strict requires the baseline be dropped (decompose it). " + f"Owner: {entry.owner_ticket or 'UNASSIGNED'}." + for key, entry in baseline.items() + if key in scanned_keys + ] + + +# --------------------------------------------------------------------------- # +# CLI +# --------------------------------------------------------------------------- # +def _resolve_repo_root(arg: str | None) -> Path: + if arg: + return Path(arg).resolve() + # Default: repo root = four parents up from this file + # (.../src/omnibase_infra/nodes/node_architecture_validator/validators/) + return Path(__file__).resolve().parents[5] + + +def main(argv: list[str] | None = None) -> int: + """CLI entrypoint for the ARCH-004 ratchet.""" + parser = argparse.ArgumentParser( + prog="imperative-orchestrator-ratchet", + description=( + "ARCH-004 imperative-orchestrator ratchet (OMN-13472). Scans node " + "directories for contract-declared-but-unbound orchestrator FSMs " + "driven imperatively by a handler." + ), + ) + mode = parser.add_mutually_exclusive_group(required=True) + mode.add_argument( + "--check-all", + action="store_true", + help="Scan every node directory in the repo.", + ) + mode.add_argument( + "--check-changed", + action="store_true", + help="Scan only node directories containing the given --files.", + ) + parser.add_argument( + "--report", + action="store_true", + help="Print one report row per hard-fail (does not block on its own).", + ) + parser.add_argument( + "--ratchet", + action="store_true", + help="Fail on a new/worsened/untracked finding for a touched node.", + ) + parser.add_argument( + "--strict", + action="store_true", + help="Drop-baseline mode: a baselined node that still hard-fails fails.", + ) + parser.add_argument( + "--write-baseline", + action="store_true", + help="Write/refresh the baseline YAML from a --check-all scan.", + ) + parser.add_argument( + "--repo-root", + default=None, + help="Repo root (defaults to this repo).", + ) + parser.add_argument( + "--repo-name", + default="omnibase_infra", + help="Repo name recorded in baseline entries.", + ) + parser.add_argument( + "files", + nargs="*", + help="Changed files (relative to repo root) for --check-changed.", + ) + args = parser.parse_args(argv) + + repo_root = _resolve_repo_root(args.repo_root) + baseline_path = repo_root / BASELINE_RELATIVE_PATH + baseline = load_baseline(baseline_path) + + if args.check_all: + node_dirs = discover_node_dirs(repo_root) + else: + node_dirs = node_dirs_for_changed_files(repo_root, args.files) + + result = scan_node_dirs(args.repo_name, node_dirs, repo_root=repo_root) + + if args.report or args.check_all: + print( + f"ARCH-004 imperative-orchestrator scan ({args.repo_name}): " + f"{len(result.hard_fails)} hard-fail node(s) across " + f"{len(node_dirs)} node dir(s)." + ) + for line in result.report_lines: + print(line) + + if args.write_baseline: + write_baseline(baseline_path, args.repo_name, result.hard_fails) + print(f"Wrote baseline: {baseline_path}") + return 0 + + failures: list[str] = [] + if args.strict: + failures.extend(strict_violations(result.hard_fails, baseline)) + if args.ratchet: + failures.extend(ratchet_violations(result.hard_fails, baseline)) + + if failures: + print("\nARCH-004 ratchet FAILED:", file=sys.stderr) + for f in failures: + print(f" - {f}", file=sys.stderr) + return 1 + + return 0 + + +__all__ = [ + "BaselineEntry", + "ScanResult", + "discover_node_dirs", + "node_dirs_for_changed_files", + "scan_node_dirs", + "load_baseline", + "render_baseline_yaml", + "write_baseline", + "ratchet_violations", + "strict_violations", + "main", + "BASELINE_RELATIVE_PATH", + "DELEGATION_OWNER_TICKET", +] + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/src/omnibase_infra/nodes/node_architecture_validator/validators/validator_contract_declared_orchestrator_workflow.py b/src/omnibase_infra/nodes/node_architecture_validator/validators/validator_contract_declared_orchestrator_workflow.py new file mode 100644 index 0000000000..21d7de90b8 --- /dev/null +++ b/src/omnibase_infra/nodes/node_architecture_validator/validators/validator_contract_declared_orchestrator_workflow.py @@ -0,0 +1,794 @@ +# SPDX-FileCopyrightText: 2025 OmniNode.ai Inc. +# SPDX-License-Identifier: MIT +"""Validator for ARCH-004: Contract-Declared Orchestrator Workflow Must Be Bound To An Executor. + +This validator detects the *imperative-orchestrator* anti-pattern that ARCH-003 +is structurally incapable of catching (see ``validator_no_orchestrator_fsm.py``). + +Why ARCH-003 misses it (proven by ``node_delegation_orchestrator``): + - ARCH-003 is a single-file AST visitor that only inspects classes whose + *name* contains ``"Orchestrator"``. The delegation FSM lives in class + ``HandlerDelegationWorkflow`` (no "Orchestrator" in the name) → never + inspected. + - ARCH-003's method match-set is ``{transition, can_transition, + apply_transition, get_current_state, set_state}``. The delegation handler + drives state via ``_transition`` (leading underscore, and predominantly a + *call* not a *def*) → never matched. + - ARCH-003 never joins ``contract.yaml`` + ``handler_routing`` + handler + source. It cannot tell "declared FSM table but no executor binds to it, + while a handler drives the transitions itself." + +ARCH-004 is therefore a **cross-file, node-directory** rule. Per node directory +it joins four sources: + 1. ``contract.yaml`` — ``fsm:`` / ``workflow_coordination`` / node_type / name. + 2. ``handler_routing.routing_strategy`` — ``payload_type_match`` catchall. + 3. Handler source under ``handlers/`` — state-driving symbols. + 4. Node type/name — orchestrator-like classification + reducer exemption. + +Canonical truth (verified 2026-06-22, OMN-13472): + ``ModelContractOrchestrator`` + (``omnibase_core/models/contracts/model_contract_orchestrator.py``) has only + ``workflow_coordination: ModelWorkflowConfig`` — **no typed + ``fsm``/``state_machine``/``transition`` field** — and no traverser consumes + orchestrator ``fsm.transitions``. So a declared ``fsm:`` block on an + orchestrator contract is *decorative*: nothing executes it. FSM execution + exists only for reducers (``node_*_fsm_reducer`` + contract ``state_machine:`` + + a pure transition executor) — which are EXEMPT. + +Hard-fail signals (for changed/new orchestrator-like nodes): + H1. Decorative-FSM-with-handler-driven-transitions: contract declares an + ``fsm:``/workflow-state set, no executor binds to it, AND a handler drives + the transitions itself (``_transition(`` etc. or an ``Enum*State`` whose + members ≈ the contract ``fsm.states``). + H2. ``handler_routing.routing_strategy: payload_type_match`` funnels 3+ + event/payload types into a single workflow handler. + H3. One handler both selects the next workflow state AND constructs + terminal/compat events (the OMN-13408 footgun). + +Warning / baseline-score signals: + W1. Handler > 750 lines for an orchestrator-like node. + W2. Handler > 125 branch/control markers. + W3. >= 10 publish/event-construction markers (regex reported in + ``EVENT_MARKER_REGEX``). + W4. Declared states/transitions are not a typed, executor-bound schema. + +Related: + - Ticket: OMN-13472 (ARCH-004 ratchet — Workstream B) + - Epic: OMN-13471 (decompose imperative delegation orchestrator) + - OMN-12550 (wire ARCH-001/002/003 as blocking gates — ARCH-004 rides it) + - OMN-13325 (ratchet-enforcement epic) + - Plan: docs/plans/2026-06-22-imperative-orchestrator-ratchet-and-recovery-plan-verified.md §4 + - Audit: docs/audits/2026-06-22-imperative-orchestrator-audit.md +""" + +from __future__ import annotations + +import ast +import re +from pathlib import Path +from typing import TYPE_CHECKING + +import yaml + +from omnibase_infra.enums import EnumValidationSeverity +from omnibase_infra.nodes.node_architecture_validator.models.model_architecture_violation import ( + ModelArchitectureViolation, +) +from omnibase_infra.nodes.node_architecture_validator.models.model_validation_result import ( + ModelFileValidationResult, +) + +if TYPE_CHECKING: + from omnibase_infra.nodes.node_architecture_validator.models import ( + ModelRuleCheckResult, + ) + +RULE_ID = "ARCH-004" +RULE_NAME = "Contract-Declared Orchestrator Workflow Must Be Bound To An Executor" + +# --- Detection thresholds (warning / baseline score) ---------------------- + +#: A handler larger than this for an orchestrator-like node is debt (W1). +HANDLER_LINE_WARN_THRESHOLD = 750 +#: More branch/control markers than this in a single handler is debt (W2). +BRANCH_MARKER_WARN_THRESHOLD = 125 +#: At least this many publish/event-construction markers is debt (W3). +EVENT_MARKER_WARN_THRESHOLD = 10 +#: payload_type_match funneling this many distinct payload types into one +#: workflow handler is a hard fail (H2). +PAYLOAD_FANIN_HARD_THRESHOLD = 3 + +# --- Regexes (reported in violation details so they can be re-derived) ----- + +#: Publish / event-construction markers (W3 / H3). Reported verbatim. +EVENT_MARKER_REGEX = r"\.publish\(|ModelEventEnvelope|\.emit\(|[A-Za-z_]+Event\(" +#: Branch / control-flow markers (W2). Reported verbatim. +BRANCH_MARKER_REGEX = r"\b(if|elif|for|while|except|case|and|or)\b" + +#: Handler symbols that prove the handler is itself driving an FSM (H1). +#: +#: These are the *human-readable* symbol names reported in violation detail and +#: used by the synthetic-fixture documentation. Detection is performed by +#: :data:`HANDLER_STATE_DRIVE_REGEX` below, which is precise (boundary-anchored) +#: to avoid the substring traps that the audit flagged with ``0*`` — e.g. +#: ``record_phase_transition(`` must NOT match ``_transition(``, and a log-format +#: string ``"current_state=%s"`` must NOT match ``current_state``. +#: +#: NOTE: ``_transition`` (leading underscore) is included precisely because +#: ARCH-003 omits it. +HANDLER_STATE_DRIVE_SYMBOLS: tuple[str, ...] = ( + "_transition(", + "transition(", + "can_transition", + "apply_transition", + "set_state", + "current_state", + "next_state", +) + +#: Boundary-anchored detection for the H1 state-drive symbols. +#: +#: * A ``transition(``/``_transition(`` call must be a *method call* — preceded +#: by ``self.`` / ``cls.`` / an identifier-dot, with a word boundary at the +#: start of the method name so ``record_phase_transition(`` (matches as +#: ``record_phase_`` + ``transition(``) does NOT count. We anchor on +#: ``(?:self|cls)\.(?:_)?transition\(`` and ``\b_transition\(`` is excluded +#: when preceded by another word char. +#: * ``can_transition`` / ``apply_transition`` / ``set_state`` are method names; +#: require a leading ``.`` or ``def ``. +#: * ``current_state`` / ``next_state`` count only as an *assignment target* or +#: attribute access (``self.current_state``, ``.current_state =``), never +#: inside a string literal. +#: +#: Each entry maps a reported symbol name -> its precise regex. +HANDLER_STATE_DRIVE_REGEX: dict[str, str] = { + "_transition(": r"(?:self|cls)\.\s*_transition\s*\(", + "transition(": r"(?:self|cls)\.\s*transition\s*\(", + "can_transition": r"(?:\.|def\s+)can_transition\b", + "apply_transition": r"(?:\.|def\s+)apply_transition\b", + "set_state": r"(?:\.|def\s+)set_state\b", + "current_state": r"(?:self|cls)\.current_state\s*=", + "next_state": r"(?:self|cls)\.next_state\s*=", +} + +#: Names on ``ModelContractOrchestrator`` (and any contract field) that would +#: indicate a *typed, executor-bound* workflow schema. Presence in the contract +#: under one of these keys means the FSM is bound (not decorative) → not H1. +EXECUTOR_BOUND_CONTRACT_KEYS: tuple[str, ...] = ( + "state_machine", + "execution_graph", + "workflow_dag", + "execution_dag", +) + +#: Fraction of contract fsm.states that must appear as Enum*State members for +#: the enum to be treated as the handler's parallel FSM (H1, approximate match). +ENUM_STATE_MATCH_RATIO = 0.6 + + +# --------------------------------------------------------------------------- # +# Node-directory analysis result (internal, not a public model) +# --------------------------------------------------------------------------- # +class _NodeAnalysis: + """Mutable accumulator for one node directory's cross-file analysis. + + This is an internal scratch object; the public surface is the list of + :class:`ModelArchitectureViolation` it produces via :meth:`violations`. + """ + + def __init__(self, node_dir: Path, contract_path: Path) -> None: + self.node_dir = node_dir + self.contract_path = contract_path + self.node_name: str = node_dir.name + self.node_type: str = "" + self.is_orchestrator_like: bool = False + self.is_reducer: bool = False + self.has_fsm_block: bool = False + self.fsm_states: tuple[str, ...] = () + self.has_executor_bound_workflow: bool = False + self.routing_strategy: str = "" + self.payload_type_count: int = 0 + self.max_handler_path: Path | None = None + self.max_handler_lines: int = 0 + self.handler_drives_state: bool = False + self.handler_state_drive_symbol: str = "" + self.enum_state_match: bool = False + self.handler_branch_markers: int = 0 + self.handler_event_markers: int = 0 + self.handler_selects_state_and_emits: bool = False + # Accumulated finding codes (e.g. "H1", "W3") for the baseline. + self.finding_codes: list[str] = [] + self._violations: list[ModelArchitectureViolation] = [] + + # -- violation helpers -------------------------------------------------- # + def _add( + self, + *, + code: str, + severity: EnumValidationSeverity, + message: str, + suggestion: str, + location: Path | None = None, + ) -> None: + self.finding_codes.append(code) + loc = str(location) if location is not None else str(self.contract_path) + self._violations.append( + ModelArchitectureViolation( + rule_id=RULE_ID, + rule_name=RULE_NAME, + severity=severity, + target_type="orchestrator_node", + target_name=self.node_name, + message=f"[{code}] {message}", + location=loc, + suggestion=suggestion, + details={ + "finding_code": code, + "node_dir": str(self.node_dir), + "node_type": self.node_type, + "routing_strategy": self.routing_strategy, + "event_marker_regex": EVENT_MARKER_REGEX, + "branch_marker_regex": BRANCH_MARKER_REGEX, + }, + ) + ) + + def violations(self) -> list[ModelArchitectureViolation]: + return self._violations + + @property + def has_hard_fail(self) -> bool: + return any(v.severity == EnumValidationSeverity.ERROR for v in self._violations) + + @property + def risk_score(self) -> int: + """Transparent additive score mirroring the audit doc. + + ``fsm:+2 · payload_type_match:+2 · handler_drives_state:+3 · + handler>1500:+2 (>750:+1) · event_markers>=10:+1``. + """ + score = 0 + if self.has_fsm_block: + score += 2 + if self.routing_strategy == "payload_type_match": + score += 2 + if self.handler_drives_state: + score += 3 + if self.max_handler_lines > 1500: + score += 2 + elif self.max_handler_lines > HANDLER_LINE_WARN_THRESHOLD: + score += 1 + if self.handler_event_markers >= EVENT_MARKER_WARN_THRESHOLD: + score += 1 + return score + + +# --------------------------------------------------------------------------- # +# Helpers +# --------------------------------------------------------------------------- # +def _is_orchestrator_like(node_name: str, node_type: str) -> bool: + """Classify whether a node directory is orchestrator-like. + + Joins node directory name AND declared ``node_type`` so a misnamed + orchestrator (or a generic orchestrator type) is still caught. + """ + name_signals = "orchestrator" in node_name.lower() + type_signals = "ORCHESTRATOR" in node_type.upper() + return name_signals or type_signals + + +def _is_reducer(node_name: str, node_type: str) -> bool: + """Reducers (``node_*_fsm_reducer`` / ``REDUCER_*``) are EXEMPT.""" + name_signals = "reducer" in node_name.lower() + type_signals = "REDUCER" in node_type.upper() + return name_signals or type_signals + + +def _contract_has_executor_bound_workflow(contract: dict[str, object]) -> bool: + """Return True if the contract declares a typed, executor-bound workflow. + + A reducer-style ``state_machine:`` block, or an explicit execution-graph / + workflow DAG key, counts as executor-bound. A bare orchestrator ``fsm:`` + block does NOT — ``ModelContractOrchestrator`` has no typed field that binds + it, so nothing executes it. + """ + for key in EXECUTOR_BOUND_CONTRACT_KEYS: + if contract.get(key): + return True + # workflow_coordination with an explicit, non-empty execution graph / steps + wc = contract.get("workflow_coordination") + if isinstance(wc, dict): + for key in ("execution_graph", "workflow_dag", "steps", "dag"): + if wc.get(key): + return True + return False + + +def _extract_fsm_states(contract: dict[str, object]) -> tuple[bool, tuple[str, ...]]: + """Extract (has_fsm_block, states) from the contract. + + Detects a declared workflow-state set under either ``fsm:`` (orchestrator + convention) — a ``state_machine:`` block is handled as executor-bound and + does not count here. + """ + fsm = contract.get("fsm") + if not isinstance(fsm, dict): + return False, () + states = fsm.get("states") + if isinstance(states, list): + return True, tuple(str(s) for s in states) + return True, () + + +def _enum_states_match_fsm(handler_source: str, fsm_states: tuple[str, ...]) -> bool: + """Return True if an ``Enum*State`` class in the handler ≈ contract fsm.states. + + Approximate match: an Enum whose name contains "State" defines >= 60% of the + contract's declared fsm states as members. This catches the delegation + ``EnumDelegationState`` (8 members == the 8 contract states) even when the + enum lives in a sibling ``enums.py`` re-exported into the handler. + """ + if not fsm_states: + return False + try: + tree = ast.parse(handler_source) + except SyntaxError: + return False + state_set = {s.upper() for s in fsm_states} + threshold = max(1, int(len(state_set) * ENUM_STATE_MATCH_RATIO)) + for node in ast.walk(tree): + if isinstance(node, ast.ClassDef) and "State" in node.name: + members = { + t.id.upper() + for item in node.body + if isinstance(item, ast.Assign) + for t in item.targets + if isinstance(t, ast.Name) + } + if len(state_set & members) >= threshold: + return True + return False + + +def _count_payload_types(handler_routing: dict[str, object]) -> int: + """Count distinct ``event_model`` payload types in the routing table.""" + handlers = handler_routing.get("handlers") + if not isinstance(handlers, list): + return 0 + payload_names: set[str] = set() + for entry in handlers: + if not isinstance(entry, dict): + continue + event_model = entry.get("event_model") + if isinstance(event_model, dict): + name = event_model.get("name") + if isinstance(name, str): + payload_names.add(name) + elif isinstance(event_model, str): + payload_names.add(event_model) + return len(payload_names) + + +def _largest_handler(handlers_dir: Path) -> tuple[Path | None, int, str]: + """Return (path, line_count, source) of the largest handler module.""" + best_path: Path | None = None + best_lines = 0 + best_source = "" + if not handlers_dir.is_dir(): + return None, 0, "" + for py in sorted(handlers_dir.glob("handler_*.py")): + try: + source = py.read_text(encoding="utf-8") + except (OSError, UnicodeDecodeError): + continue + lines = source.count("\n") + 1 + if lines > best_lines: + best_path, best_lines, best_source = py, lines, source + return best_path, best_lines, best_source + + +def _detect_state_drive_symbol(source: str) -> str: + """Return the first precise state-drive symbol matched in ``source``, or "". + + Uses :data:`HANDLER_STATE_DRIVE_REGEX` (boundary-anchored) so substring + traps such as ``record_phase_transition(`` or a ``"current_state=%s"`` log + format string do not produce false positives. + """ + for symbol, pattern in HANDLER_STATE_DRIVE_REGEX.items(): + if re.search(pattern, source) is not None: + return symbol + return "" + + +def _handler_selects_state_and_emits(source: str) -> bool: + """H3: one handler both drives a state transition AND constructs an event. + + Detected as the *co-occurrence*, in one handler module, of a (precisely + matched) state-driving symbol and an event-construction marker. + """ + drives = _detect_state_drive_symbol(source) != "" + emits = re.search(EVENT_MARKER_REGEX, source) is not None + return drives and emits + + +# --------------------------------------------------------------------------- # +# Core node-directory analysis +# --------------------------------------------------------------------------- # +def analyze_node_directory(node_dir: Path) -> _NodeAnalysis | None: + """Run the ARCH-004 cross-file join over a single node directory. + + Args: + node_dir: Path to a ``nodes//`` directory containing a + ``contract.yaml``. String callers (CLI / the protocol rule) convert + to :class:`~pathlib.Path` before calling. + + Returns: + A populated :class:`_NodeAnalysis` (with any violations), or ``None`` if + the directory is not an analyzable node (no contract) or is exempt + (reducer / not orchestrator-like). + """ + node_path = Path(node_dir) + contract_path = node_path / "contract.yaml" + if not contract_path.is_file(): + return None + + try: + raw = yaml.safe_load(contract_path.read_text(encoding="utf-8")) + except (yaml.YAMLError, OSError, UnicodeDecodeError): + return None + if not isinstance(raw, dict): + return None + + analysis = _NodeAnalysis(node_path, contract_path) + analysis.node_type = str(raw.get("node_type", "")) + analysis.node_name = str(raw.get("name", node_path.name)) + + analysis.is_reducer = _is_reducer(analysis.node_name, analysis.node_type) + analysis.is_orchestrator_like = _is_orchestrator_like( + analysis.node_name, analysis.node_type + ) + + # Reducers are EXEMPT (they legitimately own state machines). + if analysis.is_reducer: + return None + # Only orchestrator-like nodes are in scope. + if not analysis.is_orchestrator_like: + return None + + # --- contract facts ---------------------------------------------------- # + analysis.has_fsm_block, analysis.fsm_states = _extract_fsm_states(raw) + analysis.has_executor_bound_workflow = _contract_has_executor_bound_workflow(raw) + + handler_routing = raw.get("handler_routing") + if isinstance(handler_routing, dict): + analysis.routing_strategy = str(handler_routing.get("routing_strategy", "")) + analysis.payload_type_count = _count_payload_types(handler_routing) + + # --- handler facts ----------------------------------------------------- # + handlers_dir = node_path / "handlers" + handler_path, handler_lines, handler_source = _largest_handler(handlers_dir) + analysis.max_handler_path = handler_path + analysis.max_handler_lines = handler_lines + + if handler_source: + matched_symbol = _detect_state_drive_symbol(handler_source) + if matched_symbol: + analysis.handler_drives_state = True + analysis.handler_state_drive_symbol = matched_symbol + analysis.enum_state_match = _enum_states_match_fsm( + handler_source, analysis.fsm_states + ) + analysis.handler_branch_markers = len( + re.findall(BRANCH_MARKER_REGEX, handler_source) + ) + analysis.handler_event_markers = len( + re.findall(EVENT_MARKER_REGEX, handler_source) + ) + analysis.handler_selects_state_and_emits = _handler_selects_state_and_emits( + handler_source + ) + # An Enum*State that mirrors the contract states also counts as the handler + # owning a parallel FSM (covers the case where transitions are not literal + # ``_transition(`` calls but the enum drives state). + if analysis.enum_state_match: + analysis.handler_drives_state = True + if not analysis.handler_state_drive_symbol: + analysis.handler_state_drive_symbol = "Enum*State≈fsm.states" + + _evaluate_signals(analysis) + return analysis + + +def _evaluate_signals(a: _NodeAnalysis) -> None: + """Apply the hard-fail and warning signals to a populated analysis.""" + handler_loc = a.max_handler_path or a.contract_path + + # H1: decorative FSM + handler-driven transitions. + if a.has_fsm_block and not a.has_executor_bound_workflow and a.handler_drives_state: + a._add( + code="H1", + severity=EnumValidationSeverity.ERROR, + message=( + f"Orchestrator declares an fsm: state set " + f"({len(a.fsm_states)} states) but no runtime executor binds to " + f"it (no typed fsm/state_machine field on ModelContractOrchestrator, " + f"no traverser consumes its transitions), WHILE the handler drives " + f"the transitions itself " + f"(symbol '{a.handler_state_drive_symbol}'). The declared " + f"transition table is decorative; the handler owns a parallel " + f"imperative FSM." + ), + suggestion=( + "Bind the contract fsm to an executor (migrate to a typed, " + "executor-bound workflow/DAG field per OMN-12835), or decompose " + "the handler into per-step thin handlers driven by a reducer. " + "See OMN-13471 (delegation decomposition epic)." + ), + location=handler_loc, + ) + + # H2: payload_type_match fan-in. + if ( + a.routing_strategy == "payload_type_match" + and a.payload_type_count >= PAYLOAD_FANIN_HARD_THRESHOLD + ): + a._add( + code="H2", + severity=EnumValidationSeverity.ERROR, + message=( + f"handler_routing.routing_strategy=payload_type_match funnels " + f"{a.payload_type_count} distinct payload/event types into a " + f"single workflow handler (>= {PAYLOAD_FANIN_HARD_THRESHOLD})." + ), + suggestion=( + "Split the catchall into per-event/per-step handlers " + "(operation_match or one handler per payload type)." + ), + location=handler_loc, + ) + + # H3: one handler selects next state AND constructs terminal/compat events. + if a.handler_selects_state_and_emits: + a._add( + code="H3", + severity=EnumValidationSeverity.ERROR, + message=( + "A single handler both selects the next workflow state and " + "constructs terminal/compat events (the OMN-13408 footgun): " + "dual terminal/compat event construction co-located with state " + "selection clobbers projection fields." + ), + suggestion=( + "Move terminal/compat event construction behind one " + "contract-owned builder/reducer; keep state selection separate " + "from event emission." + ), + location=handler_loc, + ) + + # W1: oversized handler. + if a.max_handler_lines > HANDLER_LINE_WARN_THRESHOLD: + a._add( + code="W1", + severity=EnumValidationSeverity.WARNING, + message=( + f"Orchestrator max handler is {a.max_handler_lines} lines " + f"(> {HANDLER_LINE_WARN_THRESHOLD}); orchestrator handlers should " + f"be thin per-step coordinators." + ), + suggestion="Decompose into per-step handlers.", + location=handler_loc, + ) + + # W2: branch/control density. + if a.handler_branch_markers > BRANCH_MARKER_WARN_THRESHOLD: + a._add( + code="W2", + severity=EnumValidationSeverity.WARNING, + message=( + f"Handler has {a.handler_branch_markers} branch/control markers " + f"(> {BRANCH_MARKER_WARN_THRESHOLD}; regex {BRANCH_MARKER_REGEX!r})." + ), + suggestion="Reduce branching by extracting decision compute nodes.", + location=handler_loc, + ) + + # W3: event-construction density. + if a.handler_event_markers >= EVENT_MARKER_WARN_THRESHOLD: + a._add( + code="W3", + severity=EnumValidationSeverity.WARNING, + message=( + f"Handler has {a.handler_event_markers} publish/event-construction " + f"markers (>= {EVENT_MARKER_WARN_THRESHOLD}; regex " + f"{EVENT_MARKER_REGEX!r})." + ), + suggestion="Centralize event construction in a contract-owned builder.", + location=handler_loc, + ) + + # W4: declared-but-not-executor-bound workflow states. + if a.has_fsm_block and not a.has_executor_bound_workflow: + a._add( + code="W4", + severity=EnumValidationSeverity.WARNING, + message=( + "Contract declares workflow states/transitions that are not a " + "typed, executor-bound schema (decorative fsm:)." + ), + suggestion=( + "Adopt the typed executor-bound workflow/DAG field (OMN-12835)." + ), + location=a.contract_path, + ) + + +def validate_contract_declared_orchestrator_workflow( + node_dir: str, +) -> ModelFileValidationResult: + """Validate a single node directory against ARCH-004. + + Args: + node_dir: Path to a ``nodes//`` directory. + + Returns: + ``ModelFileValidationResult`` with ``valid=False`` only when a HARD-fail + (ERROR) signal fires. Warnings do not flip ``valid`` (they feed the + baseline score). + """ + analysis = analyze_node_directory(Path(node_dir)) + if analysis is None: + return ModelFileValidationResult( + valid=True, + violations=[], + files_checked=0, + rules_checked=[RULE_ID], + ) + violations = analysis.violations() + return ModelFileValidationResult( + valid=not analysis.has_hard_fail, + violations=violations, + files_checked=1, + rules_checked=[RULE_ID], + ) + + +class RuleContractDeclaredOrchestratorWorkflow: + """Protocol-compliant rule: contract-declared orchestrator workflow must be bound. + + This rule implements :class:`ProtocolArchitectureRule`. Unlike the per-file + ARCH-001/002/003 rules, ARCH-004's ``target`` is a **node directory** (the + parent of a ``contract.yaml``). It accepts either the directory path or a + path to a ``contract.yaml`` (the directory parent is used). + + Thread Safety: + Stateless; safe for concurrent use. + """ + + @property + def rule_id(self) -> str: + """Return the canonical rule ID matching contract.yaml.""" + return RULE_ID + + @property + def name(self) -> str: + """Return human-readable rule name.""" + return RULE_NAME + + @property + def description(self) -> str: + """Return detailed rule description.""" + return ( + "An orchestrator-like node that declares an fsm:/workflow-state set " + "must bind it to a runtime executor. It must not leave the contract " + "table decorative while a handler drives the transitions itself " + "(_transition(...), an Enum*State mirroring the contract states, or a " + "payload_type_match catchall funneling 3+ payloads into one handler). " + "Reducers are exempt." + ) + + @property + def severity(self) -> EnumValidationSeverity: + """Return severity level for violations of this rule.""" + return EnumValidationSeverity.ERROR + + def check(self, target: object) -> ModelRuleCheckResult: + """Check a node directory (or its contract.yaml) against ARCH-004. + + Args: + target: A node directory path, or a path to a ``contract.yaml``. + Other types return ``skipped=True``. + + Returns: + ``ModelRuleCheckResult`` indicating pass/fail. When multiple + violations exist, the first ERROR (or first violation) is surfaced; + ``details["total_violations"]`` and ``details["finding_codes"]`` + carry the full picture. + """ + from omnibase_infra.nodes.node_architecture_validator.models import ( + ModelRuleCheckResult, + ) + + node_dir = self._resolve_node_dir(target) + if node_dir is None: + return ModelRuleCheckResult( + passed=True, + rule_id=self.rule_id, + skipped=True, + reason="Target is not a node directory or contract.yaml path", + ) + + analysis = analyze_node_directory(node_dir) + if analysis is None: + return ModelRuleCheckResult( + passed=True, + rule_id=self.rule_id, + skipped=True, + reason="Directory is not an in-scope orchestrator node (or is exempt)", + ) + + if not analysis.has_hard_fail: + return ModelRuleCheckResult( + passed=True, + rule_id=self.rule_id, + details={ + "finding_codes": list(analysis.finding_codes), + "risk_score": analysis.risk_score, + "max_handler_lines": analysis.max_handler_lines, + }, + ) + + # Surface the first ERROR-severity violation. + error_violations = [ + v + for v in analysis.violations() + if v.severity == EnumValidationSeverity.ERROR + ] + violation = error_violations[0] + return ModelRuleCheckResult( + passed=False, + rule_id=self.rule_id, + message=violation.message, + details={ + "target_name": violation.target_name, + "target_type": violation.target_type, + "location": violation.location, + "suggestion": violation.suggestion, + "finding_codes": list(analysis.finding_codes), + "risk_score": analysis.risk_score, + "max_handler_lines": analysis.max_handler_lines, + "max_handler_path": ( + str(analysis.max_handler_path) + if analysis.max_handler_path + else None + ), + "total_violations": len(analysis.violations()), + }, + ) + + @staticmethod + def _resolve_node_dir(target: object) -> Path | None: + """Resolve ``target`` to a node directory, or None if not applicable.""" + if not isinstance(target, str | Path): + return None + path = Path(target) + if path.name == "contract.yaml": + return path.parent + if path.is_dir() and (path / "contract.yaml").is_file(): + return path + return None + + +__all__ = [ + "validate_contract_declared_orchestrator_workflow", + "analyze_node_directory", + "RuleContractDeclaredOrchestratorWorkflow", + "RULE_ID", + "RULE_NAME", + "EVENT_MARKER_REGEX", + "BRANCH_MARKER_REGEX", + "HANDLER_LINE_WARN_THRESHOLD", + "BRANCH_MARKER_WARN_THRESHOLD", + "EVENT_MARKER_WARN_THRESHOLD", + "PAYLOAD_FANIN_HARD_THRESHOLD", +] diff --git a/tests/unit/nodes/node_architecture_validator/test_validator_contract_declared_orchestrator_workflow.py b/tests/unit/nodes/node_architecture_validator/test_validator_contract_declared_orchestrator_workflow.py new file mode 100644 index 0000000000..bebf521a7f --- /dev/null +++ b/tests/unit/nodes/node_architecture_validator/test_validator_contract_declared_orchestrator_workflow.py @@ -0,0 +1,634 @@ +# SPDX-FileCopyrightText: 2025 OmniNode.ai Inc. +# SPDX-License-Identifier: MIT +"""Tests for ARCH-004: Contract-Declared Orchestrator Workflow Must Be Bound To An Executor. + +ARCH-004 is the cross-file, node-directory rule that catches the imperative +orchestrator anti-pattern ARCH-003 structurally cannot (OMN-13472, epic +OMN-13471). + +The HEADLINE proof (``test_arch003_passes_but_arch004_fails_delegation_shape``): +a synthetic node whose contract declares an fsm: table and whose monolithic +handler drives transitions via ``self._transition(...)`` — where the handler +class is NOT named ``*Orchestrator`` — is PASSED by ARCH-003 (single-file, +class-name-gated AST) but FAILED by ARCH-004 (cross-file join). + +All fixtures are synthetic, vendored mini node directories under ``tmp_path``; +no test depends on a sibling-repo path. +""" + +from __future__ import annotations + +from pathlib import Path + +import pytest + +from omnibase_infra.enums import EnumValidationSeverity +from omnibase_infra.nodes.node_architecture_validator.validators.scanner_imperative_orchestrator_ratchet import ( + BaselineEntry, + discover_node_dirs, + load_baseline, + ratchet_violations, + render_baseline_yaml, + scan_node_dirs, + strict_violations, + write_baseline, +) +from omnibase_infra.nodes.node_architecture_validator.validators.validator_contract_declared_orchestrator_workflow import ( + EVENT_MARKER_REGEX, + RULE_ID, + RuleContractDeclaredOrchestratorWorkflow, + analyze_node_directory, + validate_contract_declared_orchestrator_workflow, +) +from omnibase_infra.nodes.node_architecture_validator.validators.validator_no_orchestrator_fsm import ( + validate_no_orchestrator_fsm, +) + +pytestmark = pytest.mark.unit + + +# --------------------------------------------------------------------------- # +# Synthetic node-directory fixture builder +# --------------------------------------------------------------------------- # +def _make_node_dir( + root: Path, + *, + node_name: str, + contract_yaml: str, + handler_name: str = "handler_workflow", + handler_source: str | None = None, + enums_source: str | None = None, +) -> Path: + """Create a synthetic ``nodes//`` directory under ``root``. + + Returns the node directory path. + """ + node_dir = root / "src" / "pkg" / "nodes" / node_name + handlers_dir = node_dir / "handlers" + handlers_dir.mkdir(parents=True, exist_ok=True) + (node_dir / "contract.yaml").write_text(contract_yaml, encoding="utf-8") + if handler_source is not None: + (handlers_dir / f"{handler_name}.py").write_text( + handler_source, encoding="utf-8" + ) + if enums_source is not None: + (node_dir / "enums.py").write_text(enums_source, encoding="utf-8") + return node_dir + + +# Contract fragment: orchestrator with a decorative fsm: table (no executor bind). +_DELEGATION_SHAPE_CONTRACT = """\ +name: node_synth_delegation +node_type: ORCHESTRATOR_GENERIC +fsm: + states: + - RECEIVED + - ROUTED + - EXECUTING + - INFERENCE_COMPLETED + - GATE_EVALUATED + - ESCALATING + - COMPLETED + - FAILED + initial_state: RECEIVED + terminal_states: + - COMPLETED + - FAILED + transitions: + - from: RECEIVED + to: ROUTED + trigger: accepted +handler_routing: + routing_strategy: payload_type_match + handlers: + - operation: orchestrate + handler: + name: HandlerDelegationWorkflow + module: pkg.nodes.node_synth_delegation.handlers.handler_workflow + event_model: + name: ModelDelegationRequest + module: pkg.models +""" + +# Monolithic handler whose CLASS NAME is NOT *Orchestrator, driving state via +# self._transition(...) — the exact gap ARCH-003 misses. Also constructs a +# terminal event (H3). +_DELEGATION_SHAPE_HANDLER = '''\ +"""Synthetic monolithic delegation workflow handler.""" +from __future__ import annotations + + +class HandlerDelegationWorkflow: + """Handler whose name lacks "Orchestrator" but owns an imperative FSM.""" + + def _transition(self, workflow, new_state): + workflow.state = new_state + + def handle(self, envelope): + workflow = envelope.payload + self._transition(workflow, "ROUTED") + self._transition(workflow, "EXECUTING") + self._transition(workflow, "COMPLETED") + # H3: constructs a terminal event in the same handler that drives state. + return ModelTaskDelegatedEvent(state=workflow.state) +''' + + +def test_arch003_passes_but_arch004_fails_delegation_shape(tmp_path: Path) -> None: + """HEADLINE: ARCH-003 passes the delegation-shape handler; ARCH-004 fails it. + + (a) Synthetic node: contract fsm table + monolithic handler using + ``self._transition(`` whose class is NOT named ``*Orchestrator``. + """ + node_dir = _make_node_dir( + tmp_path, + node_name="node_synth_delegation", + contract_yaml=_DELEGATION_SHAPE_CONTRACT, + handler_source=_DELEGATION_SHAPE_HANDLER, + ) + handler_file = node_dir / "handlers" / "handler_workflow.py" + + # ARCH-003 (single-file, class-name-gated AST) PASSES the handler: + # class is HandlerDelegationWorkflow (no "Orchestrator"); the method match + # set does not include the *call* self._transition(...). + arch003 = validate_no_orchestrator_fsm(str(handler_file)) + assert arch003.valid, ( + "ARCH-003 is expected to PASS the delegation-shape handler (it cannot " + "see _transition calls in a non-*Orchestrator class) — this is the gap." + ) + assert len(arch003.violations) == 0 + + # ARCH-004 (cross-file node-directory join) FAILS the node. + arch004 = validate_contract_declared_orchestrator_workflow(str(node_dir)) + assert not arch004.valid, "ARCH-004 must FAIL the delegation-shape node." + codes = { + str(v.details["finding_code"]) + for v in arch004.violations + if v.details and "finding_code" in v.details + } + assert "H1" in codes, "Expected H1 (decorative FSM + handler-driven transitions)." + assert all(v.rule_id == RULE_ID for v in arch004.violations) + # H3: one handler selects state and constructs an event. + assert "H3" in codes, "Expected H3 (state selection + terminal event construction)." + # This fixture has a single declared payload type, so H2 (fan-in >= 3) does + # NOT fire here; the 3+-payload fan-in case is proven in the full-audit test. + assert "H2" not in codes + + +def test_payload_type_match_three_plus_payloads_fails(tmp_path: Path) -> None: + """(b) payload_type_match routing 3+ payloads to one workflow handler -> fail.""" + contract = """\ +name: node_synth_fanin +node_type: ORCHESTRATOR_GENERIC +handler_routing: + routing_strategy: payload_type_match + handlers: + - operation: a + handler: + name: HandlerWorkflow + module: pkg.nodes.node_synth_fanin.handlers.handler_workflow + event_model: + name: ModelEventA + module: pkg.models + - operation: b + handler: + name: HandlerWorkflow + module: pkg.nodes.node_synth_fanin.handlers.handler_workflow + event_model: + name: ModelEventB + module: pkg.models + - operation: c + handler: + name: HandlerWorkflow + module: pkg.nodes.node_synth_fanin.handlers.handler_workflow + event_model: + name: ModelEventC + module: pkg.models +""" + handler = '''\ +"""Synthetic catchall handler (no FSM, no event construction).""" + + +class HandlerWorkflow: + def handle(self, envelope): + return envelope +''' + node_dir = _make_node_dir( + tmp_path, + node_name="node_synth_fanin", + contract_yaml=contract, + handler_source=handler, + ) + result = validate_contract_declared_orchestrator_workflow(str(node_dir)) + assert not result.valid, "payload_type_match with 3 payloads must hard-fail." + codes = { + str(v.details["finding_code"]) + for v in result.violations + if v.details and "finding_code" in v.details + } + assert "H2" in codes + # No fsm + no handler-driven transitions => no H1. + assert "H1" not in codes + + +def test_reducer_with_state_machine_passes(tmp_path: Path) -> None: + """(c) reducer with state_machine: + pure transition executor -> pass (exempt).""" + contract = """\ +name: node_synth_fsm_reducer +node_type: REDUCER_GENERIC +state_machine: + states: + - CREATED + - PROCESSING + - DONE + transitions: + - from: CREATED + to: PROCESSING + - from: PROCESSING + to: DONE +""" + handler = '''\ +"""Pure transition executor for a reducer (legitimately owns the FSM).""" + + +class ReducerStateMachine: + def reduce(self, state, event): + # Pure: returns next state, no I/O. + return self._transition(state, event) + + def _transition(self, state, event): + return state +''' + node_dir = _make_node_dir( + tmp_path, + node_name="node_synth_fsm_reducer", + contract_yaml=contract, + handler_source=handler, + ) + # Reducers are EXEMPT: analysis returns None (out of scope). + analysis = analyze_node_directory(node_dir) + assert analysis is None, "Reducers must be exempt from ARCH-004." + result = validate_contract_declared_orchestrator_workflow(str(node_dir)) + assert result.valid, "Reducer with state_machine must PASS (exempt)." + assert len(result.violations) == 0 + + +def test_executor_bound_orchestrator_passes(tmp_path: Path) -> None: + """(d) executor-bound orchestrator workflow DAG + thin handlers -> pass.""" + contract = """\ +name: node_synth_bound_orchestrator +node_type: ORCHESTRATOR_GENERIC +workflow_coordination: + execution_graph: + - step: route + next: infer + - step: infer + next: emit +handler_routing: + routing_strategy: operation_match + handlers: + - operation: route + handler: + name: HandlerRouteStep + module: pkg.nodes.node_synth_bound_orchestrator.handlers.handler_route +""" + handler = '''\ +"""Thin per-step handler; no FSM ownership, no terminal event construction.""" + + +class HandlerRouteStep: + def handle(self, envelope): + return envelope +''' + node_dir = _make_node_dir( + tmp_path, + node_name="node_synth_bound_orchestrator", + contract_yaml=contract, + handler_name="handler_route", + handler_source=handler, + ) + result = validate_contract_declared_orchestrator_workflow(str(node_dir)) + assert result.valid, ( + "Executor-bound orchestrator with thin handlers must PASS ARCH-004." + ) + assert len(result.violations) == 0 + + +def test_full_audit_detects_delegation_shaped_fixture(tmp_path: Path) -> None: + """(e) full-audit mode detects a delegation-shaped fixture. + + Mirrors the real contract.yaml: fsm: table + payload_type_match (3+ payloads) + + a handler using self._transition(...). + """ + contract = _DELEGATION_SHAPE_CONTRACT.replace( + """\ + - operation: orchestrate + handler: + name: HandlerDelegationWorkflow + module: pkg.nodes.node_synth_delegation.handlers.handler_workflow + event_model: + name: ModelDelegationRequest + module: pkg.models +""", + """\ + - operation: orchestrate + handler: + name: HandlerDelegationWorkflow + module: pkg.nodes.node_synth_delegation.handlers.handler_workflow + event_model: + name: ModelDelegationRequest + module: pkg.models + - operation: lifecycle + handler: + name: HandlerDelegationWorkflow + module: pkg.nodes.node_synth_delegation.handlers.handler_workflow + event_model: + name: ModelAgentLifecycle + module: pkg.models + - operation: gate + handler: + name: HandlerDelegationWorkflow + module: pkg.nodes.node_synth_delegation.handlers.handler_workflow + event_model: + name: ModelGateResult + module: pkg.models +""", + ) + _make_node_dir( + tmp_path, + node_name="node_synth_delegation", + contract_yaml=contract, + handler_source=_DELEGATION_SHAPE_HANDLER, + ) + # A clean orchestrator alongside, to prove the scan is selective. + _make_node_dir( + tmp_path, + node_name="node_synth_clean_orchestrator", + contract_yaml=( + "name: node_synth_clean_orchestrator\n" + "node_type: ORCHESTRATOR_GENERIC\n" + "handler_routing:\n" + " routing_strategy: operation_match\n" + " handlers:\n" + " - operation: go\n" + " handler:\n" + " name: HandlerGo\n" + " module: pkg.nodes.node_synth_clean_orchestrator.handlers.handler_go\n" + ), + handler_name="handler_go", + handler_source="class HandlerGo:\n def handle(self, e):\n return e\n", + ) + + node_dirs = discover_node_dirs(tmp_path) + result = scan_node_dirs("synthetic", node_dirs, repo_root=tmp_path) + hard_fail_nodes = {e.node for e in result.hard_fails} + assert "node_synth_delegation" in hard_fail_nodes, ( + "Full-audit must detect the delegation-shaped node." + ) + assert "node_synth_clean_orchestrator" not in hard_fail_nodes, ( + "Full-audit must NOT flag the clean operation_match orchestrator." + ) + deleg = next(e for e in result.hard_fails if e.node == "node_synth_delegation") + assert {"H1", "H2", "H3"}.issubset(set(deleg.finding_codes)) + # Repo-relative handler path (never a machine-absolute path). + assert not deleg.max_handler_path.startswith("/") + + +# --------------------------------------------------------------------------- # +# Enum-state-match detection (the leading-underscore-free FSM) +# --------------------------------------------------------------------------- # +def test_enum_state_match_drives_h1(tmp_path: Path) -> None: + """An Enum*State mirroring contract states (no _transition calls) -> H1.""" + contract = """\ +name: node_synth_enum_fsm +node_type: ORCHESTRATOR_GENERIC +fsm: + states: + - IDLE + - RUNNING + - DONE + - FAILED +handler_routing: + routing_strategy: operation_match + handlers: + - operation: go + handler: + name: HandlerEnumFsm + module: pkg.nodes.node_synth_enum_fsm.handlers.handler_enum +""" + handler = '''\ +"""Handler that owns an Enum*State mirroring the contract states.""" +from enum import StrEnum + + +class EnumWorkflowState(StrEnum): + IDLE = "IDLE" + RUNNING = "RUNNING" + DONE = "DONE" + FAILED = "FAILED" + + +class HandlerEnumFsm: + def handle(self, envelope): + state = EnumWorkflowState.IDLE + return state +''' + node_dir = _make_node_dir( + tmp_path, + node_name="node_synth_enum_fsm", + contract_yaml=contract, + handler_name="handler_enum", + handler_source=handler, + ) + analysis = analyze_node_directory(node_dir) + assert analysis is not None + assert analysis.enum_state_match, "Enum*State must match the contract states." + assert analysis.handler_drives_state + assert "H1" in analysis.finding_codes + + +def test_substring_trap_not_a_false_positive(tmp_path: Path) -> None: + """``record_phase_transition(`` and a ``current_state=%s`` log must NOT match. + + This is the precision check that distinguishes ARCH-004 from a naive + substring scanner (the audit flagged pr_lifecycle's ``_transition`` as 0*). + """ + contract = """\ +name: node_synth_logonly +node_type: ORCHESTRATOR_GENERIC +fsm: + states: + - A + - B +handler_routing: + routing_strategy: operation_match + handlers: + - operation: go + handler: + name: HandlerLogOnly + module: pkg.nodes.node_synth_logonly.handlers.handler_log +""" + handler = '''\ +"""Handler that only LOGS transition-ish strings; does not drive state.""" +import logging + +logger = logging.getLogger(__name__) + + +def record_phase_transition(phase): + return phase + + +class HandlerLogOnly: + def handle(self, envelope): + logger.info("current_state=%s", envelope) + record_phase_transition("done") + return envelope +''' + node_dir = _make_node_dir( + tmp_path, + node_name="node_synth_logonly", + contract_yaml=contract, + handler_name="handler_log", + handler_source=handler, + ) + analysis = analyze_node_directory(node_dir) + assert analysis is not None + assert not analysis.handler_drives_state, ( + "record_phase_transition( and a current_state=%s log string must NOT " + "register as handler-driven state transitions." + ) + # 2 fsm states + a log-only handler => no H1 hard fail (only the W4 warning). + assert "H1" not in analysis.finding_codes + assert not analysis.has_hard_fail + + +# --------------------------------------------------------------------------- # +# Protocol-rule surface +# --------------------------------------------------------------------------- # +def test_rule_class_protocol_surface(tmp_path: Path) -> None: + """RuleContractDeclaredOrchestratorWorkflow implements the rule protocol.""" + rule = RuleContractDeclaredOrchestratorWorkflow() + assert rule.rule_id == "ARCH-004" + assert rule.severity == EnumValidationSeverity.ERROR + assert rule.name + assert rule.description + + # Non-path target => skipped, not a violation. + skipped = rule.check(object()) + assert skipped.skipped is True + assert skipped.passed is True + + # contract.yaml path resolves to the parent node dir. + node_dir = _make_node_dir( + tmp_path, + node_name="node_synth_delegation", + contract_yaml=_DELEGATION_SHAPE_CONTRACT, + handler_source=_DELEGATION_SHAPE_HANDLER, + ) + res = rule.check(str(node_dir / "contract.yaml")) + assert res.passed is False + assert res.details is not None + assert "H1" in (res.details.get("finding_codes") or []) + + +def test_event_marker_regex_is_reported(tmp_path: Path) -> None: + """The exact event-marker regex is surfaced in violation details (auditability).""" + node_dir = _make_node_dir( + tmp_path, + node_name="node_synth_delegation", + contract_yaml=_DELEGATION_SHAPE_CONTRACT, + handler_source=_DELEGATION_SHAPE_HANDLER, + ) + result = validate_contract_declared_orchestrator_workflow(str(node_dir)) + assert any( + v.details and v.details.get("event_marker_regex") == EVENT_MARKER_REGEX + for v in result.violations + ) + + +# --------------------------------------------------------------------------- # +# Ratchet behaviour +# --------------------------------------------------------------------------- # +def _delegation_entry() -> BaselineEntry: + return BaselineEntry( + repo="synthetic", + node="node_synth_delegation", + max_handler_path="src/pkg/nodes/node_synth_delegation/handlers/handler_workflow.py", + line_count=20, + risk_score=8, + finding_codes=("H1", "H2", "H3"), + owner_ticket="OMN-13471", + ) + + +def test_ratchet_passes_when_node_matches_baseline() -> None: + """A scanned hard-fail already in the baseline (unchanged) passes the ratchet.""" + entry = _delegation_entry() + baseline = {f"{entry.repo}::{entry.node}": entry} + assert ratchet_violations([entry], baseline) == [] + + +def test_ratchet_fails_on_new_untracked_node() -> None: + """A scanned hard-fail not in the baseline fails the ratchet.""" + entry = _delegation_entry() + failures = ratchet_violations([entry], baseline={}) + assert len(failures) == 1 + assert "NEW imperative-orchestrator hard-fail" in failures[0] + + +def test_ratchet_fails_on_worsened_risk_and_growth() -> None: + """A scanned hard-fail that worsens risk / grows the handler fails the ratchet.""" + base = _delegation_entry() + worse = BaselineEntry( + repo=base.repo, + node=base.node, + max_handler_path=base.max_handler_path, + line_count=base.line_count + 500, + risk_score=base.risk_score + 2, + finding_codes=base.finding_codes + ("W1",), + owner_ticket=base.owner_ticket, + ) + failures = ratchet_violations([worse], baseline={f"{base.repo}::{base.node}": base}) + joined = "\n".join(failures) + assert "risk score WORSENED" in joined + assert "max handler GREW" in joined + assert "NEW finding codes" in joined + + +def test_ratchet_fails_when_baseline_entry_has_no_owner_ticket() -> None: + """Every accepted hard-fail must cite an owner ticket.""" + base = BaselineEntry( + repo="synthetic", + node="node_synth_delegation", + max_handler_path="x.py", + line_count=20, + risk_score=8, + finding_codes=("H1",), + owner_ticket="", + ) + failures = ratchet_violations([base], baseline={f"{base.repo}::{base.node}": base}) + assert any("no owner_ticket" in f for f in failures) + + +def test_strict_mode_fails_on_still_baselined_node() -> None: + """--strict: a baselined node that still hard-fails must be remediated.""" + entry = _delegation_entry() + baseline = {f"{entry.repo}::{entry.node}": entry} + failures = strict_violations([entry], baseline) + assert len(failures) == 1 + assert "still hard-fails ARCH-004 but is baselined" in failures[0] + + +def test_baseline_roundtrip(tmp_path: Path) -> None: + """render -> write -> load round-trips entries, with no absolute paths.""" + entry = _delegation_entry() + rendered = render_baseline_yaml("omnibase_infra", [entry]) + assert "/Users/" not in rendered and "/Volumes/" not in rendered + baseline_path = tmp_path / "architecture-handshakes" / "baseline.yaml" + write_baseline(baseline_path, "omnibase_infra", [entry]) + loaded = load_baseline(baseline_path) + key = f"{entry.repo}::{entry.node}" + assert key in loaded + assert loaded[key].finding_codes == entry.finding_codes + assert loaded[key].owner_ticket == "OMN-13471" From ee04217c79fe8cdd292f61331f18fc0c3a064449 Mon Sep 17 00:00:00 2001 From: jonahgabriel Date: Mon, 22 Jun 2026 09:48:20 -0400 Subject: [PATCH 2/5] fix(OMN-13472): rename _NodeAnalysis -> OrchestratorNodeAnalysis (Patterns validator: no leading-underscore class names) --- ...alidator_contract_declared_orchestrator_workflow.py | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/src/omnibase_infra/nodes/node_architecture_validator/validators/validator_contract_declared_orchestrator_workflow.py b/src/omnibase_infra/nodes/node_architecture_validator/validators/validator_contract_declared_orchestrator_workflow.py index 21d7de90b8..5cf1190b38 100644 --- a/src/omnibase_infra/nodes/node_architecture_validator/validators/validator_contract_declared_orchestrator_workflow.py +++ b/src/omnibase_infra/nodes/node_architecture_validator/validators/validator_contract_declared_orchestrator_workflow.py @@ -169,7 +169,7 @@ # --------------------------------------------------------------------------- # # Node-directory analysis result (internal, not a public model) # --------------------------------------------------------------------------- # -class _NodeAnalysis: +class OrchestratorNodeAnalysis: """Mutable accumulator for one node directory's cross-file analysis. This is an internal scratch object; the public surface is the list of @@ -414,7 +414,7 @@ def _handler_selects_state_and_emits(source: str) -> bool: # --------------------------------------------------------------------------- # # Core node-directory analysis # --------------------------------------------------------------------------- # -def analyze_node_directory(node_dir: Path) -> _NodeAnalysis | None: +def analyze_node_directory(node_dir: Path) -> OrchestratorNodeAnalysis | None: """Run the ARCH-004 cross-file join over a single node directory. Args: @@ -423,7 +423,7 @@ def analyze_node_directory(node_dir: Path) -> _NodeAnalysis | None: to :class:`~pathlib.Path` before calling. Returns: - A populated :class:`_NodeAnalysis` (with any violations), or ``None`` if + A populated :class:`OrchestratorNodeAnalysis` (with any violations), or ``None`` if the directory is not an analyzable node (no contract) or is exempt (reducer / not orchestrator-like). """ @@ -439,7 +439,7 @@ def analyze_node_directory(node_dir: Path) -> _NodeAnalysis | None: if not isinstance(raw, dict): return None - analysis = _NodeAnalysis(node_path, contract_path) + analysis = OrchestratorNodeAnalysis(node_path, contract_path) analysis.node_type = str(raw.get("node_type", "")) analysis.node_name = str(raw.get("name", node_path.name)) @@ -499,7 +499,7 @@ def analyze_node_directory(node_dir: Path) -> _NodeAnalysis | None: return analysis -def _evaluate_signals(a: _NodeAnalysis) -> None: +def _evaluate_signals(a: OrchestratorNodeAnalysis) -> None: """Apply the hard-fail and warning signals to a populated analysis.""" handler_loc = a.max_handler_path or a.contract_path From 3c7490cbb6197974391a93bc39240807332c6299 Mon Sep 17 00:00:00 2001 From: jonahgabriel Date: Mon, 22 Jun 2026 09:56:26 -0400 Subject: [PATCH 3/5] =?UTF-8?q?fix(OMN-13472):=20address=20CodeRabbit=20?= =?UTF-8?q?=E2=80=94=20derive=20repo=20key=20from=20worktree,=20guard=20--?= =?UTF-8?q?write-baseline=20behind=20--check-all?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - scripts/validate.py: derive repo_name from repo_root.name instead of hardcoding 'omnibase_infra' (correct repo::node baseline keying). - scanner: --write-baseline now requires --check-all (a baseline from a --check-changed scan would silently shrink the ratchet below true state); add test_write_baseline_requires_check_all. --- scripts/validate.py | 8 ++++++-- .../scanner_imperative_orchestrator_ratchet.py | 10 ++++++++++ ...alidator_contract_declared_orchestrator_workflow.py | 10 ++++++++++ 3 files changed, 26 insertions(+), 2 deletions(-) diff --git a/scripts/validate.py b/scripts/validate.py index 84ec529fa6..3507fc066c 100755 --- a/scripts/validate.py +++ b/scripts/validate.py @@ -574,6 +574,10 @@ def run_imperative_orchestrators( from pathlib import Path as _Path repo_root = _Path.cwd() + # Derive the repo key from the working tree (worktrees nest the repo name in + # a ticket dir, so the immediate dir name is authoritative) instead of + # hardcoding it; baseline entries are keyed by repo::node. + repo_name = repo_root.name baseline = load_baseline( repo_root / "architecture-handshakes" / "imperative-orchestrator-baseline.yaml" ) @@ -584,7 +588,7 @@ def run_imperative_orchestrators( if verbose: print("Imperative Orchestrators: SKIP (no node dirs in changeset)") return True - result = scan_node_dirs("omnibase_infra", node_dirs, repo_root=repo_root) + result = scan_node_dirs(repo_name, node_dirs, repo_root=repo_root) failures = ratchet_violations(result.hard_fails, baseline) passed = not failures if verbose or not passed: @@ -603,7 +607,7 @@ def run_imperative_orchestrators( ) node_dirs = discover_node_dirs(repo_root) - result = scan_node_dirs("omnibase_infra", node_dirs, repo_root=repo_root) + result = scan_node_dirs(repo_name, node_dirs, repo_root=repo_root) print( f"Imperative Orchestrators (full report): " f"{len(result.hard_fails)} hard-fail node(s) across {len(node_dirs)} dirs." diff --git a/src/omnibase_infra/nodes/node_architecture_validator/validators/scanner_imperative_orchestrator_ratchet.py b/src/omnibase_infra/nodes/node_architecture_validator/validators/scanner_imperative_orchestrator_ratchet.py index 63d6972e99..ae35e703b8 100644 --- a/src/omnibase_infra/nodes/node_architecture_validator/validators/scanner_imperative_orchestrator_ratchet.py +++ b/src/omnibase_infra/nodes/node_architecture_validator/validators/scanner_imperative_orchestrator_ratchet.py @@ -437,6 +437,16 @@ def main(argv: list[str] | None = None) -> int: print(line) if args.write_baseline: + # A baseline must capture the FULL current debt; writing one from a + # --check-changed scan would silently drop every hard-fail not in the + # changeset, shrinking the ratchet baseline below true state. + if not args.check_all: + print( + "--write-baseline requires --check-all (a baseline must record " + "the full current debt, not just the changed nodes).", + file=sys.stderr, + ) + return 1 write_baseline(baseline_path, args.repo_name, result.hard_fails) print(f"Wrote baseline: {baseline_path}") return 0 diff --git a/tests/unit/nodes/node_architecture_validator/test_validator_contract_declared_orchestrator_workflow.py b/tests/unit/nodes/node_architecture_validator/test_validator_contract_declared_orchestrator_workflow.py index bebf521a7f..67e14b47d6 100644 --- a/tests/unit/nodes/node_architecture_validator/test_validator_contract_declared_orchestrator_workflow.py +++ b/tests/unit/nodes/node_architecture_validator/test_validator_contract_declared_orchestrator_workflow.py @@ -27,6 +27,7 @@ class is NOT named ``*Orchestrator`` — is PASSED by ARCH-003 (single-file, BaselineEntry, discover_node_dirs, load_baseline, + main, ratchet_violations, render_baseline_yaml, scan_node_dirs, @@ -632,3 +633,12 @@ def test_baseline_roundtrip(tmp_path: Path) -> None: assert key in loaded assert loaded[key].finding_codes == entry.finding_codes assert loaded[key].owner_ticket == "OMN-13471" + + +def test_write_baseline_requires_check_all() -> None: + """--write-baseline must reject --check-changed (would shrink the baseline).""" + rc = main(["--check-changed", "--write-baseline"]) + assert rc == 1, ( + "--write-baseline from a --check-changed scan must fail: a baseline must " + "record the full current debt, not just the changed nodes." + ) From 4269756ebd4ffce94060098da71d35d207de81cc Mon Sep 17 00:00:00 2001 From: jonahgabriel Date: Mon, 22 Jun 2026 10:26:09 -0400 Subject: [PATCH 4/5] fix(OMN-13472): refresh infra sync gates --- .../0020_delegation_context_pack_hash.sql | 9 +++++++++ docker/runners/runner-image.lock.json | 4 ++-- 2 files changed, 11 insertions(+), 2 deletions(-) create mode 100644 docker/migrations/forward/nodes/node_projection_delegation/0020_delegation_context_pack_hash.sql diff --git a/docker/migrations/forward/nodes/node_projection_delegation/0020_delegation_context_pack_hash.sql b/docker/migrations/forward/nodes/node_projection_delegation/0020_delegation_context_pack_hash.sql new file mode 100644 index 0000000000..f1d69bab17 --- /dev/null +++ b/docker/migrations/forward/nodes/node_projection_delegation/0020_delegation_context_pack_hash.sql @@ -0,0 +1,9 @@ +-- OMN-13407: persist delegation context-pack identity on the canonical +-- delegation projection so context ON/OFF ROI is measurable from the +-- correlation-trace surface. Empty string is the OFF/no-context arm. + +ALTER TABLE delegation_events + ADD COLUMN IF NOT EXISTS context_pack_hash TEXT NOT NULL DEFAULT ''; + +CREATE INDEX IF NOT EXISTS idx_delegation_events_context_pack_hash + ON delegation_events (context_pack_hash); diff --git a/docker/runners/runner-image.lock.json b/docker/runners/runner-image.lock.json index 0d11c75a21..18c62a3d48 100644 --- a/docker/runners/runner-image.lock.json +++ b/docker/runners/runner-image.lock.json @@ -1,12 +1,12 @@ { "base_image_digest": "sha256:3ba65aa20f86a0fad9df2b2c259c613df006b2e6d0bfcc8a146afb8c525a9751", "gh_version": "2.67.0", - "identity_digest": "8c3208f1d0b18f6a94b6fe27a270e7cb", + "identity_digest": "0f337da65f9dbcd3018bccd77bdd1420", "image_version": 5, "kubectl_version": "1.32.1", "python_version": "3.12", "runner_version": "2.334.0", - "shared_env_digest": "a796970abe13b64c009206f2", + "shared_env_digest": "efb9011b0136952b9f7d9886", "shared_env_install_args": "--frozen --all-extras --all-groups --no-install-project", "uv_version": "0.6.14" } From aaa0de603288b1d6e83b68bd8e4969a4b212bde7 Mon Sep 17 00:00:00 2001 From: jonahgabriel Date: Mon, 22 Jun 2026 10:47:35 -0400 Subject: [PATCH 5/5] fix(OMN-13472): align runner-identity test assertions with refreshed lock Commit 4269756eb refreshed docker/runners/runner-image.lock.json (identity_digest 8c3208f1->0f337da6, shared_env_digest a796970a->efb9011b) to clear runner-image-build-smoke, but left test_release_backmerge_preserves_runner_identity_lock asserting the stale digests, failing Tests (Split 1/15). Update the assertions to the committed lock values; runner-image-build-smoke validates the lock vs real image identity. --- .../infra/test_omn_12765_release_backmerge_identity.py | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/tests/integration/infra/test_omn_12765_release_backmerge_identity.py b/tests/integration/infra/test_omn_12765_release_backmerge_identity.py index 2058aafa37..9c1c7b29ad 100644 --- a/tests/integration/infra/test_omn_12765_release_backmerge_identity.py +++ b/tests/integration/infra/test_omn_12765_release_backmerge_identity.py @@ -32,8 +32,7 @@ def test_release_backmerge_preserves_runner_identity_lock() -> None: (ROOT / "docker/runners/runner-image.lock.json").read_text(encoding="utf-8") ) - # OMN-13445: identity regenerated after advancing the omnibase-core pin to the - # Phase-1b core SHA (the pin lives in pyproject.toml + uv.lock, which are - # binding components of the runner image identity digest). - assert lock["identity_digest"] == "8c3208f1d0b18f6a94b6fe27a270e7cb" - assert lock["shared_env_digest"] == "a796970abe13b64c009206f2" + # OMN-13472: identity regenerated after refreshing the shared CI env binding + # for the current dev/runtime proof inputs. + assert lock["identity_digest"] == "0f337da65f9dbcd3018bccd77bdd1420" + assert lock["shared_env_digest"] == "efb9011b0136952b9f7d9886"