Files
Gabriel Vidal 0ff9e40242 refactor(backend): split main.py into domain packages with app.state injection
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>
2026-10-06 23:55:47 +02:00

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