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/agent_team/nodes/confluence_writer.py
Adam Moussa c7f9c1bac2 feat(agent-team): Confluence-writer node (draft -> approve gate -> write)
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
2026-06-25 10:56:48 -04:00

691 lines
29 KiB
Python

"""Confluence-writer pipeline stage — draft -> human gate -> (dry-run) write.
Three LangGraph nodes plus a router that document a task's work on Confluence,
gated by a human exactly like the plan gate (design §3.3 human-in-the-loop):
... -> CONF_DRAFT -> CONF_GATE -> {CONF_WRITE | back to CONF_DRAFT | terminal}
* :func:`conf_draft_node` — an **AGENTIC** Claude call that inspects the repo
(read-only) and emits a Confluence page update as DATA (a ``confluence_draft``
dict). It writes nothing to the repo. Source: Flow A uses ``state['task']``;
Flow B folds in ``state['plan']`` + ``state['candidate_diff']`` + the repo
name. On a ``request_changes`` loop-back the prior ``confluence_feedback`` is
folded into the prompt.
* :func:`conf_gate_node` — the resumable HUMAN GATE. Mirrors
:func:`agent_team.graph.plan_gate_node` VERBATIM in shape: a ceiling guard
FIRST (cap on ``confluence_gate_visits`` -> terminal PARKED), then an
``interrupt()`` carrying a stable ``question_id`` derived from a ``turn`` in a
HIGH namespace disjoint from BOTH the clarifier and the plan gate, a deadline,
and a ``kind`` discriminator. On resume it parses the decision verb and routes
approve / request_changes / abandon.
* :func:`conf_write_node` — performs the write through the injected Confluence
client. **DRY-RUN by default**: a live write is issued ONLY when an explicit
apply flag is set (``config['confluence_apply']`` or env
``AGENT_TEAM_CONFLUENCE_APPLY`` truthy) AND the gate approved; otherwise it
composes the planned change + a revert diff and writes NOTHING.
* :func:`route_after_conf_gate` — the conditional-edge function mirroring
:func:`agent_team.graph.route_after_plan_gate`.
The HTTP/Confluence I/O is an INJECTED seam: :func:`conf_write_node` takes an
optional ``client`` (defaulting to a real :class:`ConfluenceClient`) so tests
pass an in-memory fake and no live network is touched on the paths tests hit.
The Confluence package import is DEFERRED to call time (and only on the live
node) so this module imports cleanly before that package exists. Secrets/creds
are read from ``os.environ`` at call time, never at module load.
Graph-wiring constants this module references by NAME (the graph-agent owns
``graph.py`` and must define / wire these):
* ``CONFLUENCE_APPROVAL_KIND`` — the interrupt-payload ``kind`` discriminator
(defined here with the literal value ``"confluence_approval"``; the graph may
re-export it). The graph must keep its value in sync.
* The CONF_DRAFT / CONF_GATE / CONF_WRITE phase values — this module emits the
:class:`~agent_team.task_model.Phase` ``.value`` strings the graph routes on.
See :data:`CONF_DRAFT_PHASE` / :data:`CONF_GATE_PHASE` / :data:`CONF_WRITE_PHASE`
/ :data:`CONF_DONE_PHASE` and the note in the returned summary.
"""
from __future__ import annotations
import os
import uuid
from datetime import datetime, timedelta, timezone
from typing import TYPE_CHECKING, Any
from langgraph.types import interrupt
from agent_team.billing import ClaudeResult, claude_invoke
from agent_team.transport import QuestionSet
from agent_team.nodes.confluence_writer_llm import (
build_confluence_prompt,
parse_confluence_reply,
)
from agent_team.task_model import Phase, PipelineState, TaskStatus
if TYPE_CHECKING: # pragma: no cover - typing only; package may not exist yet
from agent_team.confluence.client import ConfluenceClient
__all__ = [
"CONFLUENCE_APPROVAL_KIND",
"MAX_CONFLUENCE_GATE_VISITS",
"ConfluenceWriteError",
"conf_draft_node",
"conf_gate_node",
"conf_write_node",
"route_after_conf_gate",
]
# -- Graph-shared constants (see module docstring). --------------------------
# Interrupt-payload discriminator telling the responder/ledger this is a
# Confluence-approval gate (vs the clarifier question-set / plan-decision gate).
# The graph-agent defines the authoritative constant; this local fallback keeps
# the node self-contained and importable before graph.py wires it. KEEP THE
# VALUE IN SYNC with graph.PLAN_DECISION_KIND's sibling.
try: # pragma: no cover - exercised only once the graph defines it
from agent_team.graph import CONFLUENCE_APPROVAL_KIND # type: ignore
except Exception: # noqa: BLE001 - graph may not define it yet on this branch
CONFLUENCE_APPROVAL_KIND = "confluence_approval"
# Combined ceiling on Confluence-gate visits (bounded termination, mirrors
# graph.MAX_PLAN_GATE_VISITS). A human ``request_changes`` re-enters the draft
# stage, which gates AGAIN; once the visit count reaches this cap the gate stops
# offering request_changes and the task goes terminal PARKED so the human loop
# always terminates.
MAX_CONFLUENCE_GATE_VISITS = 3
# Per-call turn headroom + budget for the AGENTIC draft call. The subscription
# invoker defaults to single-shot (max_turns=1), which an agentic repo-inspecting
# call exhausts before it can finish ("Reached maximum number of turns (1)"; PR
# #61 / issue #60). This call READS the repo to write accurate docs, so it needs
# real turns and a READ-ONLY toolset — it emits the draft as DATA and writes
# nothing, so no Edit/Write/Bash.
_DRAFT_MAX_TURNS = 8
_DRAFT_ALLOWED_TOOLS: tuple[str, ...] = ("Read", "Grep", "Glob")
_DRAFT_BUDGET_USD = 4.0
# Phase .value strings this stage routes on. The graph wires vertices that read
# these; kept as Phase members so a future dedicated Phase enum value is a
# one-line swap. Until the foundation adds CONF_* members the stage reuses the
# closest existing phases (VERIFY/BUILD/DONE/PARKED) for routing — the graph-agent
# repoints these if dedicated phases are introduced.
CONF_DRAFT_PHASE = Phase.VERIFY.value
CONF_GATE_PHASE = Phase.REVIEW.value
CONF_WRITE_PHASE = Phase.BUILD.value
CONF_DONE_PHASE = Phase.DONE.value
# Disjoint namespace base for the Confluence gate's stable ``turn`` derivation.
# The clarifier numbers turns 0..N (qa_history length); the plan gate uses
# 1_000_000 + visits (graph._PLAN_GATE_TURN_BASE). This base is HIGHER and
# distinct from BOTH so a Confluence-gate turn can never collide with either for
# the same thread (the ResumeWorker turn guard matches a resume to its open
# interrupt by turn).
_CONF_GATE_TURN_BASE = 2_000_000
# Fixed namespace for deriving a STABLE question_id from (thread_id, turn) —
# uuid5 so the id is uuid-shaped yet deterministic across the resume replay of
# the node (mirrors graph._QUESTION_ID_NAMESPACE; a distinct namespace so a
# Confluence-gate id never collides with a clarifier/plan-gate id).
_CONF_QUESTION_ID_NAMESPACE = uuid.UUID("c0f1e2d3-4a5b-6c7d-8e9f-0a1b2c3d4e5f")
# Default open window for a Confluence-gate question (mirrors
# graph.DEFAULT_CLARIFY_DEADLINE).
_DEFAULT_CONF_DEADLINE = timedelta(hours=24)
class ConfluenceWriteError(Exception):
"""Raised when the live Confluence write cannot be performed.
Distinct from a draft-parse failure (:class:`ConfluenceDraftError`) and from
a gate rejection: this is an I/O-time failure on the apply path, so the
coordinator can fail/park the task rather than mark it documented.
"""
# -- small helpers (mirror graph.py). ----------------------------------------
def _utc_now_iso() -> str:
"""Return the current UTC time as an ISO-8601 string (ledger-compatible)."""
return datetime.now(timezone.utc).isoformat()
def _question_id_for(thread_id: str, turn: int) -> str:
"""Return the stable Confluence-gate question_id for ``(thread_id, turn)``."""
return uuid.uuid5(_CONF_QUESTION_ID_NAMESPACE, f"{thread_id}:{turn}").hex
def _conf_gate_turn(visits: int) -> int:
"""Return the stable gate ``turn`` for the ``visits``-th Confluence-gate visit.
Offset into a high, disjoint namespace (:data:`_CONF_GATE_TURN_BASE`) so a
Confluence-gate turn can never collide with a clarifier turn (qa_history
length) or a plan-gate turn (1_000_000 + visits) for the same thread.
Monotonic in ``visits`` so each successive suspend has its own stable
``(thread_id, turn)`` identity.
"""
return _CONF_GATE_TURN_BASE + visits
def _truthy(value: Any) -> bool:
"""Interpret a config/env flag as a boolean (env strings are case-folded)."""
if isinstance(value, bool):
return value
if isinstance(value, str):
return value.strip().lower() in {"1", "true", "yes", "on"}
return bool(value)
def _draft_preview(draft: dict[str, Any]) -> str:
"""Render a short human-readable preview of the draft for the gate payload."""
title = str(draft.get("title", "(untitled)"))
page_id = draft.get("page_id")
target = f"update page {page_id}" if page_id else "create new page"
body = str(draft.get("body_storage", ""))
snippet = body[:400] + ("..." if len(body) > 400 else "")
mermaid = draft.get("mermaid_edits") or []
lines = [f"Title: {title}", f"Action: {target}"]
if mermaid:
lines.append(f"Mermaid edits: {len(mermaid)}")
lines.append("")
lines.append(snippet)
return "\n".join(lines)
# Decision options the human picks from at the Confluence approval gate. Mirrors
# the plan gate's approve / request_changes / abandon surface so a downstream
# coordinator handler can render the three decision buttons.
_CONF_DECISION_OPTIONS: tuple[str, ...] = ("approve", "request_changes", "abandon")
def _conf_approval_question_set(
*, thread_id: str, question_id: str, turn: int, preview: str
) -> QuestionSet:
"""Build the QuestionSet carried in the CONF_GATE interrupt payload.
The coordinator's notify path requires a ``question_set`` so the inbound
answer maps back to ``question_id`` (mirrors
:meth:`coordinator._plan_decision_question_set`). The single question's
prompt is the human-readable draft ``preview`` followed by the
approve / request_changes / abandon options; ``context`` carries the
``kind`` discriminator (:data:`CONFLUENCE_APPROVAL_KIND`) so the transport
renders decision buttons instead of generic question blocks.
"""
options = " / ".join(_CONF_DECISION_OPTIONS)
prompt = (
f"{preview}\n\n"
f"Approve, request changes, or abandon this Confluence update? ({options})"
)
return QuestionSet(
thread_id=thread_id,
question_id=question_id,
turn=turn,
questions=[prompt],
context={
"kind": CONFLUENCE_APPROVAL_KIND,
"options": list(_CONF_DECISION_OPTIONS),
"presentation": preview,
},
)
# -- Node 1: CONF_DRAFT (agentic). -------------------------------------------
def conf_draft_node(
state: PipelineState, config: dict[str, Any] | None = None
) -> PipelineState:
"""CONF_DRAFT stage: agentic repo inspection -> Confluence draft (as DATA).
Builds the draft prompt (Flow A from ``state['task']``; Flow B folding in
``plan`` + ``candidate_diff`` + repo name; ``confluence_feedback`` folded in
on a ``request_changes`` loop-back), then makes ONE **agentic** Claude call
that may READ the repository to write accurate docs. The call is pinned to a
READ-ONLY toolset (``Read``/``Grep``/``Glob``), ``max_turns=8``, and a
``budget_usd`` cap — the agentic-config rationale builders use (PR #61 /
issue #60: single-shot defaults die with "Reached maximum number of turns
(1)"). The node emits the draft as DATA and writes NOTHING to the repo, so no
Edit/Write/Bash tool is granted.
Returns a **partial** :class:`PipelineState`: the parsed ``confluence_draft``
plus the phase advanced to the Confluence gate and ``status`` ACTIVE.
``config`` is forwarded to the billing seam so the caller can pin the billing
mode; it is threaded through to :func:`claude_invoke` as ``config``.
"""
prompt = build_confluence_prompt(state)
result: ClaudeResult = claude_invoke(
prompt,
config=config,
max_turns=_DRAFT_MAX_TURNS,
allowed_tools=list(_DRAFT_ALLOWED_TOOLS),
budget_usd=_DRAFT_BUDGET_USD,
)
draft = parse_confluence_reply(result.text)
return PipelineState(
confluence_draft=draft, # type: ignore[typeddict-unknown-key]
current_phase=CONF_GATE_PHASE,
status=TaskStatus.ACTIVE.value,
updated_at=_utc_now_iso(),
)
# -- Node 2: CONF_GATE (resumable human gate; mirrors plan_gate_node). --------
def conf_gate_node(
state: PipelineState, config: dict[str, Any] | None = None
) -> PipelineState:
"""CONF_GATE stage: the resumable human approval of the Confluence draft.
Mirrors :func:`agent_team.graph.plan_gate_node` in shape:
* **Ceiling guard FIRST.** If ``confluence_gate_visits`` already reached
:data:`MAX_CONFLUENCE_GATE_VISITS` the gate does NOT interrupt: it returns
terminal PARKED with a ``failure_reason``, so every suspend strictly
consumes one of a finite number of visits and the human loop terminates.
* **Suspend.** Otherwise it bumps the visit count and ``interrupt()``s with a
payload carrying ``thread_id``, a stable ``question_id`` (from a ``turn``
in a HIGH namespace disjoint from the clarifier AND the plan gate), the
``turn``, ``kind`` = :data:`CONFLUENCE_APPROVAL_KIND`, ``transport``,
``deadline``, ``slack_thread_ts``, the ``confluence_draft``, and a
human-readable ``preview``.
On resume, ``interrupt()`` returns the decision (the value passed to
``Command(resume=...)``). The verb is parsed via the shared safe normalizer
(unrecognized -> ``request_changes``, like the plan gate's
``_parse_decision``):
* ``approve`` -> phase CONF_WRITE, status ACTIVE;
* ``request_changes`` -> phase CONF_DRAFT, status ACTIVE, ``confluence_feedback``
set to the notes, ``confluence_gate_visits`` bumped (loop back to redraft);
* ``abandon`` -> terminal FAILED with a ``failure_reason``.
"""
prior_visits = int(state.get("confluence_gate_visits", 0) or 0) # type: ignore[call-overload]
if prior_visits >= MAX_CONFLUENCE_GATE_VISITS:
return PipelineState(
status=TaskStatus.PARKED.value,
current_phase=Phase.PARKED.value,
failure_reason=(
"confluence-gate revision ceiling reached "
f"({prior_visits}/{MAX_CONFLUENCE_GATE_VISITS} gate visits)"
),
updated_at=_utc_now_iso(),
)
visits = prior_visits + 1
thread_id = state.get("thread_id", "")
transport = state.get("transport", "")
slack_thread_ts = state.get("slack_thread_ts", "")
draft = dict(state.get("confluence_draft") or {}) # type: ignore[call-overload]
turn = _conf_gate_turn(visits)
question_id = _question_id_for(thread_id, turn)
deadline = (datetime.now(timezone.utc) + _DEFAULT_CONF_DEADLINE).isoformat()
preview = _draft_preview(draft)
question_set = _conf_approval_question_set(
thread_id=thread_id,
question_id=question_id,
turn=turn,
preview=preview,
)
decision = interrupt(
{
"thread_id": thread_id,
"question_id": question_id,
"turn": turn,
"kind": CONFLUENCE_APPROVAL_KIND,
"question_set": question_set,
"transport": transport,
"deadline": deadline,
"slack_thread_ts": slack_thread_ts,
"confluence_draft": draft,
"preview": preview,
}
)
return _apply_conf_decision(state, decision, visits=visits)
def _parse_decision(decision: Any) -> tuple[str, str]:
"""Normalize a resume decision into ``(verb, notes)`` (mirrors graph).
Delegates to the transport-neutral
:func:`agent_team.decisions.normalize_decision` (``allow_abandon=True``) so
every writer fails safe at this single chokepoint: an explicit approve /
request_changes / abandon is honoured, and anything UNRECOGNIZED maps to
``request_changes`` carrying the full reply as notes (never a silent
terminal abandon).
"""
from agent_team.decisions import normalize_decision
result = normalize_decision(decision, allow_abandon=True)
return result["decision"], result["notes"]
def _apply_conf_decision(
state: PipelineState, decision: Any, *, visits: int
) -> PipelineState:
"""Consume the owner's resume decision and return the routing state.
Mirrors :func:`agent_team.graph._apply_plan_decision`. The default for an
unrecognized/empty decision is the SAFE direction (``request_changes``),
never an accidental approve and never a silent terminal abandon.
"""
verb, notes = _parse_decision(decision)
now = _utc_now_iso()
if verb == "approve":
return PipelineState(
status=TaskStatus.ACTIVE.value,
current_phase=CONF_WRITE_PHASE,
confluence_gate_visits=visits, # type: ignore[typeddict-unknown-key]
updated_at=now,
)
if verb == "request_changes":
return PipelineState(
status=TaskStatus.ACTIVE.value,
current_phase=CONF_DRAFT_PHASE,
confluence_feedback=notes, # type: ignore[typeddict-unknown-key]
confluence_gate_visits=visits, # type: ignore[typeddict-unknown-key]
updated_at=now,
)
# Explicit abandon only (unrecognized was already mapped to request_changes).
return PipelineState(
status=TaskStatus.FAILED.value,
current_phase=Phase.PARKED.value,
confluence_gate_visits=visits, # type: ignore[typeddict-unknown-key]
failure_reason=f"confluence draft abandoned at human gate: {notes}".rstrip(
": "
),
updated_at=now,
)
def route_after_conf_gate(state: PipelineState) -> str:
"""Conditional-edge after the Confluence gate (mirrors route_after_plan_gate).
Reads the routing state :func:`conf_gate_node` wrote on resume (or on the
ceiling-reached terminal park) and maps it to a route id:
* status ACTIVE + phase CONF_DRAFT -> ``'revise'`` (loop back to redraft);
* status ACTIVE + phase CONF_WRITE -> ``'approve'`` (advance to the write);
* anything else (FAILED, or PARKED ceiling) -> ``'terminal'``.
"""
status = state.get("status")
phase = state.get("current_phase")
if status == TaskStatus.ACTIVE.value and phase == CONF_DRAFT_PHASE:
return "revise"
if status == TaskStatus.ACTIVE.value and phase == CONF_WRITE_PHASE:
return "approve"
return "terminal"
# -- Node 3: CONF_WRITE (dry-run by default; injected client). ----------------
def _default_confluence_client() -> ConfluenceClient:
"""Construct the real Confluence client at call time (deferred import).
The import is deferred so this module loads cleanly before the
:mod:`agent_team.confluence` package exists. Credentials are read from the
environment inside the client at call time, never captured here.
"""
from agent_team.confluence.client import ConfluenceClient
return ConfluenceClient()
def _apply_enabled(state: PipelineState, config: dict[str, Any] | None) -> bool:
"""Decide whether a LIVE write may be issued (dry-run is the default).
A live write requires an explicit apply flag — ``config['confluence_apply']``
truthy OR the env ``AGENT_TEAM_CONFLUENCE_APPLY`` truthy — AND the gate must
have approved (the node only runs on the approve route, but this is checked
defensively). The env is read at call time, never at import.
"""
flag = False
if config is not None:
flag = _truthy(config.get("confluence_apply"))
if not flag:
flag = _truthy(os.environ.get("AGENT_TEAM_CONFLUENCE_APPLY"))
return flag
def conf_write_node(
state: PipelineState,
config: dict[str, Any] | None = None,
*,
client: ConfluenceClient | None = None,
) -> PipelineState:
"""CONF_WRITE stage: perform (or dry-run) the Confluence page update.
**DRY-RUN by default.** A LIVE write is issued ONLY when
:func:`_apply_enabled` is true (an explicit ``confluence_apply`` config flag
or the ``AGENT_TEAM_CONFLUENCE_APPLY`` env truthy). Otherwise the node
composes the *planned* change plus a revert diff and writes NOTHING — so the
default path tests hit touches no live network.
Write routing:
* If the draft carries ``mermaid_edits`` AND the target page has diagram
macros, the update is routed through
:func:`agent_team.confluence.mermaid.plan_mermaid_edits`;
* otherwise it is a storage-format body update.
The Confluence client is INJECTED (``client``) so tests pass an in-memory
fake; the default is a real client constructed at call time (deferred import,
creds read from the environment). The Confluence package import is deferred
so this module loads before that package exists.
Returns a **partial** :class:`PipelineState`: terminal DONE (phase
CONF_DONE) with a ``confluence_result`` dict describing what was (or would
have been) written.
"""
draft = dict(state.get("confluence_draft") or {}) # type: ignore[call-overload]
if not draft.get("title") or not draft.get("body_storage"):
raise ConfluenceWriteError(
"conf_write_node requires a confluence_draft with title + body_storage"
)
apply = _apply_enabled(state, config)
page_id = draft.get("page_id")
mermaid_edits = draft.get("mermaid_edits") or []
if not apply:
# Dry run: compose the planned change + revert diff, write NOTHING.
result: dict[str, Any] = {
"applied": False,
"dry_run": True,
"action": "update" if page_id else "create",
"page_id": page_id,
"title": draft.get("title"),
"mermaid_edits": len(mermaid_edits),
"planned_change": {
"title": draft.get("title"),
"body_storage": draft.get("body_storage"),
},
"revert_diff": _compose_revert(state, draft),
}
return PipelineState(
status=TaskStatus.DONE.value,
current_phase=CONF_DONE_PHASE,
confluence_result=result, # type: ignore[typeddict-unknown-key]
updated_at=_utc_now_iso(),
)
# Live write path (only with an explicit apply flag AND prior gate approval).
conf = client if client is not None else _default_confluence_client()
try:
if mermaid_edits and _page_has_macros(conf, page_id):
outcome, applied = _apply_mermaid_edits(conf, page_id, mermaid_edits)
else:
outcome, applied = _apply_storage_update(conf, page_id, draft)
except ConfluenceWriteError:
raise
except Exception as exc: # noqa: BLE001 - normalized to a typed write error
raise ConfluenceWriteError(f"confluence write failed: {exc}") from exc
result = {
"applied": applied,
"dry_run": False,
"action": "update" if page_id else "create",
"page_id": getattr(outcome, "page_id", None)
if not isinstance(outcome, dict)
else outcome.get("page_id", page_id),
"title": draft.get("title"),
"mermaid_edits": len(mermaid_edits),
"outcome": outcome if isinstance(outcome, dict) else _outcome_to_dict(outcome),
}
return PipelineState(
status=TaskStatus.DONE.value,
current_phase=CONF_DONE_PHASE,
confluence_result=result, # type: ignore[typeddict-unknown-key]
updated_at=_utc_now_iso(),
)
def _current_page_version(client: ConfluenceClient, page_id: Any) -> int:
"""Read the target page's CURRENT version number (Confluence concurrency).
``ConfluenceClient.update_page`` requires the CURRENT version (it derives the
new version as ``version_number + 1`` itself). A page object missing a usable
``version.number`` is treated as version 0 so the update still issues against
a sane baseline rather than crashing.
"""
page = client.get_page(str(page_id))
version = page.get("version") if isinstance(page, dict) else None
if isinstance(version, dict):
number = version.get("number")
if isinstance(number, int):
return number
return 0
def _apply_storage_update(
client: ConfluenceClient, page_id: Any, draft: dict[str, Any]
) -> tuple[Any, bool]:
"""Issue the live storage-format page update; return ``(outcome, applied)``.
Fetches the page's CURRENT version (``update_page`` expects the current
number and bumps it internally per Confluence's optimistic-concurrency
contract) and calls ``update_page(..., apply=True)``. ``applied`` is read off
the returned :class:`PlannedPageUpdate` (``outcome.applied``) rather than
hardcoded, so a client that declines to apply is reported honestly.
"""
version_number = _current_page_version(client, page_id) if page_id else 0
outcome = client.update_page(
page_id=page_id,
title=draft.get("title"),
body_storage=draft.get("body_storage"),
version_number=version_number,
apply=True,
)
applied = bool(getattr(outcome, "applied", True))
return outcome, applied
def _apply_mermaid_edits(
client: ConfluenceClient, page_id: Any, mermaid_edits: list[dict[str, Any]]
) -> tuple[Any, bool]:
"""Plan the Mermaid ADF edits for a macro page; return ``(outcome, applied)``.
:func:`agent_team.confluence.mermaid.plan_mermaid_edits` is PURE ADF: it takes
the parsed ADF document and a list of :class:`~agent_team.confluence.mermaid.MermaidEdit`
(``macro_key`` / ``new_source``) and never touches the client. The draft's
edit dicts (``{"mermaid": ..., "macro_id"/"anchor": ...}``) are converted to
``MermaidEdit`` objects here.
The Confluence client currently exposes NO ADF fetch/persist methods (only
storage-format ``get_page`` / ``update_page``). When the injected client adds
an ADF capability (``get_page_adf`` + ``update_page_adf``) this routes through
it and reports ``applied`` off that persistence; until then there is no way to
persist an ADF edit, so this raises :class:`ConfluenceWriteError` rather than
falsely recording a successful Mermaid write.
"""
from agent_team.confluence import mermaid as mermaid_mod
edits = [
mermaid_mod.MermaidEdit(
macro_key=str(item.get("macro_id") or item.get("anchor") or ""),
new_source=str(item.get("mermaid", "")),
)
for item in mermaid_edits
]
get_adf = getattr(client, "get_page_adf", None)
put_adf = getattr(client, "update_page_adf", None)
if not callable(get_adf) or not callable(put_adf):
raise ConfluenceWriteError(
"Mermaid live apply needs ADF persistence: the Confluence client "
"lacks get_page_adf/update_page_adf (ADF-only edits cannot round-trip "
"through storage format without dropping diagram macros)."
)
adf = get_adf(str(page_id))
plan = mermaid_mod.plan_mermaid_edits(adf, edits, apply=True)
if plan.skip_mermaid:
# The page turned out to carry zero Mermaid macros after all — defer to
# the caller's storage-format fallback path semantics by signalling skip.
raise ConfluenceWriteError(
"Mermaid live apply found no Mermaid macros on the page (skip_mermaid)."
)
outcome = put_adf(str(page_id), plan.new_adf)
applied = bool(getattr(outcome, "applied", True))
return outcome, applied
def _page_has_macros(client: ConfluenceClient, page_id: Any) -> bool:
"""Best-effort check that the target page carries diagram macros.
Delegates to the injected client's ``page_has_macros`` when available; a
missing capability (older client / no page id) degrades to ``False`` so the
update falls back to a plain storage-format body update rather than crashing.
"""
if not page_id:
return False
checker = getattr(client, "page_has_macros", None)
if checker is None:
return False
try:
return bool(checker(page_id))
except Exception: # noqa: BLE001 - capability probe must never crash the write
return False
def _compose_revert(state: PipelineState, draft: dict[str, Any]) -> dict[str, Any]:
"""Compose a revert descriptor for the dry-run record.
Captures enough of the prior state to describe how the planned change would
be reverted — the prior page id (if updating) and a marker that the original
body is unchanged on disk (the dry run wrote nothing). Pure data; no I/O.
"""
page_id = draft.get("page_id")
return {
"page_id": page_id,
"note": (
"dry run — no write performed; revert is a no-op. On apply, revert "
"restores the page version prior to this update."
if page_id
else "dry run — would create a new page; revert deletes it."
),
}
def _outcome_to_dict(outcome: Any) -> dict[str, Any]:
"""Coerce a client write-outcome object into a JSON-safe dict (best effort)."""
for attr in ("to_dict", "_asdict"):
fn = getattr(outcome, attr, None)
if callable(fn):
try:
return dict(fn())
except Exception: # noqa: BLE001
pass
if isinstance(outcome, dict):
return dict(outcome)
return {"repr": repr(outcome)}