feat: OpenMesh 基础平台与 MD/PDF 转换技能
- 后端: coworker 智能体框架, WS API, 文件上传, 附件处理 - 前端: Open WebUI, 文件全量走 upload API (含 MD/TXT/JSON 等文本类) - 技能: md-to-office (pandoc + wkhtmltopdf) - 修复: 上传文件路径丢失, Agent 搜索浪费, 输出文件跑到 uploads/ - 打包: PyInstaller one-dir, 预打包 pandoc/wkhtmltopdf/chromium
This commit is contained in:
128
coworker/automation/scheduler.py
Normal file
128
coworker/automation/scheduler.py
Normal file
@@ -0,0 +1,128 @@
|
||||
"""The scheduler loop — runs in the always-on server.
|
||||
|
||||
Policy (agreed): **run-once-catch-up** for runs missed while down (due tasks fire once on
|
||||
startup, then resume), and **skip-on-overlap** (don't stack a run if the previous is still
|
||||
going). The actual execution is injected as `runner(task, trigger) -> TaskRun` so this stays
|
||||
independent of the engine/manager.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import Awaitable, Callable, Optional
|
||||
|
||||
from .models import ScheduledTask, TaskRun
|
||||
from .store import TaskStore
|
||||
|
||||
logger = logging.getLogger("coworker.automation")
|
||||
|
||||
Runner = Callable[[ScheduledTask, str], Awaitable[TaskRun]]
|
||||
|
||||
|
||||
class Scheduler:
|
||||
def __init__(
|
||||
self,
|
||||
store: TaskStore,
|
||||
runner: Runner,
|
||||
*,
|
||||
tick_seconds: float = 30.0,
|
||||
extra_tick: Optional[Callable[[], Awaitable[None]]] = None,
|
||||
) -> None:
|
||||
self.store = store
|
||||
self.runner = runner
|
||||
self.tick_seconds = tick_seconds
|
||||
# An extra per-tick coroutine (self-wake resumption: resume sessions whose wakes are due).
|
||||
self.extra_tick = extra_tick
|
||||
self._task: Optional[asyncio.Task] = None
|
||||
self._running_ids: set[str] = set() # overlap guard
|
||||
self._spawned: set[asyncio.Task] = set() # keep spawned runs referenced
|
||||
|
||||
def start(self) -> None:
|
||||
if self._task is None:
|
||||
self._task = asyncio.create_task(self._loop())
|
||||
|
||||
async def stop(self) -> None:
|
||||
if self._task is not None:
|
||||
self._task.cancel()
|
||||
try:
|
||||
await self._task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._task = None
|
||||
# In-flight runs died with the loop before they were spawned; keep that shutdown
|
||||
# contract now that they're independent tasks (a suspended run must not outlive us).
|
||||
for spawned in list(self._spawned):
|
||||
spawned.cancel()
|
||||
try:
|
||||
await spawned
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._spawned.clear()
|
||||
|
||||
async def _loop(self) -> None:
|
||||
# First pass = run-once-catch-up for anything missed while the server was down.
|
||||
try:
|
||||
await self._tick(trigger="catchup")
|
||||
except Exception:
|
||||
logger.exception("scheduler catch-up failed")
|
||||
while True:
|
||||
await asyncio.sleep(self.tick_seconds)
|
||||
try:
|
||||
await self._tick(trigger="schedule")
|
||||
except Exception:
|
||||
logger.exception("scheduler tick failed")
|
||||
|
||||
async def _tick(self, *, trigger: str) -> None:
|
||||
for task in self.store.due():
|
||||
# Spawn, don't await: a run can suspend on a parked approval (standing
|
||||
# scoped approvals, §25) and one blocked automation must never stall the
|
||||
# scheduler loop, other due tasks, or self-wake resumption. The overlap
|
||||
# guard must be claimed *here*, before the spawn: this due() snapshot
|
||||
# goes stale, and if the in-flight run finishes before a spawned
|
||||
# duplicate gets its first step, a guard checked inside the spawn is
|
||||
# already clear — the task runs twice.
|
||||
if not self._claim(task.id):
|
||||
continue
|
||||
spawned = asyncio.create_task(self._run_claimed(task, trigger=trigger))
|
||||
self._spawned.add(spawned)
|
||||
spawned.add_done_callback(self._spawned.discard)
|
||||
if self.extra_tick is not None:
|
||||
try:
|
||||
await self.extra_tick()
|
||||
except Exception:
|
||||
logger.exception("scheduler extra_tick (wake resume) failed")
|
||||
|
||||
def _claim(self, task_id: str) -> bool:
|
||||
if task_id in self._running_ids: # skip-on-overlap
|
||||
logger.info("skipping %s — previous run still going", task_id)
|
||||
return False
|
||||
self._running_ids.add(task_id)
|
||||
return True
|
||||
|
||||
async def run_task(self, task: ScheduledTask, *, trigger: str) -> Optional[TaskRun]:
|
||||
if not self._claim(task.id):
|
||||
return None
|
||||
return await self._run_claimed(task, trigger=trigger)
|
||||
|
||||
async def _run_claimed(
|
||||
self, task: ScheduledTask, *, trigger: str
|
||||
) -> Optional[TaskRun]:
|
||||
try:
|
||||
run = await self.runner(task, trigger)
|
||||
except Exception as exc:
|
||||
logger.exception("task %s run failed", task.id)
|
||||
run = TaskRun(
|
||||
task_id=task.id, status="error", error=str(exc), trigger=trigger
|
||||
)
|
||||
self.store.add_run(run)
|
||||
finally:
|
||||
self._running_ids.discard(task.id)
|
||||
# advance the task (run_count/last_run) → save recomputes next_run.
|
||||
fresh = self.store.get(task.id)
|
||||
if fresh is not None:
|
||||
fresh.run_count += 1
|
||||
fresh.last_run = run.started_at if run else None
|
||||
fresh.last_status = run.status if run else "error"
|
||||
self.store.save(fresh)
|
||||
return run
|
||||
Reference in New Issue
Block a user