feat: add api key auth, devii agent, openai gateway, and admin service management
This commit introduces a comprehensive set of new features including API key authentication with CLI management commands (get, reset, backfill), a Devii agentic assistant with WebSocket terminal and session bootstrap, an OpenAI-compatible LLM gateway service, and an admin service management panel. It also adds Playwright browser automation for bot support, configures internal gateway URLs, refactors content editing/deletion to support JSON API responses, and updates documentation across AGENTS.md, README.md, and the developer docs site.
This commit is contained in:
@@ -0,0 +1,8 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
from .actions import TASK_ACTIONS
|
||||
from .controller import TaskController
|
||||
from .scheduler import Scheduler
|
||||
from .store import TaskStore
|
||||
|
||||
__all__ = ["TASK_ACTIONS", "TaskController", "Scheduler", "TaskStore"]
|
||||
@@ -0,0 +1,128 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from ..actions.spec import Action, Param
|
||||
|
||||
|
||||
def field(name: str, description: str, required: bool = False, kind: str = "string") -> Param:
|
||||
return Param(name=name, location="body", description=description, required=required, type=kind)
|
||||
|
||||
|
||||
SCHEDULE_FIELDS: tuple[Param, ...] = (
|
||||
field(
|
||||
"kind",
|
||||
"Schedule type: 'once' (run a single time), 'interval' (repeat every N seconds), "
|
||||
"or 'cron' (5-field cron expression).",
|
||||
required=True,
|
||||
),
|
||||
field(
|
||||
"run_at",
|
||||
"For kind=once: absolute UTC time in ISO 8601, e.g. 2026-06-09T14:30:00.",
|
||||
),
|
||||
field("delay_seconds", "For kind=once: run this many seconds from now.", kind="integer"),
|
||||
field("every_seconds", "For kind=interval: number of seconds between runs.", kind="integer"),
|
||||
field(
|
||||
"start_at",
|
||||
"For kind=interval: optional ISO 8601 UTC time of the first run. "
|
||||
"Defaults to one interval from now.",
|
||||
),
|
||||
field(
|
||||
"cron",
|
||||
"For kind=cron: a 5-field cron expression 'minute hour day-of-month month day-of-week'. "
|
||||
"Supports *, ranges (1-5), lists (1,3,5) and steps (*/15).",
|
||||
),
|
||||
field(
|
||||
"max_runs",
|
||||
"Optional maximum number of executions; the task disables itself afterwards.",
|
||||
kind="integer",
|
||||
),
|
||||
)
|
||||
|
||||
TASK_ACTIONS: tuple[Action, ...] = (
|
||||
Action(
|
||||
name="create_task",
|
||||
method="LOCAL",
|
||||
path="",
|
||||
summary="Queue a prompt to run autonomously on a schedule with full tool access",
|
||||
description=(
|
||||
"The prompt is executed later by a fresh agent that has every tool available, "
|
||||
"running inside the current authenticated session."
|
||||
),
|
||||
handler="task",
|
||||
requires_auth=False,
|
||||
params=(
|
||||
field(
|
||||
"prompt",
|
||||
"A complete, self-contained instruction the future agent will execute when the "
|
||||
"task fires (it has no other context). Write an imperative command, not a "
|
||||
"description: e.g. 'Create a post in the fun topic about raccoons with a short "
|
||||
"funny text, then verify it appears.'",
|
||||
required=True,
|
||||
),
|
||||
field("label", "Optional short human-readable label for the task."),
|
||||
*SCHEDULE_FIELDS,
|
||||
),
|
||||
),
|
||||
Action(
|
||||
name="list_tasks",
|
||||
method="LOCAL",
|
||||
path="",
|
||||
summary="List queued tasks and their schedules and last results",
|
||||
handler="task",
|
||||
requires_auth=False,
|
||||
params=(
|
||||
field("enabled_only", "Only return enabled tasks.", kind="boolean"),
|
||||
field("status", "Filter by status: pending, running, done, error, disabled."),
|
||||
),
|
||||
),
|
||||
Action(
|
||||
name="get_task",
|
||||
method="LOCAL",
|
||||
path="",
|
||||
summary="Get a single task including its full last result",
|
||||
handler="task",
|
||||
requires_auth=False,
|
||||
params=(field("uid", "Uid of the task.", required=True),),
|
||||
),
|
||||
Action(
|
||||
name="update_task",
|
||||
method="LOCAL",
|
||||
path="",
|
||||
summary="Update a task's prompt, label, enabled state, or schedule",
|
||||
description="Provide schedule fields together with kind to reschedule the task.",
|
||||
handler="task",
|
||||
requires_auth=False,
|
||||
params=(
|
||||
field("uid", "Uid of the task.", required=True),
|
||||
field("prompt", "New prompt."),
|
||||
field("label", "New label."),
|
||||
field("enabled", "Enable or disable the task.", kind="boolean"),
|
||||
field("kind", "New schedule type when rescheduling."),
|
||||
field("run_at", "New absolute run time for kind=once."),
|
||||
field("delay_seconds", "New relative delay for kind=once.", kind="integer"),
|
||||
field("every_seconds", "New interval for kind=interval.", kind="integer"),
|
||||
field("start_at", "New first-run time for kind=interval."),
|
||||
field("cron", "New cron expression for kind=cron."),
|
||||
field("max_runs", "New maximum number of executions.", kind="integer"),
|
||||
),
|
||||
),
|
||||
Action(
|
||||
name="delete_task",
|
||||
method="LOCAL",
|
||||
path="",
|
||||
summary="Delete a queued task",
|
||||
handler="task",
|
||||
requires_auth=False,
|
||||
params=(field("uid", "Uid of the task.", required=True),),
|
||||
),
|
||||
Action(
|
||||
name="run_task_now",
|
||||
method="LOCAL",
|
||||
path="",
|
||||
summary="Trigger a task to execute immediately on the next scheduler tick",
|
||||
handler="task",
|
||||
requires_auth=False,
|
||||
params=(field("uid", "Uid of the task.", required=True),),
|
||||
),
|
||||
)
|
||||
@@ -0,0 +1,211 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import uuid
|
||||
from typing import Any
|
||||
|
||||
from pydantic import ValidationError
|
||||
|
||||
from ..errors import ToolInputError
|
||||
from .schedule import Schedule, next_run, now_utc, to_iso
|
||||
from .store import TaskStore
|
||||
|
||||
logger = logging.getLogger("devii.tasks.controller")
|
||||
|
||||
SCHEDULE_KEYS = ("kind", "run_at", "delay_seconds", "every_seconds", "start_at", "cron", "max_runs")
|
||||
RESULT_PREVIEW_CHARS = 500
|
||||
TRUTHY = {"1", "true", "yes", "on"}
|
||||
|
||||
|
||||
def _as_bool(value: Any, default: bool = True) -> bool:
|
||||
if isinstance(value, bool):
|
||||
return value
|
||||
if value is None:
|
||||
return default
|
||||
return str(value).strip().lower() in TRUTHY
|
||||
|
||||
|
||||
def _serialize(row: dict[str, Any], preview: bool) -> dict[str, Any]:
|
||||
view = {
|
||||
"uid": row.get("uid"),
|
||||
"label": row.get("label"),
|
||||
"prompt": row.get("prompt"),
|
||||
"enabled": bool(row.get("enabled")),
|
||||
"status": row.get("status"),
|
||||
"kind": row.get("kind"),
|
||||
"next_run_at": row.get("next_run_at"),
|
||||
"last_run_at": row.get("last_run_at"),
|
||||
"run_count": row.get("run_count"),
|
||||
"max_runs": row.get("max_runs"),
|
||||
"every_seconds": row.get("every_seconds"),
|
||||
"cron": row.get("cron"),
|
||||
"run_at": row.get("run_at"),
|
||||
}
|
||||
last_error = row.get("last_error")
|
||||
if last_error:
|
||||
view["last_error"] = last_error
|
||||
result = row.get("last_result")
|
||||
if result:
|
||||
view["last_result"] = result[:RESULT_PREVIEW_CHARS] if preview else result
|
||||
return view
|
||||
|
||||
|
||||
class TaskController:
|
||||
def __init__(self, store: TaskStore) -> None:
|
||||
self._store = store
|
||||
|
||||
async def dispatch(self, name: str, arguments: dict[str, Any]) -> str:
|
||||
handlers = {
|
||||
"create_task": self.create_task,
|
||||
"list_tasks": self.list_tasks,
|
||||
"get_task": self.get_task,
|
||||
"update_task": self.update_task,
|
||||
"delete_task": self.delete_task,
|
||||
"run_task_now": self.run_task_now,
|
||||
}
|
||||
handler = handlers.get(name)
|
||||
if handler is None:
|
||||
raise ToolInputError(f"Unknown task tool: {name}")
|
||||
return handler(arguments)
|
||||
|
||||
def create_task(self, arguments: dict[str, Any]) -> str:
|
||||
prompt = str(arguments.get("prompt", "")).strip()
|
||||
if not prompt:
|
||||
raise ToolInputError("create_task requires a non-empty prompt.")
|
||||
schedule = self._build_schedule(arguments)
|
||||
reference = now_utc()
|
||||
first = schedule.first_run(reference)
|
||||
record: dict[str, Any] = {
|
||||
"uid": uuid.uuid4().hex,
|
||||
"label": (arguments.get("label") or "").strip() or None,
|
||||
"prompt": prompt,
|
||||
"enabled": True,
|
||||
"status": "pending",
|
||||
"created_at": to_iso(reference),
|
||||
"next_run_at": to_iso(first),
|
||||
"last_run_at": None,
|
||||
"run_count": 0,
|
||||
"last_result": None,
|
||||
"last_error": None,
|
||||
**schedule.columns(),
|
||||
}
|
||||
self._store.create(record)
|
||||
return json.dumps(
|
||||
{"status": "created", "task": _serialize(record, preview=True)},
|
||||
ensure_ascii=False,
|
||||
)
|
||||
|
||||
def list_tasks(self, arguments: dict[str, Any]) -> str:
|
||||
rows = self._store.list(
|
||||
enabled_only=_as_bool(arguments.get("enabled_only"), default=False),
|
||||
status=(arguments.get("status") or None),
|
||||
)
|
||||
rows.sort(key=lambda row: row.get("next_run_at") or "")
|
||||
return json.dumps(
|
||||
{"count": len(rows), "tasks": [_serialize(row, preview=True) for row in rows]},
|
||||
ensure_ascii=False,
|
||||
)
|
||||
|
||||
def get_task(self, arguments: dict[str, Any]) -> str:
|
||||
row = self._require_task(arguments)
|
||||
return json.dumps({"task": _serialize(row, preview=False)}, ensure_ascii=False)
|
||||
|
||||
def update_task(self, arguments: dict[str, Any]) -> str:
|
||||
row = self._require_task(arguments)
|
||||
changes: dict[str, Any] = {}
|
||||
|
||||
if "prompt" in arguments and arguments["prompt"] is not None:
|
||||
new_prompt = str(arguments["prompt"]).strip()
|
||||
if not new_prompt:
|
||||
raise ToolInputError("prompt cannot be empty.")
|
||||
changes["prompt"] = new_prompt
|
||||
if "label" in arguments:
|
||||
changes["label"] = (arguments.get("label") or "").strip() or None
|
||||
if "enabled" in arguments:
|
||||
changes["enabled"] = _as_bool(arguments.get("enabled"))
|
||||
|
||||
if any(key in arguments and arguments[key] is not None for key in SCHEDULE_KEYS):
|
||||
merged = {key: row.get(key) for key in SCHEDULE_KEYS}
|
||||
for key in SCHEDULE_KEYS:
|
||||
if key in arguments and arguments[key] is not None:
|
||||
merged[key] = arguments[key]
|
||||
schedule = self._build_schedule(merged)
|
||||
changes.update(schedule.columns())
|
||||
changes["next_run_at"] = to_iso(schedule.first_run(now_utc()))
|
||||
changes["status"] = "pending"
|
||||
|
||||
if changes.get("enabled") and row.get("status") in ("done", "disabled", "error"):
|
||||
changes.setdefault("status", "pending")
|
||||
if not changes.get("next_run_at") and not row.get("next_run_at"):
|
||||
schedule = self._build_schedule({key: row.get(key) for key in SCHEDULE_KEYS})
|
||||
changes["next_run_at"] = to_iso(schedule.first_run(now_utc()))
|
||||
if changes.get("enabled") is False:
|
||||
changes["status"] = "disabled"
|
||||
|
||||
if not changes:
|
||||
raise ToolInputError("No updatable fields supplied.")
|
||||
|
||||
self._store.update(row["uid"], changes)
|
||||
return json.dumps(
|
||||
{"status": "updated", "task": _serialize(self._store.get(row["uid"]), preview=True)},
|
||||
ensure_ascii=False,
|
||||
)
|
||||
|
||||
def delete_task(self, arguments: dict[str, Any]) -> str:
|
||||
uid = self._require_uid(arguments)
|
||||
deleted = self._store.delete(uid)
|
||||
if not deleted:
|
||||
raise ToolInputError(f"No task found with uid {uid}.")
|
||||
return json.dumps({"status": "deleted", "uid": uid}, ensure_ascii=False)
|
||||
|
||||
def run_task_now(self, arguments: dict[str, Any]) -> str:
|
||||
row = self._require_task(arguments)
|
||||
self._store.update(
|
||||
row["uid"],
|
||||
{"enabled": True, "status": "pending", "next_run_at": to_iso(now_utc())},
|
||||
)
|
||||
return json.dumps(
|
||||
{"status": "queued", "uid": row["uid"], "note": "Will execute on the next scheduler tick."},
|
||||
ensure_ascii=False,
|
||||
)
|
||||
|
||||
def _build_schedule(self, source: dict[str, Any]) -> Schedule:
|
||||
payload = {key: source.get(key) for key in SCHEDULE_KEYS if source.get(key) is not None}
|
||||
try:
|
||||
return Schedule(**payload)
|
||||
except ValidationError as exc:
|
||||
raise ToolInputError(f"Invalid schedule: {exc.errors()[0]['msg']}") from exc
|
||||
except ValueError as exc:
|
||||
raise ToolInputError(f"Invalid schedule: {exc}") from exc
|
||||
|
||||
def _require_uid(self, arguments: dict[str, Any]) -> str:
|
||||
uid = str(arguments.get("uid", "")).strip()
|
||||
if not uid:
|
||||
raise ToolInputError("This task tool requires a uid.")
|
||||
return uid
|
||||
|
||||
def _require_task(self, arguments: dict[str, Any]) -> dict[str, Any]:
|
||||
uid = self._require_uid(arguments)
|
||||
row = self._store.get(uid)
|
||||
if row is None:
|
||||
raise ToolInputError(f"No task found with uid {uid}.")
|
||||
return row
|
||||
|
||||
|
||||
def compute_followup(row: dict[str, Any], reference: Any) -> tuple[dict[str, Any], bool]:
|
||||
run_count = int(row.get("run_count") or 0) + 1
|
||||
max_runs = row.get("max_runs")
|
||||
upcoming = next_run(row.get("kind"), row.get("every_seconds"), row.get("cron"), reference)
|
||||
|
||||
changes: dict[str, Any] = {"run_count": run_count, "last_run_at": to_iso(reference)}
|
||||
if upcoming is None or (max_runs is not None and run_count >= int(max_runs)):
|
||||
changes["status"] = "done"
|
||||
changes["enabled"] = False
|
||||
changes["next_run_at"] = None
|
||||
else:
|
||||
changes["status"] = "pending"
|
||||
changes["next_run_at"] = to_iso(upcoming)
|
||||
return changes, changes["status"] == "done"
|
||||
@@ -0,0 +1,147 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Literal, Optional
|
||||
|
||||
from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator
|
||||
|
||||
ISO_FORMAT = "%Y-%m-%dT%H:%M:%S"
|
||||
CRON_MINUTES_LIMIT = 525600 * 4
|
||||
|
||||
ScheduleKind = Literal["once", "interval", "cron"]
|
||||
|
||||
|
||||
def now_utc() -> datetime:
|
||||
return datetime.now(timezone.utc).replace(microsecond=0)
|
||||
|
||||
|
||||
def to_iso(moment: datetime) -> str:
|
||||
if moment.tzinfo is None:
|
||||
moment = moment.replace(tzinfo=timezone.utc)
|
||||
return moment.astimezone(timezone.utc).strftime(ISO_FORMAT)
|
||||
|
||||
|
||||
def from_iso(value: str) -> datetime:
|
||||
return datetime.strptime(value, ISO_FORMAT).replace(tzinfo=timezone.utc)
|
||||
|
||||
|
||||
def _parse_field(field: str, low: int, high: int) -> set[int]:
|
||||
values: set[int] = set()
|
||||
for part in field.split(","):
|
||||
token, _, step_text = part.partition("/")
|
||||
step = int(step_text) if step_text else 1
|
||||
if step < 1:
|
||||
raise ValueError(f"Invalid step in cron field: {part}")
|
||||
if token == "*":
|
||||
start, end = low, high
|
||||
elif "-" in token:
|
||||
start_text, end_text = token.split("-", 1)
|
||||
start, end = int(start_text), int(end_text)
|
||||
else:
|
||||
start = end = int(token)
|
||||
if start < low or end > high or start > end:
|
||||
raise ValueError(f"Cron value out of range: {part}")
|
||||
values.update(range(start, end + 1, step))
|
||||
return values
|
||||
|
||||
|
||||
def cron_next(expr: str, after: datetime) -> datetime:
|
||||
fields = expr.split()
|
||||
if len(fields) != 5:
|
||||
raise ValueError("Cron expression must have 5 fields: minute hour dom month dow")
|
||||
minutes = _parse_field(fields[0], 0, 59)
|
||||
hours = _parse_field(fields[1], 0, 23)
|
||||
days = _parse_field(fields[2], 1, 31)
|
||||
months = _parse_field(fields[3], 1, 12)
|
||||
weekdays_raw = _parse_field(fields[4], 0, 7)
|
||||
weekdays = {0 if value == 7 else value for value in weekdays_raw}
|
||||
day_restricted = fields[2] != "*"
|
||||
weekday_restricted = fields[4] != "*"
|
||||
|
||||
candidate = after.replace(second=0, microsecond=0) + timedelta(minutes=1)
|
||||
for _ in range(CRON_MINUTES_LIMIT):
|
||||
if (
|
||||
candidate.minute in minutes
|
||||
and candidate.hour in hours
|
||||
and candidate.month in months
|
||||
):
|
||||
cron_weekday = (candidate.weekday() + 1) % 7
|
||||
day_ok = candidate.day in days
|
||||
weekday_ok = cron_weekday in weekdays
|
||||
if day_restricted and weekday_restricted:
|
||||
match_day = day_ok or weekday_ok
|
||||
elif day_restricted:
|
||||
match_day = day_ok
|
||||
elif weekday_restricted:
|
||||
match_day = weekday_ok
|
||||
else:
|
||||
match_day = True
|
||||
if match_day:
|
||||
return candidate
|
||||
candidate += timedelta(minutes=1)
|
||||
raise ValueError("No cron match found within four years")
|
||||
|
||||
|
||||
class Schedule(BaseModel):
|
||||
model_config = ConfigDict(extra="forbid")
|
||||
|
||||
kind: ScheduleKind
|
||||
run_at: Optional[datetime] = None
|
||||
delay_seconds: Optional[int] = Field(default=None, ge=1)
|
||||
every_seconds: Optional[int] = Field(default=None, ge=1)
|
||||
start_at: Optional[datetime] = None
|
||||
cron: Optional[str] = None
|
||||
max_runs: Optional[int] = Field(default=None, ge=1)
|
||||
|
||||
@field_validator("run_at", "start_at")
|
||||
@classmethod
|
||||
def _ensure_utc(cls, value: Optional[datetime]) -> Optional[datetime]:
|
||||
if value is None:
|
||||
return None
|
||||
if value.tzinfo is None:
|
||||
return value.replace(tzinfo=timezone.utc)
|
||||
return value.astimezone(timezone.utc)
|
||||
|
||||
@model_validator(mode="after")
|
||||
def _validate_kind(self) -> "Schedule":
|
||||
if self.kind == "once" and self.run_at is None and self.delay_seconds is None:
|
||||
raise ValueError("kind=once requires run_at or delay_seconds")
|
||||
if self.kind == "interval" and self.every_seconds is None:
|
||||
raise ValueError("kind=interval requires every_seconds")
|
||||
if self.kind == "cron":
|
||||
if not self.cron:
|
||||
raise ValueError("kind=cron requires a cron expression")
|
||||
cron_next(self.cron, now_utc())
|
||||
return self
|
||||
|
||||
def first_run(self, reference: datetime) -> datetime:
|
||||
if self.kind == "once":
|
||||
if self.run_at is not None:
|
||||
return self.run_at
|
||||
return reference + timedelta(seconds=self.delay_seconds or 0)
|
||||
if self.kind == "interval":
|
||||
if self.start_at is not None:
|
||||
return self.start_at
|
||||
return reference + timedelta(seconds=self.every_seconds or 0)
|
||||
return cron_next(self.cron or "", reference)
|
||||
|
||||
def columns(self) -> dict[str, object]:
|
||||
return {
|
||||
"kind": self.kind,
|
||||
"run_at": to_iso(self.run_at) if self.run_at else None,
|
||||
"delay_seconds": self.delay_seconds,
|
||||
"every_seconds": self.every_seconds,
|
||||
"start_at": to_iso(self.start_at) if self.start_at else None,
|
||||
"cron": self.cron,
|
||||
"max_runs": self.max_runs,
|
||||
}
|
||||
|
||||
|
||||
def next_run(kind: str, every_seconds: Optional[int], cron: Optional[str], reference: datetime) -> Optional[datetime]:
|
||||
if kind == "interval" and every_seconds:
|
||||
return reference + timedelta(seconds=int(every_seconds))
|
||||
if kind == "cron" and cron:
|
||||
return cron_next(str(cron), reference)
|
||||
return None
|
||||
@@ -0,0 +1,86 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import Any, Awaitable, Callable
|
||||
|
||||
from .controller import compute_followup
|
||||
from .schedule import now_utc, to_iso
|
||||
from .store import TaskStore
|
||||
|
||||
logger = logging.getLogger("devii.tasks.scheduler")
|
||||
|
||||
PromptExecutor = Callable[[str], Awaitable[str]]
|
||||
EventCallback = Callable[[str, dict[str, Any], str], None]
|
||||
|
||||
DEFAULT_TICK_SECONDS = 1.0
|
||||
|
||||
|
||||
class Scheduler:
|
||||
def __init__(
|
||||
self,
|
||||
store: TaskStore,
|
||||
executor: PromptExecutor,
|
||||
on_event: EventCallback,
|
||||
tick_seconds: float = DEFAULT_TICK_SECONDS,
|
||||
) -> None:
|
||||
self._store = store
|
||||
self._executor = executor
|
||||
self._on_event = on_event
|
||||
self._tick_seconds = tick_seconds
|
||||
self._loop_task: asyncio.Task[None] | None = None
|
||||
|
||||
def start(self) -> None:
|
||||
if self._loop_task is None:
|
||||
self._store.recover_running()
|
||||
self._loop_task = asyncio.create_task(self._loop())
|
||||
logger.info("Scheduler started")
|
||||
|
||||
async def stop(self) -> None:
|
||||
if self._loop_task is None:
|
||||
return
|
||||
self._loop_task.cancel()
|
||||
try:
|
||||
await self._loop_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._loop_task = None
|
||||
logger.info("Scheduler stopped")
|
||||
|
||||
async def _loop(self) -> None:
|
||||
while True:
|
||||
try:
|
||||
await self._tick()
|
||||
except Exception: # noqa: BLE001 - the scheduler must never die
|
||||
logger.exception("Scheduler tick failed")
|
||||
await asyncio.sleep(self._tick_seconds)
|
||||
|
||||
async def _tick(self) -> None:
|
||||
now = now_utc()
|
||||
for row in self._store.due(to_iso(now)):
|
||||
await self._execute(row)
|
||||
|
||||
async def _execute(self, row: dict[str, Any]) -> None:
|
||||
uid = row["uid"]
|
||||
self._store.update(uid, {"status": "running"})
|
||||
self._on_event("start", row, "")
|
||||
logger.info("Executing task uid=%s", uid)
|
||||
|
||||
try:
|
||||
result = await self._executor(row["prompt"])
|
||||
changes, finished = compute_followup(row, now_utc())
|
||||
changes["last_result"] = result
|
||||
changes["last_error"] = None
|
||||
except Exception as exc: # noqa: BLE001 - surfaced into the task record
|
||||
logger.exception("Task uid=%s crashed", uid)
|
||||
self._store.update(
|
||||
uid,
|
||||
{"status": "error", "last_error": str(exc), "last_run_at": to_iso(now_utc())},
|
||||
)
|
||||
self._on_event("error", row, str(exc))
|
||||
return
|
||||
|
||||
self._store.update(uid, changes)
|
||||
self._on_event("done" if not finished else "finished", row, result)
|
||||
@@ -0,0 +1,80 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import Any
|
||||
|
||||
import dataset
|
||||
|
||||
logger = logging.getLogger("devii.tasks.store")
|
||||
|
||||
TABLE = "devii_tasks"
|
||||
INDEXED_COLUMNS = (["owner_kind", "owner_id"], ["uid"], ["enabled"], ["next_run_at"], ["status"])
|
||||
|
||||
|
||||
def memory_db() -> Any:
|
||||
return dataset.connect("sqlite:///:memory:")
|
||||
|
||||
|
||||
class TaskStore:
|
||||
def __init__(self, db: Any, owner_kind: str, owner_id: str) -> None:
|
||||
self._db = db
|
||||
self._owner_kind = owner_kind
|
||||
self._owner_id = owner_id
|
||||
self._ensure_indexes()
|
||||
|
||||
def _ensure_indexes(self) -> None:
|
||||
if TABLE not in self._db.tables:
|
||||
return
|
||||
table = self._db[TABLE]
|
||||
for columns in INDEXED_COLUMNS:
|
||||
table.create_index(columns)
|
||||
|
||||
@property
|
||||
def _table(self) -> Any:
|
||||
return self._db[TABLE]
|
||||
|
||||
@property
|
||||
def _scope(self) -> dict[str, str]:
|
||||
return {"owner_kind": self._owner_kind, "owner_id": self._owner_id}
|
||||
|
||||
def create(self, record: dict[str, Any]) -> None:
|
||||
self._table.insert({**record, **self._scope})
|
||||
logger.info("Task created uid=%s owner=%s/%s", record.get("uid"), self._owner_kind, self._owner_id)
|
||||
|
||||
def get(self, uid: str) -> dict[str, Any] | None:
|
||||
return self._table.find_one(uid=uid, **self._scope)
|
||||
|
||||
def list(self, enabled_only: bool = False, status: str | None = None) -> list[dict[str, Any]]:
|
||||
criteria: dict[str, Any] = dict(self._scope)
|
||||
if enabled_only:
|
||||
criteria["enabled"] = True
|
||||
if status:
|
||||
criteria["status"] = status
|
||||
return list(self._table.find(**criteria))
|
||||
|
||||
def update(self, uid: str, changes: dict[str, Any]) -> None:
|
||||
changes = {**changes, "uid": uid, **self._scope}
|
||||
self._table.update(changes, ["uid", "owner_kind", "owner_id"])
|
||||
logger.debug("Task updated uid=%s changes=%s", uid, list(changes))
|
||||
|
||||
def delete(self, uid: str) -> bool:
|
||||
deleted = self._table.delete(uid=uid, **self._scope)
|
||||
logger.info("Task deleted uid=%s ok=%s", uid, deleted)
|
||||
return bool(deleted)
|
||||
|
||||
def recover_running(self) -> int:
|
||||
if TABLE not in self._db.tables:
|
||||
return 0
|
||||
stuck = list(self._table.find(status="running", **self._scope))
|
||||
for row in stuck:
|
||||
self._table.update({"uid": row["uid"], "status": "pending", **self._scope},
|
||||
["uid", "owner_kind", "owner_id"])
|
||||
if stuck:
|
||||
logger.info("Recovered %d task(s) stuck in running", len(stuck))
|
||||
return len(stuck)
|
||||
|
||||
def due(self, now_iso: str) -> list[dict[str, Any]]:
|
||||
rows = self._table.find(enabled=True, status="pending", **self._scope)
|
||||
return [row for row in rows if row.get("next_run_at") and row["next_run_at"] <= now_iso]
|
||||
Reference in New Issue
Block a user