Files
ai-agent/backend/conversations/catalog.py
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

275 lines
11 KiB
Python

"""The conversation cards every list view reads, and the session-id lookups."""
import json
import pathlib
import threading
import time
from fastapi import HTTPException
import conversations
from agents.runs import _agent_runs
from conversations.cards import _conv_meta, _conv_title
from conversations.tree import (
_conv_id,
_is_deploying,
_is_sidechain_path,
_rollup_children,
_root_path_of,
_running_agents,
)
from core.config import TRANSCRIPTS_DIR
from core.memo import _Memo
from core.state import AppState
_skills_by_path_cache: tuple[int, dict[str, dict[str, int]]] | None = None
def _skills_by_path(deps: AppState) -> dict[str, dict[str, int]]:
"""transcript path → ``{skill: calls}``, decoded once per store version.
The indexer keeps a transcript's skill map in its **own column** (it feeds
the catalog rollup), not inside the summary blob — so a stored summary has
no `skills` key at all, and this is the join that puts it back.
"""
global _skills_by_path_cache
hit = _skills_by_path_cache
if hit and hit[0] == deps.store.summaries_version:
return hit[1]
out: dict[str, dict[str, int]] = {}
for path, raw in deps.store.all_transcript_skill_rows():
try:
contrib = json.loads(raw or "{}")
except ValueError:
continue
counts = {name: d.get("count") or 0 for name, d in contrib.items()}
if counts:
out[path] = counts
_skills_by_path_cache = (deps.store.summaries_version, out)
return out
def _skills_used(deps: AppState, paths: list[str], extra: dict | None = None) -> list[str]:
"""Skill names used across `paths`, busiest first.
Callers pass a conversation *and its subagent transcripts*: a `Skill(...)`
run by a spawned Explore is work this conversation caused, and the child is
never listed on its own — so hiding its skills would lose them entirely.
`extra` folds in a freshly parsed ``{skill: {count, …}}`` map for a
transcript the store may not have re-read yet (the detail view parses live).
"""
by_path = _skills_by_path(deps, )
counts: dict[str, int] = {}
for p in paths:
for name, n in by_path.get(p, {}).items():
counts[name] = counts.get(name, 0) + n
for name, d in (extra or {}).items():
counts[name] = counts.get(name, 0) + (d.get("count") or 0)
return sorted(counts, key=lambda n: (-counts[n], n))
def _build_agents_by_conversation(deps: AppState) -> dict[str, list[str]]:
counts: dict[str, dict[str, int]] = {}
for r in _agent_runs(deps, ):
cid, name = r.get("conversationId"), r.get("agent")
if not cid or not name:
continue
c = counts.setdefault(cid, {})
c[name] = c.get(name, 0) + 1
return {cid: sorted(c, key=lambda n: (-c[n], n)) for cid, c in counts.items()}
_agents_by_conv_memo = _Memo(_build_agents_by_conversation)
def _agents_by_conversation(deps: AppState) -> dict[str, list[str]]:
"""conversation id → the agent types that ran in it, busiest first.
Built from the same run list the agents catalog counts (`_agent_runs`) —
not from a second pass over the `Task` calls — so a conversation card can
never disagree with the agent page. That list already resolves both origins:
a subagent run is keyed to its *parent* conversation (where its Task card
lives), and a cron/manual agent session to the conversation it *is*.
"""
return _agents_by_conv_memo.get(
(deps.store.summaries_version, deps.meta_store.version), deps)
def _conv_card(deps: AppState, path: str, s: dict) -> dict:
"""Compact conversation reference used when cross-linking from a memory."""
return {
"id": _conv_id(path),
"sessionId": s.get("sessionId"),
"title": _conv_title(deps, s),
"model": s.get("model"),
"endedAt": s.get("endedAt"),
"cost": s.get("cost"),
"tokens": s.get("tokens"),
"meta": _conv_meta(deps, s),
}
# Card fields every list view needs. `usage`/`byToolSub` are heavy per-card and
# only the analytics dashboard reads them — it asks with ``full=1``.
_CARD_KEYS = ("sessionId", "project", "gitBranch", "model", "models", "efforts",
"title", "userTurns", "assistantTurns", "messages", "startedAt", "endedAt",
"tokens", "cost", "byTool", "subagentCount",
"subagentTokens", "subagentCost", "tasks")
_CARD_KEYS_FULL = _CARD_KEYS + ("usage", "byToolSub")
def _dedup_by_session(summaries: list[tuple[str, dict]]) -> list[tuple[str, dict]]:
"""Collapse transcript files that share a ``sessionId`` into one.
Resuming a session in a different working directory (a worktree, or
``claude --resume`` from another dir) makes Claude Code write the
continuation to a *new* cwd-encoded folder under the SAME sessionId — so
one logical conversation lands on disk as several ``<cwd>/<sid>.jsonl``
files. Keyed on the file path, each became its own card and the
conversation showed up two (or more) times. Keep only the file with the
most recent activity (latest ``endedAt``, then most messages, then longest
path as a stable tiebreak) — that's the fully-resumed transcript. Files
without a sessionId are never collapsed (nothing links them)."""
best: dict[str, tuple[str, dict]] = {}
passthrough: list[tuple[str, dict]] = []
for path, s in summaries:
sid = s.get("sessionId")
if not sid:
passthrough.append((path, s))
continue
cur = best.get(sid)
if cur is None or (
(s.get("endedAt") or "", s.get("messages") or 0, path)
> (cur[1].get("endedAt") or "", cur[1].get("messages") or 0, cur[0])
):
best[sid] = (path, s)
return list(best.values()) + passthrough
# A card's `state` also depends on the clock (`_auto_state`'s recency window,
# 10 min), so the memo key carries a coarse time bucket too: an idle list flips
# a stale "running" badge within this many seconds even with no write.
_CARDS_TIME_BUCKET_S = 15
def _conversation_cards(deps: AppState, full: bool) -> list[dict]:
"""Every non-sidechain conversation as a card dict, unsorted — memoized per
(store, sidecar, time bucket) and shared across concurrent requests.
The returned list and its cards are shared: callers must copy before
sorting/filtering in place, and never mutate a card."""
key = (deps.store.summaries_version, deps.meta_store.version,
int(time.time() // _CARDS_TIME_BUCKET_S))
return _cards_memo[bool(full)].get(key, deps, bool(full))
def _build_conversation_cards(deps: AppState, full: bool) -> list[dict]:
"""Every non-sidechain conversation as a card dict, unsorted.
Summaries come from the store's shared decode cache and are **never
mutated** here — rollups copy, and per-card extras land on the card."""
summaries = deps.store.all_summaries()
# Group subagent summaries under their parent so we can roll their usage up.
children: dict[str, list[dict]] = {}
child_paths: dict[str, list[str]] = {}
for path, s in summaries:
if s.get("isSidechain"):
pp = _root_path_of(path) # every depth rolls into the root card
if pp:
children.setdefault(pp, []).append(s)
child_paths.setdefault(pp, []).append(path)
keys = _CARD_KEYS_FULL if full else _CARD_KEYS
# A session resumed across working dirs has one transcript file per cwd —
# collapse those so the conversation shows up once, not once per copy.
out = []
agents_by_conv = _agents_by_conversation(deps, )
for path, s in _dedup_by_session(
[(p, s) for p, s in summaries if not s.get("isSidechain")]
):
if not s.get("messages"):
continue # skip empty transcripts
skills_used = _skills_used(deps, [path, *child_paths.get(path, [])])
s = _rollup_children(s, children.get(path, []))
m = _conv_meta(deps, s)
cid = _conv_id(path)
out.append({"id": cid, **{k: s.get(k) for k in keys},
"title": _conv_title(deps, s),
"runningAgents": _running_agents(deps, s, m.get("state")),
"deploying": _is_deploying(deps, s, m.get("state")),
"skillsUsed": skills_used,
"agentsUsed": agents_by_conv.get(cid, []),
"meta": m})
return out
_cards_memo = {f: _Memo(_build_conversation_cards) for f in (False, True)}
# ── diff view: per-conversation commits + commit / conversation diffs ────────
def _resolve_transcript(id: str) -> pathlib.Path:
"""Validate a conversation id and return its archived transcript path."""
if not TRANSCRIPTS_DIR:
raise HTTPException(404, "transcripts not mounted")
rel = pathlib.PurePosixPath(id)
if rel.is_absolute() or ".." in rel.parts or rel.suffix != ".jsonl":
raise HTTPException(400, "invalid id")
ap = (TRANSCRIPTS_DIR / id).resolve()
try:
ap.relative_to(TRANSCRIPTS_DIR)
except ValueError:
raise HTTPException(400, "id escapes transcripts")
if not ap.is_file():
raise HTTPException(404, "not found")
return ap
def _summary_and_meta(deps: AppState, id: str) -> tuple[str, dict, dict]:
"""(sessionId, summary, meta sidecar) for a conversation id."""
ap = _resolve_transcript(id)
summary = deps.store.get_summary(str(ap)) or conversations.parse_conversation(ap, full=False)
sid = summary.get("sessionId") or ""
return sid, summary, deps.meta_store.get(sid)
def _meta_key(raw: str) -> str:
"""Accept a bare session id or a transcript id; store under the session id."""
name = pathlib.PurePosixPath(raw).name
return name[:-6] if name.endswith(".jsonl") else name
# ── cron jobs: scheduled agent sessions (see cron.py) ────────────────────────
# sessionId → conversation id, for linking a job's run history to the viewer.
# Memoized on the summaries version (rebuilding walks ~900 summaries).
_sid_map_cache: tuple[int, dict] | None = None
_sid_map_lock = threading.Lock()
def _session_conv_map(deps: AppState) -> dict:
global _sid_map_cache
v = deps.store.summaries_version
with _sid_map_lock:
if _sid_map_cache and _sid_map_cache[0] == v:
return _sid_map_cache[1]
m: dict[str, str] = {}
for path, s in deps.store.all_summaries():
sid = s.get("sessionId")
if sid and not _is_sidechain_path(path):
m[sid] = _conv_id(path)
with _sid_map_lock:
_sid_map_cache = (v, m)
return m
def _conv_ids_by_session(deps: AppState) -> dict[str, str]:
"""sessionId → transcript id (newest summary wins) for conversation links."""
best: dict[str, tuple[str, str]] = {}
for path, s in deps.store.all_summaries():
sid = s.get("sessionId")
if not sid or not s.get("messages"):
continue
ended = s.get("endedAt") or ""
prev = best.get(sid)
if prev is None or ended > prev[0]:
best[sid] = (ended, _conv_id(path))
return {sid: cid for sid, (_, cid) in best.items()}