Files
ai-agent/backend/cron/service.py
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

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"})