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

378 lines
16 KiB
Python

"""Conversation list / detail / search / scaffold / metadata routes."""
import copy
import datetime
import pathlib
import threading
from fastapi import APIRouter, HTTPException
from fastapi.responses import PlainTextResponse
from pydantic import BaseModel
import accounts as accounts_mod
import conversations
import memories as memories_mod
import plans as plans_mod
import schemas
from conversations import scaffold as scaffold_mod
from conversations.cards import _conv_meta, _conv_title
from conversations.catalog import (
_CARD_KEYS,
_agents_by_conversation,
_conversation_cards,
_meta_key,
_resolve_transcript,
_skills_used,
_summary_and_meta,
)
from conversations.tree import (
_attach_subagents,
_conv_id,
_is_deploying,
_parent_path_of,
_rollup_children,
_root_path_of,
_running_agents,
_sidechain_state,
_subagent_descendants,
_subagent_meta_cached,
_with_child_models,
)
from core.http import _json_in_worker, _r
from core.state import State
from indexer import INPUT_PRICE_PER_MTOK, MODEL
from projects import artefacts as artefacts_mod
router = APIRouter()
@router.get("/api/conversations", responses=_r(schemas.ConversationsResponse))
@_json_in_worker
def conversations_list(deps: State, limit: int = 50, before: str | None = None,
archived: str = "false", full: bool = False,
sessions: str = "", account: str | None = None):
"""List parsed transcripts (newest first) for the conversation viewer.
Paginated: ``limit`` caps the page (0 = everything), ``before`` is the
``endedAt`` cursor of the previous page's last card, and ``archived``
filters hidden conversations server-side (``false`` — the default view —
excludes them; ``all`` includes them). ``full=1`` adds the per-card
``usage``/``byToolSub`` blocks the analytics dashboard aggregates.
``sessions`` (comma-separated session ids) resolves a specific handful of
conversations instead of a page — what the home page's "waiting for
feedback" cards use to join a pending form/ask back to the session that
posted it, without pulling the whole (~MB) list for one title. An explicit
id lookup ignores the archived filter and the cursor.
``account`` keeps one Claude account's conversations (``totals`` = the
ones counted in totals). Unlike the cost endpoints, a *list* is unfiltered
by default — every conversation is listed; only the sums leave work out."""
cards = list(_conversation_cards(deps, full)) # shared memo — sorted below
if sessions:
want = {s.strip() for s in sessions.split(",") if s.strip()}
cards = [c for c in cards if c.get("sessionId") in want]
return {"conversations": cards, "count": len(cards),
"nextBefore": None,
"pricing": {"model": MODEL,
"inputPerMTok": INPUT_PRICE_PER_MTOK}}
if archived != "all":
cards = [c for c in cards if not c["meta"].get("archived")]
if account:
sel = accounts_mod.selector(account)
cards = [c for c in cards if sel(c["meta"].get("account"))]
cards.sort(key=lambda c: c.get("endedAt") or "", reverse=True)
total = len(cards)
if before:
cards = [c for c in cards if (c.get("endedAt") or "") < before]
next_before = None
if limit and limit > 0 and len(cards) > limit:
cards = cards[:limit]
next_before = cards[-1].get("endedAt")
return {"conversations": cards, "count": total, "nextBefore": next_before,
"pricing": {"model": MODEL, "inputPerMTok": INPUT_PRICE_PER_MTOK}}
@router.get("/api/subagents", responses=_r(schemas.SubagentsResponse))
@_json_in_worker
def subagents_list(deps: State):
"""Every subagent (sidechain) transcript as a flat lite ref carrying its
parent conversation id. The top-level conversation list hides sidechains;
this is the one place they're enumerated — the graph view's node source.
All from the in-process summaries cache (+ cached meta.json reads)."""
out = []
open_by_parent: dict[str, set] = {}
for path, s in deps.store.all_summaries():
if not s.get("isSidechain") or not s.get("messages"):
continue
pp = _parent_path_of(path)
rp = _root_path_of(path)
if not pp or not rp:
continue
# One open-Task lookup per parent, shared by all its children. A nested
# agent's parent is a sidechain whose state comes from *its* parent's
# open Task calls (`_sidechain_state`), not from the shared meta state.
if pp not in open_by_parent:
psum = deps.store.get_summary(pp) or {}
pstate = _sidechain_state(deps, pathlib.Path(pp))
open_by_parent[pp] = {a.get("toolUseId")
for a in _running_agents(deps, psum, pstate)}
mj = _subagent_meta_cached(pathlib.Path(path))
tuid = mj.get("toolUseId")
out.append({
"id": _conv_id(path),
"parentId": _conv_id(pp),
"rootId": _conv_id(rp),
"depth": mj.get("spawnDepth") or (1 if pp == rp else 2),
"agentType": mj.get("agentType"),
"description": mj.get("description"),
"title": s.get("title"),
"tokens": s.get("tokens"),
"cost": s.get("cost"),
"state": "running" if tuid and tuid in open_by_parent[pp]
else "finished",
})
return {"subagents": out}
@router.get("/api/conversation-summary",
responses=_r(schemas.ConversationSummary))
def conversation_summary(deps: State, id: str):
"""One conversation's list card — what an SSE `transcript` ping patches
into the cached list, so a live turn costs a few KB instead of a full
list refetch."""
ap = _resolve_transcript(id)
s = deps.store.get_summary(str(ap))
if not s:
raise HTTPException(404, "summary not indexed yet")
children: list[dict] = []
paths = [str(ap)]
for child in _subagent_descendants(ap): # the whole tree, any depth
paths.append(str(child))
cs = deps.store.get_summary(str(child))
if cs:
children.append(cs)
skills_used = _skills_used(deps, paths)
s = _rollup_children(s, children)
m = _conv_meta(deps, s)
return {"id": id, **{k: s.get(k) for k in _CARD_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_conversation(deps, ).get(id, []),
"meta": m}
@router.get("/api/search-index", responses=_r(schemas.SearchIndexResponse))
def search_index(deps: State):
"""Full-text search corpus: every conversation's user/assistant text blocks.
Shipped once to the client, which runs a fuzzy search (Fuse.js) over it. Each
message carries its ordinal `i`, which matches the `mi` on the detail thread's
text items, so a hit can deep-link to `/conversation/<id>?m=<i>` and scroll to
the exact message."""
out = []
for path, s in deps.store.all_summaries():
if not s.get("messages"):
continue
if s.get("isSidechain"):
continue # subagent messages are reachable via the parent, not listed
msgs = s.get("searchMsgs") or []
if not msgs:
continue
out.append({
"id": _conv_id(path),
"title": _conv_title(deps, s),
"project": s.get("project"),
"model": s.get("model"),
"endedAt": s.get("endedAt"),
# serve-time clip: search + snippets don't need more, and this
# corpus ships to the client in one payload
"msgs": [{**m, "text": (m.get("text") or "")[:500]} for m in msgs],
})
out.sort(key=lambda c: c.get("endedAt") or "", reverse=True)
return {"conversations": out, "count": len(out)}
# Recently parsed full threads, keyed by path and validated on (mtime, size).
# A parse of a large transcript costs hundreds of ms; re-opening a conversation
# (or a meta-only SSE ping refetching the open one) shouldn't pay it again.
# Served as a deep copy — the endpoint decorates the dict in place.
_detail_cache: dict[str, tuple[tuple[float, int], dict]] = {}
_DETAIL_CACHE_CAP = 4
_detail_lock = threading.Lock()
def _parse_full_cached(ap: pathlib.Path) -> dict:
key = str(ap)
try:
st = ap.stat()
sig = (st.st_mtime, st.st_size)
except OSError:
sig = None
with _detail_lock:
hit = _detail_cache.get(key)
if hit and sig and hit[0] == sig:
_detail_cache[key] = _detail_cache.pop(key) # LRU bump
return copy.deepcopy(hit[1])
data = conversations.parse_conversation(ap, full=True)
if sig:
with _detail_lock:
_detail_cache[key] = (sig, copy.deepcopy(data))
while len(_detail_cache) > _DETAIL_CACHE_CAP:
_detail_cache.pop(next(iter(_detail_cache)))
return data
@router.get("/api/conversation", responses=_r(schemas.ConversationDetail))
@_json_in_worker
def conversation_detail(deps: State, id: str):
"""Full parsed thread (tools collapsible client-side) + per-turn usage."""
ap = _resolve_transcript(id)
data = _parse_full_cached(ap)
# The raw `{skill: {count, last}}` map is indexer-internal; the thread only
# needs the names it used, and those roll the subagents' calls in. This
# parse is fresher than the store (a live turn re-parses here first), so the
# conversation's own map comes off `data` and only the children are joined.
kid_paths = [str(c) for c in _subagent_descendants(ap)] # any depth
data["skillsUsed"] = _skills_used(deps, kid_paths, extra=data.get("skills"))
data["agentsUsed"] = _agents_by_conversation(deps, ).get(id, [])
data["models"] = _with_child_models(
data.get("models"), [c for c in map(deps.store.get_summary, kid_paths) if c])
data.pop("skills", None)
data["id"] = id
data["title"] = _conv_title(deps, data)
data["meta"] = _conv_meta(deps, data)
sub_state = _sidechain_state(deps, ap)
if sub_state:
data["meta"]["state"] = sub_state
# Resolve the state first: which Task calls count as running depends on it.
data["runningAgents"] = _running_agents(deps, data, data["meta"].get("state"))
data["deploying"] = _is_deploying(deps, data, data["meta"].get("state"))
_attach_subagents(deps, ap, data)
# Cross-link the memories this conversation read (recalled) and created.
read_slugs = set(data.get("memoriesRead") or [])
data["memoriesRead"] = [m for m in memories_mod.list_memories()
if m["slug"] in read_slugs]
data["memoriesCreated"] = memories_mod.by_origin(data.get("sessionId"))
# Plans this conversation authored (the `plan` skill stamps the planning
# session id / conversationIds into each plan's frontmatter).
sid = data.get("sessionId")
data["plans"] = ([p for p in plans_mod.list_plans()
if sid in plans_mod.linked_session_ids(p)] if sid else [])
# Generated artefacts (screenshots, renders…) the conversation's projects/
# services hold for this session — <dir>/.ai/artefacts/<date>/<sid>/<file>.
data["artefacts"] = artefacts_mod.list_for_conversation(
sid or "", data["meta"].get("projects") or [],
data["meta"].get("services") or [])
return data
# ── conversation scaffold + published summary (see scaffold.py) ─────────────
# The algorithmic half of the `complete` skill: everything about a conversation
# can be aggregated rather than reasoned about, handed to the model as markdown
# with two sections left blank. `conv-scaffold` (scripts/) is the CLI wrapper.
@router.get("/api/conversations/{id:path}/scaffold",
response_class=PlainTextResponse)
def conversation_scaffold(deps: State, id: str):
"""The completion scaffold for one conversation, as markdown."""
sid, _, meta = _summary_and_meta(deps, id)
path = _resolve_transcript(id)
# A full parse: the merged tool sections need the thread, which the cached
# rollup summary deliberately doesn't carry.
full = conversations.parse_conversation(path, full=True)
return scaffold_mod.build(_conv_id(str(path)), sid, full, meta, deps.store)
class SummaryBody(BaseModel):
markdown: str
@router.put("/api/conversations/{id:path}/summary", responses=_r(schemas.OkResponse))
def conversation_summary_put(deps: State, id: str, body: SummaryBody):
"""Publish a filled-in scaffold as the conversation's summary. Rejects an
untouched one — an unedited scaffold is a report of nothing."""
sid, _, _ = _summary_and_meta(deps, id)
md = (body.markdown or "").strip()
if not md:
raise HTTPException(400, "markdown is required")
if not scaffold_mod.is_filled(md):
raise HTTPException(422, "the Summary section is still the scaffold's "
"placeholder — fill it in before publishing")
now = datetime.datetime.now(datetime.timezone.utc) \
.strftime("%Y-%m-%dT%H:%M:%SZ")
deps.meta_store.update(sid, {"summary": {
"markdown": scaffold_mod.clean_for_publish(md), "at": now}})
deps.hub.publish({"type": "meta", "id": sid})
return {"ok": True}
# ── conversation action metadata (edited by skills, streamed to the UI) ──────
class MetaPatch(BaseModel):
id: str # Claude Code session id (or "<proj>/<sid>.jsonl")
# Display title the session chose for itself (`conv-meta title "…"`), used
# instead of the parsed one. Lets an agent spawned with a generic prompt
# rename its card once it knows what it's actually building.
title: str | None = None
projects: list[str] | str | None = None
services: list[str] | str | None = None
state: str | None = None
committed: str | None = None
pushed: str | None = None
merged: str | None = None
deployed: str | None = None
notified: str | None = None
notification: dict | None = None
worktree: dict | None = None # a single worktree create/remove event to fold in
archived: bool | None = None # hide/unhide from the default conversation list
# Pin the conversation's Claude account (accounts.py) — the explicit stamp
# that outranks the import source and the cwd rule.
account: str | None = None
@router.get("/api/conversation-meta", responses=_r(schemas.ConversationMetaResponse))
def conversation_meta(deps: State):
"""The raw action-metadata sidecar (keyed by session id)."""
return {"meta": deps.meta_store.all()}
@router.post("/api/conversation-meta", responses=_r(schemas.MetaUpdateResult))
def update_conversation_meta(deps: State, patch: MetaPatch):
body = patch.model_dump(exclude_none=True)
cid = _meta_key(body.pop("id"))
if not cid:
raise HTTPException(400, "missing id")
if "account" in body and not accounts_mod.valid(body["account"]):
raise HTTPException(400, f"unknown account: {body['account']!r}")
merged = deps.meta_store.update(cid, body)
deps.hub.publish({"type": "meta", "id": cid})
return {"id": cid, "meta": merged}
@router.post("/api/conversations/archive-old",
responses=_r(schemas.ArchiveOldResult))
def archive_old_conversations(deps: State):
"""Archive (hide) every conversation that ended before the start of yesterday.
Bulk cleanup for the drawer: keeps today's and yesterday's transcripts
visible and folds everything older behind the "show archived" toggle."""
# endedAt is a UTC ISO timestamp; compare against a UTC start-of-yesterday so
# the day boundary lines up (lexicographic works at day granularity).
today = datetime.datetime.now(datetime.timezone.utc).date()
cutoff = datetime.datetime.combine(
today - datetime.timedelta(days=1), datetime.time.min).isoformat()
n = 0
for path, s in deps.store.all_summaries():
if not s.get("messages") or s.get("isSidechain"):
continue
ended = s.get("endedAt") or ""
if ended and ended < cutoff:
sid = s.get("sessionId") or _meta_key(_conv_id(path))
if sid and not deps.meta_store.get(sid).get("archived"):
deps.meta_store.update(sid, {"archived": True})
n += 1
deps.hub.publish({"type": "meta"})
return {"archived": n, "cutoff": cutoff}