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>
82 lines
3.4 KiB
Python
82 lines
3.4 KiB
Python
"""Firing a cron job: prompt file → spawn → history + conversation tag."""
|
|
|
|
import pathlib
|
|
|
|
from fastapi import HTTPException
|
|
|
|
import cron as cron_mod
|
|
from conversations.catalog import _session_conv_map
|
|
from core.config import WORKSPACE
|
|
from core.state import AppState
|
|
from runs.service import LOCAL_RUNNER, _send_message
|
|
from runs.sidecar_client import _detail_text
|
|
|
|
|
|
def _cron_job_out(deps: AppState, job: dict) -> dict:
|
|
"""A stored job + the computed fields the UI shows: the next fire time, the
|
|
run config read out of its prompt file's frontmatter, and the run history
|
|
enriched with conversation ids (newest first)."""
|
|
out = dict(job)
|
|
out["nextRun"] = (cron_mod.next_run(job.get("schedule") or "")
|
|
if job.get("enabled") else None)
|
|
out["runConfig"] = cron_mod.run_config(_cron_prompt_text(job))
|
|
sid2conv = _session_conv_map(deps, )
|
|
out["runs"] = [{"at": at, "sessionId": sid,
|
|
"conversationId": sid2conv.get(sid)}
|
|
for at, sid in sorted((job.get("history") or {}).items(),
|
|
reverse=True)]
|
|
return out
|
|
|
|
|
|
def _cron_prompt_path(job: dict) -> tuple[pathlib.Path, str]:
|
|
rel = cron_mod.valid_prompt_file(job.get("promptFile") or "")
|
|
return (WORKSPACE / rel), rel
|
|
|
|
|
|
def _cron_prompt_text(job: dict) -> str:
|
|
"""The job's prompt file verbatim (frontmatter included), or ``""`` if it
|
|
isn't readable — a missing file is reported when the job fires, not here."""
|
|
try:
|
|
ap, _ = _cron_prompt_path(job)
|
|
return ap.read_text(encoding="utf-8")
|
|
except (OSError, ValueError):
|
|
return ""
|
|
|
|
|
|
def _fire_cron_job(deps: AppState, job: dict) -> dict:
|
|
"""Run a cron job once: read its prompt file, spawn a session on the
|
|
harness/model/parameters its frontmatter declares — on the job's
|
|
``worker`` (unset ⇒ the lab's own sidecar, never auto-routed) — record the
|
|
run in the job's history and tag the conversation with the job."""
|
|
_, rel = _cron_prompt_path(job)
|
|
raw = _cron_prompt_text(job)
|
|
cfg = cron_mod.run_config(raw)
|
|
prompt = cron_mod.strip_frontmatter(raw).strip()
|
|
if not prompt:
|
|
deps.cron_store.record_run(job["id"], "",
|
|
status=f"error: prompt file {rel} missing/empty")
|
|
deps.hub.publish({"type": "cron"})
|
|
raise HTTPException(409, f"prompt file {rel} is missing or empty")
|
|
# Frontmatter wins; the store's `model` is the pre-frontmatter fallback.
|
|
out = _send_message(deps, prompt, model=cfg["model"] or job.get("model"),
|
|
harness=cfg["harness"], thinking=cfg["thinking"],
|
|
effort=cfg["effort"], account=cfg.get("account"),
|
|
worker=job.get("worker") or LOCAL_RUNNER)
|
|
sid = out["sessionId"]
|
|
deps.meta_store.update(sid, {"cron": {"id": job["id"],
|
|
"name": job.get("name")}})
|
|
deps.cron_store.record_run(job["id"], sid, status="ok")
|
|
deps.hub.publish({"type": "cron"})
|
|
deps.hub.publish({"type": "meta", "id": sid})
|
|
return out
|
|
|
|
|
|
def _scheduler_fire(deps: AppState, job: dict) -> None:
|
|
"""The scheduler-thread wrapper: outcomes land in the store (surfaced as
|
|
the job's lastStatus), never raise into the loop."""
|
|
try:
|
|
_fire_cron_job(deps, job)
|
|
except HTTPException as e:
|
|
deps.cron_store.record_run(job["id"], "", status=f"error: {_detail_text(e)}")
|
|
deps.hub.publish({"type": "cron"})
|