from __future__ import annotations

import io
import json
import re
import runpy
import sqlite3
import subprocess
import tempfile
from pathlib import Path

import pytest
import yaml

from adw_modules import kubernetes_job, repository_registry, result_mailbox, tracer

IMAGE = kubernetes_job.APPROVED_IMAGE
_IMMUTABLE_FACTORY_IMAGE = re.compile(
    r"ghcr\.io/alexcodeplace/overdeck-agent-sandbox@sha256:[0-9a-f]{64}"
)
_KEY_DIRECTORY = tempfile.TemporaryDirectory(prefix="factory-test-signing-")
_PRIVATE_KEY_PATH = Path(_KEY_DIRECTORY.name) / "private.pem"
_PUBLIC_KEY_PATH = Path(_KEY_DIRECTORY.name) / "public.pem"
subprocess.run(["openssl", "genpkey", "-algorithm", "ED25519", "-out", str(_PRIVATE_KEY_PATH)], check=True, capture_output=True)
subprocess.run(["openssl", "pkey", "-in", str(_PRIVATE_KEY_PATH), "-pubout", "-out", str(_PUBLIC_KEY_PATH)], check=True, capture_output=True)
KEY = _PRIVATE_KEY_PATH.read_bytes()
PUBLIC_KEY = _PUBLIC_KEY_PATH.read_bytes()
EXEC = "refs/heads/factory-exec/" + "1" * 24
RESULT = "refs/heads/factory-result/" + "1" * 24


def registry() -> repository_registry.Registry:
    return repository_registry.load()


def envelope(repository_id: str = "overdeck", invocation: tuple[str, ...] = ("build", "request")) -> repository_registry.AttemptEnvelope:
    return repository_registry.issue_envelope(
        registry(),
        repository_id,
        repository_registry.WORKFLOW,
        KEY,
        attempt_id="1" * 24,
        execution_ref=EXEC,
        result_ref=RESULT,
        input_commit="0" * 40,
        invocation=invocation,
        upload_sha256=result_mailbox.target_sha256(upload()),
        ttl_seconds=600,
    )


def upload() -> result_mailbox.UploadTarget:
    return result_mailbox.UploadTarget(
        f"https://{kubernetes_job.APPROVED_RESULT_HOST}:4443/v1/results/" + "1" * 24,
        "t" * 43,
        "-----BEGIN CERTIFICATE-----\ntest\n-----END CERTIFICATE-----\n",
    )


def job(repository_id: str = "overdeck") -> dict:
    return kubernetes_job.manifest(
        job_name="factory-1111111111111111",
        namespace="overdeck-factory",
        image=IMAGE,
        secret=envelope(repository_id).credential_handle,
        upload_secret="factory-upload-1111111111111111",
        service_account="overdeck-factory-runner",
        envelope=envelope(repository_id),
        execution_ref=EXEC,
        result_ref=RESULT,
        invocation=["build", "request"],
        upload=upload(),
        timeout=600,
    )


def test_registered_repository_builds_secret_free_restricted_job() -> None:
    manifest = job()
    encoded = json.dumps(manifest)
    assert "hostPath" not in encoded and "nodeSelector" not in encoded and "ConfigMap" not in encoded
    pod = manifest["spec"]["template"]["spec"]
    containers = {item["name"]: item for item in pod["containers"]}
    container = containers["factory"]
    assert pod["automountServiceAccountToken"] is False
    assert pod["imagePullSecrets"] == [{"name": "overdeck-ghcr-pull"}]
    assert pod["hostAliases"] == [{"ip": "10.43.200.8", "hostnames": ["overdeck-factory-egress"]}]
    assert container["imagePullPolicy"] == "Always"
    volumes = {volume["name"]: volume for volume in pod["volumes"]}
    assert volumes["repository-key"]["secret"] == {"secretName": "factory-git-overdeck", "defaultMode": 256}
    assert volumes["runtime-credentials"]["secret"] == {"secretName": "overdeck-factory-runtime", "defaultMode": 288}
    mounts = {mount["mountPath"]: mount for mount in container["volumeMounts"]}
    assert mounts["/home/factory/.ssh/id_ed25519"] == {"name": "repository-key", "mountPath": "/home/factory/.ssh/id_ed25519", "subPath": "id_ed25519", "readOnly": True}
    assert "/home/factory/.netrc" not in mounts
    assert container["securityContext"] == {"allowPrivilegeEscalation": False, "readOnlyRootFilesystem": True, "capabilities": {"drop": ["ALL"]}}
    assert container["resources"] == {
        "requests": {"cpu": "1", "memory": "2Gi"},
        "limits": {"cpu": "4", "memory": "8Gi"},
    }
    environment = {item["name"]: item for item in container["env"]}
    assert environment["FACTORY_REPOSITORY_ID"]["value"] == "overdeck"
    assert environment["FACTORY_CREDENTIAL_HANDLE"]["value"] == "factory-git-overdeck"
    assert environment["HTTPS_PROXY"]["value"] == kubernetes_job.APPROVED_PROXY_URL
    assert environment["HTTP_PROXY"]["value"] == kubernetes_job.APPROVED_PROXY_URL
    assert environment["UV_NO_SYNC"]["value"] == "1"
    assert manifest["spec"]["template"]["metadata"]["labels"] == manifest["metadata"]["labels"]
    assert "https://github.com" not in environment["FACTORY_ATTEMPT_ENVELOPE"]["value"]
    assert environment["FACTORY_EXECUTION_REF"]["value"] == EXEC and environment["FACTORY_RESULT_REF"]["value"] == RESULT
    assert environment["FACTORY_RESULT_UPLOAD_JSON"] == {
        "name": "FACTORY_RESULT_UPLOAD_JSON",
        "valueFrom": {"secretKeyRef": {"name": "factory-upload-1111111111111111", "key": "upload-json"}},
    }
    assert set(containers) == {"factory", "result-reader"}
    assert mounts["/result"] == {"name": "result", "mountPath": "/result"}
    assert containers["result-reader"]["securityContext"] == container["securityContext"]
    assert containers["result-reader"]["resources"] == {
        "requests": {"cpu": "1m", "memory": "8Mi"},
        "limits": {"cpu": "10m", "memory": "16Mi"},
    }
    assert containers["result-reader"]["volumeMounts"] == [{"name": "result", "mountPath": "/result", "readOnly": True}]
    assert volumes["result"] == {"name": "result", "emptyDir": {"sizeLimit": "256Mi"}}
    assert volumes["workspace"] == {"name": "workspace", "emptyDir": {"sizeLimit": "8Gi"}}
    assert volumes["home"] == {"name": "home", "emptyDir": {"sizeLimit": "2Gi"}}
    assert volumes["tmp"] == {"name": "tmp", "emptyDir": {"sizeLimit": "256Mi"}}


def test_upload_secret_is_immutable_attempt_owned_and_not_controller_readable() -> None:
    target_json = result_mailbox.serialize_target(upload())
    secret = kubernetes_job.upload_secret_manifest(
        "overdeck-factory",
        "factory-upload-1111111111111111",
        "factory-1111111111111111",
        "12345678-1234-1234-1234-123456789abc",
        target_json,
    )
    assert secret["immutable"] is True and set(secret["data"]) == {"upload-json"}
    assert secret["metadata"]["ownerReferences"] == [{
        "apiVersion": "batch/v1",
        "kind": "Job",
        "name": "factory-1111111111111111",
        "uid": "12345678-1234-1234-1234-123456789abc",
        "controller": True,
        "blockOwnerDeletion": True,
    }]
    assert upload().token not in json.dumps(job())


def test_mailbox_policy_is_attempt_owned_and_allows_only_exact_endpoint() -> None:
    policy = kubernetes_job.mailbox_network_policy(
        "overdeck-factory",
        "factory-mailbox-1111111111111111",
        "factory-1111111111111111",
        "12345678-1234-1234-1234-123456789abc",
        upload(),
    )
    assert policy["metadata"]["ownerReferences"] == [{
        "apiVersion": "batch/v1",
        "kind": "Job",
        "name": "factory-1111111111111111",
        "uid": "12345678-1234-1234-1234-123456789abc",
        "controller": True,
        "blockOwnerDeletion": True,
    }]
    assert policy["spec"] == {
        "podSelector": {"matchLabels": {"overdeck.dev/factory-job": "factory-1111111111111111"}},
        "policyTypes": ["Egress"],
        "egress": [{
            "to": [{"ipBlock": {"cidr": f"{kubernetes_job.APPROVED_RESULT_HOST}/32"}}],
            "ports": [{"protocol": "TCP", "port": 4443}],
        }],
    }


def test_envelope_uses_asymmetric_attempt_key_and_binds_upload() -> None:
    source = envelope()
    assert source.upload_sha256 == result_mailbox.target_sha256(upload())
    assert repository_registry.verify_envelope(source.serialize(), registry(), PUBLIC_KEY) == source
    with pytest.raises(ValueError, match="integrity"):
        repository_registry.verify_envelope(source.serialize(), registry(), b"not a public key")
    with pytest.raises(ValueError, match="private signing key"):
        repository_registry.issue_envelope(
            registry(), "overdeck", repository_registry.WORKFLOW, PUBLIC_KEY,
            attempt_id="1" * 24, execution_ref=EXEC, result_ref=RESULT,
            input_commit="0" * 40, invocation=("build",),
            upload_sha256=result_mailbox.target_sha256(upload()), ttl_seconds=60,
        )


def test_unknown_repository_is_rejected_before_submission(tmp_path: Path, monkeypatch) -> None:
    calls: list[list[str]] = []
    monkeypatch.setattr(kubernetes_job, "workspace_snapshot", lambda _: (_ for _ in ()).throw(AssertionError("submission started")))
    monkeypatch.setattr(kubernetes_job, "run", lambda argv, **kwargs: calls.append(list(argv)))
    with pytest.raises(ValueError, match="not registered"):
        kubernetes_job.submit(tmp_path, ["build"], "unknown", key=KEY)
    assert calls == []


@pytest.mark.parametrize("url", [
    "http://github.com/alexcodeplace/overdeck.git",
    "https://github.com/AlexCodePlace/overdeck.git",
    "https://github.com/alexcodeplace/overdeck",
    "https://github.com/alexcodeplace/overdeck.git/",
    "https://github.com/alexcodeplace/overdeck.git?redirect=1",
    "https://github.com/alexcodeplace/overdeck.git#redirect",
    "https://github.com/alexcodeplace/overdeck.git/../attacker.git",
    "https://user:password@github.com/alexcodeplace/overdeck.git",
    "https://github.com/avi-ezra/Press.Zone-Works.git",
])
def test_origin_aliases_redirects_credentials_and_cross_repository_fail(url: str) -> None:
    with pytest.raises(ValueError, match="origin URL"):
        repository_registry.authorize_origin(registry(), "overdeck", repository_registry.WORKFLOW, url)


