#!/usr/bin/env python3
from __future__ import annotations

import importlib.util
import io
import contextlib
import hashlib
import json
import re
import sqlite3
import subprocess
import tempfile
import threading
import time
from datetime import datetime, timedelta, timezone
from pathlib import Path

ROOT = Path(__file__).resolve().parents[1]
SCRIPTS = ROOT / "skills" / "agent-orchestration-skill" / "scripts"


def load(name: str):
    path = SCRIPTS / f"{name}.py"
    spec = importlib.util.spec_from_file_location(name, path)
    if spec is None or spec.loader is None:
        raise RuntimeError(f"cannot import {path}")
    module = importlib.util.module_from_spec(spec)
    spec.loader.exec_module(module)
    return module


def run(*args: str, cwd: Path) -> None:
    subprocess.run(args, cwd=cwd, check=True, capture_output=True, text=True)


def output(*args: str, cwd: Path) -> str:
    return subprocess.run(args, cwd=cwd, check=True, capture_output=True, text=True).stdout


def assert_equal(actual, expected, label: str) -> None:
    if actual != expected:
        raise AssertionError(f"{label}: expected {expected!r}, got {actual!r}")


def capture_project_report(usage_ledger, summary: dict) -> str:
    stream = io.StringIO()
    with contextlib.redirect_stdout(stream):
        usage_ledger.print_project_report(summary)
    return stream.getvalue()


def gui_metric_value(html: str, label: str) -> str:
    matches = re.findall(r'<div class="metric"><b>([^<]*)</b><span>' + re.escape(label) + r'</span>', html)
    assert_equal(len(matches), 1, f"GUI {label} metric binding")
    return matches[0]


def test_route_and_config(workflow_config) -> None:
    assert_equal(workflow_config.normalize_route(None), "light", "default route")
    assert_equal(workflow_config.normalize_route("median"), "medium", "median alias")
    assert_equal(workflow_config.normalize_route("heavy"), "hard", "heavy alias")
    try:
        workflow_config.normalize_route("large")
    except ValueError:
        pass
    else:
        raise AssertionError("task-size label must not be accepted as a route")
    assert_equal(workflow_config.resolve_user_route("Please use route hard for this task"), "hard", "explicit hard route")
    assert_equal(workflow_config.resolve_user_route("Documentation says `route hard`, but just answer this"), "light", "inline quote cannot select route")
    assert_equal(workflow_config.resolve_user_route("> route hard\nPlease explain the code"), "light", "quoted block cannot select route")
    assert_equal(workflow_config.resolve_user_route("continue", current_route="medium", continuation=True), "medium", "continuation keeps route")
    assert_equal(workflow_config.resolve_user_route("new unrelated task", current_route="hard", continuation=False), "light", "unrelated task resets to light")
    try:
        workflow_config.resolve_user_route("route medium, then route hard")
    except ValueError:
        pass
    else:
        raise AssertionError("conflicting explicit routes must be rejected")

    config = workflow_config.default_config()
    assert_equal(config["routes"]["default"], "light", "configured default")
    assert_equal(config["stages"]["planning"]["model_name"], "gpt-5.6-sol", "planning model")
    assert_equal(config["stages"]["implementation"]["model_name"], "gpt-5.6-luna", "implementation model")
    assert_equal(config["stages"]["reviewer"]["model_name"], "gpt-5.6-sol", "reviewer model")
    exact_stage_models = {
        "planning": "gpt-5.6-sol",
        "implementation": "gpt-5.6-luna",
        "reviewer": "gpt-5.6-sol",
    }
    for stage, expected_model in exact_stage_models.items():
        swapped = json.loads(json.dumps(config))
        swapped["stages"][stage]["model_name"] = (
            "gpt-5.6-luna" if expected_model == "gpt-5.6-sol" else "gpt-5.6-sol"
        )
        if not any(f"stages.{stage}.model_name" in error for error in workflow_config.validate_config(swapped)):
            raise AssertionError(f"swapped model must be rejected for {stage}")
    assert_equal(workflow_config.select_effort(config, "implementation", "high"), "high", "allowed effort")
    try:
        workflow_config.select_effort(config, "reviewer", "max")
    except ValueError:
        pass
    else:
        raise AssertionError("reviewer max effort must be rejected")


def test_shared_project_store(workflow_config, base: Path) -> None:
    repo = base / "repo"
    linked = base / "linked"
    repo.mkdir()
    run("git", "init", "-q", cwd=repo)
    run("git", "config", "user.email", "aow@example.invalid", cwd=repo)
    run("git", "config", "user.name", "AOW Test", cwd=repo)
    (repo / "tracked.txt").write_text("base\n", encoding="utf-8")
    run("git", "add", "tracked.txt", cwd=repo)
    run("git", "commit", "-qm", "base", cwd=repo)
    run("git", "worktree", "add", "-q", str(linked), cwd=repo)
    assert_equal(
        workflow_config.project_store_root(repo),
        workflow_config.project_store_root(linked),
        "linked worktrees share a usage store",
    )


def test_usage_accounting(usage_ledger, aow_gui, base: Path) -> None:
    repo = base / "usage-repo"
    repo.mkdir()
    run("git", "init", "-q", cwd=repo)

    fixture = usage_ledger.normalize_record(
        {
            "schema_version": 2,
            "event_id": "evt-arithmetic",
            "project_id": "project-a",
            "worktree_id": "root",
            "run_id": "run-a",
            "session_id": "session-a",
            "route": "hard",
            "stage": "implementation",
            "model": "gpt-5.6-luna",
            "effort": "high",
            "input_tokens": 1000,
            "cached_input_tokens": 600,
            "output_tokens": 200,
            "reasoning_output_tokens": 150,
        }
    )
    assert_equal(fixture["total_tokens"], 1200, "cached/reasoning subsets are not double counted")
    assert_equal(fixture["uncached_input_tokens"], 400, "uncached input")
    assert_equal(fixture["visible_output_tokens"], 50, "visible output")

    first = usage_ledger.record_usage_event(repo, fixture)
    replay = usage_ledger.record_usage_event(repo, fixture)
    assert first["indexed"] is True
    assert replay["indexed"] is False
    summary = usage_ledger.project_usage_summary(repo)
    assert_equal(summary["totals"]["total_tokens"], 1200, "event replay is idempotent")
    assert_equal(summary["cursor"], 1, "cursor advances only for unique events")

    assert_equal(usage_ledger.usage_coverage({}, {"records": 0}), "unknown", "no source coverage")
    registered = {"registered_sources": 1, "sources": [{"status": "collecting", "attribution_status": "confirmed"}]}
    assert_equal(usage_ledger.usage_coverage(registered, {"records": 0}), "registered_no_events", "registered empty coverage")
    assert_equal(usage_ledger.usage_coverage(registered, {"records": 1}), "observed", "observed coverage")
    assert_equal(
        usage_ledger.usage_coverage(
            {"registered_sources": 1, "sources": [{"status": "collecting"}]},
            {"records": 1},
        ),
        "partial",
        "missing attribution status is a legacy coverage gap",
    )
    assert_equal(usage_ledger.usage_coverage({**registered, "worker_manifests": [{"telemetry_status": "unknown"}]}, {"records": 1}), "partial", "manifest gap coverage")
    assert_equal(usage_ledger.usage_coverage({"registered_sources": 1, "sources": [{"counter_coverage": "partial"}]}, {"records": 1}), "partial", "last delta coverage")
    assert "Observed tokens: unknown" in capture_project_report(usage_ledger, {"coverage_status": "unknown", "totals": {"cost_usd": None, "pricing_status": "unknown"}})
    assert "Cost: partial" in capture_project_report(usage_ledger, {"coverage_status": "observed", "totals": {"cost_usd": 1, "pricing_status": "partial_unknown"}})


    common = {
        "schema_version": 2,
        "project_id": "project-a",
        "worktree_id": "root",
        "run_id": "run-a",
        "session_id": "cumulative-session",
        "source": "codex-runtime",
        "route": "medium",
        "stage": "root",
        "model": "gpt-5.6-sol",
        "effort": "medium",
        "counter_mode": "cumulative",
    }
    usage_ledger.record_usage_event(repo, {**common, "event_id": "cum-1", "input_tokens": 1000, "output_tokens": 200})
    usage_ledger.record_usage_event(repo, {**common, "event_id": "cum-2", "input_tokens": 1200, "output_tokens": 260})
    usage_ledger.record_usage_event(repo, {**common, "event_id": "cum-3", "input_tokens": 100, "output_tokens": 20})
    summary = usage_ledger.project_usage_summary(repo, session_id="cumulative-session")
    assert_equal(summary["totals"]["input_tokens"], 1300, "cumulative input snapshots are differenced and resets handled")
    assert_equal(summary["totals"]["output_tokens"], 280, "cumulative output snapshots are differenced and resets handled")

    # Detail fields use their own availability markers.  A missing first
    # snapshot cannot be differenced against a later newly-known value; the
    # following snapshot establishes a safe baseline, and only the next one
    # becomes measurable.
    detail_repo = base / "cumulative-detail-repo"
    detail_repo.mkdir()
    run("git", "init", "-q", cwd=detail_repo)
    detail_common = {**common, "session_id": "cumulative-detail"}
    first_detail = usage_ledger.record_usage_event(detail_repo, {
        **detail_common, "event_id": "detail-1", "input_tokens": 100,
        "output_tokens": 20,
    })
    baseline_detail = usage_ledger.record_usage_event(detail_repo, {
        **detail_common, "event_id": "detail-2", "input_tokens": 120,
        "cached_input_tokens": 60, "output_tokens": 25,
        "reasoning_output_tokens": 5,
    })
    measured_detail = usage_ledger.record_usage_event(detail_repo, {
        **detail_common, "event_id": "detail-3", "input_tokens": 130,
        "cached_input_tokens": 70, "output_tokens": 30,
        "reasoning_output_tokens": 6,
    })
    assert_equal(first_detail["cached_input_tokens"], None, "first missing cache detail stays unknown")
    assert_equal(baseline_detail["cached_input_tokens"], None, "newly-known cache detail waits for a baseline")
    assert_equal(measured_detail["cached_input_tokens"], 10, "cache detail is measurable after a safe baseline")
    assert_equal(measured_detail["reasoning_output_tokens"], 1, "reasoning detail is measurable after a safe baseline")
    assert_equal(usage_ledger.project_usage_summary(detail_repo, session_id="cumulative-detail")["totals"]["total_tokens"], 160, "detail transition does not double-count cumulative totals")

    old_ts = (datetime.now(timezone.utc) - timedelta(hours=6)).isoformat()
    usage_ledger.record_usage_event(
        repo,
        {
            **common,
            "event_id": "old-window",
            "session_id": "old-session",
            "counter_mode": "delta",
            "ts": old_ts,
            "input_tokens": 900,
            "output_tokens": 100,
        },
    )
    window = usage_ledger.project_usage_summary(repo, window="5h")
    if window["totals"]["total_tokens"] >= usage_ledger.project_usage_summary(repo)["totals"]["total_tokens"]:
        raise AssertionError("five-hour window did not exclude an older event")

    db = usage_ledger.usage_db_path(repo)
    log = usage_ledger.usage_v2_jsonl(repo)
    if not db.is_file() or not log.is_file():
        raise AssertionError("usage JSONL and SQLite index were not created")
    if "prompt" in log.read_text(encoding="utf-8").lower():
        raise AssertionError("usage evidence must not persist prompt text")
    dashboard = aow_gui.collect(repo)
    if "discovery" in dashboard:
        raise AssertionError("v2 dashboard must not expose global Codex discovery")
    assert_equal(dashboard["usage_cursor"], usage_ledger.latest_usage_cursor(repo), "dashboard usage cursor")
    assert_equal(dashboard["telemetry"]["model_calls"], 0, "monitoring is model-free")

    pricing_repo = base / "pricing-repo"
    pricing_repo.mkdir()
    (pricing_repo / ".orca").mkdir()
    (pricing_repo / ".orca" / "pricing.json").write_text(
        json.dumps({
            "schema": "aoc.pricing.v1",
            "models": {
                "gpt-5.6-luna": {
                    "input_per_million": 2.0,
                    "cached_input_per_million": 0.5,
                    "output_per_million": 8.0,
                }
            },
        }),
        encoding="utf-8",
    )
    priced = usage_ledger.record_usage_event(pricing_repo, {
        **fixture,
        "event_id": "priced-event",
        "input_tokens": 1_000_000,
        "cached_input_tokens": 600_000,
        "output_tokens": 200_000,
        "reasoning_output_tokens": 150_000,
        "cost_usd": None,
    })
    assert_equal(priced["pricing_status"], "estimated", "catalog pricing status")
    assert_equal(round(priced["cost_usd"], 2), 2.70, "uncached/cached/output price formula")


