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

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