def test_cross_repository_credentials_fail() -> None:
    source = envelope("overdeck")
    tampered = repository_registry.AttemptEnvelope(
        **{**source.payload(), "credential_handle": "factory-git-press-zone-works", "credential_profile": "press-zone-works", "signature": source.signature}
    )
    with pytest.raises(ValueError, match="integrity"):
        repository_registry.verify_envelope(tampered.serialize(), registry(), PUBLIC_KEY)


def test_envelope_tampering_and_policy_drift_fail() -> None:
    source = envelope()
    tampered = json.loads(source.serialize())
    tampered["repository_id"] = "press-zone-works"
    with pytest.raises(ValueError, match="integrity"):
        repository_registry.verify_envelope(json.dumps(tampered), registry(), PUBLIC_KEY)
    drifted = repository_registry.Registry(registry().version, "0" * 64, registry().repositories)
    with pytest.raises(ValueError, match="registry drift"):
        repository_registry.verify_envelope(source.serialize(), drifted, PUBLIC_KEY)
    policy = {**source.payload(), "resource_policy": {**source.resource_policy, "cpu_limit": "8"}}
    signature = repository_registry._sign_payload(policy, KEY)
    with pytest.raises(ValueError, match="policy drift"):
        repository_registry.verify_envelope(json.dumps({**policy, "signature": signature}), registry(), PUBLIC_KEY)


def test_envelope_binds_attempt_refs_input_invocation_and_expiry() -> None:
    source = envelope()
    assert source.attempt_id == "1" * 24
    assert source.execution_ref == EXEC and source.result_ref == RESULT
    assert source.input_commit == "0" * 40
    assert source.invocation_sha256 == repository_registry.invocation_sha256(("build", "request"))
    mismatched = {**source.payload(), "execution_ref": "refs/heads/factory-exec/" + "2" * 24}
    signature = repository_registry._sign_payload(mismatched, KEY)
    with pytest.raises(ValueError, match="refs do not match"):
        repository_registry.verify_envelope(json.dumps({**mismatched, "signature": signature}), registry(), PUBLIC_KEY)
    expired = repository_registry.issue_envelope(
        registry(), "overdeck", repository_registry.WORKFLOW, KEY,
        attempt_id="1" * 24, execution_ref=EXEC, result_ref=RESULT,
        input_commit="0" * 40, invocation=("build",), upload_sha256=result_mailbox.target_sha256(upload()), ttl_seconds=60, now=1000,
    )
    with pytest.raises(ValueError, match="expired or is not yet valid"):
        repository_registry.verify_envelope(expired.serialize(), registry(), PUBLIC_KEY, now=1060)
    with pytest.raises(ValueError, match="expired or is not yet valid"):
        repository_registry.verify_envelope(expired.serialize(), registry(), PUBLIC_KEY, now=939)
    with pytest.raises(ValueError, match="does not match the Job transport"):
        kubernetes_job.manifest(
            job_name="factory-1111111111111111", namespace="overdeck-factory", image=IMAGE,
            secret="factory-git-overdeck", upload_secret="factory-upload-1111111111111111", service_account="runner", envelope=source,
            execution_ref=EXEC, result_ref=RESULT, invocation=["different"],
            upload=upload(), timeout=600,
        )


def test_registry_removal_revokes_new_execution_and_preserves_recovery(tmp_path: Path) -> None:
    document = json.loads(repository_registry.REGISTRY_PATH.read_text())
    document["repositories"] = [item for item in document["repositories"] if item["id"] != "overdeck"]
    revoked_path = tmp_path / "repositories.v1.json"
    revoked_path.write_text(json.dumps(document))
    revoked = repository_registry.load(revoked_path)
    recovery_ref = "refs/heads/factory-result/" + "a" * 24
    with pytest.raises(ValueError, match="not registered"):
        repository_registry.resolve(revoked, "overdeck", repository_registry.WORKFLOW)
    assert recovery_ref.startswith("refs/heads/factory-result/")


def test_manifest_rejects_cross_repository_secret() -> None:
    with pytest.raises(ValueError, match="repository-bound credential Secret"):
        kubernetes_job.manifest(job_name="factory-1111111111111111", namespace="overdeck-factory", image=IMAGE, secret="factory-git-press-zone-works", upload_secret="factory-upload-1111111111111111", service_account="runner", envelope=envelope("overdeck"), execution_ref=EXEC, result_ref=RESULT, invocation=["build"], upload=upload(), timeout=600)


def test_manifest_rejects_mutable_image_or_unauthorized_workflow() -> None:
    with pytest.raises(ValueError, match="immutable GHCR digest"):
        kubernetes_job.manifest(job_name="factory-1111111111111111", namespace="overdeck-factory", image="ghcr.io/alexcodeplace/overdeck-agent-sandbox:latest", secret="factory-git-overdeck", upload_secret="factory-upload-1111111111111111", service_account="runner", envelope=envelope(), execution_ref=EXEC, result_ref=RESULT, invocation=["build"], upload=upload(), timeout=600)
    unauthorized = repository_registry.AttemptEnvelope(**{**envelope().payload(), "workflow": "other", "signature": envelope().signature})
    with pytest.raises(ValueError, match="authorized workflow"):
        kubernetes_job.manifest(job_name="factory-1111111111111111", namespace="overdeck-factory", image=IMAGE, secret="factory-git-overdeck", upload_secret="factory-upload-1111111111111111", service_account="runner", envelope=unauthorized, execution_ref=EXEC, result_ref=RESULT, invocation=["build"], upload=upload(), timeout=600)


def init_repo(path: Path) -> None:
    subprocess.run(["git", "init", "-q", str(path)], check=True)
    subprocess.run(["git", "-C", str(path), "config", "user.email", "test@example.invalid"], check=True)
    subprocess.run(["git", "-C", str(path), "config", "user.name", "Test"], check=True)
    (path / ".gitignore").write_text("ignored\n")
    (path / "tracked").write_text("before\n")
    subprocess.run(["git", "-C", str(path), "add", ".gitignore", "tracked"], check=True)
    subprocess.run(["git", "-C", str(path), "commit", "-qm", "base"], check=True)


def fake_log_collection(argv, path: Path, **kwargs) -> tuple[bool, str | None]:
    path.write_text("worker log\n")
    return False, None


@pytest.fixture
def live_trace_follower(tmp_path: Path, monkeypatch):
    trace_db = tmp_path / "canonical" / "sssf.db"
    trace_events = tmp_path / "canonical" / "events.jsonl"
    monkeypatch.setattr(
        kubernetes_job, "canonical_trace_paths", lambda: (trace_db, trace_events),
    )

    class FakeFollower:
        def __init__(self, argv, path, mirror, **kwargs):
            self.path = path
            self.mirror = mirror
            self.truncated = False
            records = [
                {
                    "version": 1, "kind": "session_start", "attempt_id": mirror.attempt_id,
                    "ts": "2026-08-15T10:00:00Z", "adw_id": f"adw-{mirror.attempt_id[:8]}",
                    "engineer": "factory", "host": "worker-pod",
                    "started_at": "2026-08-15T10:00:00Z",
                },
                {
                    "version": 1, "kind": "session_finish", "attempt_id": mirror.attempt_id,
                    "ts": "2026-08-15T10:00:01Z", "adw_id": f"adw-{mirror.attempt_id[:8]}",
                    "status": "success", "ended_at": "2026-08-15T10:00:01Z",
                },
            ]
            lines = [
                tracer.KUBERNETES_TRACE_PREFIX + json.dumps(record, separators=(",", ":"))
                for record in records
            ]
            path.write_text("worker log\n" + "\n".join(lines) + "\n")
            for line in lines:
                mirror.consume(line)

        def finish(self, timeout=10):
            return False, None

        def abort(self):
            return None

    monkeypatch.setattr(kubernetes_job, "JobLogFollower", FakeFollower)
    return trace_db


def test_bounded_output_records_truncation_and_preserves_small_output(tmp_path: Path) -> None:
    small = tmp_path / "small.log"
    assert kubernetes_job.run_bounded_output(
        ["/usr/bin/python3", "-c", "print('small')"], small, limit=1024, timeout=10
    ) == (False, None)
    assert small.read_bytes() == b"small\n"

    large = tmp_path / "large.log"
    assert kubernetes_job.run_bounded_output(
        ["/usr/bin/python3", "-c", "import os; os.write(1, b'x' * 4096)"],
        large,
        limit=1024,
        timeout=10,
    ) == (True, None)
    assert large.stat().st_size == 1024
    assert large.read_bytes().endswith(kubernetes_job.LOG_TRUNCATION_MARKER)

    failed = tmp_path / "failed.log"
    assert kubernetes_job.run_bounded_output(
        ["/usr/bin/python3", "-c", "import sys; print('log error', file=sys.stderr); raise SystemExit(2)"],
        failed,
        limit=1024,
        timeout=10,
    ) == (False, "log error")


def _trace_line(attempt_id: str, kind: str, **fields) -> bytes:
    record = {
        "version": 1, "kind": kind, "attempt_id": attempt_id,
        "ts": "2026-08-15T10:00:00Z", **fields,
    }
    return (tracer.KUBERNETES_TRACE_PREFIX + json.dumps(record, separators=(",", ":")) + "\n").encode()


class _FollowerProcess:
    def __init__(self, output: bytes):
        self.stdout = io.BytesIO(output)
        self.returncode = 0
        self.terminated = False

    def wait(self, timeout=None):
        return self.returncode

    def poll(self):
        return self.returncode

    def terminate(self):
        self.terminated = True
        self.returncode = -15

    def kill(self):
        self.returncode = -9


def test_live_follower_mirrors_real_session_level_events(
    tmp_path: Path, monkeypatch, capsys,
) -> None:
    attempt = "7" * 24
    monkeypatch.setenv("FACTORY_ATTEMPT_ENVELOPE", json.dumps({"attempt_id": attempt}))
    producer = tracer.Tracer(tmp_path / "producer.db", tmp_path / "producer.jsonl")
    try:
        producer.session_start("live-run", "factory", repo=tmp_path)
        producer.event(tracer.EventRecord(
            adw_id="live-run", type="log", name="console", payload={"level": "info"},
        ))
        producer.session_finish("live-run", ok=True)
    finally:
        producer.conn.close()

    output = capsys.readouterr().out.encode()
    records = [
        json.loads(line.removeprefix(tracer.KUBERNETES_TRACE_PREFIX.encode()))
        for line in output.splitlines()
    ]
    session_event = next(record for record in records if record["kind"] == "event")
    assert "phase_id" not in session_event

    mirror = tracer.KubernetesTraceMirror(
        tmp_path / "mirror.db", tmp_path / "mirror.jsonl", attempt,
        {"cluster": "k3s", "job": "factory-job", "pod": "factory-pod"},
    )
    process = _FollowerProcess(output)
    follower = kubernetes_job.JobLogFollower(
        ["kubectl", "logs"], tmp_path / "job.log", mirror,
        popen=lambda *args, **kwargs: process,
    )
    truncated, error = follower.finish()
    assert truncated is False and error is None
    mirror.complete()
    assert mirror.tracer.conn.execute(
        "SELECT status FROM sessions WHERE adw_id='live-run'",
    ).fetchone() == ("success",)
    assert mirror.tracer.conn.execute(
        "SELECT type,name FROM events WHERE adw_id='live-run' AND type='log'",
    ).fetchone() == ("log", "console")
    mirror.close()


