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.
397 lines
15 KiB
Python
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)
|