def test_usage_coverage_render_matrix(usage_ledger, usage_collector, aow_gui, aow_tui, base: Path) -> None:
    repo = base / "coverage-render-matrix"
    repo.mkdir()
    source = base / "coverage-render.jsonl"
    source.write_text(json.dumps({"type": "event_msg", "payload": {"type": "token_count", "info": {"total_token_usage": {"input_tokens": 8, "cached_input_tokens": 0, "output_tokens": 2, "reasoning_output_tokens": 0, "total_tokens": 10}}}}) + "\n", encoding="utf-8")
    usage_collector.register_source(
        repo,
        source,
        session_id="target",
        run_id="target",
        model="gpt-5.6-sol",
        effort="medium",
    )
    usage_collector.collect_registered_sources(repo)
    assert_equal(usage_ledger.project_usage_summary(repo, run_id="target")["coverage_status"], "observed", "independent observed baseline")
    matching = repo / ".orchestration" / "runs" / "target" / "launches"
    matching.mkdir(parents=True)
    (matching / "worker.json").write_text(json.dumps({"telemetry_status": "unknown"}), encoding="utf-8")
    assert_equal(usage_ledger.project_usage_summary(repo, run_id="target")["coverage_status"], "partial", "matching manifest makes coverage partial")
    (matching / "worker.json").unlink()
    unrelated = repo / ".orchestration" / "runs" / "other" / "launches"
    unrelated.mkdir(parents=True)
    (unrelated / "worker.json").write_text(json.dumps({"telemetry_status": "unavailable"}), encoding="utf-8")
    assert_equal(usage_ledger.project_usage_summary(repo, run_id="target")["coverage_status"], "observed", "unrelated manifest is excluded")
    (unrelated / "worker.json").unlink()
    empty = base / "registered-no-events"
    empty.mkdir()
    empty_source = base / "empty.jsonl"
    empty_source.write_text("", encoding="utf-8")
    usage_collector.register_source(
        empty,
        empty_source,
        session_id="empty",
        model="gpt-5.6-sol",
        effort="medium",
    )
    empty_summary = usage_ledger.project_usage_summary(empty)
    assert_equal(empty_summary["coverage_status"], "registered_no_events", "registered no events")
    report = capture_project_report(usage_ledger, empty_summary)
    assert "Observed tokens: 0" in report and "$0.0000" not in report
    empty_tui = "\n".join(aow_tui.lines_usage(empty, None))
    assert "observed tokens: 0" in empty_tui and "coverage: registered_no_events" in empty_tui
    assert "cost: unknown" in empty_tui and "$0.0000" not in empty_tui
    gui_html = "".join(aow_gui.render_fragment(aow_gui.collect(empty), "usage"))
    assert_equal(gui_metric_value(gui_html, "Observed tokens"), "0", "GUI registered no events tokens")
    assert_equal(gui_metric_value(gui_html, "Coverage"), "registered_no_events", "GUI registered no events coverage")
    assert_equal(gui_metric_value(gui_html, "Cost"), "unknown", "GUI registered no events cost")
    assert "$0.0000" not in gui_html

    observed_tui = "\n".join(aow_tui.lines_usage(repo, None))
    assert "observed tokens: 10" in observed_tui and "coverage: observed" in observed_tui
    assert "cost: unknown" in observed_tui and "$0.0000" not in observed_tui
    observed_report = capture_project_report(usage_ledger, usage_ledger.project_usage_summary(repo))
    assert "Observed tokens: 10" in observed_report and "Cost: unknown" in observed_report
    observed_gui = "".join(aow_gui.render_fragment(aow_gui.collect(repo), "usage"))
    assert_equal(gui_metric_value(observed_gui, "Observed tokens"), "10", "GUI observed tokens")
    assert_equal(gui_metric_value(observed_gui, "Coverage"), "observed", "GUI observed coverage")
    assert_equal(gui_metric_value(observed_gui, "Cost"), "unknown", "GUI all-unpriced cost")
    assert "$0.0000" not in observed_gui

    usage_ledger.record_usage_event(repo, {
        "event_id": "priced-render-event",
        "run_id": "target",
        "session_id": "target",
        "input_tokens": 4,
        "cached_input_tokens": 0,
        "output_tokens": 1,
        "reasoning_output_tokens": 0,
        "total_tokens": 5,
        "cost_usd": 0.25,
    })
    mixed_tui = "\n".join(aow_tui.lines_usage(repo, None))
    assert "observed tokens: 15" in mixed_tui and "cost: partial" in mixed_tui
    mixed_report = capture_project_report(usage_ledger, usage_ledger.project_usage_summary(repo))
    assert "Observed tokens: 15" in mixed_report and "Cost: partial" in mixed_report
    mixed_gui = "".join(aow_gui.render_fragment(aow_gui.collect(repo), "usage"))
    assert_equal(gui_metric_value(mixed_gui, "Observed tokens"), "15", "GUI mixed tokens")
    assert_equal(gui_metric_value(mixed_gui, "Coverage"), "observed", "GUI mixed coverage")
    assert_equal(gui_metric_value(mixed_gui, "Cost"), "partial", "GUI mixed cost")
    assert "$0.0000" not in mixed_gui

    partial_dir = repo / ".orchestration" / "runs" / "target" / "launches"
    partial_dir.mkdir(parents=True, exist_ok=True)
    (partial_dir / "gap.json").write_text(json.dumps({"telemetry_status": "unknown"}), encoding="utf-8")
    partial_tui = "\n".join(aow_tui.lines_usage(repo, None))
    assert "coverage: partial" in partial_tui
    partial_gui = "".join(aow_gui.render_fragment(aow_gui.collect(repo), "usage"))
    assert_equal(gui_metric_value(partial_gui, "Coverage"), "partial", "GUI partial coverage")


def test_public_cli(base: Path) -> None:
    cli = ROOT / "bin" / "aow.mjs"
    repo = base / "cli-repo"
    repo.mkdir()
    run("git", "init", "-q", cwd=repo)
    help_text = output("node", str(cli), "--help", cwd=repo)
    primary = help_text
    for command in ("aow", "aow gui", "aow config", "aow usage", "aow doctor", "aow uninstall"):
        if command not in primary:
            raise AssertionError(f"primary help is missing {command}")
    for legacy in ("sessions", "import", "budget", "init"):
        if f"\n  {legacy} " in primary or f"\n  aow {legacy} " in primary:
            raise AssertionError(f"legacy command {legacy} leaked into primary help")

    preview = output("node", str(cli), "config", "--repo", str(repo), "--json", cwd=repo)
    assert_equal(json.loads(preview)["status"], "default-preview", "config preview is non-mutating")
    if (repo / ".orca" / "agents.yaml").exists():
        raise AssertionError("aow config show unexpectedly wrote a file")
    created = output("node", str(cli), "config", "init", "--repo", str(repo), "--json", cwd=repo)
    assert_equal(json.loads(created)["status"], "created", "config init")
    checked = output("node", str(cli), "config", "--check", "--repo", str(repo), "--json", cwd=repo)
    assert_equal(json.loads(checked)["status"], "valid", "config check")
    report = json.loads(output("node", str(cli), "usage", "--window", "5h", "--repo", str(repo), "--json", cwd=repo))
    assert_equal(report["schema"], "aoc.usage_summary.v2", "public usage command uses project index")
    doctor = json.loads(output("node", str(cli), "doctor", "--repo", str(repo), "--json", cwd=repo))
    assert_equal(doctor["schema"], "aoc.doctor.v2", "doctor checks v2 components")
    if doctor["status"] not in {"PASS", "WARN"}:
        raise AssertionError(f"doctor unexpectedly failed: {doctor}")


