refactor: no comments left - one-line module docstrings, contracts on public fields only; jobs/job.py; example config and README
This commit is contained in:
@@ -0,0 +1 @@
|
||||
"""Cron, webhook and event jobs, deferred injects, the subscription budget."""
|
||||
|
||||
@@ -0,0 +1,126 @@
|
||||
"""What a job is: its triggers, the budget it respects, and what a run may do."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timedelta
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import Awaitable, Callable, Coroutine
|
||||
|
||||
from beaver_gateway.conversations.distill import LineCap
|
||||
from beaver_gateway.conversations.injects import Priority
|
||||
from beaver_gateway.conversations.service import Conversations, DistillResult
|
||||
from beaver_gateway.jobs.scheduler import Scheduler
|
||||
from beaver_gateway.storage.models import Conversation
|
||||
|
||||
__all__ = ["Budget", "Job", "JobRun"]
|
||||
|
||||
_log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class Job:
|
||||
name: str
|
||||
run: Callable[[JobRun], Awaitable[None]]
|
||||
cron: str | None = None
|
||||
webhook: bool = False
|
||||
events: tuple[str, ...] = ()
|
||||
critical: bool = True
|
||||
dedupe: bool = True
|
||||
|
||||
@property
|
||||
def entrypoint(self) -> str:
|
||||
return f"job:{self.name}"
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class Budget:
|
||||
threshold: float = 0.7
|
||||
tokens: int | None = None
|
||||
window: timedelta = timedelta(hours=5)
|
||||
limit_window: str = "five_hour"
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class JobRun:
|
||||
job: Job
|
||||
trigger: str
|
||||
payload: dict[str, Any]
|
||||
scheduler: Scheduler
|
||||
|
||||
@property
|
||||
def conversations(self) -> Conversations:
|
||||
return self.scheduler.conversations
|
||||
|
||||
async def master(self) -> Conversation | None:
|
||||
masters = await self.conversations.find(kind="master", status="open", limit=1)
|
||||
return masters[0] if masters else None
|
||||
|
||||
async def inject_master(
|
||||
self, text: str, *, urgency: Priority = "normal", origin: str | None = None
|
||||
) -> bool:
|
||||
master = await self.master()
|
||||
if master is None:
|
||||
_log.error(
|
||||
"job %s: no open master, inject lost: %s", self.job.name, text[:200]
|
||||
)
|
||||
return False
|
||||
await self.conversations.inject(
|
||||
master, text, urgency=urgency, origin=origin or self.job.name
|
||||
)
|
||||
return True
|
||||
|
||||
async def spawn_job(
|
||||
self,
|
||||
*,
|
||||
agent: str,
|
||||
text: str,
|
||||
title: str | None = None,
|
||||
line_cap: LineCap | None = None,
|
||||
) -> Conversation:
|
||||
"""A headless job turn; ``line_cap`` bounces a rewrite past the cap."""
|
||||
return await self.conversations.spawn(
|
||||
kind="job",
|
||||
agent=agent,
|
||||
seed="brief",
|
||||
text=text,
|
||||
title=title,
|
||||
origin="job",
|
||||
flags={"line_cap": line_cap.as_flags()} if line_cap else None,
|
||||
)
|
||||
|
||||
async def close_idle(
|
||||
self,
|
||||
*,
|
||||
kind: str = "deep",
|
||||
days: int = 2,
|
||||
limit: int = 3,
|
||||
since: datetime | None = None,
|
||||
) -> list[DistillResult]:
|
||||
"""Close chats quiet for ``days``, at most ``limit`` per run."""
|
||||
out: list[DistillResult] = []
|
||||
for conv in await self.conversations.idle(
|
||||
kind=kind, days=days, since=since, limit=limit
|
||||
):
|
||||
try:
|
||||
out.append(
|
||||
await self.conversations.distill(conv, reason=f"idle {days}d")
|
||||
)
|
||||
except Exception: # noqa: BLE001
|
||||
_log.exception("closing idle %s failed", conv.external_id)
|
||||
return out
|
||||
|
||||
async def retry_in(self, delay: timedelta) -> None:
|
||||
await self.scheduler.trigger(
|
||||
self.job, self.payload, delay=delay, trigger=self.trigger
|
||||
)
|
||||
|
||||
async def rotate(self) -> list[Conversation]:
|
||||
rotation = self.scheduler.rotation
|
||||
return await rotation.tick() if rotation is not None else []
|
||||
|
||||
def background(self, coro: Coroutine[Any, Any, Any]) -> None:
|
||||
self.scheduler.background(coro)
|
||||
@@ -1,13 +1,7 @@
|
||||
"""Jobs and deferred injects on pgqueuer (§3.6, §4.5).
|
||||
"""Jobs and deferred injects on pgqueuer.
|
||||
|
||||
A job is a name, a handler and its triggers: a cron expression, the
|
||||
webhook ``/hooks/<name>``, gateway bus events. The executor is pgqueuer on
|
||||
the gateway's own Postgres (a dedicated autocommit connection for
|
||||
LISTEN/NOTIFY), so cron ticks, webhook deliveries and the ``schedule``
|
||||
tool's one-off injects all live in one table and survive a restart. A
|
||||
handler only queues work for a conversation and returns; the turn itself
|
||||
runs in the conversation's worker. Non-critical jobs step aside while the
|
||||
subscription window is past its threshold.
|
||||
A job is a name, a handler and its triggers (cron, webhook, bus events); a
|
||||
handler only queues work, and the turn itself runs in the conversation's worker.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -34,6 +28,7 @@ from starlette.responses import JSONResponse
|
||||
from starlette.routing import Route
|
||||
|
||||
from beaver_gateway.conversations.service import parse_at
|
||||
from beaver_gateway.jobs.job import Budget, Job, JobRun
|
||||
from beaver_gateway.storage.models import JobRunRecord
|
||||
|
||||
if TYPE_CHECKING:
|
||||
@@ -44,10 +39,9 @@ if TYPE_CHECKING:
|
||||
from pgqueuer.ports.driver import Driver
|
||||
from starlette.requests import Request
|
||||
|
||||
from beaver_gateway.conversations.distill import LineCap
|
||||
from beaver_gateway.conversations.injects import Priority
|
||||
from beaver_gateway.conversations.rotation import Rotation
|
||||
from beaver_gateway.conversations.service import Conversations, DistillResult
|
||||
from beaver_gateway.conversations.service import Conversations
|
||||
from beaver_gateway.storage.models import Conversation
|
||||
|
||||
__all__ = ["INJECT", "Budget", "Job", "JobRun", "LocalCron", "Scheduler", "next_run"]
|
||||
@@ -73,111 +67,6 @@ class LocalCron(ScheduleExecutor):
|
||||
return next_run(self.parameters.expression, self.tz)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class Job:
|
||||
name: str
|
||||
run: Callable[[JobRun], Awaitable[None]]
|
||||
cron: str | None = None
|
||||
webhook: bool = False
|
||||
events: tuple[str, ...] = ()
|
||||
critical: bool = True
|
||||
dedupe: bool = True
|
||||
|
||||
@property
|
||||
def entrypoint(self) -> str:
|
||||
return f"job:{self.name}"
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class Budget:
|
||||
threshold: float = 0.7
|
||||
tokens: int | None = None
|
||||
window: timedelta = timedelta(hours=5)
|
||||
limit_window: str = "five_hour"
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class JobRun:
|
||||
job: Job
|
||||
trigger: str
|
||||
payload: dict[str, Any]
|
||||
scheduler: Scheduler
|
||||
|
||||
@property
|
||||
def conversations(self) -> Conversations:
|
||||
return self.scheduler.conversations
|
||||
|
||||
async def master(self) -> Conversation | None:
|
||||
masters = await self.conversations.find(kind="master", status="open", limit=1)
|
||||
return masters[0] if masters else None
|
||||
|
||||
async def inject_master(
|
||||
self, text: str, *, urgency: Priority = "normal", origin: str | None = None
|
||||
) -> bool:
|
||||
master = await self.master()
|
||||
if master is None:
|
||||
_log.error(
|
||||
"job %s: no open master, inject lost: %s", self.job.name, text[:200]
|
||||
)
|
||||
return False
|
||||
await self.conversations.inject(
|
||||
master, text, urgency=urgency, origin=origin or self.job.name
|
||||
)
|
||||
return True
|
||||
|
||||
async def spawn_job(
|
||||
self,
|
||||
*,
|
||||
agent: str,
|
||||
text: str,
|
||||
title: str | None = None,
|
||||
line_cap: LineCap | None = None,
|
||||
) -> Conversation:
|
||||
"""A headless job turn; ``line_cap`` bounces a rewrite past the cap."""
|
||||
return await self.conversations.spawn(
|
||||
kind="job",
|
||||
agent=agent,
|
||||
seed="brief",
|
||||
text=text,
|
||||
title=title,
|
||||
origin="job",
|
||||
flags={"line_cap": line_cap.as_flags()} if line_cap else None,
|
||||
)
|
||||
|
||||
async def close_idle(
|
||||
self,
|
||||
*,
|
||||
kind: str = "deep",
|
||||
days: int = 2,
|
||||
limit: int = 3,
|
||||
since: datetime | None = None,
|
||||
) -> list[DistillResult]:
|
||||
"""§4.5: close chats quiet for ``days``, at most ``limit`` per run."""
|
||||
out: list[DistillResult] = []
|
||||
for conv in await self.conversations.idle(
|
||||
kind=kind, days=days, since=since, limit=limit
|
||||
):
|
||||
try:
|
||||
out.append(
|
||||
await self.conversations.distill(conv, reason=f"idle {days}d")
|
||||
)
|
||||
except Exception: # noqa: BLE001
|
||||
_log.exception("closing idle %s failed", conv.external_id)
|
||||
return out
|
||||
|
||||
async def retry_in(self, delay: timedelta) -> None:
|
||||
await self.scheduler.trigger(
|
||||
self.job, self.payload, delay=delay, trigger=self.trigger
|
||||
)
|
||||
|
||||
async def rotate(self) -> list[Conversation]:
|
||||
rotation = self.scheduler.rotation
|
||||
return await rotation.tick() if rotation is not None else []
|
||||
|
||||
def background(self, coro: Coroutine[Any, Any, Any]) -> None:
|
||||
self.scheduler.background(coro)
|
||||
|
||||
|
||||
class Scheduler:
|
||||
def __init__(
|
||||
self,
|
||||
@@ -212,8 +101,6 @@ class Scheduler:
|
||||
def job(self, name: str) -> Job | None:
|
||||
return self._jobs.get(name)
|
||||
|
||||
# ---- lifecycle -----------------------------------------------------
|
||||
|
||||
async def start(self) -> None:
|
||||
if self._driver is not None:
|
||||
queries = Queries(self._driver)
|
||||
@@ -274,8 +161,6 @@ class Scheduler:
|
||||
self._dispatch(job, trigger="event", payload=dict(event))
|
||||
)
|
||||
|
||||
# ---- dispatch ------------------------------------------------------
|
||||
|
||||
async def _dispatch(
|
||||
self, job: Job, *, trigger: str, payload: dict[str, Any]
|
||||
) -> None:
|
||||
@@ -350,8 +235,6 @@ class Scheduler:
|
||||
)
|
||||
return int(ids[0]) if ids and ids[0] is not None else None
|
||||
|
||||
# ---- deferred injects ----------------------------------------------
|
||||
|
||||
async def schedule(
|
||||
self,
|
||||
conv: Conversation,
|
||||
@@ -424,8 +307,6 @@ class Scheduler:
|
||||
origin="schedule",
|
||||
)
|
||||
|
||||
# ---- budget --------------------------------------------------------
|
||||
|
||||
async def utilization(self) -> float | None:
|
||||
now = datetime.now(UTC)
|
||||
values: list[float] = []
|
||||
@@ -446,8 +327,6 @@ class Scheduler:
|
||||
utilization = await self.utilization()
|
||||
return utilization is not None and utilization > self.budget.threshold
|
||||
|
||||
# ---- introspection -------------------------------------------------
|
||||
|
||||
async def snapshot(self) -> dict[str, Any]:
|
||||
crons: dict[str, PgSchedule] = {}
|
||||
queue: list[dict[str, Any]] = []
|
||||
@@ -534,8 +413,6 @@ class Scheduler:
|
||||
"payload": row.payload,
|
||||
}
|
||||
|
||||
# ---- http ----------------------------------------------------------
|
||||
|
||||
def app(self, authorize: Callable[[Request], Awaitable[Any]]) -> Starlette:
|
||||
async def hook(request: Request) -> JSONResponse:
|
||||
await authorize(request)
|
||||
@@ -549,8 +426,6 @@ class Scheduler:
|
||||
|
||||
return Starlette(routes=[Route("/{name}", hook, methods=["POST"])])
|
||||
|
||||
# ---- internals -----------------------------------------------------
|
||||
|
||||
def background(self, coro: Coroutine[Any, Any, Any]) -> None:
|
||||
self._spawn(coro)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user