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>
275 lines
11 KiB
Python
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()}
|