def test_incremental_collector(usage_collector, usage_ledger, workflow_config, base: Path) -> None:
    repo = base / "collector-repo"
    repo.mkdir()
    run("git", "init", "-q", cwd=repo)
    source = base / "runtime.jsonl"
    event = {
        "timestamp": datetime.now(timezone.utc).isoformat(),
        "ordinal": 7,
        "type": "event_msg",
        "payload": {
            "type": "token_count",
            "info": {
                "last_token_usage": {
                    "input_tokens": 1000,
                    "cached_input_tokens": 600,
                    "output_tokens": 200,
                    "reasoning_output_tokens": 150,
                    "total_tokens": 1200,
                }
            },
        },
    }
    complete = json.dumps(event, separators=(",", ":")) + "\n"
    source.write_text(complete + '{"partial":', encoding="utf-8")
    usage_collector.register_source(
        repo,
        source,
        session_id="collector-session",
        route="hard",
        stage="implementation",
        model="gpt-5.6-luna",
        effort="high",
    )
    result = usage_collector.collect_registered_sources(repo)
    assert_equal(result["indexed"], 1, "collector indexes complete token event")
    assert_equal(result["partial_lines"], 1, "collector defers partial final line")
    assert_equal(usage_ledger.project_usage_summary(repo)["totals"]["total_tokens"], 1200, "collector token arithmetic")
    replay = usage_collector.collect_registered_sources(repo)
    assert_equal(replay["indexed"], 0, "persisted byte offset prevents replay")

    entered = threading.Event()
    released = threading.Event()

    def hold_registry_lock() -> None:
        with usage_collector._registry_lock(repo):
            entered.set()
            released.wait(timeout=2)

    holder = threading.Thread(target=hold_registry_lock)
    holder.start()
    if not entered.wait(timeout=1):
        raise AssertionError("registry lock holder did not start")
    acquired = threading.Event()

    def wait_for_registry_lock() -> None:
        with usage_collector._registry_lock(repo):
            acquired.set()

    waiter = threading.Thread(target=wait_for_registry_lock)
    waiter.start()
    time.sleep(0.05)
    if acquired.is_set():
        raise AssertionError("registry lock did not serialize concurrent writers")
    released.set()
    holder.join(timeout=1)
    waiter.join(timeout=1)
    if not acquired.is_set():
        raise AssertionError("registry lock was not released")

    def token_event(*, last_input, last_output, total_input, total_output):
        return {
            "timestamp": datetime.now(timezone.utc).isoformat(),
            "type": "event_msg",
            "payload": {
                "type": "token_count",
                "info": {
                    "last_token_usage": {"input_tokens": last_input, "output_tokens": last_output},
                    "total_token_usage": {"input_tokens": total_input, "output_tokens": total_output},
                },
            },
        }

    # Total counters are normalized before reaching the ledger. Repeated
    # cumulative snapshots therefore do not create duplicate usage records.
    cumulative_repo = base / "cumulative-repo"
    cumulative_repo.mkdir()
    run("git", "init", "-q", cwd=cumulative_repo)
    cumulative_source = base / "cumulative.jsonl"
    cumulative_source.write_text(
        "".join(json.dumps(event, separators=(",", ":")) + "\n" for event in (
            token_event(last_input=100, last_output=20, total_input=100, total_output=20),
            token_event(last_input=100, last_output=20, total_input=100, total_output=20),
            token_event(last_input=50, last_output=10, total_input=150, total_output=30),
        )),
        encoding="utf-8",
    )
    usage_collector.register_source(cumulative_repo, cumulative_source, session_id="cumulative")
    collected = usage_collector.collect_registered_sources(cumulative_repo)
    assert_equal(collected["indexed"], 2, "repeated cumulative snapshot is ignored")
    assert_equal(usage_ledger.project_usage_summary(cumulative_repo)["totals"]["total_tokens"], 180, "total-first accounting")

    # Pending bytes are attributed under the old stage before re-registration.
    boundary_repo = base / "boundary-repo"
    boundary_repo.mkdir()
    run("git", "init", "-q", cwd=boundary_repo)
    boundary_source = base / "boundary.jsonl"
    boundary_source.write_text(json.dumps(token_event(last_input=10, last_output=2, total_input=10, total_output=2)) + "\n", encoding="utf-8")
    usage_collector.register_source(boundary_repo, boundary_source, session_id="boundary", stage="planning")
    usage_collector.register_source(boundary_repo, boundary_source, session_id="boundary", stage="implementation")
    boundary_rows = usage_ledger.project_usage_summary(boundary_repo, group_by="stage")["rows"]
    assert_equal(boundary_rows[0]["key"], "planning", "re-registration drains old stage")
    assert_equal(boundary_rows[0]["total_tokens"], 12, "pending old-stage event accounted")

    # An unavailable source defers metadata changes until the old route can be
    # drained; after restoration, pending bytes remain planning-attributed.
    unavailable_repo = base / "unavailable-repo"
    unavailable_repo.mkdir()
    run("git", "init", "-q", cwd=unavailable_repo)
    unavailable_source = base / "unavailable.jsonl"
    unavailable_source.write_text(json.dumps(token_event(last_input=11, last_output=1, total_input=11, total_output=1)) + "\n", encoding="utf-8")
    usage_collector.register_source(unavailable_repo, unavailable_source, session_id="unavailable", stage="planning")
    hidden_source = base / "unavailable.hidden"
    unavailable_source.rename(hidden_source)
    deferred = usage_collector.register_source(unavailable_repo, unavailable_source, session_id="unavailable", stage="implementation")
    assert_equal(deferred["stage"], "planning", "unavailable drain defers re-registration")
    hidden_source.rename(unavailable_source)
    usage_collector.collect_registered_sources(unavailable_repo)
    unavailable_rows = usage_ledger.project_usage_summary(unavailable_repo, group_by="stage")["rows"]
    assert_equal(unavailable_rows[0]["key"], "planning", "restored pending bytes keep old stage")
    usage_collector.register_source(unavailable_repo, unavailable_source, session_id="unavailable", stage="implementation")

    # Totals-absent last_token_usage is a delta with explicitly partial coverage.
    partial_repo = base / "partial-repo"
    partial_repo.mkdir()
    run("git", "init", "-q", cwd=partial_repo)
    partial_source = base / "partial.jsonl"
    partial_event = token_event(last_input=13, last_output=2, total_input=13, total_output=2)
    partial_event["payload"]["info"].pop("total_token_usage")
    partial_source.write_text(json.dumps(partial_event) + "\n", encoding="utf-8")
    usage_collector.register_source(partial_repo, partial_source, session_id="partial")
    usage_collector.collect_registered_sources(partial_repo)
    partial_evidence = json.loads(usage_ledger.usage_v2_jsonl(partial_repo).read_text(encoding="utf-8").splitlines()[0])
    assert_equal(partial_evidence["counter_mode"], "delta", "partial fallback remains delta")
    assert_equal(partial_evidence["counter_coverage"], "partial", "partial fallback is marked")
    partial_registry = json.loads(usage_collector.registry_path(partial_repo).read_text(encoding="utf-8"))
    partial_source_row = next(iter(partial_registry["sources"].values()))
    assert_equal(partial_source_row["counter_coverage"], "partial", "partial coverage is registered")
    assert_equal(usage_ledger.project_usage_summary(partial_repo)["coverage_status"], "partial", "last delta summary coverage")
    manifest_dir = partial_repo / ".orchestration" / "runs" / "target-run" / "launches"
    manifest_dir.mkdir(parents=True)
    (manifest_dir / "unknown.attempt1.json").write_text(json.dumps({"telemetry_status": "unknown"}), encoding="utf-8")
    assert_equal(usage_ledger.project_usage_summary(partial_repo, run_id="target-run")["coverage_status"], "partial", "launch manifest gap coverage")
    unrelated_dir = partial_repo / ".orchestration" / "runs" / "other-run" / "launches"
    unrelated_dir.mkdir(parents=True)
    (unrelated_dir / "unknown.attempt1.json").write_text(json.dumps({"telemetry_status": "unknown"}), encoding="utf-8")

    # A source registered from a linked worktree retains that identity when
    # another checkout performs collection.
    linked_repo = base / "linked-repo"
    linked_repo.mkdir()
    run("git", "init", "-q", cwd=linked_repo)
    run("git", "config", "user.email", "aow@example.invalid", cwd=linked_repo)
    run("git", "config", "user.name", "AOW Test", cwd=linked_repo)
    (linked_repo / "tracked").write_text("base\n", encoding="utf-8")
    run("git", "add", "tracked", cwd=linked_repo)
    run("git", "commit", "-qm", "base", cwd=linked_repo)
    linked_checkout = base / "linked-checkout"
    run("git", "worktree", "add", "-q", str(linked_checkout), cwd=linked_repo)
    linked_source = base / "linked.jsonl"
    linked_source.write_text(json.dumps(token_event(last_input=7, last_output=3, total_input=7, total_output=3)) + "\n", encoding="utf-8")
    usage_collector.register_source(linked_checkout, linked_source, session_id="linked")
    usage_collector.collect_registered_sources(linked_repo)
    linked_summary = usage_ledger.project_usage_summary(linked_repo, group_by="worktree_id")
    assert_equal(linked_summary["rows"][0]["key"], workflow_config.worktree_id(linked_checkout), "registration worktree attribution")

    # Counter resets produce a fresh delta and durable reset evidence.
    reset_repo = base / "reset-repo"
    reset_repo.mkdir()
    run("git", "init", "-q", cwd=reset_repo)
    reset_source = base / "reset.jsonl"
    reset_source.write_text(
        "".join(json.dumps(event) + "\n" for event in (
            token_event(last_input=100, last_output=20, total_input=100, total_output=20),
            token_event(last_input=50, last_output=10, total_input=50, total_output=10),
        )), encoding="utf-8",
    )
    usage_collector.register_source(reset_repo, reset_source, session_id="reset")
    usage_collector.collect_registered_sources(reset_repo)
    assert_equal(usage_ledger.project_usage_summary(reset_repo)["totals"]["total_tokens"], 180, "counter reset delta")
    with sqlite3.connect(usage_ledger.usage_db_path(reset_repo)) as connection:
        reset_count = connection.execute("SELECT counter_reset FROM usage_events WHERE counter_reset = 1").fetchone()
    if reset_count is None:
        raise AssertionError("counter reset evidence was not persisted")

    # Replaying after a commit-before-offset crash remains idempotent.
    crash_repo = base / "crash-repo"
    crash_repo.mkdir()
    run("git", "init", "-q", cwd=crash_repo)
    crash_source = base / "crash.jsonl"
    crash_source.write_text(json.dumps(token_event(last_input=9, last_output=1, total_input=9, total_output=1)) + "\n", encoding="utf-8")
    usage_collector.register_source(crash_repo, crash_source, session_id="crash")
    original_write = usage_collector._write_registry
    failed = {"value": False}
    def fail_once(root, registry):
        if not failed["value"]:
            failed["value"] = True
            raise RuntimeError("simulated offset persistence crash")
        return original_write(root, registry)
    usage_collector._write_registry = fail_once
    try:
        try:
            usage_collector.collect_registered_sources(crash_repo)
        except RuntimeError:
            pass
        else:
            raise AssertionError("crash simulation did not interrupt offset persistence")
    finally:
        usage_collector._write_registry = original_write
    replay = usage_collector.collect_registered_sources(crash_repo)
    assert_equal(replay["duplicates"], 1, "commit-before-offset replay deduplication")
    assert_equal(usage_ledger.project_usage_summary(crash_repo)["totals"]["total_tokens"], 10, "crash replay remains idempotent")

    # v1 nonzero offsets bootstrap the last consumed cumulative snapshot.
    migration_repo = base / "migration-repo"
    migration_repo.mkdir()
    run("git", "init", "-q", cwd=migration_repo)
    migration_source = base / "migration.jsonl"
    historical = json.dumps(token_event(last_input=100, last_output=20, total_input=100, total_output=20), separators=(",", ":")) + "\n"
    appended = json.dumps(token_event(last_input=50, last_output=10, total_input=150, total_output=30), separators=(",", ":")) + "\n"
    migration_source.write_text(historical + appended, encoding="utf-8")
    # Keep source-id construction identical to the collector.
    migration_id = hashlib.sha256(f"{migration_source.resolve()}\x1fmigration".encode()).hexdigest()[:24]
    registry_file = usage_collector.registry_path(migration_repo)
    registry_file.parent.mkdir(parents=True, exist_ok=True)
    registry_file.write_text(json.dumps({"schema": "aoc.telemetry_sources.v1", "sources": {migration_id: {
        "source_id": migration_id, "path": str(migration_source.resolve()), "session_id": "migration",
        "route": "light", "stage": "root", "model": "unknown", "effort": "unknown", "offset": len(historical),
    }}}), encoding="utf-8")
    usage_collector.collect_registered_sources(migration_repo)
    assert_equal(usage_ledger.project_usage_summary(migration_repo)["totals"]["total_tokens"], 60, "v1 cumulative baseline bootstrap")
    if json.loads(registry_file.read_text(encoding="utf-8"))["schema"] != "aoc.telemetry_sources.v2":
        raise AssertionError("v1 registry was not upgraded")

    # A pre-detail-flag SQLite source_counters row migrates with the only
    # honest default available: old numeric values are known because their
    # missing-field provenance was not persisted.  A newly missing detail is
    # still surfaced as unknown on the first post-upgrade snapshot.
    legacy_db_repo = base / "legacy-db-repo"
    legacy_db_repo.mkdir()
    run("git", "init", "-q", cwd=legacy_db_repo)
    legacy_record = {
        "project_id": "legacy-project", "worktree_id": "legacy-worktree",
        "source": "legacy", "session_id": "legacy-session", "run_id": "legacy-run",
        "stage": "root", "model": "gpt-5.6-sol", "effort": "medium",
        "counter_mode": "cumulative",
    }
    legacy_source_key = usage_ledger._source_key(legacy_record)
    legacy_db = usage_ledger.usage_db_path(legacy_db_repo)
    legacy_db.parent.mkdir(parents=True, exist_ok=True)
    with sqlite3.connect(legacy_db) as connection:
        connection.execute("""
            CREATE TABLE source_counters (
              source_key TEXT PRIMARY KEY,
              input_tokens INTEGER NOT NULL,
              cached_input_tokens INTEGER NOT NULL,
              output_tokens INTEGER NOT NULL,
              reasoning_output_tokens INTEGER NOT NULL,
              total_tokens INTEGER NOT NULL,
              cost_usd REAL,
              updated_at TEXT NOT NULL
            )
        """)
        connection.execute(
            "INSERT INTO source_counters VALUES (?,?,?,?,?,?,?,?)",
            (legacy_source_key, 100, 40, 20, 2, 120, None, "2026-01-01T00:00:00+00:00"),
        )
    migrated_connection = usage_ledger._connect(legacy_db_repo)
    migrated_flags = migrated_connection.execute(
        "SELECT cache_detail_known, reasoning_detail_known FROM source_counters WHERE source_key = ?",
        (legacy_source_key,),
    ).fetchone()
    migrated_connection.close()
    assert_equal(tuple(migrated_flags), (1, 1), "old-schema detail flags default to unrecoverably-known")
    upgraded = usage_ledger.record_usage_event(legacy_db_repo, {
        **legacy_record, "event_id": "legacy-after-upgrade", "input_tokens": 120,
        "output_tokens": 25,
    })
    assert_equal(upgraded["input_tokens"], 20, "old-schema cumulative baseline is preserved")
    assert_equal(upgraded["output_tokens"], 5, "old-schema cumulative output baseline is preserved")
    assert_equal(upgraded["cached_input_tokens"], None, "post-upgrade missing cache detail stays unknown")
    with sqlite3.connect(legacy_db) as connection:
        columns = {row[1] for row in connection.execute("PRAGMA table_info(source_counters)")}
        assert {"cache_detail_known", "reasoning_detail_known"}.issubset(columns)
        flags = connection.execute(
            "SELECT cache_detail_known, reasoning_detail_known FROM source_counters WHERE source_key = ?",
            (legacy_source_key,),
        ).fetchone()
    assert_equal(tuple(flags), (0, 0), "migration records current unknown detail without inventing precision")

    # Rotation creates a new event-id generation even when offsets are reused.
    rotation_repo = base / "rotation-repo"
    rotation_repo.mkdir()
    run("git", "init", "-q", cwd=rotation_repo)
    rotation_source = base / "rotation.jsonl"
    rotation_source.write_text(json.dumps(token_event(last_input=100, last_output=20, total_input=100, total_output=20)) + "\n", encoding="utf-8")
    usage_collector.register_source(rotation_repo, rotation_source, session_id="rotation")
    usage_collector.collect_registered_sources(rotation_repo)
    first_rotation_size = rotation_source.stat().st_size
    rotation_source.write_text(json.dumps(token_event(last_input=1, last_output=1, total_input=1, total_output=1)) + "\n", encoding="utf-8")
    if rotation_source.stat().st_size >= first_rotation_size:
        rotation_source.write_text(json.dumps(token_event(last_input=1, last_output=1, total_input=1, total_output=1), separators=(",", ":")) + "\n", encoding="utf-8")
    usage_collector.collect_registered_sources(rotation_repo)
    rotation_registry = json.loads(usage_collector.registry_path(rotation_repo).read_text(encoding="utf-8"))
    rotation_source_row = next(iter(rotation_registry["sources"].values()))
    assert_equal(rotation_source_row["generation"], 1, "source rotation generation")
    with sqlite3.connect(usage_ledger.usage_db_path(rotation_repo)) as connection:
        rotation_ids = [row[0] for row in connection.execute("SELECT event_id FROM usage_events ORDER BY cursor")]
    if len(rotation_ids) != 2 or rotation_ids[0] == rotation_ids[1]:
        raise AssertionError("rotation generation did not produce a distinct event id")

    # A consumed-prefix digest detects in-place rewrites even when the source
    # does not shrink. Without the digest check these files are mistaken for an
    # unchanged file or a normal append and their replacement events are lost.
    rewrite_repo = base / "rewrite-repo"
    rewrite_repo.mkdir()
    run("git", "init", "-q", cwd=rewrite_repo)
    rewrite_source = base / "rewrite.jsonl"
    first_rewrite = json.dumps(
        token_event(last_input=100, last_output=20, total_input=100, total_output=20),
        separators=(",", ":"),
    ) + "\n"
    same_size_rewrite = json.dumps(
        token_event(last_input=200, last_output=30, total_input=200, total_output=30),
        separators=(",", ":"),
    ) + "\n"
    assert_equal(len(same_size_rewrite), len(first_rewrite), "same-size rewrite fixture")
    rewrite_source.write_text(first_rewrite, encoding="utf-8")
    usage_collector.register_source(rewrite_repo, rewrite_source, session_id="rewrite")
    usage_collector.collect_registered_sources(rewrite_repo)

    rewrite_source.write_text(same_size_rewrite, encoding="utf-8")
    same_size_result = usage_collector.collect_registered_sources(rewrite_repo)
    assert_equal(same_size_result["rotations"], 1, "same-size in-place rewrite rotation")

    larger_rewrite = (
        json.dumps(
            token_event(last_input=300, last_output=40, total_input=300, total_output=40),
            separators=(",", ":"),
        )
        + "\n"
        + json.dumps(
            token_event(last_input=310, last_output=45, total_input=310, total_output=45),
            separators=(",", ":"),
        )
        + "\n"
    )
    if len(larger_rewrite) <= len(same_size_rewrite):
        raise AssertionError("larger rewrite fixture must grow the source")
    rewrite_source.write_text(larger_rewrite, encoding="utf-8")
    larger_result = usage_collector.collect_registered_sources(rewrite_repo)
    assert_equal(larger_result["rotations"], 1, "larger in-place rewrite rotation")
    rewrite_registry = json.loads(usage_collector.registry_path(rewrite_repo).read_text(encoding="utf-8"))
    rewrite_row = next(iter(rewrite_registry["sources"].values()))
    assert_equal(rewrite_row["generation"], 2, "each in-place rewrite starts a generation")
    with sqlite3.connect(usage_ledger.usage_db_path(rewrite_repo)) as connection:
        rewrite_events = connection.execute(
            "SELECT event_id, input_tokens, output_tokens FROM usage_events ORDER BY cursor"
        ).fetchall()
    assert_equal(
        [(int(str(row[0]).split(":")[-2]), row[1], row[2]) for row in rewrite_events],
        [(0, 100, 20), (1, 200, 30), (2, 300, 40), (2, 10, 5)],
        "rewritten generations are indexed once with cumulative deltas preserved",
    )


