diff --git a/AGENTS.md b/AGENTS.md index 4878e3a..3fc0fb3 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -71,6 +71,14 @@ so authorization, grants and scope checks are one code path. `agentdrive.scripts.apply_schema`: an empty database takes `schema.sql` alone; anything else replays pending migrations. Never edit a shipped migration; fold a new one's end state into `schema.sql` in the same change. +- **Maintenance jobs** — `jobs/`: garbage collection (`jobs/gc.py`, semantics + in `core/gc.py`) and usage maintenance (`jobs/usage_snapshot.py`), on the + cadence in `jobs/schedule.py`. A self-hosted install runs them with + `python -m agentdrive.jobs.scheduler`, which the process supervisor starts + inside the API container when `SCHEDULER_ENABLED=true` (as + `compose.selfhost.yml` sets it) — never as a separate container, so the jobs + that delete content cannot be configured apart from the API. Without it, + deleted content is never reclaimed. - **Formal model** — `specs/tla/ArtifactGC.tla`, the garbage collector's safety argument; its TLC configs run in CI. diff --git a/Makefile b/Makefile index af498bd..878bcfd 100644 --- a/Makefile +++ b/Makefile @@ -2,10 +2,12 @@ # # Everything here runs against your own checkout or install: the local stack, # the generated stylesheet, the model checker, the OpenAPI policy, the image, -# and the maintenance jobs a deployment schedules. How you schedule those jobs -# (cron, a Kubernetes CronJob, your cloud's scheduler) is your deployment's -# business; each one is `python -m agentdrive.jobs.` against the same -# DATABASE_URL and store the API uses. +# and the maintenance jobs a deployment must schedule. compose.selfhost.yml +# runs them on their cadence inside the API container +# (SCHEDULER_ENABLED=true starts `python -m agentdrive.jobs.scheduler` +# there); anything else runs the commands +# `python -m agentdrive.jobs.scheduler --list` prints, on that cadence, +# against the same settings as the API. The targets below run one job now. SHELL := /usr/bin/env bash diff --git a/README.md b/README.md index 9abc99a..6dcdef3 100644 --- a/README.md +++ b/README.md @@ -53,6 +53,27 @@ reads: - `COMPOSE_PROJECT_NAME` for a second install on the same machine; each project gets its own containers and volumes. +**Maintenance runs itself.** The API container also runs the maintenance +jobs on the hosted product's schedule (UTC): garbage collection of abandoned +uploads hourly, the full sweep — purging expired deletes and unreferenced +content — daily at 03:00, the same plus an orphan sweep on Sundays at 04:00, +and usage maintenance every 15 minutes. Without them, deleted content is +never reclaimed and storage only grows. The scheduler remembers the last +minute it handled, so a job that came due while the stack was down runs once +when it is back. See the schedule and its state, or run one job now: + +```bash +docker compose -f compose.selfhost.yml exec api python -m agentdrive.jobs.scheduler --list +docker compose -f compose.selfhost.yml exec api python -m agentdrive.jobs.scheduler --status +docker compose -f compose.selfhost.yml exec api python -m agentdrive.jobs.scheduler --run gc-daily +``` + +A job that is already running holds its lock, so a second `--run` of it +reports `"skipped": true` and does nothing. Not using this compose file? Set +`SCHEDULER_ENABLED=true` (and `SCHEDULER_STATE_FILE` to a writable path) on +the API container, or run the commands `--list` prints, on the same +schedule, from your own scheduler with the API's settings. + Data lives in the `agentdrive_pgdata` and `agentdrive_data` volumes. **Upgrade** with `git pull && docker compose -f compose.selfhost.yml up -d --build` (without `--build`, `up` keeps running the old image). `docker compose -f diff --git a/compose.selfhost.yml b/compose.selfhost.yml index d7024e8..c0f9dfc 100644 --- a/compose.selfhost.yml +++ b/compose.selfhost.yml @@ -5,8 +5,10 @@ # from this file). Three services: `postgres` (plain postgres:16 — schema.sql # uses only built-in tsvector), `migrate` (the image running apply_schema # once, after Postgres is healthy), and `api` (the image under -# AUTH_MODE=local with the filesystem object store on a named volume and the -# MCP transport served at /mcp by the sidecar inside the same container). +# AUTH_MODE=local with the filesystem object store on a named volume, the +# MCP transport served at /mcp by the sidecar inside the same container, and +# the job scheduler running garbage collection and usage maintenance there +# too — without it deleted content is never reclaimed). # # Settings come from a `.env` beside this file, which every `docker compose` # command reads: @@ -80,6 +82,12 @@ services: # loopback; the supervisor puts the sidecar into api-key mode. MCP_PROXY_URL: http://127.0.0.1:8081 PORT: "8080" + # Run the maintenance jobs on their schedule (src/agentdrive/jobs/ + # schedule.py) inside this container, so they can never be pointed at + # a different database or store than the API. The last handled minute + # persists on the data volume; a restart catches up what came due. + SCHEDULER_ENABLED: "true" + SCHEDULER_STATE_FILE: /data/.agentdrive-scheduler.json ports: - "${AGENTDRIVE_BIND:-127.0.0.1}:${AGENTDRIVE_PORT:-8080}:8080" volumes: @@ -94,8 +102,12 @@ services: timeout: 5s retries: 30 # The supervisor gives its children 10 s to exit; Docker's default grace - # period is also 10 s, so give the supervisor room to finish first. + # period is also 10 s, so give the supervisor room to finish first. A + # running job is terminated with its process group; its advisory lock is + # released with its database connection and the next run resumes. stop_grace_period: 20s + # A real PID 1 that reaps orphaned processes. + init: true restart: unless-stopped volumes: diff --git a/scripts/selfhost-smoke.sh b/scripts/selfhost-smoke.sh index 0a5adc8..0207c33 100755 --- a/scripts/selfhost-smoke.sh +++ b/scripts/selfhost-smoke.sh @@ -168,6 +168,43 @@ else fi rm -f "$AFTER_BODY" +echo "== the maintenance jobs run on this install" +# Without the scheduler, deleted content is never reclaimed. It runs inside +# the api container; its loop records the minute it last handled, so a state +# file that appears proves the loop itself is alive (it ticks every minute). +LOOP_ALIVE="" +for _ in $(seq 1 45); do + if $COMPOSE exec -T api python -m agentdrive.jobs.scheduler --status > /dev/null; then LOOP_ALIVE=1; break; fi + sleep 2 +done +if [ -n "$LOOP_ALIVE" ]; then ok "the scheduler loop is ticking"; else bad "the scheduler loop never recorded a tick"; fi +# Then each scheduled job runs once, now, against this install's database and +# store — the rows the steps above created included; the sweeps' age rules +# mean nothing that young is deleted, which the test suite proves on the same +# store. A run that met the loop's own run of the same job reports +# `"skipped": true` (lock held); that is retried, never counted as a pass. +JOBS=$($COMPOSE exec -T api python -m agentdrive.jobs.scheduler --list | grep -v '^ ' | awk '{print $1}') +RUN_OUT=$(mktemp) +for name in gc-hourly gc-daily gc-weekly usage-snapshot; do + if ! echo "$JOBS" | grep -qx "$name"; then bad "--list does not show $name"; continue; fi + result="" + for _ in 1 2 3; do + if $COMPOSE exec -T api python -m agentdrive.jobs.scheduler --run "$name" > "$RUN_OUT" 2>&1; then + if grep -q '"skipped": true' "$RUN_OUT"; then result="skipped"; sleep 5; continue; fi + if grep -q "Traceback" "$RUN_OUT"; then result="traceback"; break; fi + result="ok"; break + fi + result="failed"; break + done + if [ "$result" = "ok" ]; then + ok "scheduled job $name exits 0" + else + bad "scheduled job $name: $result" + sed 's/^/ /' "$RUN_OUT" | tail -20 + fi +done +rm -f "$RUN_OUT" + echo "== the container logs for THIS run are clean" LOGS=$($COMPOSE logs --no-color --since "$START" api 2>/dev/null) if echo "$LOGS" | grep -q "Traceback"; then bad "api log has a traceback"; else ok "api log has no traceback"; fi diff --git a/src/agentdrive/jobs/gc.py b/src/agentdrive/jobs/gc.py index 5225d9d..5a161fe 100644 --- a/src/agentdrive/jobs/gc.py +++ b/src/agentdrive/jobs/gc.py @@ -1,11 +1,13 @@ -"""Cloud Run Job entrypoint for the GC sweeper. +"""Job entrypoint for the GC sweeper. -Thin wrapper around ``GCSweeper(...).run()``. Cloud Scheduler invokes +Thin wrapper around ``GCSweeper(...).run()``, run as ``python -m agentdrive.jobs.gc`` hourly with ``--sessions-only``, daily at 03:00 UTC for the full sweep (session reconciliation + transfer cleanup + purge + CAS mark-sweep + scratch sweep), and weekly on Sunday 04:00 UTC with -``--orphan-sweep`` appended. Operator runbook lives in ``deploy/README.md``; -the sweep semantics live in ``agentdrive.core.gc``. +``--orphan-sweep`` appended. That cadence is ``agentdrive.jobs.schedule``: +a self-hosted install runs it with ``python -m agentdrive.jobs.scheduler``, +a hosted deployment with its platform's scheduler. The sweep semantics live +in ``agentdrive.core.gc``. CLI (the scheduled contract — do not change argument meanings): --dry-run Preview without persisting: all database work runs inside diff --git a/src/agentdrive/jobs/schedule.py b/src/agentdrive/jobs/schedule.py new file mode 100644 index 0000000..5aaba99 --- /dev/null +++ b/src/agentdrive/jobs/schedule.py @@ -0,0 +1,72 @@ +"""When the maintenance jobs run: the one record of the cadence. + +A deployment must run these, or it never reclaims storage: deleted artifacts +stay in the database and their bytes stay in the store, abandoned uploads keep +their reserved bytes, and usage is never finalized. A self-hosted install runs +them with `python -m agentdrive.jobs.scheduler`, which the process +supervisor starts inside the API container when `SCHEDULER_ENABLED=true` +(`compose.selfhost.yml` sets it); a deployment with its own scheduler (cron, a +Kubernetes CronJob, a cloud scheduler) runs the same commands on the same +cadence, listed by `python -m agentdrive.jobs.scheduler --list`. + +Times are UTC, in five-field cron syntax. Each job is +`python -m ` against the same settings as the API. +""" + +from __future__ import annotations + +from dataclasses import dataclass + + +@dataclass(frozen=True) +class ScheduledJob: + name: str + cron: str + module: str + args: tuple[str, ...] + # The scheduler kills a run that outlives this. The GC stops itself at + # `GCSweeper.HARD_TIMEOUT_S` (50 min), so its kill comes after that. + timeout_s: int + description: str = "" + + +SCHEDULE: tuple[ScheduledJob, ...] = ( + ScheduledJob( + name="gc-hourly", + cron="0 * * * *", + module="agentdrive.jobs.gc", + args=("--sessions-only",), + timeout_s=3600, + description=( + "Session phases only: terminalize abandoned uploads and release " + "their reserved bytes. No object-store listing." + ), + ), + ScheduledJob( + name="gc-daily", + cron="0 3 * * *", + module="agentdrive.jobs.gc", + args=(), + timeout_s=3600, + description=( + "Full sweep: purge expired soft-deletes, mark-sweep unreferenced " + "content blobs, clean up scratch." + ), + ), + ScheduledJob( + name="gc-weekly", + cron="0 4 * * 0", + module="agentdrive.jobs.gc", + args=("--orphan-sweep",), + timeout_s=3600, + description="The full sweep plus the orphan sweep over purged drives' prefixes.", + ), + ScheduledJob( + name="usage-snapshot", + cron="*/15 * * * *", + module="agentdrive.jobs.usage_snapshot", + args=("--no-notify",), + timeout_s=900, + description="Finalize usage reservations and prune bounded usage history.", + ), +) diff --git a/src/agentdrive/jobs/scheduler.py b/src/agentdrive/jobs/scheduler.py new file mode 100644 index 0000000..fea8d50 --- /dev/null +++ b/src/agentdrive/jobs/scheduler.py @@ -0,0 +1,403 @@ +"""Run the maintenance jobs on their schedule, for a self-hosted install. + + python -m agentdrive.jobs.scheduler # run forever + python -m agentdrive.jobs.scheduler --list # print the schedule + python -m agentdrive.jobs.scheduler --run NAME # run one job now; exit with its status + python -m agentdrive.jobs.scheduler --status # print the persisted state + +In `compose.selfhost.yml` this runs INSIDE the API container, started by the +process supervisor when `SCHEDULER_ENABLED=true` — one environment and one +volume for the API and the jobs that delete what it wrote, so no override can +point them at different databases or stores. The hosted product never sets +it: its platform scheduler runs the same jobs. + +The schedule is `agentdrive.jobs.schedule.SCHEDULE`. Each due job runs as a +child process, `python -m `, in its own process group, +inheriting this process's environment. Jobs run one at a time. A job whose +due time passed while another ran, or while the scheduler was down (the last +handled minute persists in `SCHEDULER_STATE_FILE`), runs once when it can — +never once per missed minute. A job that outlives its timeout is killed with +its whole process group. A failing job is logged at ERROR and the scheduler +carries on. + +Deliberately standard library only, and it never loads the application +settings: each job validates its own configuration and fails loudly. +""" + +from __future__ import annotations + +import argparse +import contextlib +import json +import logging +import os +import signal +import subprocess +import sys +import threading +import time +from collections.abc import Callable, Iterable, Sequence +from dataclasses import dataclass, field +from datetime import UTC, datetime, timedelta +from pathlib import Path + +from agentdrive.jobs.schedule import SCHEDULE, ScheduledJob + +log = logging.getLogger("agentdrive.jobs.scheduler") + +# How far back a catch-up looks. A job due at least monthly is always inside +# it, and the scan stays small after an arbitrarily long outage. +CATCH_UP_WINDOW = timedelta(days=32) +# Grace a child's process group gets between SIGTERM and SIGKILL. +KILL_GRACE_S = 10 +STATE_FILE_ENV = "SCHEDULER_STATE_FILE" + + +class CronError(ValueError): + pass + + +_FIELDS = (("minute", 0, 59), ("hour", 0, 23), ("day", 1, 31), ("month", 1, 12), ("weekday", 0, 7)) + + +def _parse_field(text: str, low: int, high: int) -> frozenset[int]: + values: set[int] = set() + for part in text.split(","): + body, slash, step_text = part.partition("/") + try: + step = int(step_text) if slash else 1 + if body == "*": + start, end = low, high + elif "-" in body: + a, b = body.split("-", 1) + start, end = int(a), int(b) + else: + start = int(body) + # `N/S` means N, N+S, … up to the field's maximum, as in cron. + end = high if slash else start + except ValueError as exc: + raise CronError(f"{text!r} is not a cron field") from exc + if step < 1 or start > end or start < low or end > high: + raise CronError(f"{text!r} is outside {low}-{high}") + values.update(range(start, end + 1, step)) + return frozenset(values) + + +@dataclass(frozen=True) +class Cron: + """Five-field cron: numbers, `*`, ranges, lists and steps; Sunday is 0 or + 7. When both day-of-month and day-of-week are restricted (neither field + starts with `*`), either may match, as in classic cron.""" + + minute: frozenset[int] + hour: frozenset[int] + day: frozenset[int] + month: frozenset[int] + weekday: frozenset[int] + day_restricted: bool + weekday_restricted: bool + + @classmethod + def parse(cls, expr: str) -> Cron: + parts = expr.split() + if len(parts) != 5: + raise CronError(f"{expr!r} must have five fields") + fields = [_parse_field(p, lo, hi) for p, (_, lo, hi) in zip(parts, _FIELDS, strict=True)] + weekday = frozenset(0 if d == 7 else d for d in fields[4]) + return cls(fields[0], fields[1], fields[2], fields[3], weekday, + not parts[2].startswith("*"), not parts[4].startswith("*")) + + def matches(self, when: datetime) -> bool: + if when.minute not in self.minute or when.hour not in self.hour: + return False + if when.month not in self.month: + return False + day_ok = when.day in self.day + weekday_ok = (when.isoweekday() % 7) in self.weekday + if self.day_restricted and self.weekday_restricted: + return day_ok or weekday_ok + return day_ok and weekday_ok + + +def _minute(when: datetime) -> datetime: + return when.astimezone(UTC).replace(second=0, microsecond=0) + + +def latest_due(cron: Cron, after: datetime, until: datetime) -> datetime | None: + """The last minute in (after, until] the cron matches, if any.""" + tick = _minute(until) + floor = max(_minute(after), tick - CATCH_UP_WINDOW) + while tick > floor: + if cron.matches(tick): + return tick + tick -= timedelta(minutes=1) + return None + + +def due_between( + jobs: Sequence[ScheduledJob], + after: datetime, + until: datetime, + started: dict[str, datetime] | None = None, +) -> list[ScheduledJob]: + """Jobs due in (after, until], each once, in schedule order — skipping a + job that already STARTED at or after its latest due minute (it ran late, + behind another job, and that run covered this due time).""" + started = started or {} + due = [] + for job in jobs: + when = latest_due(Cron.parse(job.cron), after, until) + if when is None: + continue + ran = started.get(job.name) + if ran is not None and ran >= when: + continue + due.append(job) + return due + + +def module_command(job: ScheduledJob) -> list[str]: + return [sys.executable, "-m", job.module, *job.args] + + +def display_command(job: ScheduledJob) -> str: + """The command as an operator would type it in their own scheduler.""" + return " ".join(["python", "-m", job.module, *job.args]) + + +# --- persisted state ------------------------------------------------------ + + +@dataclass +class State: + """The last minute the loop handled, and when each job last started.""" + + last_tick: datetime | None = None + started: dict[str, datetime] = field(default_factory=dict) + + @classmethod + def load(cls, path: Path | None) -> State: + if path is None or not path.exists(): + return cls() + try: + raw = json.loads(path.read_text()) + return cls( + datetime.fromisoformat(raw["last_tick"]) if raw.get("last_tick") else None, + {k: datetime.fromisoformat(v) for k, v in raw.get("started", {}).items()}, + ) + except (OSError, ValueError, KeyError, TypeError, AttributeError): + log.warning("at=scheduler.state_unreadable path=%s; starting fresh", path) + return cls() + + def save(self, path: Path | None) -> None: + if path is None: + return + payload = { + "last_tick": self.last_tick.isoformat() if self.last_tick else None, + "started": {k: v.isoformat() for k, v in self.started.items()}, + } + tmp = path.with_name(path.name + ".tmp") + try: + tmp.write_text(json.dumps(payload)) + os.replace(tmp, path) + except OSError: + log.warning("at=scheduler.state_unwritable path=%s", path) + + +# --- running jobs --------------------------------------------------------- + + +@dataclass(frozen=True) +class RunResult: + name: str + exit_code: int + duration_s: float + timed_out: bool + + +def _exit_status(code: int) -> int: + """A signal death as a shell would report it (128 + signal).""" + return 128 - code if code < 0 else code + + +class Runner: + """Runs jobs one at a time, each in its own process group.""" + + def __init__(self, command_for: Callable[[ScheduledJob], list[str]] = module_command): + self._command_for = command_for + self._child: subprocess.Popen[bytes] | None = None + self._stopping = threading.Event() + self._stop_requested_at: float | None = None + + @property + def stopping(self) -> bool: + return self._stopping.is_set() + + def stop(self) -> None: + """Ask the current job to finish and run nothing after it.""" + self._stop_requested_at = time.monotonic() + self._stopping.set() + self._signal_child(signal.SIGTERM) + + def _signal_child(self, sig: int) -> None: + child = self._child + if child is None or child.poll() is not None: + return + with contextlib.suppress(ProcessLookupError, PermissionError): + os.killpg(child.pid, sig) + + def run_all(self, jobs: Iterable[ScheduledJob]) -> list[RunResult]: + results = [] + for job in jobs: + if self.stopping: + break + results.append(self.run(job)) + return results + + def run(self, job: ScheduledJob) -> RunResult: + log.info("at=scheduler.job_started name=%s", job.name) + started = time.monotonic() + timed_out = False + self._child = subprocess.Popen(self._command_for(job), start_new_session=True) + try: + if self.stopping: # a stop that landed between the check and the start + self._signal_child(signal.SIGTERM) + term_sent_at: float | None = None + while True: + try: + code = self._child.wait(timeout=0.2) + break + except subprocess.TimeoutExpired: + pass + now = time.monotonic() + if term_sent_at is None and now - started >= job.timeout_s: + timed_out = True + term_sent_at = now + self._signal_child(signal.SIGTERM) + if self.stopping and term_sent_at is None: + term_sent_at = self._stop_requested_at or now + if term_sent_at is not None and now - term_sent_at >= KILL_GRACE_S: + self._signal_child(signal.SIGKILL) + # Take down anything the job left behind in its group. + with contextlib.suppress(ProcessLookupError, PermissionError): + os.killpg(self._child.pid, signal.SIGKILL) + finally: + self._child = None + result = RunResult(job.name, _exit_status(code), time.monotonic() - started, timed_out) + level = logging.INFO if result.exit_code == 0 else logging.ERROR + log.log( + level, + "at=scheduler.job_finished name=%s exit=%s duration_s=%.1f timed_out=%s", + result.name, result.exit_code, result.duration_s, result.timed_out, + ) + return result + + +def serve( + runner: Runner, + jobs: Sequence[ScheduledJob] = SCHEDULE, + *, + state_path: Path | None = None, + now: Callable[[], datetime] = lambda: datetime.now(UTC), + sleep: Callable[[float], None] = time.sleep, +) -> int: + for job in jobs: + log.info("at=scheduler.scheduled name=%s cron=%r command=%r", + job.name, job.cron, display_command(job)) + state = State.load(state_path) + current = _minute(now()) + # Resume from the persisted minute, so a job due while the scheduler was + # down runs once now; with no state (a first start), from this minute. + last = state.last_tick if state.last_tick and state.last_tick <= current else current + last = max(last, current - CATCH_UP_WINDOW) + first = True + while not runner.stopping: + if not first: + next_tick = last + timedelta(minutes=1) + while not runner.stopping and now() < next_tick: + sleep(min(1.0, max(0.0, (next_tick - now()).total_seconds()))) + if runner.stopping: + break + first = False + tick = _minute(now()) + if tick < last: + # The clock stepped backwards: replay nothing, resume from here. + log.warning("at=scheduler.clock_stepped_back from=%s to=%s", last, tick) + last = tick + state.last_tick = tick + state.save(state_path) + continue + due = due_between(jobs, last, tick, state.started) + for job in due: + if runner.stopping: + break + state.started[job.name] = now() + state.save(state_path) + runner.run(job) + else: + # Only mark the batch handled after every due job started. If + # shutdown interrupts it, retain the old checkpoint so a restart + # catches up pending jobs; started deduplicates the earlier ones. + last = tick + state.last_tick = tick + state.save(state_path) + log.info("at=scheduler.stopped") + return 0 + + +def _state_path() -> Path | None: + value = os.environ.get(STATE_FILE_ENV, "").strip() + return Path(value) if value else None + + +def main(argv: list[str] | None = None) -> int: + logging.basicConfig( + level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s" + ) + parser = argparse.ArgumentParser( + prog="agentdrive.jobs.scheduler", description=__doc__.split("\n\n")[0] + ) + group = parser.add_mutually_exclusive_group() + group.add_argument("--list", action="store_true", help="print the schedule and exit") + group.add_argument("--run", metavar="NAME", help="run one job now and exit with its status") + group.add_argument("--status", action="store_true", help="print the persisted state") + args = parser.parse_args(argv) + + if args.list: + for job in SCHEDULE: + print(f"{job.name:16} {job.cron:14} {display_command(job)}") + if job.description: + print(f"{'':16} {job.description}") + return 0 + + if args.status: + path = _state_path() + if path is None: + print(f"{STATE_FILE_ENV} is not set; no state is kept", file=sys.stderr) + return 1 + state = State.load(path) + print(json.dumps({ + "last_tick": state.last_tick.isoformat() if state.last_tick else None, + "started": {k: v.isoformat() for k, v in sorted(state.started.items())}, + }, indent=2)) + return 0 if state.last_tick else 1 + + runner = Runner() + def _stop(signum: int, _frame: object) -> None: + log.info("at=scheduler.stopping signal=%s", signum) + runner.stop() + + signal.signal(signal.SIGTERM, _stop) + signal.signal(signal.SIGINT, _stop) + if args.run: + job = next((j for j in SCHEDULE if j.name == args.run), None) + if job is None: + print(f"no job named {args.run!r}; see --list", file=sys.stderr) + return 2 + return runner.run(job).exit_code + + return serve(runner, state_path=_state_path()) + + +if __name__ == "__main__": # pragma: no cover + sys.exit(main()) diff --git a/src/agentdrive/process_supervisor.py b/src/agentdrive/process_supervisor.py index fb1d752..d62f373 100644 --- a/src/agentdrive/process_supervisor.py +++ b/src/agentdrive/process_supervisor.py @@ -1,4 +1,5 @@ -"""Run the API, the private MCP ingress, and the MCP sidecar as one process. +"""Run the API, the private MCP ingress, the MCP sidecar and (for a +self-hosted install) the job scheduler as one process. Cloud Run's existing service owns the canonical AgentDrive domain, so the MCP cannot be a second service without changing the domain-routing topology. @@ -19,6 +20,13 @@ Generated per boot and never persisted: it is not in Terraform state, not in Secret Manager, and not in any log. A revision restart mints a new one, and the only two processes that need it are started by this file. + +THE SCHEDULER. With `SCHEDULER_ENABLED=true` (set by `compose.selfhost.yml`, +never by the hosted deployment, whose platform scheduler runs the same jobs) +this also starts `python -m agentdrive.jobs.scheduler`, which runs garbage +collection and usage maintenance on their schedule. It lives here rather than +in a container of its own so the jobs that DELETE content can never be +configured apart from the API that wrote it: one environment, one volume. """ from __future__ import annotations @@ -96,6 +104,14 @@ def _internal_ingress_command() -> list[str]: ] +def _scheduler_command() -> list[str]: + return [sys.executable, "-m", "agentdrive.jobs.scheduler"] + + +def scheduler_enabled(source: Mapping[str, str]) -> bool: + return source.get("SCHEDULER_ENABLED", "").strip().lower() in {"1", "true", "yes"} + + def _terminate(processes: Sequence[subprocess.Popen[bytes]]) -> None: for process in processes: if process.poll() is None: @@ -188,6 +204,14 @@ def stop(_signum: int, _frame) -> None: ) log.info("started AgentDrive MCP sidecar on 127.0.0.1:%s", MCP_SIDECAR_PORT) + if scheduler_enabled(os.environ): + # No proof: the jobs talk to Postgres and the store, never to + # the ingress. + children.append( + subprocess.Popen(_scheduler_command(), env=os.environ.copy()) # noqa: S603 + ) + log.info("started AgentDrive job scheduler") + # The PUBLIC API child, which deliberately does NOT receive the proof: # it neither calls the ingress nor needs to recognise it. children.append(subprocess.Popen(_api_command(), env=os.environ.copy())) # noqa: S603 diff --git a/tests/test_jobs_scheduler.py b/tests/test_jobs_scheduler.py new file mode 100644 index 0000000..21cd01d --- /dev/null +++ b/tests/test_jobs_scheduler.py @@ -0,0 +1,454 @@ +"""The self-hosted job scheduler (`agentdrive.jobs.scheduler`). + +A hosted deployment schedules the maintenance jobs with its platform's +scheduler; a self-hosted install runs this one inside its API container. These +tests pin the three things that decide whether an install ever reclaims +storage: the schedule itself, the cron matching, and the runner's behaviour +when jobs are due together, overlap, fail or hang. The runner is exercised +with real child processes (tiny Python one-liners), never mocks. +""" + +from __future__ import annotations + +import sys +from datetime import UTC, datetime, timedelta + +import pytest + +from agentdrive.jobs import scheduler +from agentdrive.jobs.schedule import SCHEDULE, ScheduledJob +from agentdrive.jobs.scheduler import Cron, CronError, Runner, due_between + + +def _utc(*args: int) -> datetime: + return datetime(*args, tzinfo=UTC) + + +# --- the schedule -------------------------------------------------------- + + +def test_the_schedule_is_the_hosted_cadence(): + """The same jobs, arguments and times the hosted deployment runs, so a + self-hosted install behaves like the product its tests describe.""" + rows = {(j.name, j.cron, j.module, tuple(j.args)) for j in SCHEDULE} + assert rows == { + ("gc-hourly", "0 * * * *", "agentdrive.jobs.gc", ("--sessions-only",)), + ("gc-daily", "0 3 * * *", "agentdrive.jobs.gc", ()), + ("gc-weekly", "0 4 * * 0", "agentdrive.jobs.gc", ("--orphan-sweep",)), + ("usage-snapshot", "*/15 * * * *", "agentdrive.jobs.usage_snapshot", ("--no-notify",)), + } + + +def test_every_job_has_a_unique_name_and_a_timeout_inside_an_hour(): + names = [j.name for j in SCHEDULE] + assert len(names) == len(set(names)) + for job in SCHEDULE: + Cron.parse(job.cron) # parses + assert 0 < job.timeout_s <= 3600, job.name + + +def test_the_gc_timeout_leaves_room_for_the_sweepers_own_deadline(): + """The sweeper stops itself at HARD_TIMEOUT_S; the scheduler's kill must + come after that, or it interrupts a sweep that was about to finish.""" + from agentdrive.core.gc import GCSweeper + + for job in SCHEDULE: + if job.module == "agentdrive.jobs.gc": + assert job.timeout_s > GCSweeper.HARD_TIMEOUT_S + + +# --- cron ----------------------------------------------------------------- + + +@pytest.mark.parametrize( + "expr, when, expected", + [ + ("0 * * * *", _utc(2026, 9, 29, 14, 0), True), + ("0 * * * *", _utc(2026, 9, 29, 14, 1), False), + ("0 3 * * *", _utc(2026, 9, 29, 3, 0), True), + ("0 3 * * *", _utc(2026, 9, 29, 4, 0), False), + # 2026-10-04 is a Sunday; cron's Sunday is 0 (and 7). + ("0 4 * * 0", _utc(2026, 10, 4, 4, 0), True), + ("0 4 * * 7", _utc(2026, 10, 4, 4, 0), True), + ("0 4 * * 0", _utc(2026, 10, 5, 4, 0), False), + ("*/15 * * * *", _utc(2026, 9, 29, 14, 45), True), + ("*/15 * * * *", _utc(2026, 9, 29, 14, 50), False), + ("5,35 1-3 * * *", _utc(2026, 9, 29, 2, 35), True), + ("5,35 1-3 * * *", _utc(2026, 9, 29, 4, 35), False), + ("0 0 1 * *", _utc(2026, 10, 1, 0, 0), True), + ("0 0 * 10 *", _utc(2026, 9, 1, 0, 0), False), + ], +) +def test_cron_matching(expr, when, expected): + assert Cron.parse(expr).matches(when) is expected + + +@pytest.mark.parametrize( + "expr", + ["", "* * * *", "* * * * * *", "60 * * * *", "* 24 * * *", "*/0 * * * *", "a * * * *", + "5-1 * * * *", "* * 0 * *", "* * * 13 *", "* * * * 8"], +) +def test_malformed_cron_is_refused(expr): + with pytest.raises(CronError): + Cron.parse(expr) + + +def test_when_both_day_fields_are_restricted_either_may_match(): + """Classic cron: day-of-month OR day-of-week once both are restricted.""" + cron = Cron.parse("0 0 13 * 5") # the 13th, or any Friday + assert cron.matches(_utc(2026, 10, 13, 0, 0)) # a Tuesday the 13th + assert cron.matches(_utc(2026, 10, 2, 0, 0)) # a Friday the 2nd + assert not cron.matches(_utc(2026, 10, 3, 0, 0)) + + +# --- which jobs are due --------------------------------------------------- + + +def test_three_oclock_runs_the_hourly_the_daily_and_usage_in_schedule_order(): + due = due_between(SCHEDULE, _utc(2026, 9, 29, 2, 59), _utc(2026, 9, 29, 3, 0)) + assert [j.name for j in due] == ["gc-hourly", "gc-daily", "usage-snapshot"] + + +def test_minutes_missed_while_a_job_ran_are_caught_up_once_each(): + """A daily sweep that ran from 03:00 to 03:40 must not cost the 03:15 and + 03:30 usage passes their run — but it owes ONE usage pass, not two.""" + due = due_between(SCHEDULE, _utc(2026, 9, 29, 3, 0), _utc(2026, 9, 29, 3, 40)) + assert [j.name for j in due] == ["usage-snapshot"] + + +def test_nothing_is_due_retroactively_for_the_minute_already_handled(): + assert due_between(SCHEDULE, _utc(2026, 9, 29, 3, 0), _utc(2026, 9, 29, 3, 0)) == [] + + +def test_a_long_outage_catches_up_each_job_once(): + """A host asleep over a weekend still runs each job exactly once.""" + due = due_between(SCHEDULE, _utc(2026, 10, 2, 0, 0), _utc(2026, 10, 5, 0, 0)) + assert sorted(j.name for j in due) == sorted(j.name for j in SCHEDULE) + + +def test_the_catch_up_window_is_bounded(): + """Days of missed minutes are scanned without iterating forever.""" + start = _utc(2026, 1, 1, 0, 0) + assert due_between(SCHEDULE, start, start + timedelta(days=400)) + + +# --- running jobs --------------------------------------------------------- + + +def _job(name: str, code: str, timeout_s: int = 30) -> ScheduledJob: + """A job whose "module" is a Python one-liner, run through `-c`.""" + return ScheduledJob(name=name, cron="* * * * *", module="", args=("-c", code), + timeout_s=timeout_s) + + +def _runner() -> Runner: + return Runner(command_for=lambda job: [sys.executable, *job.args]) + + +def test_jobs_run_one_at_a_time_and_report_their_exit_status(tmp_path): + marker = tmp_path / "order" + first = _job("first", f"open({str(marker)!r}, 'a').write('1')") + second = _job("second", f"open({str(marker)!r}, 'a').write('2'); raise SystemExit(3)") + results = _runner().run_all([first, second]) + assert marker.read_text() == "12" + assert [(r.name, r.exit_code) for r in results] == [("first", 0), ("second", 3)] + + +def test_a_failing_job_does_not_stop_the_ones_after_it(tmp_path): + marker = tmp_path / "ran" + results = _runner().run_all([ + _job("boom", "raise SystemExit(1)"), + _job("after", f"open({str(marker)!r}, 'w').write('ok')"), + ]) + assert marker.read_text() == "ok" + assert [r.exit_code for r in results] == [1, 0] + + +def test_a_hung_job_is_killed_at_its_timeout(): + results = _runner().run_all([_job("hang", "import time; time.sleep(60)", timeout_s=1)]) + assert results[0].timed_out and results[0].exit_code != 0 + assert results[0].duration_s < 30 + + +def test_the_real_jobs_are_started_as_python_modules(): + job = next(j for j in SCHEDULE if j.name == "gc-weekly") + assert scheduler.module_command(job) == [ + sys.executable, "-m", "agentdrive.jobs.gc", "--orphan-sweep", + ] + + +def test_run_one_by_name_and_refuse_an_unknown_name(capsys): + assert scheduler.main(["--list"]) == 0 + listed = capsys.readouterr().out + for job in SCHEDULE: + assert job.name in listed and job.cron in listed + assert scheduler.display_command(job) in listed + assert scheduler.main(["--run", "no-such-job"]) == 2 + + +def test_the_service_loop_stops_promptly_when_asked(): + """`docker compose stop` sends SIGTERM; the loop must not sit out the + rest of its minute before exiting.""" + import threading + + runner = _runner() + exit_codes: list[int] = [] + thread = threading.Thread(target=lambda: exit_codes.append(scheduler.serve(runner, ()))) + thread.start() + runner.stop() + thread.join(timeout=5) + assert not thread.is_alive() and exit_codes == [0] + + +def test_stop_ends_the_running_child_and_runs_nothing_after_it(tmp_path): + import threading + + marker = tmp_path / "second" + runner = _runner() + results: list = [] + jobs = [_job("slow", "import time; time.sleep(60)"), + _job("second", f"open({str(marker)!r}, 'w').write('ran')")] + thread = threading.Thread(target=lambda: results.extend(runner.run_all(jobs))) + thread.start() + import time + + time.sleep(1) + runner.stop() + thread.join(timeout=15) + assert not thread.is_alive() + assert [r.name for r in results] == ["slow"] and results[0].exit_code != 0 + assert not marker.exists() + + +# --- review round: cron parity with classic cron --------------------------- + + +def test_a_step_on_a_single_number_runs_to_the_fields_maximum(): + assert Cron.parse("5/15 * * * *").minute == {5, 20, 35, 50} + + +def test_a_starred_day_field_is_unrestricted_even_with_a_step(): + """Classic cron ORs the day fields only when NEITHER starts with `*`: + `*/2` in day-of-month still ANDs with a restricted weekday.""" + cron = Cron.parse("0 0 */2 * 1") # odd days of the month that are Mondays + assert cron.matches(_utc(2026, 10, 5, 0, 0)) # Monday the 5th + assert not cron.matches(_utc(2026, 10, 12, 0, 0)) # Monday the 12th + assert not cron.matches(_utc(2026, 10, 7, 0, 0)) # Wednesday the 7th + + +def test_every_scheduled_job_fires_within_any_five_weeks(): + """A job that never fires would be a schedule typo nothing else notices.""" + start = _utc(2026, 1, 1, 0, 0) + for job in SCHEDULE: + cron = Cron.parse(job.cron) + minutes = range(35 * 24 * 60) + assert any(cron.matches(start + timedelta(minutes=m)) for m in minutes), job.name + + +def test_the_catch_up_window_really_bounds_the_look_back(): + """A yearly job last due more than CATCH_UP_WINDOW ago is not caught up + — and the four real jobs, all due inside it, each are.""" + yearly = ScheduledJob(name="yearly", cron="0 0 1 1 *", module="m", args=(), timeout_s=60) + assert due_between([yearly], _utc(2026, 1, 1, 0, 0) - timedelta(days=1), + _utc(2026, 2, 15, 0, 0)) == [] + everything = due_between(SCHEDULE, _utc(2025, 1, 1, 0, 0), _utc(2026, 2, 15, 0, 0)) + assert [j.name for j in everything] == [j.name for j in SCHEDULE] + + +# --- review round: the service loop, driven by a fake clock ---------------- + + +class _Clock: + def __init__(self, start: datetime): + self.t = start + + def now(self) -> datetime: + return self.t + + def sleep(self, seconds: float) -> None: + self.t += timedelta(seconds=max(seconds, 0.001)) + + +class _FakeRunner(Runner): + """Records runs, advances the clock by each job's duration, and stops the + loop once the clock passes `until`.""" + + def __init__(self, clock: _Clock, until: datetime, durations: dict[str, timedelta]): + super().__init__() + self.clock, self.until, self.durations = clock, until, durations + self.runs: list[tuple[str, datetime]] = [] + + @property + def stopping(self) -> bool: + return self.clock.t >= self.until + + def run(self, job): + self.runs.append((job.name, self.clock.t)) + self.clock.t += self.durations.get(job.name, timedelta(seconds=5)) + return scheduler.RunResult(job.name, 0, 0.0, False) + + +def _serve(start, until, durations=None, state_path=None): + clock = _Clock(start) + runner = _FakeRunner(clock, until, durations or {}) + scheduler.serve(runner, SCHEDULE, state_path=state_path, now=clock.now, sleep=clock.sleep) + return [(name, when.strftime("%H:%M")) for name, when in runner.runs] + + +def test_a_job_that_waited_behind_a_long_sweep_is_not_run_twice(): + """gc-daily takes 55 minutes; usage due at 03:00 runs once when it ends, + and the 03:15/03:30/03:45 dues it covered do not run it again.""" + runs = _serve(_utc(2026, 9, 29, 2, 59, 30), _utc(2026, 9, 29, 4, 1), + {"gc-daily": timedelta(minutes=55)}) + assert [name for name, _ in runs] == [ + "gc-hourly", "gc-daily", "usage-snapshot", # the 03:00 batch + "gc-hourly", "usage-snapshot", # 04:00 + ] + + +def test_a_restart_catches_up_what_came_due_while_it_was_down(tmp_path): + """Down across Sunday 04:00: the weekly orphan sweep runs once on start, + beside one hourly and one usage pass — not once per missed minute.""" + state = tmp_path / "state.json" + scheduler.State(last_tick=_utc(2026, 10, 4, 3, 50)).save(state) + runs = _serve(_utc(2026, 10, 4, 6, 10, 20), _utc(2026, 10, 4, 6, 11), state_path=state) + assert sorted(name for name, _ in runs) == ["gc-hourly", "gc-weekly", "usage-snapshot"] + + +def test_a_first_start_runs_nothing_retroactively(tmp_path): + runs = _serve(_utc(2026, 9, 29, 3, 0, 20), _utc(2026, 9, 29, 3, 1), + state_path=tmp_path / "state.json") + assert runs == [] + + +def test_state_persists_the_minute_and_each_start(tmp_path): + state = tmp_path / "state.json" + _serve(_utc(2026, 9, 29, 2, 59, 30), _utc(2026, 9, 29, 3, 2), state_path=state) + saved = scheduler.State.load(state) + assert saved.last_tick == _utc(2026, 9, 29, 3, 2) or saved.last_tick == _utc(2026, 9, 29, 3, 1) + assert set(saved.started) == {"gc-hourly", "gc-daily", "usage-snapshot"} + + +def test_a_clock_stepped_backwards_replays_nothing(tmp_path): + state = tmp_path / "state.json" + scheduler.State(last_tick=_utc(2026, 9, 29, 5, 0)).save(state) + runs = _serve(_utc(2026, 9, 29, 4, 0, 10), _utc(2026, 9, 29, 4, 14), state_path=state) + assert runs == [] + + +def test_an_unreadable_state_file_starts_fresh(tmp_path): + state = tmp_path / "state.json" + state.write_text("{not json") + assert _serve(_utc(2026, 9, 29, 3, 0, 20), _utc(2026, 9, 29, 3, 1), state_path=state) == [] + + +# --- review round: the runner ---------------------------------------------- + + +def test_a_timeout_kills_the_jobs_whole_process_group(tmp_path): + """A job's own children (a wrapper, a pool) die with it, TERM-deaf or not.""" + import os + import time as _time + + pidfile = tmp_path / "grandchild" + code = ( + "import subprocess, time\n" + f"p = subprocess.Popen(['sh', '-c', 'trap \"\" TERM; echo $$ > {pidfile}; sleep 60'])\n" + "time.sleep(60)\n" + ) + result = _runner().run(_job("spawner", code, timeout_s=1)) + assert result.timed_out + grandchild = int(pidfile.read_text()) + for _ in range(50): + try: + os.kill(grandchild, 0) + except ProcessLookupError: + break + _time.sleep(0.1) + else: + raise AssertionError("the job's grandchild outlived the timeout") + + +def test_stop_escalates_to_kill_for_a_job_that_ignores_sigterm(monkeypatch): + import threading + import time as _time + + monkeypatch.setattr(scheduler, "KILL_GRACE_S", 1) + runner = _runner() + deaf = _job("deaf", "import signal, time; signal.signal(signal.SIGTERM, signal.SIG_IGN); " + "time.sleep(60)") + results: list = [] + thread = threading.Thread(target=lambda: results.append(runner.run(deaf))) + started = _time.monotonic() + thread.start() + _time.sleep(1) + runner.stop() + thread.join(timeout=15) + assert not thread.is_alive() and _time.monotonic() - started < 10 + assert results[0].exit_code == 128 + 9 + + +def test_a_stop_that_lands_before_the_child_starts_still_stops_it(): + runner = _runner() + runner.stop() + result = runner.run(_job("late", "import time; time.sleep(30)")) + assert result.exit_code == 128 + 15 and result.duration_s < 10 + + +def test_a_signal_death_reports_as_a_shell_would(): + result = _runner().run(_job("sig", "import os, signal; os.kill(os.getpid(), signal.SIGTERM)")) + assert result.exit_code == 143 + + +def test_restart_preserves_unstarted_jobs_in_an_interrupted_batch(tmp_path): + state = tmp_path / "state.json" + first = _serve( + _utc(2026, 10, 4, 3, 59, 30), _utc(2026, 10, 4, 4, 0, 30), + {"gc-hourly": timedelta(minutes=1)}, state_path=state, + ) + assert first == [("gc-hourly", "04:00")] + resumed = _serve( + _utc(2026, 10, 4, 4, 1), _utc(2026, 10, 4, 4, 2), state_path=state, + ) + assert resumed == [("gc-weekly", "04:01"), ("usage-snapshot", "04:01")] + + +def test_manual_cli_terminates_its_child_on_sigterm(tmp_path): + import os + import signal + import subprocess + import time + + marker = tmp_path / "child.pid" + child_code = ( + f"import os,time; open({str(marker)!r}, 'w').write(str(os.getpid())); " + "time.sleep(60)" + ) + wrapper = ( + "import sys\n" + "from agentdrive.jobs import scheduler as s\n" + f"s.module_command = lambda job: [sys.executable, '-c', {child_code!r}]\n" + "s.Runner.__init__.__defaults__ = (s.module_command,)\n" + "raise SystemExit(s.main(['--run', 'gc-daily']))\n" + ) + process = subprocess.Popen([sys.executable, "-c", wrapper]) + child_pid = None + try: + deadline = time.monotonic() + 10 + while not marker.exists() and time.monotonic() < deadline: + time.sleep(0.05) + assert marker.exists() + child_pid = int(marker.read_text()) + process.send_signal(signal.SIGTERM) + assert process.wait(timeout=15) == 143 + with pytest.raises(ProcessLookupError): + os.kill(child_pid, 0) + finally: + if process.poll() is None: + process.kill() + process.wait() + if child_pid: + import contextlib + + with contextlib.suppress(ProcessLookupError): + os.killpg(child_pid, signal.SIGKILL) diff --git a/tests/test_selfhost_compose.py b/tests/test_selfhost_compose.py index 47954bd..e2fb45e 100644 --- a/tests/test_selfhost_compose.py +++ b/tests/test_selfhost_compose.py @@ -40,6 +40,8 @@ "SESSION_SECRET", "MCP_PROXY_URL", "PORT", + "SCHEDULER_ENABLED", + "SCHEDULER_STATE_FILE", } @@ -92,6 +94,21 @@ def test_api_is_local_mode_on_the_filesystem_store_with_the_sidecar(): assert api["stop_grace_period"] == "20s" +def test_the_api_container_runs_the_maintenance_jobs_on_its_own_store(): + """The jobs delete what the API wrote, so they run INSIDE its container + (the supervisor starts the scheduler): no override can point them at a + different database or store. Their state lives on the data volume, as a + dot-file the filesystem store never lists as an object.""" + api = _compose()["services"]["api"] + env = api["environment"] + assert env["SCHEDULER_ENABLED"] == "true" + root = env["STORAGE_FS_ROOT"] + state = env["SCHEDULER_STATE_FILE"] + assert state.startswith(root + "/.") + assert f"data:{root}" in api["volumes"] + assert api["init"] is True + + def test_the_published_port_and_the_public_base_url_cannot_come_apart(): """An operator who sets only AGENTDRIVE_PORT must get share links and an MCP discovery document on that port, not on 8080 — and must get the API @@ -169,8 +186,16 @@ def test_the_smoke_script_runs_nothing_the_readme_does_not_show(): """The other direction of parity: every `exec` the script runs against the stack is a command the README shows the reader.""" readme = " ".join(_readme_quickstart()) + whole_readme = " ".join(README.read_text().split()) for line in _smoke_as_readme_would_spell_it().splitlines(): line = line.strip() + if "exec api python -m agentdrive.jobs.scheduler" in line: + # The maintenance commands sit outside the quickstart block; the + # script runs each job by name where the README shows one. + command = line.split("exec api", 1)[1].split("|")[0].split(">")[0].split(")")[0] + command = " ".join(command.replace('"$name"', "gc-daily").split()) + assert f"exec api {command}" in whole_readme, f"the README never shows: {command}" + continue if "docker compose -f compose.selfhost.yml exec api" not in line: continue command = line.split("exec api", 1)[1].strip().split("|")[0].strip().rstrip(")") diff --git a/tests/test_supervisor_local_env.py b/tests/test_supervisor_local_env.py index 5ad8819..1ccd95a 100644 --- a/tests/test_supervisor_local_env.py +++ b/tests/test_supervisor_local_env.py @@ -44,3 +44,17 @@ def test_an_operators_explicit_sidecar_value_is_never_overwritten(): assert local_api_key_sidecar_env({"AUTH_MODE": "local", "MCP_AUTH_MODE": "jwt"}) == { "MCP_ENVIRONMENT": "local" } + + +def test_the_scheduler_starts_only_when_asked(): + """Self-hosted compose sets SCHEDULER_ENABLED; the hosted deployment never + does (its Cloud Scheduler runs the same jobs, and two schedulers would + double every run).""" + from agentdrive.process_supervisor import scheduler_enabled + + assert scheduler_enabled({"SCHEDULER_ENABLED": "true"}) + assert scheduler_enabled({"SCHEDULER_ENABLED": " TRUE "}) + assert scheduler_enabled({"SCHEDULER_ENABLED": "1"}) + assert not scheduler_enabled({}) + assert not scheduler_enabled({"SCHEDULER_ENABLED": "false"}) + assert not scheduler_enabled({"SCHEDULER_ENABLED": ""})