def test_live_follower_rejects_malformed_mismatched_and_incomplete_trace(tmp_path: Path) -> None:
    attempt = "e" * 24
    placement = {"cluster": "k3s", "job": "factory-job", "pod": "factory-pod"}

    for bad_line, expected in [
        (tracer.KUBERNETES_TRACE_PREFIX.encode() + b"{bad}\n", "malformed"),
        (_trace_line("f" * 24, "session_start", adw_id="run-1", host="pod", started_at="2026-08-15T10:00:00Z"), "contract mismatch"),
    ]:
        mirror = tracer.KubernetesTraceMirror(
            tmp_path / f"{expected}.db", tmp_path / f"{expected}.jsonl", attempt, placement,
        )
        process = _FollowerProcess(b"human log\n" + bad_line)
        follower = kubernetes_job.JobLogFollower(
            ["kubectl", "logs"], tmp_path / f"{expected}.log", mirror,
            popen=lambda *args, **kwargs: process,
        )
        truncated, error = follower.finish()
        assert truncated is False and expected in (error or "")
        assert not follower.thread.is_alive() and process.stdout.closed
        with pytest.raises(RuntimeError, match="missing session"):
            mirror.complete()
        mirror.close()

    mirror = tracer.KubernetesTraceMirror(
        tmp_path / "incomplete.db", tmp_path / "incomplete.jsonl", attempt, placement,
    )
    mirror.consume(_trace_line(
        attempt, "session_start", adw_id="run-2", engineer="factory", host="pod",
        started_at="2026-08-15T10:00:00Z",
    ).rstrip())
    with pytest.raises(RuntimeError, match="missing session start or finish"):
        mirror.complete()
    mirror.fail()
    assert mirror.tracer.conn.execute(
        "SELECT status FROM sessions WHERE adw_id='run-2'",
    ).fetchone() == ("fail",)
    mirror.close()


def test_mirror_rejects_every_record_after_session_finish(tmp_path: Path) -> None:
    attempt = "9" * 24
    mirror = tracer.KubernetesTraceMirror(
        tmp_path / "trace.db", tmp_path / "events.jsonl", attempt,
        {"cluster": "k3s", "job": "factory-job", "pod": "factory-pod"},
    )
    mirror.consume(_trace_line(
        attempt, "session_start", adw_id="closed-run", engineer="factory", host="pod",
        started_at="2026-08-15T10:00:00Z",
    ).rstrip())
    mirror.consume(_trace_line(
        attempt, "session_finish", adw_id="closed-run", status="success",
        ended_at="2026-08-15T10:00:01Z",
    ).rstrip())
    mirror.complete()
    post_finish = [
        _trace_line(attempt, "session_usage", adw_id="closed-run", tokens=99, cost=1.0),
        _trace_line(
            attempt, "event", event_id="late-event", adw_id="closed-run",
            event_type="log", name="console", metadata={"level": "info"},
            started_at="2026-08-15T10:00:02Z",
        ),
        _trace_line(
            attempt, "session_start", adw_id="closed-run", engineer="factory", host="pod",
            started_at="2026-08-15T10:00:02Z",
        ),
    ]
    for line in post_finish:
        with pytest.raises(ValueError, match="followed session finish"):
            mirror.consume(line.rstrip())
    assert mirror.tracer.conn.execute(
        "SELECT status,total_tokens,total_cost FROM sessions WHERE adw_id='closed-run'",
    ).fetchone() == ("success", 0, 0.0)
    assert mirror.tracer.conn.execute(
        "SELECT 1 FROM events WHERE event_id='late-event'",
    ).fetchone() is None
    mirror.close()


def test_complete_rejects_unfinished_phase_attempt_and_tool_call(tmp_path: Path) -> None:
    attempt = "8" * 24

    def open_mirror(name: str) -> tracer.KubernetesTraceMirror:
        mirror = tracer.KubernetesTraceMirror(
            tmp_path / f"{name}.db", tmp_path / f"{name}.jsonl", attempt,
            {"cluster": "k3s", "job": f"job-{name}", "pod": f"pod-{name}"},
        )
        mirror.consume(_trace_line(
            attempt, "session_start", adw_id=name, engineer="factory", host=f"pod-{name}",
            started_at="2026-08-15T10:00:00Z",
        ).rstrip())
        return mirror

    def phase_record(adw_id: str, status: str, ended_at: str | None = None) -> bytes:
        fields = {
            "adw_id": adw_id, "phase_id": f"{adw_id}-phase", "seq": 1,
            "name": "build", "phase_kind": "agent", "owner": "builder",
            "status": status, "phase_attempt": 0, "retries": 0,
            "started_at": "2026-08-15T10:00:00Z",
        }
        if ended_at is not None:
            fields["ended_at"] = ended_at
        return _trace_line(attempt, "phase", **fields)

    phase_mirror = open_mirror("unfinished-phase")
    phase_mirror.consume(phase_record("unfinished-phase", "running").rstrip())
    phase_mirror.consume(_trace_line(
        attempt, "session_finish", adw_id="unfinished-phase", status="success",
        ended_at="2026-08-15T10:00:03Z",
    ).rstrip())
    with pytest.raises(RuntimeError, match="unfinished phases"):
        phase_mirror.complete()
    phase_mirror.close()

    attempt_mirror = open_mirror("unfinished-attempt")
    attempt_mirror.consume(phase_record(
        "unfinished-attempt", "success", "2026-08-15T10:00:01Z",
    ).rstrip())
    attempt_mirror.consume(_trace_line(
        attempt, "agent_attempt_start", adw_id="unfinished-attempt",
        phase_id="unfinished-attempt-phase", agent_attempt_id="open-attempt",
        agent="builder", session_id="session", host="pod", model="model",
        started_at="2026-08-15T10:00:01Z",
    ).rstrip())
    attempt_mirror.consume(_trace_line(
        attempt, "session_finish", adw_id="unfinished-attempt", status="success",
        ended_at="2026-08-15T10:00:03Z",
    ).rstrip())
    with pytest.raises(RuntimeError, match="agent attempts"):
        attempt_mirror.complete()
    attempt_mirror.close()

    tool_mirror = open_mirror("unfinished-tool")
    tool_mirror.consume(phase_record(
        "unfinished-tool", "success", "2026-08-15T10:00:01Z",
    ).rstrip())
    tool_mirror.consume(_trace_line(
        attempt, "agent_attempt_start", adw_id="unfinished-tool",
        phase_id="unfinished-tool-phase", agent_attempt_id="closed-attempt",
        agent="builder", session_id="session", host="pod", model="model",
        started_at="2026-08-15T10:00:01Z",
    ).rstrip())
    tool_mirror.consume(_trace_line(
        attempt, "agent_attempt_finish", agent_attempt_id="closed-attempt",
        returncode=0, timed_out=False, tokens=1, ended_at="2026-08-15T10:00:02Z",
    ).rstrip())
    tool_mirror.consume(_trace_line(
        attempt, "event", event_id="open-tool-event", adw_id="unfinished-tool",
        phase_id="unfinished-tool-phase", event_type="tool_call_start", name="read",
        metadata={"tool_call_id": "open-tool", "attempt_id": "closed-attempt", "seq": 1, "tool": "read"},
        started_at="2026-08-15T10:00:02Z",
    ).rstrip())
    tool_mirror.consume(_trace_line(
        attempt, "session_finish", adw_id="unfinished-tool", status="success",
        ended_at="2026-08-15T10:00:03Z",
    ).rstrip())
    with pytest.raises(RuntimeError, match="tool calls"):
        tool_mirror.complete()
    tool_mirror.close()


def test_live_follower_caps_log_and_treats_truncation_as_trace_failure(tmp_path: Path) -> None:
    attempt = "a" * 24
    mirror = tracer.KubernetesTraceMirror(
        tmp_path / "trace.db", tmp_path / "events.jsonl", attempt,
        {"cluster": "k3s", "job": "factory-job", "pod": "factory-pod"},
    )
    output = (
        _trace_line(attempt, "session_start", adw_id="run-3", engineer="factory", host="pod", started_at="2026-08-15T10:00:00Z")
        + b"x" * 2048 + b"\n"
        + _trace_line(attempt, "session_finish", adw_id="run-3", status="success", ended_at="2026-08-15T10:00:01Z")
    )
    process = _FollowerProcess(output)
    log = tmp_path / "job.log"
    follower = kubernetes_job.JobLogFollower(
        ["kubectl", "logs"], log, mirror, limit=1024,
        popen=lambda *args, **kwargs: process,
    )
    truncated, error = follower.finish()
    assert truncated is True and "configured limit" in (error or "")
    assert log.stat().st_size == 1024 and log.read_bytes().endswith(kubernetes_job.LOG_TRUNCATION_MARKER)
    assert not follower.thread.is_alive() and process.stdout.closed
    with pytest.raises(RuntimeError, match="missing session start or finish"):
        mirror.complete()
    assert mirror.tracer.conn.execute(
        "SELECT status FROM sessions WHERE adw_id='run-3'",
    ).fetchone() == ("running",)
    mirror.close()


def test_live_follower_stops_mirroring_at_record_budget(tmp_path: Path) -> None:
    attempt = "b" * 24
    mirror = tracer.KubernetesTraceMirror(
        tmp_path / "trace.db", tmp_path / "events.jsonl", attempt,
        {"cluster": "k3s", "job": "factory-job", "pod": "factory-pod"},
    )
    output = b"".join([
        _trace_line(attempt, "session_start", adw_id="run-budget", engineer="factory", host="pod", started_at="2026-08-15T10:00:00Z"),
        _trace_line(attempt, "session_usage", adw_id="run-budget", tokens=1, cost=0.0),
        _trace_line(attempt, "session_usage", adw_id="run-budget", tokens=1000, cost=1000.0),
        _trace_line(attempt, "session_finish", adw_id="run-budget", status="success", ended_at="2026-08-15T10:00:01Z"),
    ])
    process = _FollowerProcess(output)
    follower = kubernetes_job.JobLogFollower(
        ["kubectl", "logs"], tmp_path / "job.log", mirror, limit=8192, record_limit=2,
        popen=lambda *args, **kwargs: process,
    )
    truncated, error = follower.finish()
    assert truncated is False and error == "Kubernetes trace record budget exceeded"
    assert follower.records_consumed == 2
    assert mirror.tracer.conn.execute(
        "SELECT total_tokens,total_cost,status FROM sessions WHERE adw_id='run-budget'",
    ).fetchone() == (1, 0.0, "running")
    with pytest.raises(RuntimeError, match="missing session start or finish"):
        mirror.complete()
    mirror.close()


