main.py (5146 lines, 106 routes) becomes an assembly only: one package per domain — core, files, settings, dashboards, conversations, diff, notifications, forms, runs, accounts, workers, models, cron, agents, projects, services, goals, memories, plans, templates — each exposing an APIRouter; the flat domain modules move into their package behind a barrel that keeps the old `import conversations` / `import projects` spellings. The shared singletons (store, meta_store, hub, indexer, …) are built once by core.state.build_state() and attached to app.state.ai; routes take them as the `deps: State` dependency and helpers as an explicit `deps: AppState`. conversations/pricing.py carries the per-model rates out of the parser. Verified: route table and OpenAPI byte-identical; 90 read endpoints golden-diffed against the monolith on a copy of the live data (identical); write routes smoke-tested; 66 backend tests pass. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
391 lines
18 KiB
Python
391 lines
18 KiB
Python
"""Subagent (sidechain) linkage: the transcript tree under ``subagents/``."""
|
|
|
|
import functools
|
|
import json
|
|
import pathlib
|
|
|
|
import conversations
|
|
from conversations.cards import _conv_meta, _conv_title
|
|
from core.config import TRANSCRIPTS_DIR
|
|
from core.state import AppState
|
|
|
|
|
|
# The path helpers below run once per summary on every list/agents build
|
|
# (thousands of calls per request) and `pathlib.relative_to` alone costs
|
|
# ~0.2 ms a call — memoized, since a transcript path maps to the same id forever.
|
|
@functools.lru_cache(maxsize=32768)
|
|
def _conv_id(path: str) -> str:
|
|
"""Transcript path → the rel id the detail endpoint / UI use."""
|
|
try:
|
|
return str(pathlib.Path(path).relative_to(TRANSCRIPTS_DIR)) \
|
|
if TRANSCRIPTS_DIR else pathlib.Path(path).name
|
|
except ValueError:
|
|
return pathlib.Path(path).name
|
|
|
|
|
|
# ── subagent (sidechain) linkage ─────────────────────────────────────────────
|
|
# Subagent transcripts live at `<convId>/subagents/agent-<id>.jsonl` next to the
|
|
# root `<convId>.jsonl`. Claude Code writes **every** nested agent into that one
|
|
# flat dir, whatever its depth: a subagent's own Task calls spawn grandchildren
|
|
# whose `agent-<id>.meta.json` names the spawner in `parentAgentId` (absent on a
|
|
# depth-1 agent, whose parent is the root conversation) and the level in
|
|
# `spawnDepth`. So the directory is a tree keyed by `parentAgentId`, and two
|
|
# "parent" notions matter: the **root** conversation (the `<convId>.jsonl` the
|
|
# dir belongs to — where usage rolls up and where the list card lives) and the
|
|
# **direct parent** (the transcript holding the Task card that spawned it —
|
|
# where the card links and the running state come from). Sidechains are parsed
|
|
# like any transcript (flagged `isSidechain`) but hidden from the top-level
|
|
# lists; each is reachable from its parent's Task/Agent card, at any depth.
|
|
@functools.lru_cache(maxsize=32768)
|
|
def _is_sidechain_path(abs_path: str) -> bool:
|
|
return "subagents" in pathlib.Path(abs_path).parts
|
|
|
|
|
|
@functools.lru_cache(maxsize=32768)
|
|
def _root_path_of(abs_path: str) -> str | None:
|
|
"""Given a subagent transcript path (any depth), the root `<convId>.jsonl`."""
|
|
parts = pathlib.Path(abs_path).parts
|
|
if "subagents" not in parts:
|
|
return None
|
|
idx = parts.index("subagents")
|
|
conv_dir = pathlib.Path(*parts[:idx]) # …/<convId>
|
|
return str(conv_dir.with_suffix(".jsonl")) # …/<convId>.jsonl
|
|
|
|
|
|
def _agent_id_of(abs_path: str) -> str | None:
|
|
"""`…/subagents/agent-<id>.jsonl` → `<id>`; None for a root transcript."""
|
|
if not _is_sidechain_path(abs_path):
|
|
return None
|
|
stem = pathlib.Path(abs_path).stem
|
|
return stem[len("agent-"):] if stem.startswith("agent-") else stem
|
|
|
|
|
|
def _parent_path_of(abs_path: str) -> str | None:
|
|
"""Given a subagent transcript path, its **direct** parent's transcript: the
|
|
sibling `agent-<parentAgentId>.jsonl` for a nested agent, else the root
|
|
`<convId>.jsonl`. Not memoised on the path alone: the meta.json may not be
|
|
written yet on the first look (the cached reader only keeps non-empty
|
|
reads), and a nested agent must not get pinned to the root."""
|
|
root = _root_path_of(abs_path)
|
|
if not root:
|
|
return None
|
|
pid = _subagent_meta_cached(pathlib.Path(abs_path)).get("parentAgentId")
|
|
if pid:
|
|
return str(pathlib.Path(abs_path).with_name(f"agent-{pid}.jsonl"))
|
|
return root
|
|
|
|
|
|
def _parent_chain_of(abs_path: str) -> list[str]:
|
|
"""Every ancestor transcript of a subagent, root first, direct parent last.
|
|
Bounded by the dir's size (a bad `parentAgentId` cycle can't loop)."""
|
|
chain: list[str] = []
|
|
seen: set[str] = set()
|
|
cur = _parent_path_of(abs_path)
|
|
while cur and cur not in seen:
|
|
chain.append(cur)
|
|
seen.add(cur)
|
|
cur = _parent_path_of(cur)
|
|
chain.reverse()
|
|
return chain
|
|
|
|
|
|
def _subagent_dir_of(ap: pathlib.Path) -> pathlib.Path | None:
|
|
"""The flat `subagents/` dir a transcript's descendants live in — beside a
|
|
root (`<convId>/subagents/`), the own dir for a sidechain."""
|
|
if _is_sidechain_path(str(ap)):
|
|
return ap.parent
|
|
d = ap.with_suffix("") / "subagents"
|
|
return d if d.is_dir() else None
|
|
|
|
|
|
def _subagent_children(ap: pathlib.Path) -> list[tuple[pathlib.Path, dict]]:
|
|
"""The **direct** children of a transcript, as `(agent jsonl, meta)` pairs
|
|
in spawn order (meta files sort by agent id, which Claude Code assigns in
|
|
order): a root's are the metas with no `parentAgentId`, a sidechain's are
|
|
the ones naming its agent id. A meta without `parentAgentId` under an
|
|
older transcript is a depth-1 agent by definition, so pre-nesting archives
|
|
keep their one-level reading."""
|
|
sub_dir = _subagent_dir_of(ap)
|
|
if sub_dir is None or not sub_dir.is_dir():
|
|
return []
|
|
own = _agent_id_of(str(ap))
|
|
out = []
|
|
for meta_f in sorted(sub_dir.glob("*.meta.json")):
|
|
stem = meta_f.name[:-len(".meta.json")] # agent-<id>
|
|
child = sub_dir / f"{stem}.jsonl"
|
|
if child == ap:
|
|
continue
|
|
mj = _subagent_meta_cached(child)
|
|
if (mj.get("parentAgentId") or None) == own:
|
|
out.append((child, mj))
|
|
return out
|
|
|
|
|
|
def _subagent_descendants(ap: pathlib.Path) -> list[pathlib.Path]:
|
|
"""Every transcript under `ap` in the agent tree (children, grandchildren…),
|
|
breadth-first. For a root that is the whole `subagents/` dir."""
|
|
out: list[pathlib.Path] = []
|
|
queue = [ap]
|
|
seen: set[str] = {str(ap)}
|
|
while queue:
|
|
cur = queue.pop(0)
|
|
for child, _mj in _subagent_children(cur):
|
|
if str(child) in seen:
|
|
continue
|
|
seen.add(str(child))
|
|
out.append(child)
|
|
queue.append(child)
|
|
return out
|
|
|
|
|
|
def _with_child_models(models: list[str] | None, children: list[dict]) -> list[str]:
|
|
"""The conversation's own models (busiest first) followed by any model only
|
|
its subagents ran on — so a thread on opus whose Explore agents ran haiku
|
|
is tagged with both. Appended, never interleaved: `models[0]` stays the
|
|
main thread's busiest model, which is what the family/footprint reads."""
|
|
out = list(models or [])
|
|
for c in children:
|
|
for m in c.get("models") or ([c["model"]] if c.get("model") else []):
|
|
if m not in out:
|
|
out.append(m)
|
|
return out
|
|
|
|
|
|
def _rollup_children(summary: dict, children: list[dict]) -> dict:
|
|
"""Fold each subagent's aggregate usage into the parent's totals and its
|
|
`agent` activity bucket, so a conversation's cost includes the sub-agents it
|
|
spawned. Returns a shallow copy; leaves the stored summary untouched."""
|
|
if not children:
|
|
return summary
|
|
s = dict(summary)
|
|
usage = dict(summary.get("usage") or conversations.zero_usage())
|
|
by_tool = {k: dict(v) for k, v in (summary.get("byTool") or {}).items()}
|
|
by_tool_sub = {k: {sk: dict(sv) for sk, sv in v.items()}
|
|
for k, v in (summary.get("byToolSub") or {}).items()}
|
|
add_tokens = add_output = 0
|
|
add_cost = 0.0
|
|
for c in children:
|
|
cu = c.get("usage") or {}
|
|
for k in usage:
|
|
usage[k] = (usage.get(k) or 0) + (cu.get(k) or 0)
|
|
add_tokens += c.get("tokens") or 0
|
|
add_cost += c.get("cost") or 0.0
|
|
add_output += cu.get("output") or 0
|
|
agent = by_tool.get("agent") or {"tokens": 0, "cost": 0.0, "output": 0}
|
|
by_tool["agent"] = {"tokens": (agent.get("tokens") or 0) + add_tokens,
|
|
"cost": (agent.get("cost") or 0.0) + add_cost,
|
|
"output": (agent.get("output") or 0) + add_output}
|
|
# Mirror the rollup into the drill-down: a "subagents" sub of the agent bucket.
|
|
asub = by_tool_sub.setdefault("agent", {})
|
|
sa = asub.get("subagents") or {"tokens": 0, "cost": 0.0, "output": 0}
|
|
asub["subagents"] = {"tokens": (sa.get("tokens") or 0) + add_tokens,
|
|
"cost": (sa.get("cost") or 0.0) + add_cost,
|
|
"output": (sa.get("output") or 0) + add_output}
|
|
s["usage"] = usage
|
|
s["byTool"] = by_tool
|
|
s["byToolSub"] = by_tool_sub
|
|
s["tokens"] = (summary.get("tokens") or 0) + add_tokens
|
|
s["cost"] = (summary.get("cost") or 0.0) + add_cost
|
|
s["models"] = _with_child_models(summary.get("models"), children)
|
|
s["subagentCount"] = len(children)
|
|
s["subagentTokens"] = add_tokens
|
|
s["subagentCost"] = add_cost
|
|
return s
|
|
|
|
|
|
def _subagent_meta(sub_jsonl: pathlib.Path) -> dict:
|
|
"""Read the `agent-<id>.meta.json` sitting next to a subagent transcript."""
|
|
meta_f = sub_jsonl.with_name(sub_jsonl.stem + ".meta.json")
|
|
try:
|
|
return json.loads(meta_f.read_text(encoding="utf-8"))
|
|
except (OSError, ValueError):
|
|
return {}
|
|
|
|
|
|
# meta.json files are written once at spawn and never change, so list endpoints
|
|
# may cache them (only non-empty reads — an unwritten file may appear later).
|
|
_SUBAGENT_META_CACHE: dict[str, dict] = {}
|
|
|
|
|
|
def _subagent_meta_cached(sub_jsonl: pathlib.Path) -> dict:
|
|
key = str(sub_jsonl)
|
|
hit = _SUBAGENT_META_CACHE.get(key)
|
|
if hit is not None:
|
|
return hit
|
|
mj = _subagent_meta(sub_jsonl)
|
|
if mj:
|
|
_SUBAGENT_META_CACHE[key] = mj
|
|
return mj
|
|
|
|
|
|
def _running_agents(deps: AppState, summary: dict, state: str | None = None) -> list[dict]:
|
|
"""The subagents a conversation still has in flight.
|
|
|
|
An open Task call only means "running" while the conversation itself is: a
|
|
run that was interrupted (or died) mid-agent leaves its call open forever,
|
|
and so does every transcript written before the parser learned to close a
|
|
background agent on its `<task-notification>`. Nothing is running there."""
|
|
ra = summary.get("runningAgents") or []
|
|
if not ra:
|
|
return []
|
|
st = state if state is not None else _conv_meta(deps, summary).get("state")
|
|
return ra if st == "running" else []
|
|
|
|
|
|
def _is_deploying(deps: AppState, summary: dict, state: str | None = None) -> bool:
|
|
"""True while this conversation has a deploy command in flight (issued, no
|
|
tool_result yet) — gated the same way as `_running_agents`: a dangling call
|
|
left open by an interrupted/dead run isn't actually deploying anymore."""
|
|
if not summary.get("deploying"):
|
|
return False
|
|
st = state if state is not None else _conv_meta(deps, summary).get("state")
|
|
return st == "running"
|
|
|
|
|
|
def _sidechain_state(deps: AppState, ap: pathlib.Path) -> str | None:
|
|
"""A subagent's own lifecycle state: it runs exactly as long as the parent's
|
|
Task call is still open. `None` if `ap` isn't a subagent transcript.
|
|
|
|
Its `_conv_meta` state can't be trusted: a sidechain's sessionId is the
|
|
*parent's*, so the sidecar lookup returns the parent's (possibly stale)
|
|
state, and the recency fallback would keep a just-finished subagent "running"
|
|
for the whole recency window."""
|
|
pp = _is_sidechain_path(str(ap)) and _parent_path_of(str(ap))
|
|
if not pp:
|
|
return None
|
|
psum = deps.store.get_summary(pp) or {}
|
|
tuid = _subagent_meta(ap).get("toolUseId")
|
|
# A nested agent's parent is itself a sidechain: resolve *its* state the
|
|
# same way (up the chain to the root), since its own meta state would be
|
|
# the root's too.
|
|
pstate = _sidechain_state(deps, pathlib.Path(pp))
|
|
open_calls = {a.get("toolUseId") for a in _running_agents(deps, psum, pstate)}
|
|
return "running" if tuid and tuid in open_calls else "finished"
|
|
|
|
|
|
def _attach_subagents(deps: AppState, ap: pathlib.Path, data: dict) -> dict:
|
|
"""Enrich a parsed thread with its subagents: link each Task/Agent card to
|
|
the **direct** child transcript it spawned (via meta.json `toolUseId`) and,
|
|
on a root conversation, roll every descendant's usage into its totals. A
|
|
subagent's own view gets the same card linking (its Task calls spawned the
|
|
grandchildren) plus a `parentConversation` back-link and the `parentChain`
|
|
breadcrumb up to the root — but no rollup: an agent page reports its own
|
|
spend, the root the whole tree's. Returns the (mutated) data."""
|
|
is_sub = _is_sidechain_path(str(ap))
|
|
if is_sub:
|
|
# Link back up: the direct parent + the exact Task card, and the chain
|
|
# of ancestors (root first) for the breadcrumb.
|
|
chain = _parent_chain_of(str(ap))
|
|
mj = _subagent_meta(ap)
|
|
if chain:
|
|
pp = chain[-1]
|
|
psum = deps.store.get_summary(pp) or {}
|
|
data["parentConversation"] = {
|
|
"id": _conv_id(pp),
|
|
"title": _conv_title(deps, psum),
|
|
"toolUseId": mj.get("toolUseId"),
|
|
"agentType": mj.get("agentType"),
|
|
"description": mj.get("description"),
|
|
}
|
|
refs = []
|
|
for anc in chain:
|
|
asum = deps.store.get_summary(anc) or {}
|
|
amj = (_subagent_meta_cached(pathlib.Path(anc))
|
|
if _is_sidechain_path(anc) else {})
|
|
refs.append({"id": _conv_id(anc), "title": _conv_title(deps, asum),
|
|
"agentType": amj.get("agentType"),
|
|
"description": amj.get("description")})
|
|
data["parentChain"] = refs
|
|
data["spawnDepth"] = mj.get("spawnDepth") or len(chain) or 1
|
|
# Mark every Task/Agent card with whether its agent is still running (the
|
|
# parser tracks the open calls as `runningAgents`) — including the ones
|
|
# whose subagent transcript hasn't been written yet, which is exactly the
|
|
# moment the card most needs to say "running".
|
|
open_calls = {a.get("toolUseId") for a in (data.get("runningAgents") or [])}
|
|
for it in data.get("thread") or []:
|
|
if it.get("name") in ("Task", "Agent") and it.get("toolUseId"):
|
|
it["agentRunning"] = it["toolUseId"] in open_calls
|
|
# Gather the direct children from the flat `subagents/` dir.
|
|
children = _subagent_children(ap)
|
|
if not children:
|
|
return data
|
|
link_by_tooluse: dict[str, dict] = {}
|
|
infos_by_id: dict[str, dict] = {}
|
|
subtree_cost: dict[str, tuple[float, int]] = {}
|
|
for child, mj in children:
|
|
csum = deps.store.get_summary(str(child)) or {}
|
|
tuid = mj.get("toolUseId")
|
|
grand = _subagent_descendants(child)
|
|
info = {"id": _conv_id(str(child)),
|
|
"agentType": mj.get("agentType"),
|
|
"description": mj.get("description"),
|
|
"title": csum.get("title"), "model": csum.get("model"),
|
|
"tokens": csum.get("tokens"), "cost": csum.get("cost"),
|
|
"messages": csum.get("messages"),
|
|
"startedAt": csum.get("startedAt"),
|
|
"endedAt": csum.get("endedAt"),
|
|
"state": "running" if tuid in open_calls else "finished",
|
|
"subagentCount": len(grand)}
|
|
infos_by_id[info["id"]] = info
|
|
if tuid:
|
|
link_by_tooluse[tuid] = info
|
|
# What the card that spawned this child actually caused: the child's
|
|
# own spend plus everything it spawned in turn (the grandchildren have
|
|
# no card on *this* page). The ref itself keeps the child's own cost —
|
|
# the number its page shows.
|
|
cost = csum.get("cost") or 0.0
|
|
toks = csum.get("tokens") or 0
|
|
for g in grand:
|
|
gs = deps.store.get_summary(str(g)) or {}
|
|
cost += gs.get("cost") or 0.0
|
|
toks += gs.get("tokens") or 0
|
|
subtree_cost[info["id"]] = (cost, toks)
|
|
# Link each Task/Agent card to its child and collect the subagents in the
|
|
# order they were spawned (thread order), so the detail view can list them.
|
|
ordered: list[dict] = []
|
|
seen: set[str] = set()
|
|
for it in data.get("thread") or []:
|
|
tuid = it.get("toolUseId")
|
|
info = tuid and link_by_tooluse.get(tuid)
|
|
if info:
|
|
it["subagent"] = info
|
|
# A subagent's spend is rolled into the parent's total and its `agent`
|
|
# bucket below, but it was billed to the *child's* transcript — so no
|
|
# parent thread item carries it, and the cost-growth chart (which sums
|
|
# thread items) would report less than the total. Credit it to the card
|
|
# that spawned it: the turn that actually caused the spend.
|
|
child_cost, child_tokens = subtree_cost.get(info["id"], (0.0, 0))
|
|
if child_cost > 0:
|
|
it["turnCost"] = (it.get("turnCost") or 0.0) + child_cost
|
|
it["turnTokens"] = (it.get("turnTokens") or 0) + child_tokens
|
|
# It's all downstream agent work, whatever the card's own turn drove.
|
|
it["turnBucket"] = "agent"
|
|
if info["id"] not in seen:
|
|
ordered.append(info)
|
|
seen.add(info["id"])
|
|
for cid, info in infos_by_id.items(): # any orphans (no matching card)
|
|
if cid not in seen:
|
|
ordered.append(info)
|
|
seen.add(cid)
|
|
# No card claimed this child's cost above, but the rollup still counts
|
|
# it — park it on an invisible carrier so the chart stays whole.
|
|
o_cost, o_tokens = subtree_cost.get(cid, (0.0, 0))
|
|
if o_cost > 0 and data.get("thread") is not None:
|
|
data["thread"].append({
|
|
"role": "assistant", "kind": "turn", "first": False, "out": 0,
|
|
"turnCost": o_cost, "turnTokens": o_tokens,
|
|
"turnBucket": "agent",
|
|
})
|
|
data["subagents"] = ordered
|
|
if is_sub:
|
|
return data # an agent page reports its own spend only
|
|
# Roll every descendant's usage into the root totals / `agent` bucket.
|
|
child_summaries = [cs for cs in (deps.store.get_summary(str(d))
|
|
for d in _subagent_descendants(ap)) if cs]
|
|
rolled = _rollup_children(data, child_summaries)
|
|
for k in ("usage", "byTool", "byToolSub", "tokens", "cost",
|
|
"subagentCount", "subagentTokens", "subagentCost"):
|
|
if k in rolled:
|
|
data[k] = rolled[k]
|
|
return data
|