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>
280 lines
14 KiB
Python
280 lines
14 KiB
Python
"""A conversation's resolved state and merged metadata (``_conv_meta``),
|
|
plus the per-account summary views."""
|
|
|
|
import datetime
|
|
import threading
|
|
|
|
import accounts as accounts_mod
|
|
import projects as projects_mod
|
|
from core.config import RUNNING_WINDOW_SECS
|
|
from core.state import AppState
|
|
|
|
|
|
def _auto_state(ended_at: str | None) -> str:
|
|
"""``running`` if the last recorded activity is recent, else ``finished``."""
|
|
if not ended_at:
|
|
return "finished"
|
|
try:
|
|
dt = datetime.datetime.fromisoformat(ended_at.replace("Z", "+00:00"))
|
|
except ValueError:
|
|
return "finished"
|
|
now = datetime.datetime.now(datetime.timezone.utc)
|
|
age = (now - dt).total_seconds()
|
|
return "running" if 0 <= age < RUNNING_WINDOW_SECS else "finished"
|
|
|
|
|
|
def _outlived_exit(exited_at: str | None, ended_at: str | None) -> bool:
|
|
"""True when the transcript's last activity is clearly *after* the recorded
|
|
process exit (>10s slack) — the session kept going in a run the sidecar
|
|
never saw (a terminal or remote-control turn), so the exit-time "finished"
|
|
stamp is stale and must not be trusted."""
|
|
if not exited_at or not ended_at:
|
|
return False
|
|
try:
|
|
exited = datetime.datetime.fromisoformat(exited_at.replace("Z", "+00:00"))
|
|
ended = datetime.datetime.fromisoformat(ended_at.replace("Z", "+00:00"))
|
|
except ValueError:
|
|
return False
|
|
return (ended - exited).total_seconds() > 10
|
|
|
|
|
|
def _resolve_state(summary: dict, m: dict, notified: str | None) -> str:
|
|
"""Decide a conversation's lifecycle state.
|
|
|
|
The ``worktree``/``select-project`` skills stamp ``state=running`` but
|
|
nothing ever clears it, so finished conversations get stuck "running". We
|
|
treat an explicit ``running`` as trustworthy only while activity is recent;
|
|
once a conversation has both sent a notification and ended its last message
|
|
with the ``DONE`` marker (or the `complete` skill's ``COMPLETED`` marker —
|
|
see ``_COMPLETED_RE``; that pass continues the same session, so its own
|
|
reply becomes the transcript's last message and overwrites ``doneMarker``),
|
|
it is finished regardless of the recency window.
|
|
The exit watcher (``_watch_sidecar_runs``) stamps ``finished`` the moment a
|
|
sidecar-launched run's process dies, so the badge doesn't sit on "running"
|
|
for the rest of the recency window when the prompt never said DONE.
|
|
A run whose transcript ends on a ``[Request interrupted by user]`` marker is
|
|
reported as ``interrupted`` (a terminal state, so it wins over recency and
|
|
over a stale ``running`` stamp); one ending on a ``No response requested.``
|
|
reply is reported as ``paused``. A run that parked itself waiting on a
|
|
trigger (a trailing ``ScheduleWakeup``/``Monitor`` — see the parser's
|
|
``waiting_trailing``) is ``paused`` too, instead of sitting on a "running"
|
|
badge while nothing happens. Other explicit states (e.g. ``review``) are
|
|
always honoured.
|
|
"""
|
|
auto = _auto_state(summary.get("endedAt")) # running if recent, else finished
|
|
# A trailing interrupt marker in the transcript is authoritative: the run
|
|
# was stopped by the user (Esc, or the interrupt button in this viewer).
|
|
if summary.get("interruptedByUser"):
|
|
return "interrupted"
|
|
# A trailing "No response requested." reply (a queued/resume prompt that had
|
|
# nothing to answer) leaves the run paused rather than finished.
|
|
if summary.get("pausedByUser"):
|
|
return "paused"
|
|
if (summary.get("doneMarker") or summary.get("completedMarker")) and notified:
|
|
return "finished"
|
|
explicit = m.get("state")
|
|
# An exit-time "finished" is trusted only while the transcript hasn't
|
|
# moved past that exit; later activity means the session continued outside
|
|
# the sidecar — fall back to the recency heuristic instead.
|
|
if explicit == "finished" and _outlived_exit(m.get("exitedAt"),
|
|
summary.get("endedAt")):
|
|
explicit = None
|
|
if explicit in (None, "running"):
|
|
# Claude scheduled a wake-up / blocked on a Monitor and the turn ended:
|
|
# the run is idle until the trigger fires, however recent its activity.
|
|
if summary.get("waitingForTrigger"):
|
|
return "paused"
|
|
return auto # stale/absent "running" → fall back to recency
|
|
return explicit
|
|
|
|
|
|
def _merge_worktrees(*sources) -> list[dict]:
|
|
"""Union worktree records from several sources, keyed by directory basename.
|
|
|
|
Each source is a list of ``{name, dir, createdAt, removedAt}`` (from the
|
|
transcript parse and/or the live sidecar). Same ``dir`` → one entry: earliest
|
|
``createdAt`` and latest ``removedAt`` win, and a real short ``name`` beats a
|
|
fallback that equals the dir. Ordered by creation (unknown-created last)."""
|
|
by_dir: dict[str, dict] = {}
|
|
for src in sources:
|
|
for e in src or []:
|
|
d = (e.get("dir") or e.get("name") or "").strip()
|
|
if not d:
|
|
continue
|
|
cur = by_dir.get(d)
|
|
if cur is None:
|
|
by_dir[d] = {"name": e.get("name") or d, "dir": d,
|
|
"createdAt": e.get("createdAt"),
|
|
"removedAt": e.get("removedAt")}
|
|
continue
|
|
nm = e.get("name")
|
|
if nm and nm != d and (not cur["name"] or cur["name"] == d):
|
|
cur["name"] = nm
|
|
for k, pick in (("createdAt", min), ("removedAt", max)):
|
|
vals = [v for v in (cur.get(k), e.get(k)) if v]
|
|
cur[k] = pick(vals) if vals else None
|
|
return sorted(by_dir.values(), key=lambda w: (w.get("createdAt") is None,
|
|
w.get("createdAt") or ""))
|
|
|
|
|
|
def _conv_title(deps: AppState, summary: dict) -> str | None:
|
|
"""The conversation's display title.
|
|
|
|
A session can rename itself (``conv-meta title "…"``) once it knows what it
|
|
is actually doing — the goal-keeper does this after picking its item, so its
|
|
card reads "Add a Work button to the Goals page" rather than the generic
|
|
spawn prompt it was launched with. That sidecar title wins over the parsed
|
|
one (Claude's ``ai-title``, else the first user message)."""
|
|
sid = summary.get("sessionId") or ""
|
|
return deps.meta_store.peek(sid).get("title") or summary.get("title")
|
|
|
|
|
|
def _conv_meta(deps: AppState, summary: dict) -> dict:
|
|
"""Merge a conversation's parsed auto fields with its action sidecar."""
|
|
sid = summary.get("sessionId") or ""
|
|
m = deps.meta_store.peek(sid) # read-only — never mutated below
|
|
projects = list(summary.get("projectsAuto") or [])
|
|
for p in (m.get("projects") or []):
|
|
if p not in projects:
|
|
projects.append(p)
|
|
# No project at all ⇒ the conversation worked on the lab itself (a service,
|
|
# a skill, a script, a doc, a question). Attribute it to the root project so
|
|
# it stops falling out of every project-shaped view. Derived, never stored.
|
|
account = accounts_mod.resolve(summary, m)
|
|
projects = projects_mod.attributed_projects(
|
|
projects, _account_home_project(account))
|
|
# Homelab services touched by the conversation (mirrors ``projects``).
|
|
services = list(summary.get("servicesAuto") or [])
|
|
for s in (m.get("services") or []):
|
|
if s not in services:
|
|
services.append(s)
|
|
# Lifecycle stamps: the live sidecar (written by the skills) is authoritative;
|
|
# where it's empty we backfill from evidence mined out of the transcript.
|
|
auto = summary.get("lifecycleAuto") or {}
|
|
stamp = lambda k: m.get(k) or auto.get(k) # noqa: E731
|
|
notified = stamp("notified")
|
|
return {
|
|
"state": _resolve_state(summary, m, notified),
|
|
# The custom sidecar title override, when set (the resolved top-level
|
|
# `title` already prefers it) — lets the UI tell a renamed card from an
|
|
# auto-titled one and offer "reset to auto title" in the rename dialog.
|
|
"title": m.get("title") or None,
|
|
"projects": projects,
|
|
"services": services,
|
|
# Git worktrees the conversation created/removed (with created/removed
|
|
# timestamps), merging the transcript parse with anything the worktree
|
|
# skill recorded live in the sidecar.
|
|
"worktrees": _merge_worktrees(summary.get("worktreesAuto"),
|
|
m.get("worktrees")),
|
|
"committed": stamp("committed"),
|
|
"pushed": stamp("pushed"),
|
|
"merged": stamp("merged"),
|
|
"deployed": stamp("deployed"),
|
|
"notified": notified,
|
|
"notifications": m.get("notifications") or [],
|
|
# Hidden-from-list flag, toggled from the conversation card context menu.
|
|
"archived": bool(m.get("archived")),
|
|
# Which agent harness runs this session ("claude"/"pi"). Stamped by
|
|
# /api/spawn; conversations predating the field (or spawned from a
|
|
# terminal) fall back to "claude" client-side.
|
|
"harness": m.get("harness"),
|
|
# The thinking / `--effort` setting the session was spawned or last
|
|
# resumed with — what the resume composer's chips seed from. Declared on
|
|
# the schema all along but never emitted, so they always read "default".
|
|
"thinking": m.get("thinking"),
|
|
"effort": m.get("effort"),
|
|
# Which Claude account the session ran on / bills ("personal"/"work").
|
|
# Always resolved (accounts.resolve: spawn stamp → import source → cwd
|
|
# rule → default), so a terminal session is tagged too.
|
|
"account": account,
|
|
# The paired worker (workers.py) the session runs on — launched there
|
|
# by the hub (``worker``) or mirrored from its terminal
|
|
# (``workerSource``). None = the host sidecar / a lab terminal.
|
|
"worker": m.get("worker") or m.get("workerSource"),
|
|
# Source conversation + message ordinal when this one is a fork.
|
|
"forkedFrom": m.get("forkedFrom"),
|
|
# The cron job that spawned this session ({id, name}), if any — the
|
|
# viewer tags the conversation and links back to Settings → Cron.
|
|
"cron": m.get("cron"),
|
|
# The agent definition this session was launched from ({name, at}) by
|
|
# the agent page's "Run now" — the manual sibling of `cron`, and what
|
|
# puts the run in that agent's history.
|
|
"agentRun": m.get("agentRun"),
|
|
# When the `complete` skill's end-of-task pass ran on this session:
|
|
# its reply signs off with COMPLETED instead of DONE (the parser's
|
|
# ``completedMarker``), so the marker on the last message *is* the
|
|
# record — no sidecar stamp to keep in sync. The viewer turns it into
|
|
# the COMPLETED seal on the card.
|
|
"completedAt": (summary.get("endedAt")
|
|
if summary.get("completedMarker") else None),
|
|
# The conversation's published summary, written by that same pass
|
|
# ({markdown, at}) and rendered under the thread (see scaffold.py).
|
|
"summary": m.get("summary") or None,
|
|
}
|
|
|
|
|
|
# ── accounts: which Claude subscription each run bills (accounts.py) ─────────
|
|
# Every aggregate the dashboards read (activity, projects, skills, agents) is
|
|
# split by account, and the default ("totals") view leaves out accounts flagged
|
|
# ``excludeFromTotals`` — the work seat isn't the owner's money. Resolving an
|
|
# account is a meta peek + a cwd split per summary; memoized on the (summaries,
|
|
# sidecar) versions since both feed it.
|
|
_acct_cache: tuple[tuple[int, int], list] | None = None
|
|
_acct_lock = threading.Lock()
|
|
|
|
|
|
def _account_of(deps: AppState, summary: dict) -> str:
|
|
"""A run's account. A subagent shares its parent's sessionId (so its
|
|
metadata) and cwd, so it resolves to the parent's account."""
|
|
return accounts_mod.resolve(
|
|
summary, deps.meta_store.peek(summary.get("sessionId") or ""))
|
|
|
|
|
|
def _account_home_project(account: str | None) -> str | None:
|
|
"""The project a project-less conversation of ``account`` is attributed to
|
|
(its first repo) — None for the default account, which keeps the root."""
|
|
a = accounts_mod.get(account)
|
|
return (a.get("projects") or [None])[0] if a and not a.get("default") else None
|
|
|
|
|
|
def _account_project_fallback(deps: AppState, summary: dict) -> str | None:
|
|
return _account_home_project(_account_of(deps, summary))
|
|
|
|
|
|
def _summaries_with_account(deps: AppState) -> list[tuple[str, dict, str]]:
|
|
"""``store.all_summaries()`` with each run's account attached (same order)."""
|
|
global _acct_cache
|
|
key = (deps.store.summaries_version, deps.meta_store.version)
|
|
with _acct_lock:
|
|
if _acct_cache and _acct_cache[0] == key:
|
|
return _acct_cache[1]
|
|
out = [(p, s, _account_of(deps, s)) for p, s in deps.store.all_summaries()]
|
|
with _acct_lock:
|
|
_acct_cache = (key, out)
|
|
return out
|
|
|
|
|
|
def _summaries_for(deps: AppState, account: str | None) -> list[tuple[str, dict]]:
|
|
"""The summaries an ``?account=`` view covers (see ``accounts.selector``:
|
|
unset = the aggregated totals, i.e. without the excluded accounts)."""
|
|
sel = accounts_mod.selector(account)
|
|
return [(p, s) for p, s, a in _summaries_with_account(deps, ) if sel(a)]
|
|
|
|
|
|
def _by_account(rows) -> dict[str, dict]:
|
|
"""``{account: {cost, tokens, conversations}}`` over ``(path, summary,
|
|
account)`` rows — every run's spend (subagents included, as separate
|
|
summaries), conversations counted top-level only. Every profile appears,
|
|
zeros included, so the UI never has to guess whether a key is missing."""
|
|
out = {a: {"cost": 0.0, "tokens": 0, "conversations": 0}
|
|
for a in accounts_mod.ids()}
|
|
for _p, s, a in rows:
|
|
if not s.get("messages"):
|
|
continue
|
|
e = out.setdefault(a, {"cost": 0.0, "tokens": 0, "conversations": 0})
|
|
e["cost"] += float(s.get("cost") or 0.0)
|
|
e["tokens"] += int(s.get("tokens") or 0)
|
|
if not s.get("isSidechain"):
|
|
e["conversations"] += 1
|
|
return out
|