def test_mirror_rejects_cross_attempt_identifiers_and_parent_conflicts(tmp_path: Path) -> None:
    attempt = "c" * 24
    db = tmp_path / "trace.db"
    events = tmp_path / "events.jsonl"
    seed = tracer.Tracer(db, events, emit_kubernetes_trace=False)
    try:
        seed.session_start("other-run", "factory")
        seed.conn.execute(
            "INSERT INTO phases (phase_id,adw_id,seq,name,kind,owner,status,attempt,retries)"
            " VALUES ('shared-phase','other-run',1,'other','agent','builder','running',0,0)",
        )
        seed.conn.execute(
            "INSERT INTO agent_attempts (attempt_id,adw_id,phase_id,agent,started_at)"
            " VALUES ('shared-attempt','other-run','shared-phase','builder','2026-08-15T09:00:00Z')",
        )
        seed.conn.execute(
            "INSERT INTO events (event_id,adw_id,phase_id,type,name,payload_json,started_at)"
            " VALUES ('shared-event','other-run','shared-phase','agent_start','builder','{}','2026-08-15T09:00:00Z')",
        )
    finally:
        seed.conn.close()
    mirror = tracer.KubernetesTraceMirror(
        db, events, attempt, {"cluster": "k3s", "job": "job-own", "pod": "pod-own"},
    )
    mirror.consume(_trace_line(
        attempt, "session_start", adw_id="own-run", engineer="factory", host="pod-own",
        started_at="2026-08-15T10:00:00Z",
    ).rstrip())
    phase = {
        "adw_id": "own-run", "phase_id": "shared-phase", "seq": 1, "name": "build",
        "phase_kind": "agent", "owner": "builder", "status": "running",
        "phase_attempt": 0, "retries": 0, "started_at": "2026-08-15T10:00:00Z",
    }
    with pytest.raises(ValueError, match="phase identifier conflict"):
        mirror.consume(_trace_line(attempt, "phase", **phase).rstrip())
    assert mirror.tracer.conn.execute(
        "SELECT adw_id,status FROM phases WHERE phase_id='shared-phase'",
    ).fetchone() == ("other-run", "running")

    own_phase = {**phase, "phase_id": "own-phase"}
    mirror.consume(_trace_line(attempt, "phase", **own_phase).rstrip())
    with pytest.raises(ValueError, match="session ownership mismatch"):
        mirror.consume(_trace_line(
            attempt, "phase", **{**own_phase, "adw_id": "other-run", "status": "success"},
        ).rstrip())
    with pytest.raises(ValueError, match="agent attempt identifier conflict"):
        mirror.consume(_trace_line(
            attempt, "agent_attempt_start", adw_id="own-run", phase_id="own-phase",
            agent_attempt_id="shared-attempt", agent="builder", session_id="session-own",
            host="pod-own", model="model", started_at="2026-08-15T10:00:01Z",
        ).rstrip())
    assert mirror.tracer.conn.execute(
        "SELECT adw_id,phase_id FROM agent_attempts WHERE attempt_id='shared-attempt'",
    ).fetchone() == ("other-run", "shared-phase")
    with pytest.raises(ValueError, match="event ownership mismatch"):
        mirror.consume(_trace_line(
            attempt, "agent_attempt_start", adw_id="own-run", phase_id="own-phase",
            agent_attempt_id="own-attempt", parent_id="shared-event", agent="builder",
            session_id="session-own", host="pod-own", model="model",
            started_at="2026-08-15T10:00:01Z",
        ).rstrip())
    assert mirror.tracer.conn.execute(
        "SELECT 1 FROM agent_attempts WHERE attempt_id='own-attempt'",
    ).fetchone() is None
    mirror.close()


def test_mirror_persists_sanitized_lifecycle_and_tool_records(tmp_path: Path) -> None:
    attempt = "d" * 24
    mirror = tracer.KubernetesTraceMirror(
        tmp_path / "trace.db", tmp_path / "events.jsonl", attempt,
        {"cluster": "k3s", "job": "job-own", "pod": "pod-own"},
    )
    mirror.consume(_trace_line(
        attempt, "session_start", adw_id="tool-run", engineer="factory", host="pod-own",
        started_at="2026-08-15T10:00:00Z",
    ).rstrip())
    mirror.consume(_trace_line(
        attempt, "phase", adw_id="tool-run", phase_id="tool-phase", seq=1, name="build",
        phase_kind="agent", owner="builder", status="running", phase_attempt=0, retries=0,
        started_at="2026-08-15T10:00:00Z",
    ).rstrip())
    mirror.consume(_trace_line(
        attempt, "event", event_id="evt-agent", adw_id="tool-run", phase_id="tool-phase",
        event_type="agent_start", name="builder", metadata={"model": "safe-model"},
        started_at="2026-08-15T10:00:00Z",
    ).rstrip())
    mirror.consume(_trace_line(
        attempt, "agent_attempt_start", adw_id="tool-run", phase_id="tool-phase",
        agent_attempt_id="agent-attempt", parent_id="evt-agent", agent="builder",
        session_id="session-own", host="pod-own", model="safe-model",
        started_at="2026-08-15T10:00:00Z",
    ).rstrip())
    mirror.consume(_trace_line(
        attempt, "event", event_id="evt-tool-start", adw_id="tool-run", phase_id="tool-phase",
        event_type="tool_call_start", name="read",
        metadata={"tool_call_id": "tool-call", "attempt_id": "agent-attempt", "seq": 1, "tool": "read"},
        started_at="2026-08-15T10:00:01Z",
    ).rstrip())
    mirror.consume(_trace_line(
        attempt, "event", event_id="evt-tool-end", adw_id="tool-run", phase_id="tool-phase",
        event_type="tool_call", name="read",
        metadata={"tool_call_id": "tool-call", "attempt_id": "agent-attempt", "duration_ms": 5, "ok": False},
        started_at="2026-08-15T10:00:01Z", ended_at="2026-08-15T10:00:02Z",
    ).rstrip())
    mirror.consume(_trace_line(
        attempt, "event", event_id="evt-error", adw_id="tool-run", phase_id="tool-phase",
        event_type="error", name="builder", metadata={"error": True},
        started_at="2026-08-15T10:00:02Z",
    ).rstrip())
    rows = mirror.tracer.conn.execute(
        "SELECT type,name,payload_json FROM events WHERE event_id LIKE 'evt-%' ORDER BY started_at,event_id",
    ).fetchall()
    tool = mirror.tracer.conn.execute(
        "SELECT attempt_id,seq,tool_name,args_json,duration_ms,ok,result_excerpt"
        " FROM tool_calls WHERE tool_call_id='tool-call'",
    ).fetchone()
    encoded = json.dumps({"events": rows, "tool": tool})
    assert "SECRET" not in encoded
    assert tool == ("agent-attempt", 1, "read", None, 5, 0, None)
    assert {row[0] for row in rows} == {"agent_start", "tool_call_start", "tool_call", "error"}
    mirror.close()


@pytest.mark.parametrize("result_retained", [False, True])
def test_expired_attempt_reconciliation_deletes_job_and_leased_refs(tmp_path: Path, monkeypatch, result_retained: bool) -> None:
    repo = tmp_path / "repo"
    repo.mkdir()
    init_repo(repo)
    output_root = tmp_path / "attempts"
    attempt = "d" * 24
    directory = output_root / f"factory-{attempt[:16]}"
    directory.mkdir(parents=True)
    receipt = {
        "version": 1,
        "attempt_id": attempt,
        "job_name": f"factory-{attempt[:16]}",
        "upload_secret": f"factory-upload-{attempt[:16]}",
        "mailbox_policy": f"factory-mailbox-{attempt[:16]}",
        "namespace": "overdeck-factory",
        "execution_ref": f"refs/heads/factory-exec/{attempt}",
        "execution_commit": "0" * 40,
        "result_ref": f"refs/heads/factory-result/{attempt}",
        "expires_at": 100,
        "state": "prepared",
    }
    if result_retained:
        receipt["result_retained"] = True
    kubernetes_job.write_receipt(directory / "attempt.json", receipt)
    commands: list[list[str]] = []
    refs: list[str] = []
    monkeypatch.setattr(kubernetes_job, "run", lambda argv, **kwargs: (commands.append(list(argv)) or subprocess.CompletedProcess(argv, 0, "", "")))
    monkeypatch.setattr(kubernetes_job, "delete_remote_ref", lambda _repo, ref: (refs.append(ref) or None))
    kubernetes_job.reconcile_attempts(repo, output_root, None, now=101)
    assert any("delete" in command and "job" in command for command in commands)
    expected_refs = [receipt["execution_ref"]] if result_retained else [receipt["result_ref"], receipt["execution_ref"]]
    assert refs == expected_refs
    assert json.loads((directory / "attempt.json").read_text())["state"] == "cleaned"


def test_reconciliation_recovers_retained_trace_failure_after_cleanup_error(
    tmp_path: Path, monkeypatch,
) -> None:
    output_root = tmp_path / "attempts"
    attempt = "e" * 24
    directory = output_root / f"factory-{attempt[:16]}"
    directory.mkdir(parents=True)
    receipt = {
        "version": 1,
        "attempt_id": attempt,
        "job_name": f"factory-{attempt[:16]}",
        "upload_secret": f"factory-upload-{attempt[:16]}",
        "mailbox_policy": f"factory-mailbox-{attempt[:16]}",
        "namespace": "overdeck-factory",
        "execution_ref": f"refs/heads/factory-exec/{attempt}",
        "execution_commit": "0" * 40,
        "result_ref": f"refs/heads/factory-result/{attempt}",
        "expires_at": 100,
        "state": "prepared",
        "result_retained": True,
        "trace_error": "Kubernetes trace is missing session start or finish",
    }
    receipt_path = directory / "attempt.json"
    kubernetes_job.write_receipt(receipt_path, receipt)
    cleanup_fails = True
    deleted_refs: list[str] = []

    def cleanup(argv, **kwargs):
        return subprocess.CompletedProcess(
            argv, 1 if cleanup_fails else 0, "", "initial cleanup failure" if cleanup_fails else "",
        )

    monkeypatch.setattr(kubernetes_job, "run", cleanup)
    monkeypatch.setattr(
        kubernetes_job, "delete_remote_ref",
        lambda _repo, ref: (deleted_refs.append(ref) or None),
    )
    with pytest.raises(RuntimeError, match="initial cleanup failure"):
        kubernetes_job.reconcile_attempts(tmp_path, output_root, None, now=101)
    assert json.loads(receipt_path.read_text())["state"] == "prepared"
    assert deleted_refs == [receipt["execution_ref"]]

    cleanup_fails = False
    kubernetes_job.reconcile_attempts(tmp_path, output_root, None, now=101)
    recovered = json.loads(receipt_path.read_text())
    assert recovered["state"] == "cleaned"
    assert recovered["trace_error"] == receipt["trace_error"]
    assert deleted_refs == [receipt["execution_ref"], receipt["execution_ref"]]
    assert receipt["result_ref"] not in deleted_refs