def test_orca_adapter(orca_adapter, workflow_config, aow_doctor, base: Path) -> None:
    config = workflow_config.default_config()
    original_which = orca_adapter.shutil.which
    original_preflight_run = orca_adapter.subprocess.run
    class PreflightResult:
        returncode = 0
        stdout = "worker-start --model --effort --json --display-name"
    def fake_preflight_run(argv, **kwargs):
        result = PreflightResult()
        if argv[2:4] == ["worker-release", "--help"]:
            result.stdout = "worker-release --json"
        return result
    orca_adapter.shutil.which = lambda name: "/fake/orca"
    orca_adapter.subprocess.run = fake_preflight_run
    preflight_repo = base / "preflight-repo"
    preflight_repo.mkdir()
    run("git", "init", "-q", cwd=preflight_repo)
    workflow_config.write_default(preflight_repo)
    try:
        preflight = orca_adapter.preflight(preflight_repo)
    finally:
        orca_adapter.shutil.which = original_which
        orca_adapter.subprocess.run = original_preflight_run
    if not preflight["checks"]["release_status"]:
        raise AssertionError("preflight must check worker-release capability")
    if "sender_context" not in preflight["checks"]:
        raise AssertionError("preflight must report sender-context readiness")
    assert_equal(orca_adapter.terminal_name(config, "implementation", 1), "implementation1", "stage terminal name")
    current = orca_adapter.build_worker_command(
        config,
        route="hard",
        stage="implementation", effort="medium", task_id="task-1", run_id="run-1",
        ordinal=1, worktree="current",
    )
    if "--display-name" in current:
        raise AssertionError("current worktree command must not use creation-only display name")
    created = orca_adapter.build_worker_command(
        config, route="hard", stage="implementation", effort="medium", task_id="task-2", run_id="run-1",
        ordinal=2, worktree="new-child",
    )
    assert_equal(created[created.index("--display-name") + 1], "implementation2", "created worktree display name")
    command = orca_adapter.build_worker_command(
        config, route="hard", stage="implementation", effort="high", task_id="task-123", run_id="run-123",
        ordinal=1,
        worktree="new-top-level",
    )
    joined = " ".join(command)
    for expected in ("gpt-5.6-luna", "--effort high", "--display-name implementation1"):
        if expected not in joined:
            raise AssertionError(f"Orca command missing {expected}: {joined}")
    try:
        orca_adapter.build_worker_command(config, route="medium", stage="implementation", effort="medium", task_id="t", run_id="r", ordinal=1)
    except ValueError:
        pass
    else:
        raise AssertionError("Medium route must be blocked from worker launch")
    try:
        orca_adapter.build_worker_command(config, route="hard", stage="reviewer", effort="max", task_id="t", run_id="r", ordinal=1)
    except ValueError:
        pass
    else:
        raise AssertionError("adapter must enforce reviewer effort allowlist")
    repo = base / "adapter-repo"
    repo.mkdir()
    run("git", "init", "-q", cwd=repo)
    workflow_config.write_default(repo)
    active_dir = workflow_config.project_store_root(repo) / "runs" / "run-old" / "launches"
    active_dir.mkdir(parents=True)
    (active_dir / "task-old.attempt1.json").write_text(json.dumps({"status": "ready", "resource_status": "owned", "owned_paths": ["src/payments"]}), encoding="utf-8")
    try:
        orca_adapter.validate_launch_admission(repo, config, "task-new", "attempt1", ["src/payments/api.py"])
    except ValueError as exc:
        if "ownership" not in str(exc):
            raise
    else:
        raise AssertionError("adapter must reject overlapping active write ownership")
    lifecycle_cases = [
        ({"status": "completed", "resource_status": "owned"}, True),
        ({"status": "interrupted_unknown", "resource_status": "unknown"}, True),
        ({"status": "blocked", "resource_status": "owned"}, True),
        ({"status": "completed", "resource_status": "released"}, False),
        ({"status": "failed", "resource_status": "never_started"}, False),
    ]
    lifecycle_config = {**config, "runtime": {**config["runtime"], "max_parallel_workers": 10}}
    for index, (manifest, active) in enumerate(lifecycle_cases):
        (active_dir / f"lifecycle-{index}.attempt1.json").write_text(
            json.dumps({**manifest, "owned_paths": [f"lifecycle/{index}"]}), encoding="utf-8"
        )
        admitted = orca_adapter.validate_launch_admission(repo, lifecycle_config, f"new-{index}", "attempt1", ["docs"])
        expected = 1 + int(active)
        assert_equal(admitted["active_workers"], expected, f"resource state admission {index}")
        (active_dir / f"lifecycle-{index}.attempt1.json").unlink()
    (active_dir / "task-two.attempt1.json").write_text(json.dumps({"status": "running", "owned_paths": ["tests"]}), encoding="utf-8")
    try:
        orca_adapter.validate_launch_admission(repo, config, "task-three", "attempt1", ["docs"])
    except ValueError as exc:
        if "parallel" not in str(exc):
            raise
    else:
        raise AssertionError("adapter must enforce max_parallel_workers")
    marked = orca_adapter.mark_launch(
        repo,
        "run-old",
        "task-old",
        "attempt1",
        "completed",
        "tests passed",
        observed_model="gpt-5.6-luna",
        observed_effort="high",
    )
    assert_equal(marked["status"], "completed", "task outcome is recorded independently")
    assert_equal(marked["resource_status"], "owned", "task outcome does not settle resource")
    assert_equal(marked["evidence"], "tests passed", "worker completion evidence")
    assert_equal(marked["observed_model"], "gpt-5.6-luna", "confirmed worker model is recorded separately")
    assert_equal(marked["observed_effort"], "high", "confirmed worker effort is recorded separately")
    old_manifest = orca_adapter.launch_manifest_path(repo, "run-old", "task-old", "attempt1")
    old_payload = json.loads(old_manifest.read_text(encoding="utf-8"))
    old_payload["resource_status"] = "released"
    old_manifest.write_text(json.dumps(old_payload), encoding="utf-8")
    admitted = orca_adapter.validate_launch_admission(repo, config, "task-three", "attempt1", ["src/payments/api.py"])
    assert_equal(admitted["active_workers"], 1, "completed launch no longer counts as active")

    launch_calls = []
    worker_response = {"dispatch_id": "dispatch-1", "terminal": "terminal-1"}

    class FakeResult:
        returncode = 0
        stderr = ""

        def __init__(self):
            self.stdout = json.dumps(worker_response)

    original_run = orca_adapter.subprocess.run

    def fake_run(argv, **kwargs):
        if argv[:2] == ["git", "rev-parse"]:
            return original_run(argv, **kwargs)
        launch_calls.append(argv)
        return FakeResult()

    orca_adapter.subprocess.run = fake_run
    original_register_source = orca_adapter.usage_collector.register_source
    original_ensure_daemon = orca_adapter.usage_collector.ensure_daemon
    telemetry_calls = []
    daemon_calls = []
    runtime_source = repo / "runtime-session.jsonl"
    runtime_source.write_text(json.dumps({"type": "session"}) + "\n", encoding="utf-8")
    previous_source = repo / "previous-session.jsonl"
    previous_source.write_text(json.dumps({"type": "session"}) + "\n", encoding="utf-8")

    split_identity = {"sessionId": "split-session", "runtime": {"sessionPath": str(runtime_source)}}
    if orca_adapter.extract_runtime_source(split_identity) is not None:
        raise AssertionError("session and path from split branches must remain unconfirmed")
    current_identity = {
        "history": [{"sessionId": "old-session", "sessionPath": str(previous_source)}],
        "worker": {"session_id": "current-session", "rolloutPath": str(runtime_source)},
    }
    current_source = orca_adapter.extract_runtime_source(current_identity)
    assert_equal(current_source["session_id"] if current_source else None, "current-session", "current identity wins over history")
    for historical_key in ("archivedSessions", "pastRuns", "completedWorkers"):
        historical_identity = {historical_key: [{"sessionId": "historical", "sessionPath": str(previous_source)}]}
        if orca_adapter.extract_runtime_source(historical_identity) is not None:
            raise AssertionError(f"{historical_key} must not provide current runtime identity")
    ambiguous_identity = {
        "worker": [
            {"sessionId": "session-a", "sessionPath": str(runtime_source)},
            {"sessionId": "session-b", "sessionPath": str(previous_source)},
        ],
    }
    if orca_adapter.extract_runtime_source(ambiguous_identity) is not None:
        raise AssertionError("two complete runtime identities must remain ambiguous")
    same_identity = {"runtime": {"sessionId": "same-session", "sessionPath": str(runtime_source)}}
    same_source = orca_adapter.extract_runtime_source(same_identity)
    assert_equal(same_source["session_id"] if same_source else None, "same-session", "same-object identity is confirmed")

    def fake_register_source(root, source_path, **kwargs):
        telemetry_calls.append((root, source_path, kwargs))
        return {"source_id": "source-confirmed", "status": "registered", "path": str(source_path)}

    def fake_ensure_daemon(root, **kwargs):
        daemon_calls.append(root)
        return {"status": "running", "pid": 123, "started": True}

    orca_adapter.usage_collector.register_source = fake_register_source
    orca_adapter.usage_collector.ensure_daemon = fake_ensure_daemon
    try:
        orca_adapter.launch(repo, route="hard", stage="implementation", effort="medium",
                            task_id="task-launch", run_id="run-launch", ordinal=1)
        try:
            orca_adapter.launch(repo, route="hard", stage="implementation", effort="medium",
                                task_id="task-launch", run_id="run-launch", ordinal=1)
        except FileExistsError:
            pass
        else:
            raise AssertionError("duplicate launch attempt must be rejected")
    finally:
        orca_adapter.subprocess.run = original_run
    launched_manifest = json.loads(
        orca_adapter.launch_manifest_path(repo, "run-launch", "task-launch", "attempt1").read_text(encoding="utf-8")
    )
    assert_equal(launched_manifest["resource_status"], "owned", "successful worker launch owns resource")
    assert_equal(launched_manifest["dispatch_id"], "dispatch-1", "worker dispatch identity is retained")
    assert_equal(launched_manifest["terminal"], "terminal-1", "worker terminal identity is retained")
    assert_equal(launched_manifest["telemetry_status"], "unknown", "terminal-only launch has unknown telemetry")
    if telemetry_calls or daemon_calls:
        raise AssertionError("terminal-only launch must not register guessed telemetry")
    worker_starts = [call for call in launch_calls if call[:3] == ["orca", "orchestration", "worker-start"]]
    assert_equal(len(worker_starts), 1, "duplicate launch does not invoke Orca twice")

    launch_payload_for_confirmed = json.loads(
        orca_adapter.launch_manifest_path(repo, "run-launch", "task-launch", "attempt1").read_text(encoding="utf-8")
    )
    launch_payload_for_confirmed["resource_status"] = "stopped"
    orca_adapter.launch_manifest_path(repo, "run-launch", "task-launch", "attempt1").write_text(
        json.dumps(launch_payload_for_confirmed), encoding="utf-8"
    )
    (active_dir / "task-two.attempt1.json").unlink()
    worker_response = {
        "dispatch_id": "dispatch-confirmed", "terminal": "terminal-confirmed",
        "worker": {"sessionId": "session-confirmed", "sessionPath": str(runtime_source)},
    }
    orca_adapter.subprocess.run = fake_run
    confirmed = orca_adapter.launch(
        repo, route="hard", stage="implementation", effort="medium",
        task_id="task-confirmed", run_id="run-confirmed", ordinal=2, worktree="new-child",
    )
    assert_equal(confirmed["telemetry_status"], "confirmed", "launch attaches confirmed telemetry")
    assert_equal(confirmed["telemetry_source_id"], "source-confirmed", "launch persists telemetry source")
    assert_equal(telemetry_calls[-1][2]["run_id"], "run-confirmed", "worker run is preserved")
    assert_equal(telemetry_calls[-1][2]["stage"], "implementation", "worker stage is preserved")
    assert_equal(telemetry_calls[-1][2]["model"], "unknown", "unconfirmed worker model remains unknown")
    assert_equal(telemetry_calls[-1][2]["effort"], "unknown", "unconfirmed worker effort remains unknown")
    assert_equal(telemetry_calls[-1][2]["agent"], "codex", "worker agent is preserved")
    assert_equal(telemetry_calls[-1][2]["worktree"], "new-child", "worker worktree is preserved")
    assert_equal(len(daemon_calls), 1, "confirmed registration starts daemon once")

    confirmed_response = {"worker": {"sessionId": "session-confirmed", "sessionPath": str(runtime_source)}}
    def unavailable_register(*args, **kwargs):
        raise OSError("registry unavailable")

    orca_adapter.usage_collector.register_source = unavailable_register
    unavailable_registration = orca_adapter.attach_worker_telemetry(repo, confirmed, confirmed_response)
    assert_equal(unavailable_registration["telemetry_status"], "unavailable", "registration failure is unavailable")
    orca_adapter.usage_collector.register_source = fake_register_source

    def unavailable_daemon(*args, **kwargs):
        raise RuntimeError("daemon unavailable")

    orca_adapter.usage_collector.ensure_daemon = unavailable_daemon
    unavailable_daemon_result = orca_adapter.attach_worker_telemetry(repo, confirmed, confirmed_response)
    assert_equal(unavailable_daemon_result["telemetry_status"], "unavailable", "daemon failure is unavailable")
    orca_adapter.usage_collector.ensure_daemon = fake_ensure_daemon

    def stop_confirmed_resource() -> None:
        confirmed_path = orca_adapter.launch_manifest_path(repo, "run-confirmed", "task-confirmed", "attempt1")
        confirmed_manifest = json.loads(confirmed_path.read_text(encoding="utf-8"))
        confirmed_manifest["resource_status"] = "stopped"
        confirmed_path.write_text(json.dumps(confirmed_manifest), encoding="utf-8")

    worker_response = {
        "dispatch_id": "dispatch-registration-failure", "terminal": "terminal-registration-failure",
        "worker": {"sessionId": "session-registration-failure", "sessionPath": str(runtime_source)},
    }
    stop_confirmed_resource()
    orca_adapter.usage_collector.register_source = unavailable_register
    registration_failure = orca_adapter.launch(
        repo, route="hard", stage="implementation", effort="medium",
        task_id="task-registration-failure", run_id="run-registration-failure", ordinal=3,
    )
    assert_equal(registration_failure["telemetry_status"], "unavailable", "launch persists registration failure")
    if "registration unavailable" not in registration_failure["telemetry_gap"]:
        raise AssertionError("launch registration gap was not persisted")
    assert_equal(registration_failure["resource_status"], "owned", "registration failure preserves resource ownership")
    assert_equal(registration_failure["dispatch_id"], "dispatch-registration-failure", "registration failure preserves dispatch")
    assert_equal(registration_failure["telemetry_source_id"], None, "registration failure has no source")
    registration_failure_path = orca_adapter.launch_manifest_path(
        repo, "run-registration-failure", "task-registration-failure", "attempt1"
    )
    registration_failure_disk = json.loads(registration_failure_path.read_text(encoding="utf-8"))
    assert_equal(registration_failure_disk["telemetry_status"], "unavailable", "disk registration failure status")
    if "registration unavailable" not in registration_failure_disk["telemetry_gap"]:
        raise AssertionError("disk registration failure gap was not persisted")
    assert_equal(registration_failure_disk["telemetry_source_id"], None, "disk registration failure source")
    assert_equal(registration_failure_disk["dispatch_id"], "dispatch-registration-failure", "disk registration failure dispatch")
    assert_equal(registration_failure_disk["resource_status"], "owned", "disk registration failure resource")

    stop_confirmed_resource()
    orca_adapter.usage_collector.register_source = fake_register_source
    daemon_failure_calls = []

    def unavailable_daemon_launch(*args, **kwargs):
        daemon_failure_calls.append(args[0] if args else None)
        raise RuntimeError("daemon unavailable")

    orca_adapter.usage_collector.ensure_daemon = unavailable_daemon_launch
    worker_response = {
        "dispatch_id": "dispatch-daemon-failure", "terminal": "terminal-daemon-failure",
        "worker": {"sessionId": "session-daemon-failure", "sessionPath": str(runtime_source)},
    }
    daemon_failure = orca_adapter.launch(
        repo, route="hard", stage="implementation", effort="medium",
        task_id="task-daemon-failure", run_id="run-daemon-failure", ordinal=4,
    )
    assert_equal(daemon_failure["telemetry_status"], "unavailable", "launch persists daemon failure")
    if "daemon unavailable" not in daemon_failure["telemetry_gap"]:
        raise AssertionError("launch daemon gap was not persisted")
    assert_equal(daemon_failure["resource_status"], "owned", "daemon failure preserves resource ownership")
    assert_equal(daemon_failure["telemetry_source_id"], "source-confirmed", "daemon failure preserves source")
    assert_equal(daemon_failure["telemetry_daemon"]["status"], "unavailable", "daemon failure persists daemon state")
    assert_equal(len(daemon_calls), 1, "confirmed launch starts daemon exactly once")
    assert_equal(len(daemon_failure_calls), 1, "daemon failure launch calls daemon once")
    daemon_failure_path = orca_adapter.launch_manifest_path(
        repo, "run-daemon-failure", "task-daemon-failure", "attempt1"
    )
    daemon_failure_disk = json.loads(daemon_failure_path.read_text(encoding="utf-8"))
    assert_equal(daemon_failure_disk["telemetry_status"], "unavailable", "disk daemon failure status")
    if "daemon unavailable" not in daemon_failure_disk["telemetry_gap"]:
        raise AssertionError("disk daemon failure gap was not persisted")
    assert_equal(daemon_failure_disk["telemetry_source_id"], "source-confirmed", "disk daemon failure source")
    assert_equal(daemon_failure_disk["telemetry_daemon"]["status"], "unavailable", "disk daemon failure state")
    assert_equal(daemon_failure_disk["dispatch_id"], "dispatch-daemon-failure", "disk daemon failure dispatch")
    assert_equal(daemon_failure_disk["resource_status"], "owned", "disk daemon failure resource")
    orca_adapter.usage_collector.register_source = original_register_source
    orca_adapter.usage_collector.ensure_daemon = original_ensure_daemon

    if orca_adapter.extract_runtime_source({"dispatchId": "dispatch-only", "terminal": "terminal-only"}) is not None:
        raise AssertionError("dispatch/terminal identity must not produce a telemetry source")

    coverage_repo = base / "coverage-repo"
    coverage_repo.mkdir()
    run("git", "init", "-q", cwd=coverage_repo)
    workflow_config.write_default(coverage_repo)
    coverage_dir = workflow_config.project_store_root(coverage_repo) / "runs" / "coverage-run" / "launches"
    coverage_dir.mkdir(parents=True)
    (coverage_dir / "confirmed.attempt1.json").write_text(json.dumps({"status": "completed", "telemetry_status": "confirmed"}), encoding="utf-8")
    (coverage_dir / "unknown.attempt1.json").write_text(json.dumps({"status": "running", "telemetry_status": "unknown"}), encoding="utf-8")
    coverage = aow_doctor.check(coverage_repo)
    telemetry_check = next(item for item in coverage["checks"] if item["name"] == "telemetry")
    assert_equal(telemetry_check["confirmed_workers"], 1, "doctor confirmed worker telemetry count")
    assert_equal(telemetry_check["unknown_workers"], 1, "doctor unknown worker telemetry count")
    if coverage["status"] != "WARN":
        raise AssertionError("doctor must warn for active worker without confirmed telemetry")

    launch_manifest_for_interleave = orca_adapter.launch_manifest_path(repo, "run-launch", "task-launch", "attempt1")
    launch_payload_for_interleave = json.loads(launch_manifest_for_interleave.read_text(encoding="utf-8"))
    launch_payload_for_interleave["resource_status"] = "stopped"
    launch_manifest_for_interleave.write_text(json.dumps(launch_payload_for_interleave), encoding="utf-8")
    for failure_run, failure_task in (("run-registration-failure", "task-registration-failure"), ("run-daemon-failure", "task-daemon-failure")):
        failure_path = orca_adapter.launch_manifest_path(repo, failure_run, failure_task, "attempt1")
        failure_manifest = json.loads(failure_path.read_text(encoding="utf-8"))
        failure_manifest["resource_status"] = "stopped"
        failure_path.write_text(json.dumps(failure_manifest), encoding="utf-8")
    original_attach = orca_adapter.attach_worker_telemetry
    original_launch_run = orca_adapter.subprocess.run

    def interleaving_attach(root, manifest, response):
        orca_adapter.mark_launch(root, "run-interleaved", "task-interleaved", "attempt1", "completed", "concurrent")
        interleaved_path = orca_adapter.launch_manifest_path(root, "run-interleaved", "task-interleaved", "attempt1")
        interleaved_manifest = json.loads(interleaved_path.read_text(encoding="utf-8"))
        interleaved_manifest["resource_status"] = "stopped"
        interleaved_path.write_text(json.dumps(interleaved_manifest), encoding="utf-8")
        return original_attach(root, manifest, response)

    orca_adapter.attach_worker_telemetry = interleaving_attach
    orca_adapter.subprocess.run = fake_run
    orca_adapter.usage_collector.ensure_daemon = fake_ensure_daemon
    try:
        orca_adapter.launch(repo, route="hard", stage="implementation", effort="medium",
                            task_id="task-interleaved", run_id="run-interleaved", ordinal=1)
    finally:
        orca_adapter.attach_worker_telemetry = original_attach
        orca_adapter.subprocess.run = original_launch_run
        orca_adapter.usage_collector.ensure_daemon = original_ensure_daemon
    interleaved = json.loads(
        orca_adapter.launch_manifest_path(repo, "run-interleaved", "task-interleaved", "attempt1").read_text(encoding="utf-8")
    )
    assert_equal(interleaved["status"], "completed", "launch merge preserves concurrent task status")
    assert_equal(interleaved["resource_status"], "stopped", "launch merge preserves concurrent resource status")

    action_path = orca_adapter.launch_manifest_path(repo, "run-launch", "task-launch", "attempt1")
    action_manifest = json.loads(action_path.read_text(encoding="utf-8"))
    action_manifest.update({"dispatch_id": "dispatch-1", "resource_status": "owned"})
    action_path.write_text(json.dumps(action_manifest), encoding="utf-8")
    action_results = iter([
        (0, {"status": "stopped"}),
        (0, {"status": "retained"}),
        (1, {"status": "error"}),
    ])

    def fake_action_run(argv, **kwargs):
        if argv[:2] == ["git", "rev-parse"]:
            return original_run(argv, **kwargs)
        returncode, response = next(action_results)
        class ActionResult:
            pass
        result = ActionResult()
        result.returncode = returncode
        result.stdout = json.dumps(response)
        result.stderr = ""
        launch_calls.append(argv)
        return result

    orca_adapter.subprocess.run = fake_action_run
    try:
        stopped = orca_adapter.resource_action(repo, "run-launch", "task-launch", "attempt1", "stop")
        assert_equal(stopped["resource_status"], "stopped", "explicit stop settles resource")
        pending = orca_adapter.resource_action(repo, "run-launch", "task-launch", "attempt1", "release")
        assert_equal(pending["resource_status"], "release_pending", "retained release remains active")
        try:
            orca_adapter.resource_action(repo, "run-launch", "task-launch", "attempt1", "release")
        except RuntimeError:
            pass
        else:
            raise AssertionError("nonzero resource action must fail")
    finally:
        orca_adapter.subprocess.run = original_run
    persisted_action = json.loads(action_path.read_text(encoding="utf-8"))
    assert_equal(persisted_action["resource_status"], "unknown", "nonzero resource action is conservative")
    if not persisted_action.get("resource_response") or not persisted_action.get("resource_observed_at"):
        raise AssertionError("resource action evidence was not persisted")
    resource_calls = [call for call in launch_calls if "--dispatch" in call and ("worker-stop" in call or "worker-release" in call)]
    if any(call[call.index("--dispatch") + 1] != "dispatch-1" for call in resource_calls):
        raise AssertionError("resource actions must use the recorded dispatch ID")

    def set_owned_resource() -> None:
        current = json.loads(action_path.read_text(encoding="utf-8"))
        current.update({"dispatch_id": "dispatch-1", "resource_status": "owned"})
        action_path.write_text(json.dumps(current), encoding="utf-8")

    regression_results = iter([
        (0, {"status": "retained", "history": [{"status": "released"}]}),
        (0, {"status": "retained", "resource": {"status": "released"}}),
        (0, "not-json"),
        (0, {"result": {"state": "released"}}),
    ])

    def fake_regression_run(argv, **kwargs):
        if argv[:2] == ["git", "rev-parse"]:
            return original_run(argv, **kwargs)
        returncode, response = next(regression_results)
        class RegressionResult:
            pass
        result = RegressionResult()
        result.returncode = returncode
        result.stdout = response if isinstance(response, str) else json.dumps(response)
        result.stderr = ""
        return result

    orca_adapter.subprocess.run = fake_regression_run
    try:
        set_owned_resource()
        ambiguous = orca_adapter.resource_action(repo, "run-launch", "task-launch", "attempt1", "stop")
        assert_equal(ambiguous["resource_status"], "stop_pending", "nested history is not lifecycle evidence")
        set_owned_resource()
        conflicting = orca_adapter.resource_action(repo, "run-launch", "task-launch", "attempt1", "release")
        assert_equal(conflicting["resource_status"], "release_pending", "conflicting lifecycle evidence is unknown")
        set_owned_resource()
        unparseable = orca_adapter.resource_action(repo, "run-launch", "task-launch", "attempt1", "release")
        assert_equal(unparseable["resource_status"], "unknown", "unparseable action output is conservative")
        set_owned_resource()
        nested_released = orca_adapter.resource_action(repo, "run-launch", "task-launch", "attempt1", "release")
        assert_equal(nested_released["resource_status"], "released", "current result envelope settles release")
    finally:
        orca_adapter.subprocess.run = original_run

    set_owned_resource()
    def timeout_run(argv, **kwargs):
        if argv[:2] == ["git", "rev-parse"]:
            return original_run(argv, **kwargs)
        raise subprocess.TimeoutExpired(argv, 30)
    orca_adapter.subprocess.run = timeout_run
    try:
        timed_out = orca_adapter.resource_action(repo, "run-launch", "task-launch", "attempt1", "stop")
        assert_equal(timed_out["resource_status"], "unknown", "timeout is conservative")
    finally:
        orca_adapter.subprocess.run = original_run
    released_after_timeout = json.loads(action_path.read_text(encoding="utf-8"))
    released_after_timeout["resource_status"] = "released"
    action_path.write_text(json.dumps(released_after_timeout), encoding="utf-8")

    legacy_identity = active_dir / "legacy-completed.attempt1.json"
    legacy_identity.write_text(json.dumps({"status": "completed", "terminal": "terminal-old", "owned_paths": ["legacy"]}), encoding="utf-8")
    legacy_admission = orca_adapter.validate_launch_admission(
        repo, lifecycle_config, "legacy-check", "attempt1", ["docs"]
    )
    if legacy_admission["active_workers"] < 1:
        raise AssertionError("legacy completed manifest with identity must remain active")
    legacy_identity.unlink()

    if orca_adapter._response_proves_no_resource({"resource_created": False, "status": "owned"}):
        raise AssertionError("conflicting current resource evidence must remain unknown")
    if orca_adapter._response_proves_no_resource({"resourceCreated": False, "resourceState": "owned"}):
        raise AssertionError("conflicting camelCase resource evidence must remain unknown")
    for no_resource_error in ("selector_not_found", "no_active_sender_terminal"):
        response = {"ok": False, "error": {"code": no_resource_error}}
        if not orca_adapter._response_proves_no_resource(response):
            raise AssertionError(f"top-level {no_resource_error} must prove no resource")
        if orca_adapter._response_proves_no_resource({**response, "dispatch_id": "dispatch-conflict"}):
            raise AssertionError(f"{no_resource_error} with dispatch identity must remain unknown")
        if orca_adapter._response_proves_no_resource({**response, "terminal": "terminal-conflict"}):
            raise AssertionError(f"{no_resource_error} with terminal identity must remain unknown")

    failed_launch_results = iter([(1, {"error": "launch failed"}), (0, {"status": "stopped"})])
    def failed_launch_run(argv, **kwargs):
        if argv[:2] == ["git", "rev-parse"]:
            return original_run(argv, **kwargs)
        returncode, response = next(failed_launch_results)
        class FailedResult:
            pass
        result = FailedResult()
        result.returncode = returncode
        result.stdout = json.dumps(response)
        result.stderr = ""
        return result
    orca_adapter.subprocess.run = failed_launch_run
    try:
        try:
            orca_adapter.launch(repo, route="hard", stage="implementation", effort="medium",
                                task_id="task-failed", run_id="run-failed", ordinal=1)
        except RuntimeError:
            pass
        else:
            raise AssertionError("failed worker launch must fail")
    finally:
        orca_adapter.subprocess.run = original_run
    failed_manifest = json.loads(
        orca_adapter.launch_manifest_path(repo, "run-failed", "task-failed", "attempt1").read_text(encoding="utf-8")
    )
    assert_equal(failed_manifest["resource_status"], "unknown", "failed launch without proof is unknown")

    no_dispatch_path = orca_adapter.launch_manifest_path(repo, "run-no-resource", "task-no-resource", "attempt1")
    no_dispatch_path.parent.mkdir(parents=True, exist_ok=True)
    no_dispatch_response = {"ok": False, "error": {"code": "no_active_sender_terminal"}}
    no_dispatch_path.write_text(json.dumps({
        "schema": "aoc.worker_launch_request.v1",
        "status": "failed",
        "resource_status": "unknown",
        "response": no_dispatch_response,
    }), encoding="utf-8")
    settled_no_dispatch = orca_adapter.resource_action(
        repo, "run-no-resource", "task-no-resource", "attempt1", "release"
    )
    assert_equal(settled_no_dispatch["resource_status"], "never_started", "known no-resource error settles manifest")
    assert_equal(settled_no_dispatch["resource_action"], "release", "known no-resource action is recorded")
    assert_equal(settled_no_dispatch["resource_response"], no_dispatch_response, "known no-resource evidence is persisted")
    if not settled_no_dispatch.get("resource_observed_at"):
        raise AssertionError("known no-resource settlement must persist an observation time")

    set_owned_resource()
    def interleaved_run(argv, **kwargs):
        if argv[:2] == ["git", "rev-parse"]:
            return original_run(argv, **kwargs)
        orca_adapter.mark_launch(repo, "run-launch", "task-launch", "attempt1", "completed", "interleaved")
        class InterleavedResult:
            returncode = 0
            stdout = json.dumps({"status": "stopped"})
            stderr = ""
        return InterleavedResult()
    orca_adapter.subprocess.run = interleaved_run
    try:
        interleaved = orca_adapter.resource_action(repo, "run-launch", "task-launch", "attempt1", "stop")
    finally:
        orca_adapter.subprocess.run = original_run
    assert_equal(interleaved["status"], "completed", "resource action preserves concurrent task outcome")
    assert_equal(interleaved["evidence"], "interleaved", "resource action preserves concurrent task evidence")


