Add a Confluence documentation lane to the Plane-2 pipeline, flag-gated behind AGENT_TEAM_CONFLUENCE_ENABLED (default off; daemon behavior unchanged when off). - confluence/client.py: OAuth 2LO + Basic REST client, dry-run-default writes - confluence/mermaid.py: vendored ADF-only Mermaid editor (macro-count + revert-diff guards, dry-run default) - nodes/confluence_writer.py(+_llm): conf_draft -> conf_gate -> conf_write, both direct (task_kind=confluence) and post-build documentation flows - task_model/graph/coordinator: new phases, state channels, route_after_intake, CONFLUENCE_APPROVAL_KIND gate delivery, task_kind forwarding - db schema v5: widen pending_questions kind CHECK (atomic rebuild) - tests for client, mermaid, writer node, ledger v5, coordinator gate, e2e
401 lines
16 KiB
Python
401 lines
16 KiB
Python
"""End-to-end Confluence-stage tests through a compiled graph + SQLite checkpointer.
|
|
|
|
These drive the Confluence-writer sub-pipeline through the WHOLE graph wiring
|
|
composed by :func:`agent_team.graph.build_graph` (with ``confluence=True``), over
|
|
a REAL on-disk SQLite checkpointer so the durable suspend/resume can be proven
|
|
across a *fresh graph instance* that shares the same checkpoint DB. Only the
|
|
``billing.claude_invoke`` seam and the Confluence client are faked — no live
|
|
model, no Confluence, no network.
|
|
|
|
Coverage:
|
|
|
|
* **Flow A** — ``start_task(..., task_kind="confluence")`` routes INTAKE straight
|
|
to CONF_DRAFT (NOT the clarifier), suspends at CONF_GATE; an ``approve`` resume
|
|
drives CONF_WRITE -> END. The suspend persists across a brand-new graph object
|
|
sharing the checkpointer DB, and a DUPLICATE resume against the now-settled
|
|
thread is a no-op (the terminal state is unchanged).
|
|
* **Flow B** — the post-build documentation edge: when the build/verify subgraph
|
|
is wired, the verifier's approved (PR) terminus feeds CONF_DRAFT. Exercised by
|
|
graph STRUCTURE (the edge exists) rather than driving the heavy P3 stack.
|
|
* **Regression** — ``build_graph()`` with ``confluence=False`` wires NO conf_*
|
|
nodes (the node set is unchanged), proving the flag gates the whole stage.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
from pathlib import Path
|
|
from typing import Any
|
|
|
|
import pytest
|
|
|
|
from agent_team import billing
|
|
from agent_team.billing import BillingMode, ClaudeResult
|
|
from agent_team.graph import (
|
|
CLARIFY,
|
|
CONF_DRAFT_NODE,
|
|
CONF_GATE,
|
|
CONF_WRITE_NODE,
|
|
INTAKE,
|
|
PLAN,
|
|
build_graph,
|
|
build_sqlite_checkpointer,
|
|
get_pipeline_state,
|
|
pending_question,
|
|
resume_task,
|
|
start_task,
|
|
)
|
|
from agent_team.task_model import Phase, PipelineState, TaskStatus
|
|
|
|
# A well-formed draft the faked Claude clarifier/drafter returns; conf_draft_node
|
|
# parses this JSON into the confluence_draft dict.
|
|
_VALID_DRAFT = {
|
|
"title": "AWS Architecture Map — agent-team",
|
|
"body_storage": "<p>The coordinator daemon runs on the R720.</p>",
|
|
"page_id": "1540098",
|
|
}
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Fixtures / fakes (mirror tests/test_confluence_writer.py + test_graph.py).
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
def _restore_invoker():
|
|
"""Restore the billing module invoker after each test (mirrors siblings)."""
|
|
original = billing._invoker
|
|
yield
|
|
billing._invoker = original
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
def _clear_apply_env(monkeypatch):
|
|
"""Ensure the live-write apply flag never leaks from the host environment."""
|
|
monkeypatch.delenv("AGENT_TEAM_CONFLUENCE_APPLY", raising=False)
|
|
|
|
|
|
def _bind_invoker(reply: str) -> None:
|
|
"""Bind a fake Claude invoker returning ``reply`` (no live model)."""
|
|
|
|
def _fake(prompt: str, *, mode: BillingMode, **kw: Any) -> ClaudeResult:
|
|
return ClaudeResult(text=reply, mode=mode, usage={"input_tokens": 1})
|
|
|
|
billing.set_invoker(_fake)
|
|
|
|
|
|
@pytest.fixture()
|
|
def sqlite_db(tmp_path: Path) -> Path:
|
|
"""A tmp SQLite checkpoint DB path (one per test)."""
|
|
return tmp_path / "state" / "checkpoints.sqlite"
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Flow A — direct Confluence doc task: INTAKE -> CONF_DRAFT (skip clarify).
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_flow_a_routes_intake_to_conf_draft_and_suspends_at_gate(
|
|
sqlite_db: Path,
|
|
) -> None:
|
|
"""task_kind='confluence' skips the clarifier and suspends at the CONF_GATE.
|
|
|
|
A direct documentation task routes INTAKE straight to CONF_DRAFT (NOT
|
|
CLARIFY), drafts the page (faked Claude), and suspends on the Confluence
|
|
human gate's ``interrupt()`` — proven by the pending interrupt's
|
|
``confluence_approval`` kind, NOT a clarifier question-set.
|
|
"""
|
|
_bind_invoker(json.dumps(_VALID_DRAFT))
|
|
cm = build_sqlite_checkpointer(sqlite_db)
|
|
with cm as saver:
|
|
graph = build_graph(checkpointer=saver, confluence=True)
|
|
thread_id, state = start_task(
|
|
graph, transport="slack", task="document the stack", task_kind="confluence"
|
|
)
|
|
|
|
# Suspended at the Confluence gate, not the clarifier.
|
|
assert "__interrupt__" in state
|
|
payload = pending_question(graph, thread_id=thread_id)
|
|
assert payload is not None
|
|
# The gate payload carries the confluence-approval kind + the draft. It
|
|
# ALSO carries a ``question_set`` (the approve/request_changes/abandon
|
|
# decision surface the coordinator delivers — mirroring the plan-decision
|
|
# gate) whose context discriminator marks it as a confluence approval, so
|
|
# it is NOT mistaken for a clarifier question-set.
|
|
assert payload["kind"] == "confluence_approval"
|
|
assert payload["question_set"].context.get("kind") == "confluence_approval"
|
|
assert payload["confluence_draft"]["title"] == _VALID_DRAFT["title"]
|
|
|
|
# The draft was reached WITHOUT visiting the clarifier (no qa_history).
|
|
live = get_pipeline_state(graph, thread_id=thread_id)
|
|
assert live.get("qa_history", []) == []
|
|
assert live["task_kind"] == "confluence"
|
|
|
|
|
|
def test_flow_a_default_kind_still_routes_to_clarify(sqlite_db: Path) -> None:
|
|
"""The default (empty) task_kind still goes to the clarifier, not CONF_DRAFT.
|
|
|
|
Proves route_after_intake only diverts ``confluence`` tasks — every other
|
|
task keeps the original INTAKE -> CLARIFY behaviour even with the flag on.
|
|
"""
|
|
_bind_invoker(json.dumps(_VALID_DRAFT))
|
|
cm = build_sqlite_checkpointer(sqlite_db)
|
|
with cm as saver:
|
|
graph = build_graph(checkpointer=saver, confluence=True)
|
|
thread_id, state = start_task(graph, transport="slack", task="do a thing")
|
|
|
|
assert "__interrupt__" in state
|
|
payload = pending_question(graph, thread_id=thread_id)
|
|
assert payload is not None
|
|
# A clarifier suspend: a question-set, no confluence kind.
|
|
assert "question_set" in payload
|
|
assert payload.get("kind") != "confluence_approval"
|
|
|
|
|
|
def test_flow_a_durable_suspend_resume_across_fresh_graph(sqlite_db: Path) -> None:
|
|
"""The Flow-A suspend persists across a NEW graph object sharing the DB.
|
|
|
|
A first graph instance starts the task and suspends at the gate. A brand-new
|
|
graph instance — same on-disk SQLite checkpointer, fresh compile — resumes the
|
|
SAME thread with an ``approve`` decision, drives CONF_DRAFT -> CONF_GATE ->
|
|
CONF_WRITE -> END, and the task settles DONE (dry-run write, nothing applied).
|
|
This is the durability proof: only the shared checkpoint DB carries the state.
|
|
"""
|
|
_bind_invoker(json.dumps(_VALID_DRAFT))
|
|
|
|
# Instance #1: start + suspend at the gate, then drop the graph entirely.
|
|
cm1 = build_sqlite_checkpointer(sqlite_db)
|
|
with cm1 as saver1:
|
|
graph1 = build_graph(checkpointer=saver1, confluence=True)
|
|
thread_id, state1 = start_task(
|
|
graph1, transport="slack", task="doc it", task_kind="confluence"
|
|
)
|
|
assert "__interrupt__" in state1
|
|
assert pending_question(graph1, thread_id=thread_id) is not None
|
|
|
|
# Instance #2: a FRESH checkpointer connection + freshly compiled graph over
|
|
# the SAME DB file. It must see the suspended thread and resume it.
|
|
cm2 = build_sqlite_checkpointer(sqlite_db)
|
|
with cm2 as saver2:
|
|
graph2 = build_graph(checkpointer=saver2, confluence=True)
|
|
|
|
# The pending gate is visible purely from the shared checkpoint DB.
|
|
resumed_payload = pending_question(graph2, thread_id=thread_id)
|
|
assert resumed_payload is not None
|
|
assert resumed_payload["kind"] == "confluence_approval"
|
|
|
|
final = resume_task(
|
|
graph2, thread_id=thread_id, answer={"decision": "approve", "notes": ""}
|
|
)
|
|
|
|
# CONF_WRITE ran (dry-run by default) and the task is terminal DONE.
|
|
assert final["status"] == TaskStatus.DONE.value
|
|
assert final["confluence_result"]["applied"] is False
|
|
assert final["confluence_result"]["dry_run"] is True
|
|
assert pending_question(graph2, thread_id=thread_id) is None
|
|
|
|
|
|
def test_flow_a_duplicate_resume_is_noop(sqlite_db: Path) -> None:
|
|
"""A DUPLICATE approve resume against the settled thread changes nothing.
|
|
|
|
Once the gate is approved and the task settles DONE, replaying the same
|
|
``Command(resume=...)`` must NOT re-open the gate or mutate the terminal
|
|
state — there is no open interrupt to consume, so it is an idempotent no-op.
|
|
"""
|
|
_bind_invoker(json.dumps(_VALID_DRAFT))
|
|
cm = build_sqlite_checkpointer(sqlite_db)
|
|
with cm as saver:
|
|
graph = build_graph(checkpointer=saver, confluence=True)
|
|
thread_id, _ = start_task(
|
|
graph, transport="slack", task="doc it", task_kind="confluence"
|
|
)
|
|
|
|
first = resume_task(
|
|
graph, thread_id=thread_id, answer={"decision": "approve", "notes": ""}
|
|
)
|
|
assert first["status"] == TaskStatus.DONE.value
|
|
first_result = first["confluence_result"]
|
|
assert pending_question(graph, thread_id=thread_id) is None
|
|
|
|
# Replay the same resume: no open gate -> no state change.
|
|
second = resume_task(
|
|
graph, thread_id=thread_id, answer={"decision": "approve", "notes": ""}
|
|
)
|
|
assert second["status"] == TaskStatus.DONE.value
|
|
assert second["confluence_result"] == first_result
|
|
assert pending_question(graph, thread_id=thread_id) is None
|
|
|
|
|
|
def test_flow_a_request_changes_loops_back_to_draft(sqlite_db: Path) -> None:
|
|
"""A request_changes decision loops the gate back to CONF_DRAFT (redraft).
|
|
|
|
Proves the gate's revise route is wired CONF_GATE -> CONF_DRAFT: the task
|
|
re-drafts (faked Claude again) and re-suspends at the gate, with the human
|
|
feedback folded into state and the visit count bumped.
|
|
"""
|
|
_bind_invoker(json.dumps(_VALID_DRAFT))
|
|
cm = build_sqlite_checkpointer(sqlite_db)
|
|
with cm as saver:
|
|
graph = build_graph(checkpointer=saver, confluence=True)
|
|
thread_id, _ = start_task(
|
|
graph, transport="slack", task="doc it", task_kind="confluence"
|
|
)
|
|
|
|
state = resume_task(
|
|
graph,
|
|
thread_id=thread_id,
|
|
answer={"decision": "request_changes", "notes": "tighten the intro"},
|
|
)
|
|
|
|
# Looped back through CONF_DRAFT and re-suspended at the gate.
|
|
assert "__interrupt__" in state
|
|
payload = pending_question(graph, thread_id=thread_id)
|
|
assert payload is not None
|
|
assert payload["kind"] == "confluence_approval"
|
|
live = get_pipeline_state(graph, thread_id=thread_id)
|
|
assert live["confluence_feedback"] == "tighten the intro"
|
|
assert live["confluence_gate_visits"] >= 1
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Flow B — post-build doc edge: approved-build terminus feeds CONF_DRAFT.
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_flow_b_approved_build_terminus_routes_into_conf_draft(
|
|
restore_review_invoker_b,
|
|
) -> None:
|
|
"""When build_verify is wired with confluence=True, the verifier's approved
|
|
(PR) terminus is repointed from END into CONF_DRAFT (Flow B).
|
|
|
|
Exercised via graph STRUCTURE rather than driving the heavy P3 stack: we
|
|
assert the VERIFY_NODE -> CONF_DRAFT edge exists (and that without
|
|
confluence the same wiring lands VERIFY at END instead), which is the
|
|
load-bearing Flow-B claim.
|
|
"""
|
|
from agent_team.graph import VERIFY_NODE
|
|
from agent_team.nodes import review_loop
|
|
from agent_team.nodes.build_verify_subgraph import (
|
|
make_build_node,
|
|
make_verify_node,
|
|
route_after_verify,
|
|
)
|
|
from agent_team.nodes.verifier import VerifierConfig
|
|
|
|
review_loop.set_review_invoker(lambda prompt, **kw: "VERDICT: APPROVE\nok")
|
|
build_verify = (
|
|
make_build_node(diff_builder=lambda *, plan, config: "diff"),
|
|
make_verify_node(VerifierConfig(expected_run_id="r1")),
|
|
route_after_verify,
|
|
)
|
|
|
|
def _plan_to_review(state: PipelineState) -> PipelineState:
|
|
return PipelineState(
|
|
plan={"summary": "p"},
|
|
current_phase=Phase.REVIEW.value,
|
|
status=TaskStatus.ACTIVE.value,
|
|
)
|
|
|
|
graph = build_graph(
|
|
live_plan_node=_plan_to_review,
|
|
review_node=review_loop.bind_review_node(),
|
|
route_review=review_loop.route_after_review,
|
|
build_verify=build_verify,
|
|
confluence=True,
|
|
)
|
|
|
|
g = graph.get_graph()
|
|
nodes = set(g.nodes)
|
|
assert {VERIFY_NODE, CONF_DRAFT_NODE, CONF_GATE, CONF_WRITE_NODE} <= nodes
|
|
|
|
edges = {(e.source, e.target) for e in g.edges}
|
|
# Flow B: the verifier's approved terminus feeds the Confluence draft so a
|
|
# shipped change documents itself before the graph ends.
|
|
assert (VERIFY_NODE, CONF_DRAFT_NODE) in edges
|
|
|
|
|
|
def test_flow_b_without_confluence_verify_terminus_is_not_conf_draft(
|
|
restore_review_invoker_b,
|
|
) -> None:
|
|
"""With confluence=False, the same P3 wiring lands VERIFY at END, not the
|
|
Confluence draft — the Flow-B repoint is gated on the flag."""
|
|
from agent_team.graph import VERIFY_NODE
|
|
from agent_team.nodes import review_loop
|
|
from agent_team.nodes.build_verify_subgraph import (
|
|
make_build_node,
|
|
make_verify_node,
|
|
route_after_verify,
|
|
)
|
|
from agent_team.nodes.verifier import VerifierConfig
|
|
|
|
review_loop.set_review_invoker(lambda prompt, **kw: "VERDICT: APPROVE\nok")
|
|
build_verify = (
|
|
make_build_node(diff_builder=lambda *, plan, config: "diff"),
|
|
make_verify_node(VerifierConfig(expected_run_id="r1")),
|
|
route_after_verify,
|
|
)
|
|
|
|
def _plan_to_review(state: PipelineState) -> PipelineState:
|
|
return PipelineState(
|
|
plan={"summary": "p"},
|
|
current_phase=Phase.REVIEW.value,
|
|
status=TaskStatus.ACTIVE.value,
|
|
)
|
|
|
|
graph = build_graph(
|
|
live_plan_node=_plan_to_review,
|
|
review_node=review_loop.bind_review_node(),
|
|
route_review=review_loop.route_after_review,
|
|
build_verify=build_verify,
|
|
confluence=False,
|
|
)
|
|
|
|
g = graph.get_graph()
|
|
nodes = set(g.nodes)
|
|
# No Confluence nodes exist, so the approved terminus cannot feed CONF_DRAFT.
|
|
assert CONF_DRAFT_NODE not in nodes
|
|
edges = {(e.source, e.target) for e in g.edges}
|
|
assert (VERIFY_NODE, CONF_DRAFT_NODE) not in edges
|
|
|
|
|
|
@pytest.fixture()
|
|
def restore_review_invoker_b():
|
|
"""Save/restore the review-loop module-global invoker around Flow-B tests."""
|
|
from agent_team.nodes import review_loop
|
|
|
|
saved = review_loop._review_invoker
|
|
yield
|
|
review_loop._review_invoker = saved
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Regression — confluence=False wires NO conf_* nodes (flag-gating).
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_confluence_disabled_wires_no_conf_nodes() -> None:
|
|
"""build_graph() with confluence off (the default) has NO conf_* vertices.
|
|
|
|
The node set is byte-for-byte the P1 set (intake/clarify/plan); none of the
|
|
Confluence stage exists, proving the entire sub-pipeline is flag-gated.
|
|
"""
|
|
graph = build_graph()
|
|
nodes = set(graph.get_graph().nodes)
|
|
|
|
# The plain P1 nodes are present.
|
|
assert {INTAKE, CLARIFY, PLAN} <= nodes
|
|
# NONE of the Confluence vertices are wired.
|
|
assert CONF_DRAFT_NODE not in nodes
|
|
assert CONF_GATE not in nodes
|
|
assert CONF_WRITE_NODE not in nodes
|
|
assert not ({CONF_DRAFT_NODE, CONF_GATE, CONF_WRITE_NODE} & nodes)
|
|
|
|
|
|
def test_confluence_enabled_adds_exactly_the_conf_nodes() -> None:
|
|
"""Turning the flag on adds precisely the three Confluence vertices to the
|
|
otherwise-unchanged P1 node set (no other topology drift)."""
|
|
base_nodes = set(build_graph().get_graph().nodes)
|
|
conf_nodes = set(build_graph(confluence=True).get_graph().nodes)
|
|
|
|
added = conf_nodes - base_nodes
|
|
assert added == {CONF_DRAFT_NODE, CONF_GATE, CONF_WRITE_NODE}
|