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>
181 lines
6.6 KiB
Python
181 lines
6.6 KiB
Python
"""The application state: every store, hub, watcher and mirror, built once
|
|
by :func:`build_state` and attached to ``app.state.ai``.
|
|
|
|
Routes take it as a FastAPI dependency (``deps: State``) and the helpers they
|
|
call take it explicitly (``deps: AppState``) — nothing reads a module-level
|
|
singleton, so a test can build a throwaway state and the import order of the
|
|
domain packages never matters. The domain functions this wiring needs
|
|
(creating a notification, firing a cron job) are imported *inside*
|
|
``build_state`` on purpose: those modules import ``AppState`` from here.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import dataclasses
|
|
import os
|
|
import pathlib
|
|
import typing
|
|
|
|
from fastapi import Depends, Request
|
|
|
|
from core.config import (
|
|
ACCOUNT_SOURCE_DIRS,
|
|
ARCHIVE_DIR,
|
|
CRON_PATH,
|
|
DB_PATH,
|
|
FORMS_PATH,
|
|
MAIL_TRIGGER_SENDERS_PATH,
|
|
META_PATH,
|
|
NOTIFY_PATH,
|
|
SOURCE_DIRS,
|
|
UI_STATE_PATH,
|
|
WORKERS_PATH,
|
|
WORKERS_SOURCE_DIR,
|
|
WORKSPACE,
|
|
)
|
|
from cron.scheduler import CronScheduler, CronStore
|
|
from db import Store
|
|
from deploy_status import DeployWatcher
|
|
from events import Hub, Watcher
|
|
from forms.store import FormStore
|
|
from indexer import Indexer
|
|
from meta import MetaStore
|
|
from notifications.read_ledger import NotifReadStore
|
|
from notifications.store import NotifyStore
|
|
from settings.mail_trigger import MailTriggerStore
|
|
from settings.ui_state import UiStateStore
|
|
from workers.outbox_mirror import OutboxMirror
|
|
from workers.remotefeed import FeedMirror
|
|
from workers.store import WorkerMonitor, WorkerStore
|
|
|
|
# Server-side "read" ledger for the notification feed — lets the phone fetch
|
|
# *unread* notifications (listening ≠ reading; only the UI marks them read).
|
|
NOTIF_READ_PATH = os.environ.get("NOTIF_READ_PATH", "/data/notif-read.json")
|
|
|
|
|
|
@dataclasses.dataclass
|
|
class AppState:
|
|
store: Store
|
|
meta_store: MetaStore
|
|
ui_state: UiStateStore
|
|
cron_store: CronStore
|
|
notify_store: NotifyStore
|
|
form_store: FormStore
|
|
mail_trigger_store: MailTriggerStore
|
|
notif_read_store: NotifReadStore
|
|
hub: Hub
|
|
worker_store: WorkerStore
|
|
worker_monitor: WorkerMonitor
|
|
indexer: Indexer
|
|
watcher: Watcher
|
|
feed_mirror: FeedMirror | None
|
|
outbox_mirror: OutboxMirror
|
|
deploy_watcher: DeployWatcher
|
|
cron_scheduler: CronScheduler
|
|
|
|
|
|
def _record_source_account(deps: "AppState", rel: str, account: str) -> None:
|
|
"""Stamp ``meta.accountSource`` on the session a transcript imported from
|
|
``account``'s source dir belongs to (``<proj>/<sid>.jsonl``, or a subagent
|
|
under ``<proj>/<sid>/subagents/``). Once per session: a no-op after."""
|
|
parts = pathlib.PurePosixPath(rel).parts
|
|
if len(parts) < 2:
|
|
return
|
|
sid = parts[1][:-6] if parts[1].endswith(".jsonl") else parts[1]
|
|
if deps.meta_store.peek(sid).get("accountSource") == account:
|
|
return
|
|
deps.meta_store.update(sid, {"accountSource": account})
|
|
|
|
|
|
def _record_worker_files(deps: "AppState", w: dict, rels: list[str]) -> None:
|
|
"""The mirror landed these transcripts from worker ``w``: stamp which
|
|
worker each session lives on (resume/fork/interrupt route on it) and the
|
|
worker's account (``accountSource`` — the import-source rule in
|
|
accounts.resolve). One meta save per pass, however many sessions."""
|
|
patches: dict[str, dict] = {}
|
|
for rel in rels:
|
|
parts = pathlib.PurePosixPath(rel).parts
|
|
if len(parts) < 2:
|
|
continue
|
|
sid = parts[1][:-6] if parts[1].endswith(".jsonl") else parts[1]
|
|
patches[sid] = {"workerSource": w["id"],
|
|
"accountSource": w.get("account")}
|
|
if patches and deps.meta_store.stamp_many(patches):
|
|
deps.hub.publish({"type": "meta"})
|
|
|
|
|
|
def build_state() -> AppState:
|
|
"""Construct every store and background worker (nothing is started here —
|
|
see ``main._startup``)."""
|
|
from files.discovery import _iter_files, _read_text
|
|
|
|
store = Store(DB_PATH)
|
|
meta_store = MetaStore(META_PATH)
|
|
hub = Hub()
|
|
worker_store = WorkerStore(WORKERS_PATH)
|
|
|
|
|
|
notify_store = NotifyStore(NOTIFY_PATH)
|
|
form_store = FormStore(FORMS_PATH)
|
|
cron_store = CronStore(CRON_PATH)
|
|
|
|
deps = AppState(
|
|
store=store,
|
|
meta_store=meta_store,
|
|
ui_state=UiStateStore(UI_STATE_PATH),
|
|
cron_store=cron_store,
|
|
notify_store=notify_store,
|
|
form_store=form_store,
|
|
mail_trigger_store=MailTriggerStore(MAIL_TRIGGER_SENDERS_PATH),
|
|
notif_read_store=NotifReadStore(NOTIF_READ_PATH),
|
|
hub=hub,
|
|
worker_store=worker_store,
|
|
worker_monitor=WorkerMonitor(worker_store, hub.publish),
|
|
indexer=typing.cast(Indexer, None),
|
|
watcher=typing.cast(Watcher, None),
|
|
feed_mirror=None,
|
|
outbox_mirror=typing.cast(OutboxMirror, None),
|
|
deploy_watcher=DeployWatcher(hub),
|
|
cron_scheduler=typing.cast(CronScheduler, None),
|
|
)
|
|
|
|
deps.indexer = Indexer(store, WORKSPACE, SOURCE_DIRS, ARCHIVE_DIR,
|
|
_iter_files, _read_text,
|
|
source_accounts={d: a for a, d in ACCOUNT_SOURCE_DIRS.items()},
|
|
on_account=lambda rel, a: _record_source_account(deps, rel, a))
|
|
deps.watcher = Watcher(SOURCE_DIRS, meta_store, deps.indexer, hub)
|
|
if WORKERS_SOURCE_DIR:
|
|
deps.feed_mirror = FeedMirror(
|
|
worker_store, WORKERS_SOURCE_DIR,
|
|
on_files=lambda w, rels: _record_worker_files(deps, w, rels))
|
|
|
|
# The two workers below call back into domain code with the finished state.
|
|
from cron.service import _scheduler_fire
|
|
from forms.service import _create_form
|
|
from notifications.service import _create_ask, _create_notify
|
|
|
|
# Workers' outboxes: the notifications / asks / forms a session on a remote
|
|
# worker queued in *its* sidecar (it never calls the hub). Pulled here, created
|
|
# through the same functions the routes use, answers pushed back.
|
|
deps.outbox_mirror = OutboxMirror(
|
|
worker_store,
|
|
create={"notify": lambda d: _create_notify(deps, d),
|
|
"ask": lambda d: _create_ask(deps, d),
|
|
"form": lambda d: _create_form(deps, d)},
|
|
status={"ntf": lambda i: notify_store.get(i),
|
|
"ask": lambda i: notify_store.ask_status(i),
|
|
"frm": lambda i: form_store.get(i)},
|
|
cancel={"frm": lambda i: form_store.cancel(i)},
|
|
publish=hub.publish)
|
|
deps.cron_scheduler = CronScheduler(
|
|
cron_store, lambda job: _scheduler_fire(deps, job))
|
|
return deps
|
|
|
|
|
|
def get_state(request: Request) -> AppState:
|
|
return request.app.state.ai
|
|
|
|
|
|
# The dependency every route declares: ``def route(deps: State, ...)``.
|
|
State = typing.Annotated[AppState, Depends(get_state)]
|