def test_typed_terminal_effect_identity(orca_adapter, workflow_config, base: Path) -> None:
    """The observed result.effects terminal identity is deterministic evidence."""
    repo = base / "typed-terminal-effect-repo"
    repo.mkdir()
    run("git", "init", "-q", cwd=repo)
    workflow_config.write_default(repo)

    typed_effect = {
        "kind": "terminal",
        "role": "agent",
        "action": "created",
        "id": "term_effect_exact",
    }
    worker_response = {
        "ok": True,
        "result": {
            "dispatch_id": "dispatch_effect_exact",
            "effects": [typed_effect],
        },
    }
    calls: list[list[str]] = []
    original_run = orca_adapter.subprocess.run

    class Result:
        returncode = 0
        stderr = ""

        def __init__(self, response: dict):
            self.stdout = json.dumps(response)

    def fake_run(argv, **kwargs):
        calls.append(argv)
        if argv[:3] == ["orca", "orchestration", "worker-start"]:
            return Result(worker_response)
        if argv[:3] == ["orca", "terminal", "rename"]:
            return Result({"ok": True, "result": {"renamed": True}})
        return original_run(argv, **kwargs)

    try:
        orca_adapter.subprocess.run = fake_run
        launched = orca_adapter.launch(
            repo,
            route="hard",
            stage="planning",
            effort="medium",
            task_id="task-effect",
            run_id="run-effect",
            ordinal=1,
        )
    finally:
        orca_adapter.subprocess.run = original_run

    manifest_path = orca_adapter.launch_manifest_path(repo, "run-effect", "task-effect", "attempt1")
    persisted = json.loads(manifest_path.read_text(encoding="utf-8"))
    assert_equal(launched["terminal"], "term_effect_exact", "typed effect terminal handle")
    assert_equal(persisted["terminal"], "term_effect_exact", "manifest captures typed effect handle")
    assert_equal(persisted["terminal_rename_confirmed"], True, "typed effect rename confirmation")
    rename_calls = [call for call in calls if call[:3] == ["orca", "terminal", "rename"]]
    assert_equal(len(rename_calls), 1, "typed effect invokes one terminal rename")
    assert_equal(rename_calls[0][rename_calls[0].index("--terminal") + 1], "term_effect_exact", "rename exact effect handle")

    direct_only = {"terminal": "term_direct"}
    assert_equal(orca_adapter._find_terminal_handle(direct_only), "term_direct", "documented direct terminal key")
    historical = {"history": [{"terminal": "term_historical"}]}
    assert_equal(orca_adapter._find_terminal_handle(historical), None, "nested history terminal is ignored")

    no_resource_with_effect = {
        "ok": False,
        "error": {"code": "no_active_sender_terminal"},
        "result": {"effects": [typed_effect]},
    }
    assert_equal(
        orca_adapter._find_terminal_handle(no_resource_with_effect),
        "term_effect_exact",
        "no-resource response retains typed terminal evidence",
    )
    if orca_adapter._response_proves_no_resource(no_resource_with_effect):
        raise AssertionError("known no-resource error with typed terminal must remain unknown")

    no_resource_manifest = orca_adapter.launch_manifest_path(
        repo, "run-no-resource-effect", "task-no-resource-effect", "attempt1"
    )
    no_resource_manifest.parent.mkdir(parents=True, exist_ok=True)
    no_resource_manifest.write_text(
        json.dumps({
            "schema": "aoc.worker_launch_request.v1",
            "status": "failed",
            "resource_status": "unknown",
            "response": no_resource_with_effect,
        }),
        encoding="utf-8",
    )
    try:
        orca_adapter.resource_action(
            repo,
            "run-no-resource-effect",
            "task-no-resource-effect",
            "attempt1",
            "release",
        )
    except ValueError:
        pass
    else:
        raise AssertionError("typed terminal identity must block never_started reconciliation")
    after_no_resource = json.loads(no_resource_manifest.read_text(encoding="utf-8"))
    assert_equal(after_no_resource["resource_status"], "unknown", "typed no-resource remains unknown")

    conflicting = {
        "terminal": "term_direct",
        "result": {"effects": [typed_effect]},
    }
    assert_equal(orca_adapter._find_terminal_handle(conflicting), None, "direct/effect conflict is rejected")
    multiple_effects = {
        "result": {
            "effects": [
                typed_effect,
                {**typed_effect, "id": "term_effect_other"},
            ],
        },
    }
    assert_equal(orca_adapter._find_terminal_handle(multiple_effects), None, "multiple typed effect IDs are rejected")


