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>
378 lines
16 KiB
Python
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}
|