from __future__ import annotations import hashlib import json from dataclasses import replace from datetime import UTC, datetime from types import SimpleNamespace import pytest from agent.orchestration import ( AcceptanceCriterion, ArtifactRef, ArtifactSnapshot, AttemptStatus, DecisionAction, PlannerDecision, TaskAttempt, TaskLedger, TaskRecord, TaskSpec, TaskStatus, ValidationReport, ValidationRunStatus, ValidationVerdict, ) from script_build_host.agents.validation import ( ScriptBuildValidationPolicy, precheck_business_artifact, validate_phase_two_defects, validation_layer, ) from script_build_host.application.phase_two_inputs import ( AcceptedInputResolver, ActiveFrontierResolver, PhaseTwoInputError, ) from script_build_host.domain.artifacts import ( ArtifactKind, ArtifactState, ArtifactVersion, EvidenceRecordV1, ) from script_build_host.domain.phase_two_artifacts import ( CandidateLineageV1, ParagraphArtifactV1, ScriptParagraphV1, ) from script_build_host.domain.task_contracts import ( AcceptedDecisionRef, AcceptedInput, AcceptedInputBundleV1, ScriptCriterion, ScriptIntentClass, ScriptTaskBudget, ScriptTaskContractV1, ScriptTaskKind, ) from script_build_host.infrastructure.canonical_json import canonical_sha256 DIGEST = "sha256:" + "a" * 64 SNAPSHOT_REF = "script-build://inputs/11" SCOPE = "script-build://scopes/opening" WRITE = "script-build://writes/paragraphs/1" def _contract( *, kind: ScriptTaskKind = ScriptTaskKind.PARAGRAPH, scope: str = SCOPE, write_scope: tuple[str, ...] = (WRITE,), supersedes: tuple[str, ...] = (), ) -> ScriptTaskContractV1: schemas = { ScriptTaskKind.PARAGRAPH: "paragraph-artifact/v1", ScriptTaskKind.ELEMENT_SET: "element-set-artifact/v1", ScriptTaskKind.COMPARE: "comparison-artifact/v1", ScriptTaskKind.DECODE_RETRIEVAL: "evidence-record/v1", } return ScriptTaskContractV1( task_kind=kind, scope_ref=scope, intent_class=ScriptIntentClass.EXPLORE, objective="produce one independently verifiable increment", input_decision_refs=(), base_artifact_ref=None, write_scope=write_scope, gap_ref=None, output_schema=schemas[kind], criteria=(ScriptCriterion("closed", "candidate is concretely realized"),), budget=ScriptTaskBudget(), supersedes_decision_ids=supersedes, ) def _paragraph_artifact(*, text: str = "specific opening") -> ParagraphArtifactV1: return ParagraphArtifactV1( lineage=CandidateLineageV1( scope_ref=SCOPE, input_snapshot_ref=SNAPSHOT_REF, input_closure_digest=DIGEST, write_scope=(WRITE,), ), paragraphs=( ScriptParagraphV1( paragraph_id=1, paragraph_index=1, level=1, parent_id=None, name="opening", content_range={}, description=text, ), ), ) def _task(task_id: str, *, parent: str | None, kind: str, status: TaskStatus) -> TaskRecord: return TaskRecord( task_id=task_id, goal_id=None, parent_task_id=parent, display_path=task_id, specs=[ TaskSpec( version=1, objective="bounded increment", acceptance_criteria=(AcceptanceCriterion("closed", "closed"),), context_refs=( f"script-build://task-kinds/{kind}", SNAPSHOT_REF, ), ) ], status=status, ) class _TaskStore: def __init__(self, ledger: TaskLedger) -> None: self.ledger = ledger async def load(self, root_trace_id: str) -> TaskLedger: assert root_trace_id == self.ledger.root_trace_id return self.ledger class _SnapshotStore: def __init__(self, snapshot: ArtifactSnapshot) -> None: self.snapshot = snapshot async def get(self, root_trace_id: str, snapshot_id: str) -> ArtifactSnapshot: assert root_trace_id == "root" assert snapshot_id == self.snapshot.snapshot_id return self.snapshot class _Artifacts: def __init__(self, version: ArtifactVersion) -> None: self.version = version async def read_by_ref(self, ref: ArtifactRef, **owners: object) -> ArtifactVersion: assert ref.digest == self.version.canonical_sha256 assert owners == { "script_build_id": 7, "task_id": "child", "attempt_id": "attempt-child", } return self.version class _Contracts: def __init__(self, values: dict[str, ScriptTaskContractV1]) -> None: self.values = values async def read_for_task( self, *, root_trace_id: str, task: TaskRecord, spec_version: int | None = None ) -> ScriptTaskContractV1: assert root_trace_id == "root" assert spec_version == 1 return self.values[task.task_id] def _closed_fixture() -> tuple[AcceptedInputResolver, TaskLedger, TaskAttempt]: root = _task("root-task", parent=None, kind="root", status=TaskStatus.NEEDS_REPLAN) parent = _task("parent", parent="root-task", kind="compose", status=TaskStatus.RUNNING) child = _task("child", parent="parent", kind="paragraph", status=TaskStatus.COMPLETED) parent.child_task_ids.append("child") root.child_task_ids.append("parent") ref = ArtifactRef( uri="script-build://artifact-versions/1", kind=ArtifactKind.PARAGRAPH.value, version="1", digest=DIGEST, ) child_attempt = TaskAttempt( attempt_id="attempt-child", task_id="child", spec_version=1, worker_trace_id="worker", worker_preset="script_paragraph_worker", execution_mode="new", accepted_child_decision_ids=(), status=AttemptStatus.SUBMITTED, snapshot_id="artifact-snapshot", submission=SimpleNamespace(artifact_refs=[ref], evidence_refs=[]), ) validation = ValidationReport( validation_id="validation-child", task_id="child", attempt_id="attempt-child", spec_version=1, snapshot_id="artifact-snapshot", validator_trace_id="validator", status=ValidationRunStatus.COMPLETED, verdict=ValidationVerdict.PASSED, ) decision = PlannerDecision( decision_id="decision-child", task_id="child", action=DecisionAction.ACCEPT, reason="passed", from_status=TaskStatus.AWAITING_DECISION, to_status=TaskStatus.COMPLETED, attempt_id="attempt-child", validation_id="validation-child", ) child.attempt_ids.append(child_attempt.attempt_id) child.validation_ids.append(validation.validation_id) child.decision_ids.append(decision.decision_id) parent_attempt = TaskAttempt( attempt_id="attempt-parent", task_id="parent", spec_version=1, worker_trace_id="compose-worker", worker_preset="script_compose_worker", execution_mode="new", accepted_child_decision_ids=(decision.decision_id,), ) ledger = TaskLedger( root_trace_id="root", mission="mission", root_task_id="root-task", tasks={item.task_id: item for item in (root, parent, child)}, attempts={ child_attempt.attempt_id: child_attempt, parent_attempt.attempt_id: parent_attempt, }, validations={validation.validation_id: validation}, decisions={decision.decision_id: decision}, ) normalized = { "summary": "bounded", "artifact_refs": [ { "uri": ref.uri, "kind": ref.kind, "version": ref.version, "digest": ref.digest, "summary": ref.summary, "metadata": ref.metadata, } ], "evidence_refs": [], } snapshot = ArtifactSnapshot( snapshot_id="artifact-snapshot", attempt_id=child_attempt.attempt_id, normalized_content=normalized, sha256=_snapshot_sha(normalized), artifact_refs=[ref], evidence_refs=[], ) artifact = replace(_paragraph_artifact(), canonical_sha256=DIGEST) version = ArtifactVersion( artifact_version_id=1, script_build_id=7, task_id="child", attempt_id=child_attempt.attempt_id, spec_version=1, artifact_type=ArtifactKind.PARAGRAPH, canonical_sha256=DIGEST, state=ArtifactState.FROZEN, artifact=artifact, created_at=datetime.now(UTC), frozen_at=datetime.now(UTC), ) resolver = AcceptedInputResolver( task_store=_TaskStore(ledger), artifact_store=_SnapshotStore(snapshot), artifacts=_Artifacts(version), contracts=_Contracts({"child": _contract(), "parent": _contract()}), ) return resolver, ledger, parent_attempt @pytest.mark.asyncio async def test_accepted_input_closes_decision_validation_snapshot_and_artifact() -> None: resolver, _, attempt = _closed_fixture() bundle = await resolver.resolve( root_trace_id="root", script_build_id=7, task_id="parent", attempt_id=attempt.attempt_id, contract=_contract(), input_snapshot_id="11", ) assert [item.decision_id for item in bundle.inputs] == ["decision-child"] assert bundle.inputs[0].source == "direct-child" assert bundle.input_closure_digest.startswith("sha256:") @pytest.mark.asyncio async def test_legacy_attempt_without_frozen_child_set_fails_closed() -> None: resolver, ledger, attempt = _closed_fixture() ledger.attempts[attempt.attempt_id] = replace(attempt, accepted_child_decision_ids=None) with pytest.raises(PhaseTwoInputError, match="CHILD_DECISION_INVALID"): await resolver.resolve( root_trace_id="root", script_build_id=7, task_id="parent", attempt_id=attempt.attempt_id, contract=_contract(), input_snapshot_id="11", ) @pytest.mark.asyncio async def test_mid_retrieval_replacement_freezes_new_evidence_and_old_candidate() -> None: resolver, ledger, parent_attempt = _closed_fixture() old_ref = ledger.attempts["attempt-child"].submission.artifact_refs[0] old_contract = _contract() # The old accepted Paragraph remains a completed sibling. The replacement # consumer gets new Evidence from its own direct Retrieval child and imports # the immutable old Decision explicitly; it never reads a mutable "latest". ledger.tasks["child"].parent_task_id = "root-task" ledger.tasks["root-task"].child_task_ids = ["child", "parent"] replacement_task = _task( "parent", parent="root-task", kind="paragraph", status=TaskStatus.RUNNING ) replacement_task.child_task_ids = ["retrieval"] replacement_task.attempt_ids = [parent_attempt.attempt_id] ledger.tasks["parent"] = replacement_task retrieval = _task( "retrieval", parent="parent", kind="decode-retrieval", status=TaskStatus.COMPLETED ) retrieval_ref = ArtifactRef( "script-build://artifact-versions/2", ArtifactKind.EVIDENCE.value, "2", DIGEST, ) retrieval_attempt = TaskAttempt( attempt_id="attempt-retrieval", task_id="retrieval", spec_version=1, worker_trace_id="retrieval-worker", worker_preset="script_decode_retrieval_worker", execution_mode="new", accepted_child_decision_ids=(), status=AttemptStatus.SUBMITTED, snapshot_id="retrieval-snapshot", submission=SimpleNamespace(artifact_refs=[], evidence_refs=[retrieval_ref]), ) retrieval_validation = ValidationReport( validation_id="validation-retrieval", task_id="retrieval", attempt_id=retrieval_attempt.attempt_id, spec_version=1, snapshot_id="retrieval-snapshot", validator_trace_id="retrieval-validator", status=ValidationRunStatus.COMPLETED, verdict=ValidationVerdict.PASSED, ) retrieval_decision = PlannerDecision( decision_id="decision-retrieval", task_id="retrieval", action=DecisionAction.ACCEPT, reason="new evidence passed", from_status=TaskStatus.AWAITING_DECISION, to_status=TaskStatus.COMPLETED, attempt_id=retrieval_attempt.attempt_id, validation_id=retrieval_validation.validation_id, ) retrieval.attempt_ids = [retrieval_attempt.attempt_id] retrieval.validation_ids = [retrieval_validation.validation_id] retrieval.decision_ids = [retrieval_decision.decision_id] ledger.tasks[retrieval.task_id] = retrieval ledger.attempts[retrieval_attempt.attempt_id] = retrieval_attempt ledger.validations[retrieval_validation.validation_id] = retrieval_validation ledger.decisions[retrieval_decision.decision_id] = retrieval_decision ledger.attempts[parent_attempt.attempt_id] = replace( parent_attempt, accepted_child_decision_ids=(retrieval_decision.decision_id,), ) retrieval_normalized = { "summary": "bounded retrieval evidence", "artifact_refs": [], "evidence_refs": [ { "uri": retrieval_ref.uri, "kind": retrieval_ref.kind, "version": retrieval_ref.version, "digest": retrieval_ref.digest, "summary": retrieval_ref.summary, "metadata": retrieval_ref.metadata, } ], } retrieval_snapshot = ArtifactSnapshot( snapshot_id="retrieval-snapshot", attempt_id=retrieval_attempt.attempt_id, normalized_content=retrieval_normalized, sha256=_snapshot_sha(retrieval_normalized), artifact_refs=[], evidence_refs=[retrieval_ref], ) evidence = EvidenceRecordV1( evidence_id="mid-retrieval", source_type="decode", tool_name="retrieve_decode", query={"query": "opening detail", "top_k": 3}, source_refs=("decode-index://fixture/new-evidence",), raw_artifact_ref=None, summary="One visible detail contradicts the original opening assumption.", supports=(SCOPE,), confidence="high", limitations=(), content_sha256=DIGEST, created_at=datetime.now(UTC), ) evidence_version = ArtifactVersion( artifact_version_id=2, script_build_id=7, task_id="retrieval", attempt_id=retrieval_attempt.attempt_id, spec_version=1, artifact_type=ArtifactKind.EVIDENCE, canonical_sha256=DIGEST, state=ArtifactState.FROZEN, artifact=evidence, created_at=datetime.now(UTC), frozen_at=datetime.now(UTC), ) old_snapshot = resolver._artifact_store.snapshot class SnapshotStore: async def get(self, root_trace_id: str, snapshot_id: str) -> ArtifactSnapshot: assert root_trace_id == "root" return { old_snapshot.snapshot_id: old_snapshot, retrieval_snapshot.snapshot_id: retrieval_snapshot, }[snapshot_id] old_version = resolver._artifacts.version class Artifacts: async def read_by_ref(self, ref: ArtifactRef, **owners: object) -> ArtifactVersion: version = {old_ref.uri: old_version, retrieval_ref.uri: evidence_version}[ref.uri] assert owners == { "script_build_id": 7, "task_id": version.task_id, "attempt_id": version.attempt_id, } return version replacement_contract = replace( old_contract, intent_class=ScriptIntentClass.REPLACE, input_decision_refs=( AcceptedDecisionRef( "decision-child", old_ref, SCOPE, ScriptTaskKind.PARAGRAPH, ), ), supersedes_decision_ids=("decision-child",), ) resolver._artifact_store = SnapshotStore() resolver._artifacts = Artifacts() resolver._contracts.values = { "child": old_contract, "retrieval": replace( _contract(kind=ScriptTaskKind.DECODE_RETRIEVAL), write_scope=(), ), "parent": replacement_contract, } bundle = await resolver.resolve( root_trace_id="root", script_build_id=7, task_id="parent", attempt_id=parent_attempt.attempt_id, contract=replacement_contract, input_snapshot_id="11", ) assert [item.decision_id for item in bundle.inputs] == [ "decision-retrieval", "decision-child", ] assert [item.source for item in bundle.inputs] == ["direct-child", "explicit"] assert bundle.superseded_decision_ids == ("decision-child",) assert bundle.input_closure_digest.startswith("sha256:") @pytest.mark.asyncio async def test_direct_child_is_revalidated_against_consumer_scope_and_kind() -> None: resolver, _, attempt = _closed_fixture() resolver._contracts.values["parent"] = replace( _contract(), scope_ref="script-build://scopes/ending" ) with pytest.raises(PhaseTwoInputError, match="INPUT_SCOPE_MISMATCH"): await resolver.resolve( root_trace_id="root", script_build_id=7, task_id="parent", attempt_id=attempt.attempt_id, contract=resolver._contracts.values["parent"], input_snapshot_id="11", ) resolver, _, attempt = _closed_fixture() consumer = replace( _contract(), task_kind=ScriptTaskKind.CANDIDATE_PORTFOLIO, intent_class=ScriptIntentClass.PORTFOLIO, output_schema="candidate-portfolio/v1", ) resolver._contracts.values["parent"] = consumer with pytest.raises(PhaseTwoInputError, match="INPUT_SCOPE_MISMATCH"): await resolver.resolve( root_trace_id="root", script_build_id=7, task_id="parent", attempt_id=attempt.attempt_id, contract=consumer, input_snapshot_id="11", ) @pytest.mark.asyncio async def test_explicit_input_is_revalidated_against_consumer_scope() -> None: resolver, ledger, attempt = _closed_fixture() child_ref = ledger.attempts["attempt-child"].submission.artifact_refs[0] ledger.tasks["child"].parent_task_id = "root-task" ledger.tasks["parent"].child_task_ids = [] ledger.attempts[attempt.attempt_id] = replace( attempt, accepted_child_decision_ids=(), ) consumer = replace( _contract(), scope_ref="script-build://scopes/ending", input_decision_refs=( AcceptedDecisionRef( "decision-child", child_ref, SCOPE, ScriptTaskKind.PARAGRAPH, ), ), ) resolver._contracts.values["parent"] = consumer with pytest.raises(PhaseTwoInputError, match="INPUT_SCOPE_MISMATCH"): await resolver.resolve( root_trace_id="root", script_build_id=7, task_id="parent", attempt_id=attempt.attempt_id, contract=consumer, input_snapshot_id="11", ) @pytest.mark.asyncio @pytest.mark.parametrize("mutation", ["verdict", "snapshot", "digest", "not-current"]) async def test_accepted_input_rejects_broken_causal_edges(mutation: str) -> None: resolver, ledger, attempt = _closed_fixture() if mutation == "verdict": ledger.validations["validation-child"].verdict = ValidationVerdict.FAILED elif mutation == "snapshot": ledger.validations["validation-child"].snapshot_id = "other" elif mutation == "digest": resolver._artifacts.version = replace( resolver._artifacts.version, canonical_sha256="sha256:" + "b" * 64 ) else: ledger.tasks["child"].decision_ids.append("other-decision") with pytest.raises((PhaseTwoInputError, AssertionError)): await resolver.resolve( root_trace_id="root", script_build_id=7, task_id="parent", attempt_id=attempt.attempt_id, contract=_contract(), input_snapshot_id="11", ) def test_active_frontier_replacement_cycle_scope_and_write_conflicts() -> None: ref1 = ArtifactRef("script-build://artifact-versions/1", "paragraph", "1", DIGEST) ref2 = ArtifactRef("script-build://artifact-versions/2", "paragraph", "2", DIGEST) bundle = AcceptedInputBundleV1( inputs=( AcceptedInput("d1", ref1, ScriptTaskKind.PARAGRAPH, SCOPE, "explicit"), AcceptedInput("d2", ref2, ScriptTaskKind.PARAGRAPH, SCOPE, "explicit"), ), input_closure_digest=DIGEST, ) frontier = ActiveFrontierResolver() with pytest.raises(PhaseTwoInputError, match="WRITE_SCOPE_CONFLICT"): frontier.resolve(bundle, contracts_by_decision={"d1": _contract(), "d2": _contract()}) active = frontier.resolve( bundle, contracts_by_decision={ "d1": _contract(), "d2": _contract(supersedes=("d1",)), }, ) assert [item.decision_id for item in active] == ["d2"] exact_base_patch = replace(_contract(), base_artifact_ref=ref1) active = frontier.resolve( bundle, contracts_by_decision={"d1": _contract(), "d2": exact_base_patch}, ) assert [item.decision_id for item in active] == ["d1", "d2"] stale_base_patch = replace( _contract(), base_artifact_ref=replace(ref1, digest="sha256:" + "b" * 64), ) with pytest.raises(PhaseTwoInputError, match="WRITE_SCOPE_CONFLICT"): frontier.resolve( bundle, contracts_by_decision={"d1": _contract(), "d2": stale_base_patch}, ) with pytest.raises(PhaseTwoInputError, match="SUPERSESSION_CYCLE"): frontier.resolve( bundle, contracts_by_decision={ "d1": _contract(supersedes=("d2",)), "d2": _contract(supersedes=("d1",)), }, ) incompatible = AcceptedInputBundleV1( inputs=( bundle.inputs[0], replace(bundle.inputs[1], scope_ref="script-build://scopes/ending"), ), input_closure_digest=DIGEST, ) with pytest.raises(PhaseTwoInputError, match="SUPERSESSION_SCOPE_MISMATCH"): frontier.resolve( incompatible, contracts_by_decision={ "d1": _contract(), "d2": _contract(scope="script-build://scopes/ending", supersedes=("d1",)), }, ) def test_stale_base_revision_fails_closed() -> None: artifact = replace(_paragraph_artifact(), canonical_sha256=DIGEST) version = ArtifactVersion( 1, 7, "child", "attempt-child", 1, ArtifactKind.PARAGRAPH, DIGEST, ArtifactState.FROZEN, artifact, datetime.now(UTC), datetime.now(UTC), ) lineage = replace( artifact.lineage, base_artifact_ref="script-build://artifact-versions/1", base_artifact_digest=DIGEST, base_revision=1, ) ActiveFrontierResolver.verify_base(lineage=lineage, base=version) with pytest.raises(PhaseTwoInputError, match="STALE_BASE_REVISION"): ActiveFrontierResolver.verify_base( lineage=lineage, base=replace(version, artifact_version_id=2) ) def test_four_validation_layers_defects_and_placeholder_preflight() -> None: assert validation_layer("paragraph") == "local" assert validation_layer("compare") == "compare" assert validation_layer("compose") == "global" assert validation_layer("candidate-portfolio") == "governance" assert ( ScriptBuildValidationPolicy() .plan( SimpleNamespace(task=_task("t", parent=None, kind="compose", status=TaskStatus.RUNNING)) ) .validator_preset == "script_candidate_validator" ) payload = { "defect_code": "REALIZATION_PLACEHOLDER", "criterion_id": "closed", "scope_ref": SCOPE, "observed_excerpt": "TODO", "evidence_refs": ["script-build://artifact-versions/1"], "severity": "hard", "invalidated_inputs": ["decision-child"], "recommended_action_class": "replace", } with pytest.raises(Exception, match="passed validation"): validate_phase_two_defects( [payload], criterion_ids={"closed"}, allowed_evidence_refs={"script-build://artifact-versions/1"}, passed=True, ) defects = validate_phase_two_defects( [payload], criterion_ids={"closed"}, allowed_evidence_refs={"script-build://artifact-versions/1"}, passed=False, ) assert defects[0].blocks_acceptance assert precheck_business_artifact(_paragraph_artifact(text="TODO: later")) == ( "artifact.paragraphs[0].description", ) assert precheck_business_artifact(_paragraph_artifact(text="此处描述其存在")) == ( "artifact.paragraphs[0].description", ) def test_defect_rejects_unknown_criterion_evidence_and_unbounded_excerpt() -> None: base = { "defect_code": "X", "criterion_id": "closed", "scope_ref": SCOPE, "observed_excerpt": "x", "evidence_refs": [], "severity": "warning", "invalidated_inputs": [], "recommended_action_class": "repair", } with pytest.raises(Exception, match="outside the frozen contract"): validate_phase_two_defects( [{**base, "criterion_id": "other"}], criterion_ids={"closed"}, allowed_evidence_refs=set(), passed=False, ) with pytest.raises(Exception, match="500 characters"): validate_phase_two_defects( [{**base, "observed_excerpt": "x" * 501}], criterion_ids={"closed"}, allowed_evidence_refs=set(), passed=False, ) def test_bundle_digest_changes_when_source_or_order_changes() -> None: first = canonical_sha256(["a", "b"]).wire second = canonical_sha256(["b", "a"]).wire assert first != second def _snapshot_sha(value: dict[str, object]) -> str: canonical = json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")) return hashlib.sha256(canonical.encode()).hexdigest()