This commit introduces a full-featured agent architecture including persistent memory storage using vector embeddings, multi-step planning capabilities with dynamic replanning, and an extensible tool registry supporting custom function definitions. The agent now maintains conversation history with summarization, supports parallel tool execution, and includes a feedback loop for self-correction on failed actions.
87 lines
2.7 KiB
Python
87 lines
2.7 KiB
Python
# 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)
|