def test_conservative_terminal_evidence(orca_adapter, workflow_config, base: Path) -> None:
    """Ambiguous terminal evidence blocks a no-resource conclusion everywhere."""
    cases = [
        (
            "conflicting-direct-effect",
            "selector_not_found",
            {
                "ok": False,
                "error": {"code": "selector_not_found"},
                "terminal": "term_direct",
                "result": {
                    "effects": [{
                        "kind": "terminal", "role": "agent", "action": "created", "id": "term_effect",
                    }],
                },
            },
        ),
        (
            "multiple-effects",
            "no_active_sender_terminal",
            {
                "ok": False,
                "error": {"code": "no_active_sender_terminal"},
                "result": {
                    "effects": [
                        {"kind": "terminal", "role": "agent", "action": "created", "id": "term_one"},
                        {"kind": "terminal", "role": "agent", "action": "created", "id": "term_two"},
                    ],
                },
            },
        ),
    ]

    class CommandResult:
        def __init__(self, returncode: int, stdout: str):
            self.returncode = returncode
            self.stdout = stdout
            self.stderr = ""

    original_run = orca_adapter.subprocess.run
    for label, _, response in cases:
        if orca_adapter._find_terminal_handle(response) is not None:
            raise AssertionError(f"{label} must not expose an unambiguous terminal handle")
        if orca_adapter._response_proves_no_resource(response):
            raise AssertionError(f"{label} terminal evidence must block no-resource proof")

        repo = base / f"ambiguous-{label}"
        repo.mkdir()
        run("git", "init", "-q", cwd=repo)
        workflow_config.write_default(repo)

        def failed_worker_run(argv, **kwargs):
            if argv[:3] == ["orca", "orchestration", "worker-start"]:
                return CommandResult(1, json.dumps(response))
            return original_run(argv, **kwargs)

        orca_adapter.subprocess.run = failed_worker_run
        try:
            try:
                orca_adapter.launch(
                    repo,
                    route="hard",
                    stage="planning",
                    effort="medium",
                    task_id=f"task-{label}",
                    run_id=f"run-{label}",
                    ordinal=1,
                )
            except RuntimeError:
                pass
            else:
                raise AssertionError(f"{label} failed launch must raise")
        finally:
            orca_adapter.subprocess.run = original_run

        launch_path = orca_adapter.launch_manifest_path(repo, f"run-{label}", f"task-{label}", "attempt1")
        persisted = json.loads(launch_path.read_text(encoding="utf-8"))
        assert_equal(persisted["resource_status"], "unknown", f"{label} launch remains unknown")
        assert_equal(persisted.get("terminal"), None, f"{label} launch does not guess terminal")
        assert_equal(persisted["response"], response, f"{label} launch preserves response evidence")

        reconcile_path = orca_adapter.launch_manifest_path(
            repo, f"run-reconcile-{label}", f"task-reconcile-{label}", "attempt1"
        )
        reconcile_path.parent.mkdir(parents=True, exist_ok=True)
        reconcile_path.write_text(
            json.dumps({
                "schema": "aoc.worker_launch_request.v1",
                "status": "failed",
                "resource_status": "unknown",
                "response": response,
            }),
            encoding="utf-8",
        )
        try:
            orca_adapter.resource_action(
                repo,
                f"run-reconcile-{label}",
                f"task-reconcile-{label}",
                "attempt1",
                "release",
            )
        except ValueError:
            pass
        else:
            raise AssertionError(f"{label} reconciliation must not settle never_started")
        reconciled = json.loads(reconcile_path.read_text(encoding="utf-8"))
        assert_equal(reconciled["resource_status"], "unknown", f"{label} reconciliation remains unknown")