@pytest.mark.parametrize(
    "trace_error",
    [None, False, "", "é" * (kubernetes_job.MAX_RECOVERY_TRACE_ERROR_BYTES // 2 + 1)],
)
def test_reconciliation_rejects_invalid_trace_error(
    tmp_path: Path, monkeypatch, trace_error,
) -> None:
    output_root = tmp_path / "attempts"
    attempt = "f" * 24
    directory = output_root / f"factory-{attempt[:16]}"
    directory.mkdir(parents=True)
    kubernetes_job.write_receipt(directory / "attempt.json", {
        "version": 1,
        "attempt_id": attempt,
        "job_name": f"factory-{attempt[:16]}",
        "upload_secret": f"factory-upload-{attempt[:16]}",
        "mailbox_policy": f"factory-mailbox-{attempt[:16]}",
        "namespace": "overdeck-factory",
        "execution_ref": f"refs/heads/factory-exec/{attempt}",
        "execution_commit": "0" * 40,
        "result_ref": f"refs/heads/factory-result/{attempt}",
        "expires_at": 100,
        "state": "prepared",
        "result_retained": True,
        "trace_error": trace_error,
    })
    monkeypatch.setattr(kubernetes_job, "run", lambda *args, **kwargs: pytest.fail("invalid receipt reached cleanup"))
    with pytest.raises(RuntimeError, match="invalid recovery receipt"):
        kubernetes_job.reconcile_attempts(tmp_path, output_root, None, now=101)


def test_signing_private_key_must_be_owner_only(tmp_path: Path, monkeypatch) -> None:
    path = tmp_path / "private.pem"
    path.write_bytes(KEY)
    path.chmod(0o600)
    monkeypatch.setenv("FACTORY_ATTEMPT_SIGNING_PRIVATE_KEY", str(path))
    assert kubernetes_job.signing_key() == KEY
    path.chmod(0o640)
    with pytest.raises(ValueError, match="owner-only"):
        kubernetes_job.signing_key()


def test_workspace_snapshot_includes_source_but_not_ignored_data(tmp_path: Path) -> None:
    init_repo(tmp_path)
    (tmp_path / "tracked").write_text("after\n")
    (tmp_path / "new-source").write_text("new\n")
    (tmp_path / "ignored").write_text("private\n")
    _, _, commit = kubernetes_job.workspace_snapshot(tmp_path)
    names = subprocess.check_output(["git", "-C", str(tmp_path), "ls-tree", "-r", "--name-only", commit], text=True).splitlines()
    assert "tracked" in names and "new-source" in names and "ignored" not in names


def test_uploaded_result_bundle_is_validated_and_controller_publishes(
    tmp_path: Path, monkeypatch,
) -> None:
    repo = tmp_path / "repo"
    remote = tmp_path / "remote.git"
    repo.mkdir()
    init_repo(repo)
    subprocess.run(["git", "init", "--bare", "-q", str(remote)], check=True)
    subprocess.run(["git", "-C", str(repo), "remote", "add", "origin", str(remote)], check=True)
    base = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    worker = tmp_path / "worker"
    subprocess.run(["git", "clone", "--no-local", "-q", str(repo), str(worker)], check=True)
    subprocess.run(["git", "-C", str(worker), "config", "user.email", "factory@example.com"], check=True)
    subprocess.run(["git", "-C", str(worker), "config", "user.name", "Factory"], check=True)
    (worker / "tracked").write_text("result\n")
    subprocess.run(["git", "-C", str(worker), "commit", "-qam", "result"], check=True)
    result = subprocess.check_output(["git", "-C", str(worker), "rev-parse", "HEAD"], text=True).strip()
    assert subprocess.run(["git", "-C", str(repo), "cat-file", "-e", result], check=False).returncode != 0
    bundle = tmp_path / "worker.bundle"
    subprocess.run(["git", "-C", str(worker), "bundle", "create", str(bundle), "HEAD", f"^{base}"], check=True)
    calls: list[tuple[str, ...]] = []
    run_limited_git = kubernetes_job.run_limited_git

    def capture(repo: Path, *args: str, **kwargs):
        calls.append(args)
        return run_limited_git(repo, *args, **kwargs)

    monkeypatch.setattr(kubernetes_job, "run_limited_git", capture)
    imported = kubernetes_job.import_result(repo, bundle, base)
    assert imported == result
    assert any("fetch.fsckObjects=true" in call for call in calls)
    assert any("fsck" in call and "--connectivity-only" in call for call in calls)
    result_ref = "refs/heads/factory-result/" + "b" * 24
    kubernetes_job.publish_result(repo, imported, result_ref)
    assert subprocess.check_output(["git", "--git-dir", str(remote), "rev-parse", result_ref], text=True).strip() == result


def test_result_reader_handoff_copies_exact_bytes_without_log_transport(tmp_path: Path, monkeypatch) -> None:
    expected = b"bundle\x00bytes\xff"
    seen: list[list[str]] = []

    def copy_command(argv, *, stdout, stderr, check, timeout):
        seen.append(argv)
        assert timeout == 30
        stdout.write(expected)
        return subprocess.CompletedProcess(argv, 0, b"", b"")

    monkeypatch.setattr(kubernetes_job.subprocess, "run", copy_command)
    copied = kubernetes_job.copy_result_bundle("overdeck-factory", None, "factory-pod", tmp_path)
    assert copied.read_bytes() == expected
    assert seen == [[
        "kubectl", "--namespace", "overdeck-factory", "exec", "factory-pod", "-c",
        "result-reader", "--", "/bin/cat", kubernetes_job.RESULT_BUNDLE_PATH,
    ]]


@pytest.mark.parametrize("returncode, message", [(1, "could not copy"), (0, "size is not approved")])
def test_result_reader_handoff_rejects_failed_or_empty_copy(tmp_path: Path, monkeypatch, returncode: int, message: str) -> None:
    def copy_command(argv, *, stdout, stderr, check, timeout):
        if returncode == 0:
            return subprocess.CompletedProcess(argv, 0, b"", b"")
        return subprocess.CompletedProcess(argv, returncode, b"", b"reader failed")

    monkeypatch.setattr(kubernetes_job.subprocess, "run", copy_command)
    with pytest.raises(RuntimeError, match=message):
        kubernetes_job.copy_result_bundle("overdeck-factory", None, "factory-pod", tmp_path)
    assert not (tmp_path / "result.bundle").exists()


def test_result_marker_probe_is_bounded_and_uses_reader(monkeypatch) -> None:
    calls: list[tuple[list[str], dict]] = []

    def probe(argv, **kwargs):
        calls.append((argv, kwargs))
        return subprocess.CompletedProcess(argv, 0, "", "")

    monkeypatch.setattr(kubernetes_job, "run", probe)
    assert kubernetes_job.result_ready("overdeck-factory", None, "factory-pod") is True
    assert calls[0][0][-7:] == ["factory-pod", "-c", "result-reader", "--", "/usr/bin/test", "-f", kubernetes_job.RESULT_MARKER_PATH]
    assert calls[0][1] == {"check": False, "timeout": 15}


def test_worker_uploads_bundle_without_kubernetes_result_reader() -> None:
    worker = (Path(__file__).parents[1] / "kubernetes" / "factory-kubernetes-worker").read_text()
    assert "FACTORY_RESULT_BUNDLE" not in worker
    assert "base64" not in worker
    assert "FACTORY_RESULT_BUNDLE_V1" not in worker
    assert 'Path("/tmp/factory-result.bundle")' not in worker
    assert "result_mailbox.upload_result(upload, bundle)" not in worker
    assert 'Path("/result")' in worker and 'result_dir / "complete"' in worker


def test_result_bundle_rejects_oversized_input_before_git(tmp_path: Path, monkeypatch) -> None:
    bundle = tmp_path / "oversized.bundle"
    with bundle.open("wb") as output:
        output.truncate(result_mailbox.MAX_BUNDLE_BYTES + 1)
    monkeypatch.setattr(kubernetes_job, "run", lambda *args, **kwargs: pytest.fail("oversized bundle reached Git"))
    with pytest.raises(RuntimeError, match="size is not approved"):
        kubernetes_job.import_result(tmp_path, bundle, "0" * 40)


def test_result_bundle_rejects_non_descendant_without_importing_objects(tmp_path: Path) -> None:
    repo = tmp_path / "repo"
    attacker = tmp_path / "attacker"
    repo.mkdir()
    attacker.mkdir()
    init_repo(repo)
    init_repo(attacker)
    base = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    subprocess.run(["git", "-C", str(attacker), "checkout", "-q", "--orphan", "unrelated"], check=True)
    subprocess.run(["git", "-C", str(attacker), "rm", "-q", "-rf", "."], check=True)
    (attacker / "other").write_text("unrelated\n")
    subprocess.run(["git", "-C", str(attacker), "add", "other"], check=True)
    subprocess.run(["git", "-C", str(attacker), "commit", "-qm", "unrelated"], check=True)
    unrelated = subprocess.check_output(["git", "-C", str(attacker), "rev-parse", "HEAD"], text=True).strip()
    bundle = tmp_path / "unrelated.bundle"
    subprocess.run(["git", "-C", str(attacker), "bundle", "create", str(bundle), "HEAD"], check=True)
    with pytest.raises(RuntimeError, match="does not descend"):
        kubernetes_job.import_result(repo, bundle, base)
    assert subprocess.run(["git", "-C", str(repo), "cat-file", "-e", unrelated], check=False).returncode != 0


def test_result_bundle_enforces_object_and_expanded_size_limits(tmp_path: Path, monkeypatch) -> None:
    repo = tmp_path / "repo"
    repo.mkdir()
    init_repo(repo)
    base = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    (repo / "tracked").write_text("result\n")
    subprocess.run(["git", "-C", str(repo), "commit", "-qam", "result"], check=True)
    bundle = tmp_path / "result.bundle"
    subprocess.run(["git", "-C", str(repo), "bundle", "create", str(bundle), "HEAD", f"^{base}"], check=True)

    monkeypatch.setattr(kubernetes_job, "MAX_RESULT_OBJECTS", 0)
    with pytest.raises(RuntimeError, match="object count"):
        kubernetes_job.import_result(repo, bundle, base)
    monkeypatch.setattr(kubernetes_job, "MAX_RESULT_OBJECTS", 100_000)
    monkeypatch.setattr(kubernetes_job, "MAX_RESULT_OBJECT_BYTES", 1)
    with pytest.raises(RuntimeError, match="oversized object"):
        kubernetes_job.import_result(repo, bundle, base)


def test_result_patch_is_written_to_a_bounded_file(tmp_path: Path, monkeypatch) -> None:
    repo = tmp_path / "repo"
    repo.mkdir()
    init_repo(repo)
    base = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    (repo / "tracked").write_text("result exceeding one byte\n")
    subprocess.run(["git", "-C", str(repo), "commit", "-qam", "result"], check=True)
    result = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    monkeypatch.setattr(kubernetes_job, "MAX_RESULT_PATCH_BYTES", 1)
    with pytest.raises(RuntimeError, match="git exited"):
        kubernetes_job.write_result_patch(repo, base, result, tmp_path / "result.patch")


def test_delete_remote_ref_is_idempotent(tmp_path: Path) -> None:
    repo = tmp_path / "repo"
    remote = tmp_path / "remote.git"
    repo.mkdir()
    init_repo(repo)
    subprocess.run(["git", "init", "--bare", "-q", str(remote)], check=True)
    subprocess.run(["git", "-C", str(repo), "remote", "add", "origin", str(remote)], check=True)
    ref = "refs/heads/factory-exec/" + "2" * 24
    subprocess.run(["git", "-C", str(repo), "push", "-q", "origin", f"HEAD:{ref}"], check=True)
    assert kubernetes_job.delete_remote_ref(repo, ref) is None
    assert kubernetes_job.delete_remote_ref(repo, ref) is None


@pytest.mark.parametrize("complete_trace", [True, False])
def test_failed_job_retains_validated_result_and_cleans_attempt_resources(
    tmp_path: Path, monkeypatch, live_trace_follower, complete_trace: bool,
) -> None:
    repo = tmp_path / "repo"
    remote = tmp_path / "remote.git"
    repo.mkdir()
    init_repo(repo)
    subprocess.run(["git", "init", "--bare", "-q", str(remote)], check=True)
    subprocess.run(["git", "-C", str(repo), "remote", "add", "origin", str(remote)], check=True)
    base = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    tree = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD^{tree}"], text=True).strip()
    (repo / "tracked").write_text("failed result\n")
    subprocess.run(["git", "-C", str(repo), "commit", "-qam", "failed result"], check=True)
    result = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    bundle = tmp_path / "failed.bundle"
    subprocess.run(["git", "-C", str(repo), "bundle", "create", str(bundle), "HEAD", f"^{base}"], check=True)
    subprocess.run(["git", "-C", str(repo), "reset", "--hard", "-q", base], check=True)
    nonce = "a" * 24
    result_ref = f"refs/heads/factory-result/{nonce}"
    execution_ref = f"refs/heads/factory-exec/{nonce}"
    actual_run = kubernetes_job.run

    def failed_job(argv, **kwargs):
        if argv[0] != "kubectl":
            return actual_run(argv, **kwargs)
        if "get" in argv:
            pod = {"metadata": {"name": "factory-pod"}, "status": {"containerStatuses": [{"name": "factory", "state": {"terminated": {"exitCode": 1}}}]}}
            return subprocess.CompletedProcess(argv, 0, json.dumps({"items": [pod]}), "")
        if "create" in argv and json.loads(kwargs.get("input_data", "{}") or "{}").get("kind") == "Job":
            return subprocess.CompletedProcess(argv, 0, json.dumps({"metadata": {"uid": "12345678-1234-1234-1234-123456789abc"}}), "")
        return subprocess.CompletedProcess(argv, 0, "", "")

    class FakeMailbox:
        def __init__(self, host, attempt_id, output):
            self.target = result_mailbox.UploadTarget(
                f"https://{kubernetes_job.APPROVED_RESULT_HOST}:4443/v1/results/{attempt_id}",
                "t" * 43,
                upload().ca_pem,
            )

        def start(self):
            return self.target

        def wait(self, timeout):
            assert timeout == 0
            return bundle

        def close(self):
            return None

    if not complete_trace:
        class IncompleteFollower:
            def __init__(self, argv, path, mirror, **kwargs):
                self.path = path
                self.mirror = mirror
                self.truncated = False
                path.write_text("worker log\n")
                mirror.consume(_trace_line(
                    mirror.attempt_id, "session_start", adw_id="incomplete-failed-run",
                    engineer="factory", host="factory-pod",
                    started_at="2026-08-15T10:00:00Z",
                ).rstrip())

            def finish(self, timeout=10):
                return False, None

            def abort(self):
                return None

        monkeypatch.setattr(kubernetes_job, "JobLogFollower", IncompleteFollower)

    monkeypatch.setattr(kubernetes_job, "origin", lambda _: "https://github.com/alexcodeplace/overdeck.git")
    monkeypatch.setattr(kubernetes_job, "workspace_snapshot", lambda _: (base, tree, base))
    monkeypatch.setattr(kubernetes_job, "run", failed_job)
    monkeypatch.setattr(kubernetes_job, "run_bounded_output", fake_log_collection)
    monkeypatch.setattr(kubernetes_job.result_mailbox, "ResultMailbox", FakeMailbox)
    monkeypatch.setattr(kubernetes_job, "result_ready", lambda *args: True)
    monkeypatch.setattr(kubernetes_job, "copy_result_bundle", lambda *args: bundle)
    monkeypatch.setattr(kubernetes_job.secrets, "token_hex", lambda _: nonce)
    monkeypatch.setenv("FACTORY_K8S_IMAGE", IMAGE)
    with pytest.raises(RuntimeError, match="validated result retained"):
        kubernetes_job.submit(repo, ["build", "sentinel"], "overdeck", key=KEY)
    assert subprocess.check_output(["git", "--git-dir", str(remote), "rev-parse", result_ref], text=True).strip() == result
    assert subprocess.run(["git", "--git-dir", str(remote), "show-ref", "--verify", execution_ref], check=False).returncode != 0
    output = repo / ".git" / "factory-kubernetes" / f"factory-{nonce[:16]}"
    recovery = json.loads((output / "recovery.json").read_text())
    receipt = json.loads((output / "attempt.json").read_text())
    assert recovery["result_ref"] == result_ref
    assert receipt["state"] == "cleaned" and receipt["result_retained"] is True
    assert (repo / "tracked").read_text() == "before\n"
    if complete_trace:
        assert "trace_error" not in receipt
    else:
        assert "missing session start or finish" in receipt["trace_error"]
        assert receipt["trace_error"] in recovery["error"]
        assert "validated result retained" in recovery["error"]


def test_forged_trace_success_does_not_override_failed_worker_without_upload(
    tmp_path: Path, monkeypatch, live_trace_follower,
) -> None:
    repo = tmp_path / "repo"
    remote = tmp_path / "remote.git"
    repo.mkdir()
    init_repo(repo)
    subprocess.run(["git", "init", "--bare", "-q", str(remote)], check=True)
    subprocess.run(["git", "-C", str(repo), "remote", "add", "origin", str(remote)], check=True)
    base = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    tree = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD^{tree}"], text=True).strip()
    nonce = "b" * 24
    deleted: list[str] = []
    actual_run = kubernetes_job.run

    class FakeMailbox:
        def __init__(self, host, attempt_id, output):
            self.target = result_mailbox.UploadTarget(
                f"https://{kubernetes_job.APPROVED_RESULT_HOST}:4443/v1/results/{attempt_id}",
                "t" * 43,
                upload().ca_pem,
            )

        def start(self):
            return self.target

        def wait(self, timeout):
            assert timeout == 0
            raise TimeoutError("no upload")

        def close(self):
            return None

    def failed_job(argv, **kwargs):
        if argv[0] != "kubectl":
            return actual_run(argv, **kwargs)
        if "get" in argv:
            pod = {"metadata": {"name": "factory-pod"}, "status": {"containerStatuses": [{"name": "factory", "state": {"terminated": {"exitCode": 1}}}]}}
            return subprocess.CompletedProcess(argv, 0, json.dumps({"items": [pod]}), "")
        if "create" in argv and json.loads(kwargs.get("input_data", "{}") or "{}").get("kind") == "Job":
            return subprocess.CompletedProcess(argv, 0, json.dumps({"metadata": {"uid": "12345678-1234-1234-1234-123456789abc"}}), "")
        if "delete" in argv:
            deleted.append(argv[argv.index("job") + 1])
        return subprocess.CompletedProcess(argv, 0, "", "")

    monkeypatch.setattr(kubernetes_job, "origin", lambda _: "https://github.com/alexcodeplace/overdeck.git")
    monkeypatch.setattr(kubernetes_job, "workspace_snapshot", lambda _: (base, tree, base))
    monkeypatch.setattr(kubernetes_job, "run", failed_job)
    monkeypatch.setattr(kubernetes_job, "run_bounded_output", fake_log_collection)
    monkeypatch.setattr(kubernetes_job.result_mailbox, "ResultMailbox", FakeMailbox)
    monkeypatch.setattr(kubernetes_job, "result_ready", lambda *args: True)
    monkeypatch.setattr(kubernetes_job, "copy_result_bundle", lambda *args: (_ for _ in ()).throw(RuntimeError("result marker had no bundle")))
    monkeypatch.setattr(kubernetes_job.secrets, "token_hex", lambda _: nonce)
    monkeypatch.setenv("FACTORY_K8S_IMAGE", IMAGE)
    with pytest.raises(RuntimeError, match="result marker had no bundle"):
        kubernetes_job.submit(repo, ["build", "sentinel"], "overdeck", key=KEY)
    assert f"factory-{nonce[:16]}" in deleted
    assert subprocess.run(["git", "--git-dir", str(remote), "show-ref"], capture_output=True, text=True).stdout == ""


def test_uncertain_job_create_response_still_deletes_named_job(tmp_path: Path, monkeypatch) -> None:
    repo = tmp_path / "repo"
    remote = tmp_path / "remote.git"
    repo.mkdir()
    init_repo(repo)
    subprocess.run(["git", "init", "--bare", "-q", str(remote)], check=True)
    subprocess.run(["git", "-C", str(repo), "remote", "add", "origin", str(remote)], check=True)
    base = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    tree = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD^{tree}"], text=True).strip()
    nonce = "d" * 24
    deleted: list[str] = []
    actual_run = kubernetes_job.run

    class FakeMailbox:
        def __init__(self, host, attempt_id, output):
            self.target = result_mailbox.UploadTarget(
                f"https://{kubernetes_job.APPROVED_RESULT_HOST}:4443/v1/results/{attempt_id}",
                "t" * 43,
                upload().ca_pem,
            )

        def start(self):
            return self.target

        def close(self):
            return None

    def uncertain_create(argv, **kwargs):
        if argv[0] != "kubectl":
            return actual_run(argv, **kwargs)
        if "create" in argv:
            raise subprocess.TimeoutExpired(argv, 30)
        if "delete" in argv:
            deleted.append(argv[argv.index("job") + 1])
            raise subprocess.TimeoutExpired(argv, 30)
        if "get" in argv:
            return subprocess.CompletedProcess(argv, 0, "", "")
        return subprocess.CompletedProcess(argv, 0, "", "")

    monkeypatch.setattr(kubernetes_job, "origin", lambda _: "https://github.com/alexcodeplace/overdeck.git")
    monkeypatch.setattr(kubernetes_job, "workspace_snapshot", lambda _: (base, tree, base))
    monkeypatch.setattr(kubernetes_job, "run", uncertain_create)
    monkeypatch.setattr(kubernetes_job.result_mailbox, "ResultMailbox", FakeMailbox)
    monkeypatch.setattr(kubernetes_job.secrets, "token_hex", lambda _: nonce)
    monkeypatch.setenv("FACTORY_K8S_IMAGE", IMAGE)
    with pytest.raises(subprocess.TimeoutExpired):
        kubernetes_job.submit(repo, ["build", "sentinel"], "overdeck", key=KEY)
    assert deleted == [f"factory-{nonce[:16]}"]
    assert subprocess.run(["git", "--git-dir", str(remote), "show-ref"], capture_output=True, text=True).stdout == ""


def test_successful_job_waits_for_upload_then_validates_publishes_and_applies(tmp_path: Path, monkeypatch, live_trace_follower) -> None:
    repo = tmp_path / "repo"
    remote = tmp_path / "remote.git"
    repo.mkdir()
    init_repo(repo)
    subprocess.run(["git", "init", "--bare", "-q", str(remote)], check=True)
    subprocess.run(["git", "-C", str(repo), "remote", "add", "origin", str(remote)], check=True)
    base = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD"], text=True).strip()
    tree = subprocess.check_output(["git", "-C", str(repo), "rev-parse", "HEAD^{tree}"], text=True).strip()
    (repo / "tracked").write_text("result\n")
    subprocess.run(["git", "-C", str(repo), "commit", "-qam", "result"], check=True)
    bundle = tmp_path / "result.bundle"
    subprocess.run(["git", "-C", str(repo), "bundle", "create", str(bundle), "HEAD", f"^{base}"], check=True)
    subprocess.run(["git", "-C", str(repo), "reset", "--hard", "-q", base], check=True)
    nonce = "c" * 24
    events: list[str] = []
    actual_run = kubernetes_job.run

    class FakeMailbox:
        def __init__(self, host, attempt_id, output):
            assert attempt_id == nonce
            self.target = result_mailbox.UploadTarget(
                f"https://{kubernetes_job.APPROVED_RESULT_HOST}:4443/v1/results/{attempt_id}",
                "t" * 43,
                upload().ca_pem,
            )

        def start(self):
            events.append("mailbox-started")
            return self.target

        def wait(self, timeout):
            events.append("bundle-received")
            return bundle

        def close(self):
            events.append("mailbox-closed")

    def completed_job(argv, **kwargs):
        if argv[0] != "kubectl":
            return actual_run(argv, **kwargs)
        if "get" in argv:
            pod = {
                "metadata": {"name": "factory-pod", "uid": "pod-uid-1"},
                "spec": {"nodeName": "debian2"},
                "status": {"containerStatuses": [{
                    "name": "factory", "state": {"terminated": {"exitCode": 0}},
                }]},
            }
            return subprocess.CompletedProcess(argv, 0, json.dumps({"items": [pod]}), "")
        if "create" in argv:
            document = json.loads(kwargs.get("input_data", "{}") or "{}")
            if document.get("kind") == "Job":
                events.append("job-created-suspended")
                return subprocess.CompletedProcess(argv, 0, json.dumps({"metadata": {"uid": "12345678-1234-1234-1234-123456789abc"}}), "")
            events.append("mailbox-policy-created" if document.get("kind") == "NetworkPolicy" else "upload-secret-created")
        elif "patch" in argv:
            events.append("job-started")
        elif "delete" in argv:
            events.append("job-deleted")
        return subprocess.CompletedProcess(argv, 0, "", "")

    published = kubernetes_job.publish_result

    def record_publish(repo_path, commit, ref):
        events.append("result-published")
        published(repo_path, commit, ref)

    monkeypatch.setattr(kubernetes_job, "origin", lambda _: "https://github.com/alexcodeplace/overdeck.git")
    monkeypatch.setattr(kubernetes_job, "workspace_snapshot", lambda _: (base, tree, base))
    monkeypatch.setattr(kubernetes_job, "run", completed_job)
    monkeypatch.setattr(kubernetes_job, "run_bounded_output", fake_log_collection)
    monkeypatch.setattr(kubernetes_job.result_mailbox, "ResultMailbox", FakeMailbox)
    monkeypatch.setattr(kubernetes_job, "result_ready", lambda *args: True)
    monkeypatch.setattr(kubernetes_job, "copy_result_bundle", lambda *args: bundle)
    monkeypatch.setattr(kubernetes_job, "publish_result", record_publish)
    monkeypatch.setattr(kubernetes_job.secrets, "token_hex", lambda _: nonce)
    monkeypatch.setenv("FACTORY_K8S_IMAGE", IMAGE)
    output = kubernetes_job.submit(repo, ["build", "sentinel"], "overdeck", key=KEY)
    assert output.is_dir() and (repo / "tracked").read_text() == "result\n"
    assert events == ["mailbox-started", "job-created-suspended", "mailbox-policy-created", "upload-secret-created", "job-started", "result-published", "mailbox-closed", "job-deleted"]
    with sqlite3.connect(live_trace_follower) as connection:
        session = connection.execute(
            "SELECT status,host FROM sessions WHERE adw_id=?", (f"adw-{nonce[:8]}",),
        ).fetchone()
        placement = json.loads(connection.execute(
            "SELECT payload_json FROM events WHERE adw_id=? AND type='kubernetes_placement'",
            (f"adw-{nonce[:8]}",),
        ).fetchone()[0])
    assert session == ("success", "worker-pod")
    assert placement == {
        "attempt_id": nonce, "cluster": "k3s", "namespace": "overdeck-factory",
        "job": f"factory-{nonce[:16]}",
        "job_uid": "12345678-1234-1234-1234-123456789abc",
        "pod": "factory-pod", "pod_uid": "pod-uid-1", "node": "debian2",
        "container": "factory",
    }
    assert subprocess.run(["git", "--git-dir", str(remote), "show-ref", "--verify", f"refs/heads/factory-exec/{nonce}"], check=False).returncode != 0
    assert subprocess.run(["git", "--git-dir", str(remote), "show-ref", "--verify", f"refs/heads/factory-result/{nonce}"], check=False).returncode != 0


def test_worker_materializes_approved_pnpm_cache_without_network(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None:
    root = Path(__file__).parents[1]
    worker = runpy.run_path(str(root / "kubernetes" / "factory-kubernetes-worker"))
    workspace = tmp_path / "workspace"
    seed_root = tmp_path / "seeds" / "overdeck"
    seed = seed_root / "pnpm-store"
    verification_seed = seed_root / "pnpm-cache/lockfile-verified.jsonl"
    workspace.mkdir()
    seed.mkdir(parents=True)
    verification_seed.parent.mkdir(parents=True)
    verification_seed.write_text('{"lockfile":"verified"}\n')
    (workspace / "package.json").write_text(json.dumps({"packageManager": "pnpm@11.5.2"}))
    (seed / "cached-package").write_text("immutable bytes\n")
    calls: list[tuple[list[str], Path | None]] = []
    monkeypatch.setenv("PNPM_STORE_DIR", "restore-after-test")

    def capture(args: list[str], *, cwd: Path | None = None, check=True):
        calls.append((args, cwd))
        return subprocess.CompletedProcess(args, 0)

    provision = worker["provision_dependencies"]
    provision.__globals__["WORKSPACE"] = workspace
    provision.__globals__["DEPENDENCY_SEEDS"] = tmp_path / "seeds"
    provision.__globals__["DEPENDENCY_STORE"] = tmp_path / "store"
    verification_cache = tmp_path / "xdg-cache/pnpm/lockfile-verified.jsonl"
    provision.__globals__["DEPENDENCY_VERIFICATION_CACHE"] = verification_cache
    provision.__globals__["call"] = capture
    provision("overdeck")

    store = Path(worker["os"].environ["PNPM_STORE_DIR"])
    assert (store / "cached-package").read_text() == "immutable bytes\n"
    assert verification_cache.read_text() == '{"lockfile":"verified"}\n'
    assert worker["os"].environ["XDG_CACHE_HOME"] == str(tmp_path / "xdg-cache")
    assert calls == [([
        "/usr/local/bin/pnpm", "install", "--frozen-lockfile", "--offline",
        "--fetch-retries", "0", "--store-dir", str(store),
    ], workspace)]


def test_worker_refuses_uncached_pnpm_repository(tmp_path: Path) -> None:
    root = Path(__file__).parents[1]
    worker = runpy.run_path(str(root / "kubernetes" / "factory-kubernetes-worker"))
    workspace = tmp_path / "workspace"
    workspace.mkdir()
    (workspace / "package.json").write_text(json.dumps({"packageManager": "pnpm@11.5.2"}))
    (tmp_path / "seeds/unapproved/pnpm-store").mkdir(parents=True)
    provision = worker["provision_dependencies"]
    provision.__globals__["WORKSPACE"] = workspace
    provision.__globals__["DEPENDENCY_SEEDS"] = tmp_path / "seeds"
    with pytest.raises(RuntimeError, match="approved dependency cache is unavailable"):
        provision("unapproved")


def test_worker_runtime_workflow_and_admission_are_controller_and_shape_only() -> None:
    root = Path(__file__).parents[1]
    subprocess.run(["python3", "-m", "py_compile", str(root / "kubernetes" / "factory-kubernetes-worker")], check=True)
    runtime = (root / "kubernetes" / "runtime.yaml").read_text()
    assert "hostPath" not in runtime and "pod-security.kubernetes.io/audit: restricted" in runtime
    assert "kind: ValidatingAdmissionPolicy" in runtime and "kind: ValidatingAdmissionPolicyBinding" in runtime
    assert "system:serviceaccount:overdeck-factory:overdeck-factory-controller" in runtime
    assert "FACTORY_REPOSITORY_ID" in runtime and "FACTORY_ATTEMPT_ENVELOPE" in runtime
    assert "github.com/alexcodeplace/overdeck.git" not in runtime
    assert "Press.Zone-Works" not in runtime
    assert "params.data" not in runtime
    assert "kind: RoleBinding" in runtime and "overdeck-factory-controller" in runtime
    assert 'resources: ["pods", "pods/log", "pods/exec"]' in runtime and 'verbs: ["create", "get", "list"]' in runtime and 'resources: ["pods/proxy"]' not in runtime
    assert "e.name == 'FACTORY_CREDENTIAL_HANDLE' && e.value == v.secret.secretName" in runtime
    assert "size(object.spec.template.metadata.labels) == 7" in runtime
    assert "object.spec.template.metadata.labels['overdeck.dev/factory-job'] == object.metadata.name" in runtime
    assert "object.spec.template.metadata.labels['batch.kubernetes.io/job-name'] == object.metadata.name" in runtime
    assert "object.spec.template.metadata.labels['job-name'] == object.metadata.name" in runtime
    assert "object.spec.template.metadata.labels['batch.kubernetes.io/controller-uid'] == object.metadata.uid" in runtime
    assert "object.spec.template.metadata.labels['controller-uid'] == object.metadata.uid" in runtime
    assert "quantity(v.emptyDir.sizeLimit).compareTo(quantity('8Gi')) == 0" in runtime
    assert "quantity(v.emptyDir.sizeLimit).compareTo(quantity('2Gi')) == 0" in runtime
    assert "quantity(v.emptyDir.sizeLimit).compareTo(quantity('256Mi')) == 0" in runtime
    assert f"image == '{IMAGE}'" in runtime
    assert "size(object.spec.template.spec.containers) == 2" in runtime
    assert "c.name == 'result-reader'" in runtime
    assert "size(object.spec.template.spec.volumes) == 6" in runtime
    assert "v.name == 'result'" in runtime
    assert "size(object.spec.template.spec.containers[0].volumeMounts) == 9" in runtime
    assert "size(object.spec.template.spec.containers[0].env) == 15" in runtime
    assert "FACTORY_RESULT_UPLOAD_JSON" in runtime and "UV_NO_SYNC" in runtime
    assert 'resources: ["secrets"]' in runtime and 'verbs: ["create"]' in runtime
    assert 'resources: ["networkpolicies"]' in runtime
    assert 'verbs: ["create", "delete", "get", "list", "watch", "patch"]' in runtime
    assert "reserved-upload-secret-or-controller" in runtime
    assert "reserved-mailbox-policy-or-controller" in runtime
    assert "name: live-job" in runtime
    assert 'expression: "!has(object.metadata.deletionTimestamp)"' in runtime
    assert runtime.count("request.userInfo.username == 'system:serviceaccount:overdeck-factory:overdeck-factory-controller'") == 5
    assert "object.metadata.name.startsWith('factory-upload-') ||" in runtime
    assert "object.metadata.name.startsWith('factory-mailbox-') ||" in runtime
    assert "object.spec.egress[0].to[0].ipBlock.cidr == '100.126.128.50/32'" in runtime
    assert "size(object.spec.egress[0].ports) == 1" in runtime
    assert 'verbs: ["create", "delete"]' not in runtime
    assert "kind: NetworkPolicy" in runtime and "http://10.43.200.8:8888" in runtime
    assert "name: overdeck-factory-default-deny-egress" in runtime and "podSelector: {}" in runtime
    assert "clusterIP: 10.43.200.8" in runtime
    worker_policy = runtime.split("name: overdeck-factory-workers", 1)[1].split("---", 1)[0]
    assert "kube-dns" not in worker_policy and "port: 53" not in worker_policy
    assert "100.126.128.50/32" in runtime and "FilterDefaultDeny Yes" in runtime
    assert r"^ssh\.github\.com$" in runtime
    proxy_policy = next(
        document
        for document in yaml.safe_load_all(runtime)
        if isinstance(document, dict)
        and document.get("kind") == "NetworkPolicy"
        and document.get("metadata", {}).get("name") == "overdeck-factory-egress"
    )
    proxy_peers = proxy_policy["spec"]["ingress"][0]["from"]
    assert {
        "namespaceSelector": {
            "matchLabels": {"kubernetes.io/metadata.name": "overdeck"},
        },
        "podSelector": {
            "matchLabels": {
                "app.kubernetes.io/name": "overdeck-cdx",
                "app.kubernetes.io/component": "worker",
            },
        },
    } in proxy_peers
    assert {
        "podSelector": {
            "matchLabels": {
                "app.kubernetes.io/name": "overdeck-factory",
                "app.kubernetes.io/component": "worker",
            },
        },
    } in proxy_peers
    runtime_documents = [
        document for document in yaml.safe_load_all(runtime) if isinstance(document, dict)
    ]
    cdx_proxy = next(
        document for document in runtime_documents
        if document.get("kind") == "DaemonSet"
        and document.get("metadata", {}).get("name") == "overdeck-cdx-egress"
    )
    cdx_pod = cdx_proxy["spec"]["template"]["spec"]
    node_values = cdx_pod["affinity"]["nodeAffinity"]["requiredDuringSchedulingIgnoredDuringExecution"]["nodeSelectorTerms"][0]["matchExpressions"][0]["values"]
    assert node_values == ["debian2", "debian3"]
    assert cdx_pod["automountServiceAccountToken"] is False
    assert cdx_pod["securityContext"]["runAsNonRoot"] is True
    assert cdx_pod["containers"][0]["securityContext"]["capabilities"] == {"drop": ["ALL"]}
    cdx_service = next(
        document for document in runtime_documents
        if document.get("kind") == "Service"
        and document.get("metadata", {}).get("name") == "overdeck-cdx-egress"
    )
    assert cdx_service["spec"]["type"] == "NodePort"
    assert cdx_service["spec"]["externalTrafficPolicy"] == "Local"
    assert cdx_service["spec"]["ports"][0]["nodePort"] == 31888
    cdx_policy = next(
        document for document in runtime_documents
        if document.get("kind") == "NetworkPolicy"
        and document.get("metadata", {}).get("name") == "overdeck-cdx-egress"
    )
    cdx_peer = cdx_policy["spec"]["ingress"][0]["from"][0]
    assert set(cdx_peer) == {"namespaceSelector", "podSelector"}
    assert cdx_peer["namespaceSelector"]["matchLabels"] == {
        "kubernetes.io/metadata.name": "overdeck",
    }
    assert cdx_peer["podSelector"]["matchLabels"] == {
        "app.kubernetes.io/name": "overdeck-cdx",
        "app.kubernetes.io/component": "worker",
    }
    assert {port["port"] for rule in cdx_policy["spec"]["egress"] for port in rule["ports"]} == {53, 443}
    assert "attempt-signing-public-key.pem" in runtime and "attempt-signing-key" not in runtime
    worker = (root / "kubernetes" / "factory-kubernetes-worker").read_text()
    containerfile = (root / "kubernetes" / "Containerfile").read_text()
    builder = (root / "kubernetes" / "build-candidate.sh").read_text()
    verifier = (root / "kubernetes" / "verify-candidate.sh").read_text()
    assert "podman build" in builder and "podman save --format oci-archive" in builder
    assert "verify-candidate.sh" in builder
    assert "--network none" in verifier and "--read-only" in verifier
    assert "pnpm install" in verifier and "--frozen-lockfile" in verifier and "--offline" in verifier
    assert "--fetch-retries 0" in verifier
    assert "/pnpm-cache/lockfile-verified.jsonl" in verifier
    assert "FACTORY_IMAGE_OFFLINE_INSTALL" in verifier
    assert "overdeck/factory-images" in builder and "podman push" not in builder and "login" not in builder
    assert "PYTHONPATH=/opt/factory" in containerfile
    assert "PNPM_VERSION=11.5.2" in containerfile and "BUN_VERSION=1.3.14" in containerfile
    assert "COREPACK_HOME=/opt/corepack" in containerfile
    assert "pnpm install --frozen-lockfile --lockfile-only" in containerfile
    assert "pnpm fetch --frozen-lockfile" in containerfile
    assert "--dir /workspace/repo" in containerfile
    assert "--dir /workspace/repo" in verifier
    assert "/pnpm-cache/lockfile-verified.jsonl" in containerfile
    assert "--mount=type=secret" not in containerfile
    assert "/run/secrets/" not in containerfile
    assert "github_token" not in containerfile
    assert "COPY package.json pnpm-lock.yaml pnpm-workspace.yaml .npmrc" in containerfile
    assert "openssh-client" in containerfile and "openssl" in containerfile
    assert "netcat-openbsd" in containerfile and "tinyproxy" in containerfile
    assert "python3 -c 'from adw_modules import repository_registry'" in containerfile
    assert "--check-entrypoint" not in containerfile
    assert "--check-entrypoint" not in worker
    assert "git@github.com:" in worker and "StrictHostKeyChecking=yes" in worker
    assert '"push"' not in worker and "FACTORY_RESULT_BUNDLE_V1" not in worker
    assert 'Path("/result")' in worker and 'result_dir / "complete"' in worker
    assert "FACTORY_RESULT_UPLOAD_JSON" in worker and "result_mailbox.upload_result(upload, bundle)" not in worker
    assert "envelope.invocation_sha256" in worker and "envelope.input_commit" in worker and "envelope.upload_sha256" in worker
    assert "attempt-signing-public-key.pem" in worker and "attempt-signing-key" not in worker
    assert "ProxyCommand=" in worker and 'os.environ["NO_PROXY"]' in worker
    assert "HostName=ssh.github.com -p 443 -o HostKeyAlias=github.com" in worker
    assert "execution ref does not resolve to the authorized input commit" in worker
    assert "gitea" not in (runtime + worker).lower()
    workflow = Path(__file__).parents[4] / ".github/workflows/factory-kubernetes-image.yml"
    workflow_text = workflow.read_text()
    assert "packages: write" in workflow_text and "context: ." in workflow_text
    assert "driver-opts: network=host" in workflow_text
    assert "network: host" in workflow_text and "secret-files:" not in workflow_text
    for secret_name in (
        "FACTORY_QUERY_REACT_TGZ_B64",
        "FACTORY_UI_PRIMITIVES_TGZ_B64_1",
        "FACTORY_UI_PRIMITIVES_TGZ_B64_2",
        "FACTORY_UI_TOKENS_TGZ_B64",
    ):
        assert f"secrets.{secret_name}" in workflow_text
    assert "registry-gateway/admit-package.mjs" in workflow_text
    assert "REGISTRY_HOST=100.101.104.41" in workflow_text
    assert "github_token=${{ secrets.GITHUB_TOKEN }}" not in workflow_text
    assert "- pnpm-lock.yaml" in workflow_text and "- .npmrc" in workflow_text
    preset = (root / "presets" / "k3s-spark-xhigh.yaml").read_text()
    assert preset.count("openai-codex/gpt-5.3-codex-spark") == 5
    assert preset.count("thinking: xhigh") == 5
    models = json.loads((Path(__file__).parents[3] / "workstation/pi/agent/models.json").read_text())
    spark = next(model for model in models["providers"]["openai-codex"]["models"] if model["id"] == "gpt-5.3-codex-spark")
    assert spark["contextWindow"] == 128000 and spark["input"] == ["text"]
    assert spark["thinkingLevelMap"]["xhigh"] == "xhigh"


def test_factory_image_pin_sites_match_approved_image() -> None:
    """No admission, egress, or literal test fixture pin may drift from the Job pin."""
    root = Path(__file__).parents[1]
    pin_sites: list[tuple[Path, str]] = []
    for path in root.rglob("*"):
        if not path.is_file() or ".venv" in path.parts:
            continue
        for match in _IMMUTABLE_FACTORY_IMAGE.finditer(path.read_text(errors="ignore")):
            pin_sites.append((path.relative_to(root), match.group(0)))

    assert pin_sites, "no immutable Factory image pin sites found"
    assert {image for _, image in pin_sites} == {IMAGE}, pin_sites
    assert {path for path, _ in pin_sites} >= {
        Path("adw_modules/kubernetes_job.py"),
        Path("kubernetes/runtime.yaml"),
    }
