This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
orchestrator/agent-team/tests/sim/test_p1_exit_criteria.py
Adam Moussa 0eae5dbfc3 Prove P1 exit criteria against the real LangGraph graph; fix question_id stability
Reworks the P1 sim so the four §7.1 exit criteria are demonstrated against the
ACTUAL mechanic, not a model (resolves the verifier's "sim models the ledger,
not the LangGraph integration" finding).

- New tests/sim/test_p1_graph_integration.py drives the real agent_team.graph
  StateGraph (interrupt/Command(resume)) + the real langgraph SqliteSaver
  checkpointer + the committed pending_questions compare-and-set, proving:
  (a) suspend survives a simulated restart (drop saver/conn, rebuild over the
  same checkpoint DB) and resumes; (b) duplicate answer loses the CAS and the
  graph never double-advances; (c) a post-deadline answer loses to expire and
  the task is not resumed; (d) two concurrent tasks resume to the correct
  thread, with a turn-guarded no-double-apply check.
- graph.py: derive a STABLE question_id from uuid5(thread_id, turn). The
  clarifier node replays on resume, so the prior fresh-uuid id changed between
  the delivered/ledgered question and the qa_history entry — breaking the
  §3.3.1 identity contract. Now the delivered id == ledger key == history entry
  (unit-tested in test_graph.py).
- harness._connect() now uses the committed schema.connect() (WAL + busy_timeout)
  instead of a raw sqlite3.connect, so concurrent responders genuinely serialize;
  the criterion-(d) concurrency test no longer swallows OperationalError (it
  asserts zero errors + exactly one CAS winner).
- requirements.txt: pin langgraph-checkpoint-sqlite==3.1.0 (design D9 durable
  checkpointer), now exercised by the integration test.

Full suite: 564 passed; ruff + format clean.
2026-06-17 15:16:12 -04:00

397 lines
15 KiB
Python

"""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)