def test_terminal_rename_confirmation(orca_adapter, workflow_config, base: Path) -> None:
    """A rename is confirmed only by a strict JSON success response."""
    worker_response = {"ok": True, "dispatch_id": "dispatch-rename", "terminal": "term_exact"}
    cases = [
        ("exact-success", 0, {"ok": True, "result": {"renamed": True}}, True),
        (
            "matching-identity",
            0,
            {"ok": True, "result": {"renamed": True, "terminal": "term_exact"}},
            True,
        ),
        ("invalid-json", 0, "not-json", False),
        ("empty-json", 0, "", False),
        ("ok-false", 0, {"ok": False, "result": {"renamed": True}}, False),
        ("renamed-false", 0, {"ok": True, "result": {"renamed": False}}, False),
        ("renamed-missing", 0, {"ok": True, "result": {}}, False),
        (
            "conflicting-identity",
            0,
            {"ok": True, "result": {"renamed": True, "terminal": "term_other"}},
            False,
        ),
    ]

    class CommandResult:
        def __init__(self, returncode: int, stdout: str):
            self.returncode = returncode
            self.stdout = stdout
            self.stderr = ""

    original_run = orca_adapter.subprocess.run
    for label, rename_code, rename_response, expected in cases:
        repo = base / f"rename-{label}"
        repo.mkdir()
        run("git", "init", "-q", cwd=repo)
        workflow_config.write_default(repo)
        rename_calls: list[list[str]] = []

        if isinstance(rename_response, str):
            rename_stdout = rename_response
        else:
            rename_stdout = json.dumps(rename_response)

        def fake_run(argv, **kwargs):
            if argv[:3] == ["orca", "orchestration", "worker-start"]:
                return CommandResult(0, json.dumps(worker_response))
            if argv[:3] == ["orca", "terminal", "rename"]:
                rename_calls.append(argv)
                return CommandResult(rename_code, rename_stdout)
            return original_run(argv, **kwargs)

        orca_adapter.subprocess.run = fake_run
        try:
            launched = orca_adapter.launch(
                repo,
                route="hard",
                stage="planning",
                effort="medium",
                task_id=f"task-{label}",
                run_id=f"run-{label}",
                ordinal=1,
            )
        finally:
            orca_adapter.subprocess.run = original_run

        manifest_path = orca_adapter.launch_manifest_path(repo, f"run-{label}", f"task-{label}", "attempt1")
        persisted = json.loads(manifest_path.read_text(encoding="utf-8"))
        assert_equal(launched["terminal"], "term_exact", f"{label} preserves exact terminal")
        assert_equal(persisted["terminal"], "term_exact", f"{label} persists exact terminal")
        assert_equal(persisted["terminal_rename_confirmed"], expected, f"{label} rename confirmation")
        assert_equal(len(rename_calls), 1, f"{label} invokes one rename")
        target = rename_calls[0][rename_calls[0].index("--terminal") + 1]
        assert_equal(target, "term_exact", f"{label} rename target")


