"""P1 exit-criteria simulation tests (design §7.1 P1, demonstrating §3.3.1). Phase P1 may begin only once the durable human-in-the-loop suspend/resume mechanic is *demonstrated*. §7.1 P1 lists four exit criteria; this module is the executable demonstration of each, driving the committed foundation (:mod:`agent_team.db.schema` compare-and-set helpers + the atomic, integrity-checked :mod:`agent_team.state_store`) through the :mod:`harness.SimPipeline`: * (a) kill the box mid-wait and have the task resume after restart; * (b) submit a duplicate answer and confirm it no-ops; * (c) submit an answer after the deadline expired and confirm it is rejected and the task parked; * (d) two tasks suspended concurrently resume independently to the correct thread. Each criterion has its own test (and a couple of supporting tests for the delivery/recovery edges §3.3.1 calls out). The tests assert on the *durable* state — the ledger row status and the integrity-checked task record — so they verify the real mechanic, not a harness convenience. """ from __future__ import annotations import json import threading import pytest from harness import ( PostFailingTransport, RecordingTransport, SimClock, SimPipeline, ) from agent_team.state_store import IntegrityError from agent_team.task_model import Phase, TaskStatus # --------------------------------------------------------------------------- # Baseline: a single happy-path suspend/resume cycle. # --------------------------------------------------------------------------- def test_submit_suspends_task_with_open_ledger_row( pipeline: SimPipeline, transport: RecordingTransport ) -> None: suspended = pipeline.submit(questions=["which repo?"], deadline_in=100) record = pipeline.load_record(suspended.thread_id) assert record.status is TaskStatus.WAITING_HUMAN assert record.current_phase is Phase.CLARIFY row = pipeline.ledger_row(suspended.question_id) assert row is not None assert row["status"] == "open" assert row["thread_id"] == suspended.thread_id # Delivery happened: a channel_ref was stored and it embeds the question id. assert row["channel_ref"] == f"sim:{suspended.question_id}" assert ( transport.posts and transport.posts[0]["question_id"] == suspended.question_id ) def test_first_answer_wins_and_resumes_to_plan(pipeline: SimPipeline) -> None: suspended = pipeline.submit(questions=["which repo?"], deadline_in=100) won = pipeline.submit_answer( {"question_id": suspended.question_id, "answer": "core-api", "via": "slack:U1"} ) assert won is True outcome = pipeline.resume(suspended.thread_id, suspended.question_id) assert outcome.resumed is True assert outcome.new_phase is Phase.PLAN record = pipeline.load_record(suspended.thread_id) assert record.status is TaskStatus.ACTIVE assert record.current_phase is Phase.PLAN # The won answer was durably folded into the Q&A history. assert record.qa_history == [ {"question_id": suspended.question_id, "turn": 0, "answer": "core-api"} ] assert pipeline.ledger_row(suspended.question_id)["status"] == "answered" # --------------------------------------------------------------------------- # (a) kill the box mid-wait and have the task resume after restart. # --------------------------------------------------------------------------- def test_a_restart_mid_wait_then_answer_and_resume( pipeline: SimPipeline, transport: RecordingTransport ) -> None: suspended = pipeline.submit(questions=["which repo?"], deadline_in=100) # "Kill the box": drop the in-memory pipeline; rebuild purely from disk. reopened = pipeline.reopen() # Durable state survived: ledger row still open, record still WAITING_HUMAN. row = reopened.ledger_row(suspended.question_id) assert row is not None and row["status"] == "open" record = reopened.load_record(suspended.thread_id) assert record.status is TaskStatus.WAITING_HUMAN # The human answers after the restart; the task converges via the sweep. assert reopened.submit_answer( {"question_id": suspended.question_id, "answer": "core-api", "via": "slack:U1"} ) summary = reopened.startup_sweep( questions_by_qid={suspended.question_id: ["which repo?"]} ) assert summary["resumed"] == [suspended.thread_id] resumed_record = reopened.load_record(suspended.thread_id) assert resumed_record.status is TaskStatus.ACTIVE assert resumed_record.current_phase is Phase.PLAN def test_a_restart_after_answer_recovers_via_startup_sweep( pipeline: SimPipeline, ) -> None: """An answer that won *before* the crash must still resume on restart.""" suspended = pipeline.submit(questions=["which repo?"], deadline_in=100) assert pipeline.submit_answer( {"question_id": suspended.question_id, "answer": "core-api", "via": "slack:U1"} ) # Crash before the resume worker ran; recover from disk only. reopened = pipeline.reopen() # Pre-sweep the record is still suspended (resume never ran). assert reopened.load_record(suspended.thread_id).status is TaskStatus.WAITING_HUMAN summary = reopened.startup_sweep( questions_by_qid={suspended.question_id: ["which repo?"]} ) assert summary["resumed"] == [suspended.thread_id] assert reopened.load_record(suspended.thread_id).current_phase is Phase.PLAN # --------------------------------------------------------------------------- # (b) submit a duplicate answer and confirm it no-ops. # --------------------------------------------------------------------------- def test_b_duplicate_answer_no_ops(pipeline: SimPipeline) -> None: suspended = pipeline.submit(questions=["which repo?"], deadline_in=100) raw = { "question_id": suspended.question_id, "answer": "core-api", "via": "slack:U1", } first = pipeline.submit_answer(raw) second = pipeline.submit_answer(raw) # exact redelivery / double click third = pipeline.submit_answer( { "question_id": suspended.question_id, "answer": "other-repo", "via": "github:U2", } ) # a different answer via a second channel assert first is True assert second is False assert third is False # The ledger preserved the *first* answer; later ones never overwrote it. row = pipeline.ledger_row(suspended.question_id) assert row["status"] == "answered" assert json.loads(row["answer_json"]) == "core-api" assert row["answered_via"] == "slack:U1" def test_b_resume_is_single_apply_under_redelivered_resume( pipeline: SimPipeline, ) -> None: """Even if the resume worker is invoked twice, it applies exactly once.""" suspended = pipeline.submit(questions=["which repo?"], deadline_in=100) assert pipeline.submit_answer( {"question_id": suspended.question_id, "answer": "core-api", "via": "slack:U1"} ) first = pipeline.resume(suspended.thread_id, suspended.question_id) second = pipeline.resume(suspended.thread_id, suspended.question_id) assert first.resumed is True assert second.resumed is False assert second.superseded is True # turn guard caught the stale resume # The phase advanced exactly one step; the Q&A history has one entry. record = pipeline.load_record(suspended.thread_id) assert record.current_phase is Phase.PLAN assert len(record.qa_history) == 1 assert pipeline.ledger_row(suspended.question_id)["status"] == "superseded" # --------------------------------------------------------------------------- # (c) answer after the deadline -> rejected, task parked. # --------------------------------------------------------------------------- def test_c_late_answer_rejected_and_task_parked( pipeline: SimPipeline, clock: SimClock ) -> None: suspended = pipeline.submit(questions=["which repo?"], deadline_in=50) # Time passes beyond the deadline; the timer loop expires + parks. clock.advance(51) expired = pipeline.run_deadline_sweep() assert expired == [suspended.question_id] parked = pipeline.load_record(suspended.thread_id) assert parked.status is TaskStatus.PARKED assert parked.current_phase is Phase.PARKED # A late answer loses the compare-and-set against the now-expired row. won = pipeline.submit_answer( {"question_id": suspended.question_id, "answer": "core-api", "via": "slack:U1"} ) assert won is False row = pipeline.ledger_row(suspended.question_id) assert row["status"] == "expired" assert row["answer_json"] is None # A resume attempt on the expired question does nothing. outcome = pipeline.resume(suspended.thread_id, suspended.question_id) assert outcome.resumed is False def test_c_deadline_vs_answer_race_answer_first_wins( pipeline: SimPipeline, clock: SimClock ) -> None: """If the answer lands before the sweep, the sweep must not expire it.""" suspended = pipeline.submit(questions=["which repo?"], deadline_in=50) assert pipeline.submit_answer( {"question_id": suspended.question_id, "answer": "core-api", "via": "slack:U1"} ) clock.advance(99) # well past the deadline expired = pipeline.run_deadline_sweep() # The question is already 'answered', so the sweep finds nothing to expire. assert expired == [] assert pipeline.ledger_row(suspended.question_id)["status"] == "answered" outcome = pipeline.resume(suspended.thread_id, suspended.question_id) assert outcome.resumed is True # --------------------------------------------------------------------------- # (d) two tasks suspended concurrently resume independently to the correct thread. # --------------------------------------------------------------------------- def test_d_two_concurrent_tasks_resume_to_correct_thread( pipeline: SimPipeline, ) -> None: first = pipeline.submit(questions=["repo for A?"], deadline_in=100) second = pipeline.submit(questions=["repo for B?"], deadline_in=100) assert first.thread_id != second.thread_id assert first.question_id != second.question_id # Answer the second task first, with a distinct answer. assert pipeline.submit_answer( {"question_id": second.question_id, "answer": "repo-B", "via": "slack:U2"} ) assert pipeline.submit_answer( {"question_id": first.question_id, "answer": "repo-A", "via": "slack:U1"} ) out_a = pipeline.resume(first.thread_id, first.question_id) out_b = pipeline.resume(second.thread_id, second.question_id) assert out_a.resumed and out_b.resumed rec_a = pipeline.load_record(first.thread_id) rec_b = pipeline.load_record(second.thread_id) # Each thread carries *its own* answer — no cross-contamination. assert rec_a.qa_history[0]["answer"] == "repo-A" assert rec_b.qa_history[0]["answer"] == "repo-B" assert rec_a.current_phase is Phase.PLAN assert rec_b.current_phase is Phase.PLAN def test_d_concurrent_responders_only_one_wins_per_question( pipeline: SimPipeline, ) -> None: """Two threads racing the same question: exactly one compare-and-set wins. Exercises the §3.3.1 ``BEGIN IMMEDIATE`` serialization in the committed ``answer_question`` helper under real OS threads against one SQLite file. """ suspended = pipeline.submit(questions=["which repo?"], deadline_in=100) barrier = threading.Barrier(2) results: list[bool] = [] errors: list[BaseException] = [] lock = threading.Lock() def race(via: str) -> None: barrier.wait() try: won = pipeline.submit_answer( {"question_id": suspended.question_id, "answer": via, "via": via} ) except BaseException as exc: # noqa: BLE001 - record, must be empty with lock: errors.append(exc) return with lock: results.append(won) threads = [threading.Thread(target=race, args=(f"slack:U{i}",)) for i in range(2)] for thread in threads: thread.start() for thread in threads: thread.join() # No lock error is tolerated: the committed BEGIN IMMEDIATE + busy_timeout # must *serialize* the responders, so the single winner is the compare-and # -set, not a swallowed OperationalError loser (the bug the prior version # masked). assert errors == [], f"compare-and-set raised under contention: {errors!r}" assert sum(1 for r in results if r) == 1 # exactly one winner assert results.count(False) == 1 # the other genuinely lost the CAS (rowcount 0) assert pipeline.ledger_row(suspended.question_id)["status"] == "answered" def test_d_concurrent_tasks_survive_restart_independently( pipeline: SimPipeline, ) -> None: """Two suspended tasks + a crash: each converges to its own thread.""" first = pipeline.submit(questions=["repo for A?"], deadline_in=100) second = pipeline.submit(questions=["repo for B?"], deadline_in=100) assert pipeline.submit_answer( {"question_id": first.question_id, "answer": "repo-A", "via": "slack:U1"} ) reopened = pipeline.reopen() summary = reopened.startup_sweep( questions_by_qid={ first.question_id: ["repo for A?"], second.question_id: ["repo for B?"], } ) # Only the answered task resumes; the still-open one stays suspended. assert summary["resumed"] == [first.thread_id] assert reopened.load_record(first.thread_id).current_phase is Phase.PLAN assert reopened.load_record(second.thread_id).status is TaskStatus.WAITING_HUMAN # --------------------------------------------------------------------------- # Supporting §3.3.1 edges: lost-post delivery + durable integrity. # --------------------------------------------------------------------------- def test_lost_post_leaves_open_row_then_reconcile_redelivers( tmp_path_factory: pytest.TempPathFactory, clock: SimClock, post_failing_transport: PostFailingTransport, ) -> None: pipeline = SimPipeline( tmp_path_factory.mktemp("lostpost") / "state", clock=clock, transport=post_failing_transport, ) suspended = pipeline.submit(questions=["which repo?"], deadline_in=100) # The post failed: the row is open with no channel_ref (no in-flight loss). row = pipeline.ledger_row(suspended.question_id) assert row["status"] == "open" assert row["channel_ref"] is None assert post_failing_transport.posts == [] # Reconcile retries idempotently; the second attempt succeeds. redelivered = pipeline.reconcile( questions_by_qid={suspended.question_id: ["which repo?"]} ) assert redelivered == 1 row = pipeline.ledger_row(suspended.question_id) assert row["channel_ref"] == f"sim:{suspended.question_id}" def test_durable_checkpoint_is_integrity_checked( pipeline: SimPipeline, ) -> None: """Corrupting the durable checkpoint must fail closed, not return junk.""" suspended = pipeline.submit(questions=["which repo?"], deadline_in=100) checkpoint = pipeline._checkpoint_path(suspended.thread_id) # noqa: SLF001 # Tamper with the payload after the integrity sidecar was written. checkpoint.write_bytes(checkpoint.read_bytes() + b"tampered") with pytest.raises(IntegrityError): pipeline.load_record(suspended.thread_id)