def test_final_review_regressions(workflow_config, usage_ledger, usage_collector, orca_adapter, aow_tui, aow_gui, base: Path) -> None:
    """Final review gates: conservative evidence, fidelity, attribution, scope, policy, and admission."""
    # Any affirmative resource evidence wins over a no-resource error.  The
    # persisted manifest/reconciliation path is covered below with the same
    # response shape.
    for key in ("resource_created", "resourceCreated", "resource_exists", "resourceExists"):
        response = {"ok": False, "error": {"code": "selector_not_found"}, key: True}
        if orca_adapter._response_proves_no_resource(response):
            raise AssertionError(f"affirmative {key} must block no-resource proof")
    for response in (
        {"ok": False, "error": {"code": "no_active_sender_terminal"}, "resource_created": False, "status": "owned"},
        {"ok": False, "error": {"code": "selector_not_found"}, "resourceCreated": False, "resourceState": "owned"},
    ):
        if orca_adapter._response_proves_no_resource(response):
            raise AssertionError("conflicting resource booleans/states must remain unknown")

    lifecycle_repo = base / "final-lifecycle"
    lifecycle_repo.mkdir()
    run("git", "init", "-q", cwd=lifecycle_repo)
    workflow_config.write_default(lifecycle_repo)
    lifecycle_response = {"ok": False, "error": {"code": "selector_not_found"}, "resource_created": True}
    lifecycle_path = orca_adapter.launch_manifest_path(lifecycle_repo, "run-final", "task-final", "attempt1")
    lifecycle_path.parent.mkdir(parents=True, exist_ok=True)
    lifecycle_path.write_text(json.dumps({
        "schema": "aoc.worker_launch_request.v1",
        "status": "failed",
        "resource_status": "unknown",
        "response": lifecycle_response,
    }), encoding="utf-8")
    try:
        orca_adapter.resource_action(lifecycle_repo, "run-final", "task-final", "attempt1", "release")
    except ValueError:
        pass
    else:
        raise AssertionError("affirmative launch evidence must block never_started reconciliation")
    persisted_lifecycle = json.loads(lifecycle_path.read_text(encoding="utf-8"))
    assert_equal(persisted_lifecycle["resource_status"], "unknown", "affirmative launch evidence remains unknown")

    def token_event(info: dict[str, object]) -> dict[str, object]:
        return {"type": "event_msg", "payload": {"type": "token_count", "info": info}}

    complete_info = {
        "total_token_usage": {"input_tokens": 10, "output_tokens": 4, "total_tokens": 14},
        "last_token_usage": {"input_tokens": 10, "output_tokens": 4},
    }
    missing_optional = token_event({"total_token_usage": {"input_tokens": 10, "output_tokens": 4, "total_tokens": 14}})
    extracted_missing = usage_collector._extract_usage_totals(missing_optional)
    if extracted_missing is None:
        raise AssertionError("missing optional token detail should remain a usable partial event")
    missing_usage = extracted_missing[0]
    assert_equal(missing_usage["cached_input_tokens"], None, "missing cached detail remains unknown")
    assert_equal(missing_usage["reasoning_output_tokens"], None, "missing reasoning detail remains unknown")
    for bad_info in (
        {"total_token_usage": {"input_tokens": 10, "output_tokens": 4, "cached_input_tokens": -1}},
        {"total_token_usage": {"input_tokens": 10, "output_tokens": 4, "reasoning_output_tokens": -1}},
        {"total_token_usage": {"input_tokens": 10, "output_tokens": 4, "cached_input_tokens": 11}},
        {"total_token_usage": {"input_tokens": 10, "output_tokens": 4, "reasoning_output_tokens": 5}},
        {"total_token_usage": {"input_tokens": 10, "output_tokens": "not-a-number"}},
    ):
        if usage_collector._extract_usage_totals(token_event(bad_info)) is not None:
            raise AssertionError("malformed/negative/impossible optional detail must be quarantined")
    explicit_zero = usage_collector._extract_usage_totals(token_event({
        "total_token_usage": {
            "input_tokens": 10, "cached_input_tokens": 0,
            "output_tokens": 4, "reasoning_output_tokens": 0, "total_tokens": 14,
        },
    }))
    if explicit_zero is None:
        raise AssertionError("explicit zero token detail must remain valid")
    assert_equal(explicit_zero[0]["cached_input_tokens"], 0, "explicit cached zero remains measured")

    malformed_repo = base / "final-malformed"
    malformed_repo.mkdir()
    malformed_source = base / "final-malformed.jsonl"
    malformed_source.write_text(json.dumps(token_event({
        "total_token_usage": {
            "input_tokens": 10, "cached_input_tokens": 11,
            "output_tokens": 4, "reasoning_output_tokens": 0, "total_tokens": 14,
        },
    })) + "\n", encoding="utf-8")
    usage_collector.register_source(
        malformed_repo,
        malformed_source,
        session_id="malformed",
        model="gpt-5.6-sol",
        effort="medium",
    )
    malformed_result = usage_collector.collect_registered_sources(malformed_repo)
    assert_equal(malformed_result["malformed"], 1, "invalid token event is quarantined")
    assert_equal(malformed_result["ignored"], 0, "invalid token event is not silently ignored")
    assert_equal(
        usage_ledger.project_usage_summary(malformed_repo)["coverage_status"],
        "partial",
        "quarantined token event leaves a visible coverage gap",
    )

    legacy_key_record = {
        "project_id": "project", "source": "codex-runtime", "session_id": "session",
        "run_id": "run", "task_id": "task", "attempt_id": "attempt2",
        "stage": "implementation", "model": "gpt-5.6-luna",
    }
    legacy_key = hashlib.sha256(
        "project\x1fcodex-runtime\x1fsession\x1frun\x1fimplementation\x1fgpt-5.6-luna".encode("utf-8")
    ).hexdigest()
    assert_equal(
        usage_ledger._source_key(legacy_key_record),
        legacy_key,
        "task attribution preserves the pre-upgrade cumulative source key",
    )

    registration_repo = base / "final-registration-upgrade"
    registration_repo.mkdir()
    registration_source = base / "final-registration-upgrade.jsonl"
    registration_source.write_text("", encoding="utf-8")
    usage_collector.register_source(
        registration_repo,
        registration_source,
        session_id="registration-upgrade",
        model="gpt-5.6-sol",
        effort="high",
    )
    preserved_registration = usage_collector.register_source(
        registration_repo,
        registration_source,
        session_id="registration-upgrade",
    )
    assert_equal(preserved_registration["model"], "gpt-5.6-sol", "re-registration preserves observed model")
    assert_equal(preserved_registration["effort"], "high", "re-registration preserves observed effort")
    assert_equal(preserved_registration["observed_model"], "gpt-5.6-sol", "re-registration preserves observed model field")
    assert_equal(preserved_registration["observed_effort"], "high", "re-registration preserves observed effort field")
    assert_equal(preserved_registration["attribution_status"], "confirmed", "re-registration preserves attribution")

    # A hand-written v1 registry may contain only requested policy values in
    # model/effort.  Upgrade must not promote them to observed attribution or
    # use them to fabricate model pricing.  Explicit observed provenance can
    # confirm the source later.
    legacy_repo = base / "legacy-registry-attribution"
    legacy_repo.mkdir()
    run("git", "init", "-q", cwd=legacy_repo)
    (legacy_repo / ".orca").mkdir()
    (legacy_repo / ".orca" / "pricing.json").write_text(json.dumps({"models": {
        "gpt-5.6-luna": {
            "input_per_million": 1, "cached_input_per_million": 0.5,
            "output_per_million": 2,
        },
    }}), encoding="utf-8")
    legacy_source = base / "legacy-registry-attribution.jsonl"
    legacy_source.write_text(json.dumps(token_event({
        "total_token_usage": {
            "input_tokens": 10, "cached_input_tokens": 0,
            "output_tokens": 4, "reasoning_output_tokens": 0, "total_tokens": 14,
        },
    })) + "\n", encoding="utf-8")
    legacy_source_id = hashlib.sha256(
        f"{legacy_source.resolve()}\x1flegacy-registry-session".encode("utf-8")
    ).hexdigest()[:24]
    legacy_registry = usage_collector.registry_path(legacy_repo)
    legacy_registry.parent.mkdir(parents=True, exist_ok=True)
    legacy_registry.write_text(json.dumps({"schema": "aoc.telemetry_sources.v1", "sources": {
        legacy_source_id: {
            "source_id": legacy_source_id,
            "path": str(legacy_source.resolve()),
            "session_id": "legacy-registry-session",
            "route": "hard",
            "stage": "implementation",
            "model": "gpt-5.6-luna",
            "effort": "high",
            "offset": 0,
        },
    }}), encoding="utf-8")
    legacy_status = usage_collector.status(legacy_repo)
    assert_equal(legacy_status["sources"][0]["model"], "unknown", "legacy status does not expose old model as actual")
    usage_collector.collect_registered_sources(legacy_repo)
    upgraded_registry = json.loads(legacy_registry.read_text(encoding="utf-8"))
    upgraded_source = upgraded_registry["sources"][legacy_source_id]
    assert_equal(upgraded_source["model"], "unknown", "legacy model is not promoted to observed")
    assert_equal(upgraded_source["effort"], "unknown", "legacy effort is not promoted to observed")
    assert_equal(upgraded_source["requested_model"], "gpt-5.6-luna", "legacy model retained as requested metadata")
    assert_equal(upgraded_source["requested_effort"], "high", "legacy effort retained as requested metadata")
    assert_equal(upgraded_source["attribution_status"], "partial", "legacy attribution remains partial")
    legacy_summary = usage_ledger.project_usage_summary(legacy_repo)
    assert_equal(legacy_summary["coverage_status"], "partial", "legacy attribution gap is visible")
    assert_equal(legacy_summary["totals"]["cost_usd"], None, "legacy unknown model has no fabricated cost")
    confirmed_legacy = usage_collector.register_source(
        legacy_repo,
        legacy_source,
        session_id="legacy-registry-session",
        model="gpt-5.6-sol",
        effort="high",
        observed_model="gpt-5.6-sol",
        observed_effort="high",
    )
    assert_equal(confirmed_legacy["model"], "gpt-5.6-sol", "explicit observed model confirms legacy source")
    assert_equal(confirmed_legacy["effort"], "high", "explicit observed effort confirms legacy source")
    assert_equal(confirmed_legacy["attribution_status"], "confirmed", "explicit observed attribution confirms legacy source")

    fidelity_repo = base / "final-fidelity"
    fidelity_repo.mkdir()
    run("git", "init", "-q", cwd=fidelity_repo)
    fidelity_source = base / "final-fidelity.jsonl"
    fidelity_source.write_text(json.dumps(missing_optional) + "\n", encoding="utf-8")
    usage_collector.register_source(fidelity_repo, fidelity_source, session_id="fidelity", run_id="fidelity-run", model="gpt-5.6-sol", effort="medium")
    usage_collector.collect_registered_sources(fidelity_repo)
    fidelity_summary = usage_ledger.project_usage_summary(fidelity_repo, run_id="fidelity-run")
    assert_equal(fidelity_summary["coverage_status"], "partial", "missing optional detail marks partial coverage")
    assert_equal(fidelity_summary["totals"]["cached_input_tokens"], None, "summary preserves unknown cached detail")
    report = capture_project_report(usage_ledger, fidelity_summary)
    if "Cached input: unknown" not in report or "Reasoning output: unknown" not in report:
        raise AssertionError("text renderer must preserve unknown token detail")
    tui = "\n".join(aow_tui.lines_usage(fidelity_repo, "fidelity-run"))
    if "cached input: unknown" not in tui or "reasoning output: unknown" not in tui:
        raise AssertionError("TUI renderer must preserve unknown token detail")
    gui = "".join(aow_gui.render_fragment(aow_gui.collect(fidelity_repo), "usage"))
    assert_equal(gui_metric_value(gui, "Cached input"), "unknown", "GUI cached detail")
    assert_equal(gui_metric_value(gui, "Reasoning output"), "unknown", "GUI reasoning detail")

    priced_repo = base / "final-fidelity-priced"
    priced_repo.mkdir()
    (priced_repo / ".orca").mkdir()
    (priced_repo / ".orca" / "pricing.json").write_text(json.dumps({"models": {"gpt-5.6-sol": {
        "input_per_million": 1, "cached_input_per_million": 0.5, "output_per_million": 2,
    }}}), encoding="utf-8")
    unknown_price = usage_ledger.record_usage_event(priced_repo, {
        "event_id": "unknown-split", "model": "gpt-5.6-sol", "input_tokens": 10, "output_tokens": 4,
    })
    assert_equal(unknown_price["cost_usd"], None, "unknown cache split has no invented numeric cost")
    assert_equal(unknown_price["pricing_status"], "unknown", "unknown cache split pricing status")
    zero_price = usage_ledger.record_usage_event(priced_repo, {
        "event_id": "explicit-zero-split", "model": "gpt-5.6-sol", "input_tokens": 10,
        "cached_input_tokens": 0, "output_tokens": 4, "reasoning_output_tokens": 0,
    })
    if zero_price["cost_usd"] is None:
        raise AssertionError("explicit zero cache/reasoning detail should permit pricing")

    # Confirmed source registration keeps requested policy separate from actual
    # runtime identity and carries retry identity into emitted events.
    attribution_repo = base / "final-attribution"
    attribution_repo.mkdir()
    run("git", "init", "-q", cwd=attribution_repo)
    runtime_path = base / "final-attribution.jsonl"
    runtime_path.write_text(json.dumps({"type": "session"}) + "\n", encoding="utf-8")
    registrations: list[dict[str, object]] = []
    original_register = orca_adapter.usage_collector.register_source
    original_daemon = orca_adapter.usage_collector.ensure_daemon
    def capture_register(root, source_path, **kwargs):
        registrations.append(kwargs)
        return {"source_id": "source-final", "status": "registered"}
    try:
        orca_adapter.usage_collector.register_source = capture_register
        orca_adapter.usage_collector.ensure_daemon = lambda root: {"status": "running"}
        attached = orca_adapter.attach_worker_telemetry(attribution_repo, {
            "run_id": "run-attribution", "task_id": "task-attribution", "attempt_id": "attempt2",
            "route": "hard", "stage": "implementation", "worktree": "new-child",
            "agent": "codex", "requested_model": "gpt-5.6-luna", "requested_effort": "medium",
        }, {"worker": {"sessionId": "session-attribution", "sessionPath": str(runtime_path),
                         "observedModel": "gpt-5.6-sol", "observedEffort": "high"}})
    finally:
        orca_adapter.usage_collector.register_source = original_register
        orca_adapter.usage_collector.ensure_daemon = original_daemon
    assert_equal(attached["telemetry_status"], "confirmed", "confirmed worker source status")
    assert_equal(registrations[0]["task_id"], "task-attribution", "source task identity")
    assert_equal(registrations[0]["attempt_id"], "attempt2", "source retry identity")
    assert_equal(registrations[0]["requested_model"], "gpt-5.6-luna", "requested model is separate")
    assert_equal(registrations[0]["requested_effort"], "medium", "requested effort is separate")
    assert_equal(registrations[0]["model"], "gpt-5.6-sol", "observed model is actual")
    assert_equal(registrations[0]["effort"], "high", "observed effort is actual")
    registrations.clear()
    try:
        orca_adapter.usage_collector.register_source = capture_register
        orca_adapter.usage_collector.ensure_daemon = lambda root: {"status": "running"}
        orca_adapter.attach_worker_telemetry(attribution_repo, {
            "run_id": "run-attribution", "task_id": "task-attribution", "attempt_id": "attempt3",
            "route": "hard", "stage": "implementation", "worktree": "new-child",
            "agent": "codex", "requested_model": "gpt-5.6-luna", "requested_effort": "medium",
        }, {"worker": {"sessionId": "session-attribution-2", "sessionPath": str(runtime_path)}})
    finally:
        orca_adapter.usage_collector.register_source = original_register
        orca_adapter.usage_collector.ensure_daemon = original_daemon
    assert_equal(registrations[0]["model"], "unknown", "unconfirmed actual model remains unknown")
    assert_equal(registrations[0]["effort"], "unknown", "unconfirmed actual effort remains unknown")
    source_line = usage_collector._usage_from_line({
        "source_id": "source-final", "generation": 0, "worktree_id": "wt-final", "run_id": "run-attribution",
        "session_id": "session-attribution", "route": "hard", "stage": "implementation", "agent": "codex",
        "model": "unknown", "effort": "unknown", "task_id": "task-attribution", "attempt_id": "attempt2",
    }, (json.dumps({"type": "event_msg", "payload": {"type": "token_count", "info": complete_info}}) + "\n").encode(), 0)
    assert_equal(source_line["task_id"], "task-attribution", "emitted task attribution")
    assert_equal(source_line["attempt_id"], "attempt2", "emitted retry attribution")

    # Coverage registration and source counts are scoped with the same filters
    # as usage rows.
    scope_repo = base / "final-scope"
    scope_repo.mkdir()
    run("git", "init", "-q", cwd=scope_repo)
    for run_id, session_id in (("run-a", "session-a"), ("run-b", "session-b")):
        source = base / f"{run_id}.jsonl"
        source.write_text(json.dumps(token_event({"total_token_usage": {"input_tokens": 3, "cached_input_tokens": 0, "output_tokens": 2, "reasoning_output_tokens": 0, "total_tokens": 5}})) + "\n", encoding="utf-8")
        usage_collector.register_source(
            scope_repo,
            source,
            session_id=session_id,
            run_id=run_id,
            model="gpt-5.6-sol",
            effort="medium",
        )
    usage_collector.collect_registered_sources(scope_repo)
    scoped = usage_ledger.project_usage_summary(scope_repo, run_id="run-a")
    assert_equal(scoped["source_count"], 1, "run-scoped source count")
    assert_equal(scoped["coverage_status"], "observed", "run-scoped observed coverage")
    missing_scope = usage_ledger.project_usage_summary(scope_repo, run_id="run-missing")
    assert_equal(missing_scope["source_count"], 0, "missing run source count")
    assert_equal(missing_scope["coverage_status"], "unknown", "missing run coverage stays unknown")
    session_scope = usage_ledger.project_usage_summary(scope_repo, session_id="session-a")
    assert_equal(session_scope["source_count"], 1, "session-scoped source count")
    assert_equal(usage_ledger.project_usage_summary(scope_repo, worktree="wt-does-not-exist")["coverage_status"], "unknown", "worktree isolation")

    # Config policy is fail-closed for defaults and stage model names.
    invalid_config = workflow_config.default_config()
    invalid_config["routes"]["default"] = "hard"
    invalid_config["stages"]["planning"]["model_name"] = "gpt-4o"
    config_errors = workflow_config.validate_config(invalid_config)
    if not any("routes.default" in error for error in config_errors):
        raise AssertionError("non-light config default must be rejected")
    if not any("model_name" in error and "gpt-5.6" in error for error in config_errors):
        raise AssertionError("unsupported stage model must be rejected")

    # Corrupt, unreadable, and path-unsafe manifests block admission without
    # being rewritten or quarantined.
    admission_repo = base / "final-admission"
    admission_repo.mkdir()
    run("git", "init", "-q", cwd=admission_repo)
    workflow_config.write_default(admission_repo)
    launch_root = workflow_config.project_store_root(admission_repo) / "runs" / "bad-run" / "launches"
    launch_root.mkdir(parents=True, exist_ok=True)
    malformed = launch_root / "malformed.attempt1.json"
    malformed.write_text("{not-json", encoding="utf-8")
    try:
        orca_adapter.validate_launch_admission(admission_repo, workflow_config.default_config(), "new", "attempt1", [])
    except ValueError as exc:
        if "malformed.attempt1.json" not in str(exc) or "repair" not in str(exc).lower():
            raise AssertionError(f"malformed manifest error is not actionable: {exc}")
    else:
        raise AssertionError("malformed launch manifest must block admission")
    assert_equal(malformed.read_text(encoding="utf-8"), "{not-json", "malformed evidence is not overwritten")
    malformed.unlink()
    unreadable = launch_root / "unreadable.attempt1.json"
    unreadable.mkdir()
    try:
        orca_adapter.validate_launch_admission(admission_repo, workflow_config.default_config(), "new", "attempt1", [])
    except ValueError as exc:
        if "unreadable.attempt1.json" not in str(exc):
            raise AssertionError(f"unreadable manifest path is missing: {exc}")
    else:
        raise AssertionError("unreadable launch manifest must block admission")
    unreadable.rmdir()
    outside = base / "outside-manifest.json"
    outside.write_text(json.dumps({"status": "running"}), encoding="utf-8")
    unsafe = launch_root / "unsafe.attempt1.json"
    unsafe.symlink_to(outside)
    try:
        orca_adapter.validate_launch_admission(admission_repo, workflow_config.default_config(), "new", "attempt1", [])
    except ValueError:
        pass
    else:
        raise AssertionError("symlinked launch manifest must block admission")


def main() -> None:
    workflow_config = load("workflow_config")
    usage_ledger = load("usage_ledger")
    aow_gui = load("aow_gui")
    aow_tui = load("aow_tui")
    usage_collector = load("usage_collector")
    orca_adapter = load("orca_adapter")
    aow_doctor = load("aow_doctor")
    with tempfile.TemporaryDirectory(prefix="aow-v2-") as raw:
        base = Path(raw)
        test_route_and_config(workflow_config)
        test_shared_project_store(workflow_config, base)
        test_usage_accounting(usage_ledger, aow_gui, base)
        test_usage_coverage_render_matrix(usage_ledger, usage_collector, aow_gui, aow_tui, base)
        test_public_cli(base)
        test_incremental_collector(usage_collector, usage_ledger, workflow_config, base)
        test_orca_adapter(orca_adapter, workflow_config, aow_doctor, base)
        test_typed_terminal_effect_identity(orca_adapter, workflow_config, base)
        test_conservative_terminal_evidence(orca_adapter, workflow_config, base)
        test_terminal_rename_confirmation(orca_adapter, workflow_config, base)
        test_final_review_regressions(workflow_config, usage_ledger, usage_collector, orca_adapter, aow_tui, aow_gui, base)
    print("ALL WORKFLOW V2 VALIDATION CHECKS PASSED")


if __name__ == "__main__":
    main()
