Files
svcforge/services/worker/main.py
T
Nguyen Minh Phuc 66eb6cb0ee refactor: converge the patterns multiple authors left divergent
The codebase was written by several agents and had the same concept done more
than one way. This makes it read as one voice, with no behaviour change.

Dedup, each to a single canonical form:
  - INSTANCE_COLUMNS: the 13-column instances SELECT list existed as _COLUMNS in
    instances.py and reconcile.py (byte-identical) and inlined a third time in
    the worker. One exported constant now.
  - Settings.runtime_dsn: the three entrypoints each chose between str(pg_dsn)
    and pg_dsn.unicode_string(). One property.
  - yaml_tempfile: helm._ValuesFile and k8s._ManifestFile were the same
    write-yaml-to-a-temp-dir context manager. One helper in adapters/tempyaml.py.
  - services/_runtime.py: sleep_or_stop and install_stop_signals were copied
    between the worker and reconciler loops. One module, so shutdown behaviour
    cannot drift between them.
  - k8s.ensure_namespace used MANAGED_BY_LABEL/VALUE from helm.py instead of a
    hardcoded literal, so the managed-by label has one definition.
  - SvcforgeError is now the root of every svcforge exception (CatalogError,
    IllegalTransition, BadWindow, HandlerError), keeping each stdlib base in the
    MRO, so `except SvcforgeError` means what errors.py says it does.
  - ERROR_MAX_CHARS replaces the repeated `[-2000:]` truncation feeding the same
    error columns.
  - the reconciler reads settings.metrics_port like the worker, dropping its
    duplicate DEFAULT_METRICS_PORT and redundant --metrics-port option; the
    SVCFORGE_METRICS_PORT env override still applies through pydantic.

Two smaller correctness/consistency fixes:
  - RateLimitResult.retry_after_s computed its delta against datetime.now(UTC)
    while the limiter runs on an injectable clock, so it was meaningless under a
    FakeClock and drifted by request latency in production. It now carries a
    checked_at from the same clock as reset_at.
  - handle_provision's notifier.send is wrapped like the reconciler's: a flaky
    webhook after the READY CAS would fail the task, and the retry would hit the
    READY early-return and drop the notification, turning a good provision into a
    failed one.
2026-07-21 01:46:54 +00:00

198 lines
8.1 KiB
Python

"""The claim loop.
Poll every 5 seconds. Claim while a semaphore slot is free. Run the handler. Report.
That is the whole design, and the restraint is the point: LISTEN/NOTIFY would shave the
latency, is fire-and-forget so it can never replace the poll anyway, is strictly extra
code, and does not exist on pgbouncer's transaction pooler. The poll is not a placeholder
for something better.
"""
from __future__ import annotations
import asyncio
import time
from collections.abc import Awaitable
from opentelemetry import trace
from services._runtime import install_stop_signals, sleep_or_stop
from services.worker.deps import WorkerDeps
from services.worker.handlers import HANDLERS
from svcforge_core import obs
from svcforge_core.adapters.clock import SystemClock
from svcforge_core.adapters.helm import HelmProvisioner
from svcforge_core.adapters.notify import LogNotifier
from svcforge_core.domain.catalog import load_catalog
from svcforge_core.domain.models import Task, TaskKind
from svcforge_core.repo.db import make_pool
from svcforge_core.repo.instances import InstanceRepo
from svcforge_core.repo.tasks import TaskRepo
from svcforge_core.settings import Settings, load_settings
log = obs.get_logger("svcforge.worker")
async def _report(coro: Awaitable[bool], task_id: int, what: str) -> None:
"""Run a terminal report, and never let its failure escape.
Reporting is the one thing that must not kill the worker. `_run_one` runs inside a
TaskGroup, and a TaskGroup cancels every sibling the moment one child raises — so a
DB blip during `tasks.fail()` would abort every other in-flight provision on this pod,
not just this one. The task itself is safe either way: it stays `running` and the
reconciler's lease sweep returns it to the queue. Losing the report costs one lease
interval; losing the siblings costs their work.
"""
try:
if not await coro:
# The lease was stolen while we were working: another worker owns this task
# now and is mid-run. Reporting is theirs to do, not ours.
log.warning("lease lost before report; another worker owns this task", task_id=task_id)
except Exception:
log.exception("could not report task %s (%s); lease will expire", task_id, what)
async def _run_one(deps: WorkerDeps, task: Task, sem: asyncio.Semaphore) -> None:
"""Run one task to a terminal report. Never lets an exception escape the TaskGroup."""
worker_id = deps.settings.worker_id
try:
# Every log line from here carries instance_id/task_id/team. Bound once, at claim,
# rather than passed down: the alternative is threading three arguments through
# every function that might log, and the first one anyone forgets is the one you
# need at 3am.
obs.bind_task_context(task.instance_id, task.id, team=task.team or "unknown")
log.info("task claimed", kind=task.kind.value, attempt=task.attempts)
obs.TASKS_CLAIMED.labels(kind=task.kind.value).inc()
handler = HANDLERS.get(task.kind)
if handler is None:
await _report(
deps.tasks.fail(task.id, f"no handler for {task.kind}", worker_id, max_attempts=1),
task.id,
"no-handler",
)
return
# Re-parent to the span that enqueued this task. Without the stored traceparent
# the worker's span starts a brand-new trace, and the POST that caused the work
# is in a different trace to the helm call that did it.
ctx = obs.context_from_traceparent(task.traceparent)
started = time.monotonic()
with obs.tracer().start_as_current_span(
f"task.{task.kind.value}",
context=ctx,
kind=trace.SpanKind.CONSUMER,
) as span:
span.set_attribute("task.id", task.id)
span.set_attribute("task.kind", task.kind.value)
span.set_attribute("instance.id", str(task.instance_id))
try:
await handler(task, deps)
except asyncio.CancelledError:
# CancelledError inherits from BaseException, so `except Exception` below
# would never see it. Catch it only to release the claim, then ALWAYS
# re-raise: swallowing it breaks cancellation for everyone above us.
await _report(
deps.tasks.fail(task.id, "cancelled", worker_id, deps.settings.max_attempts),
task.id,
"cancelled",
)
raise
except Exception as exc:
log.exception("task failed", kind=task.kind.value, error=str(exc))
span.record_exception(exc)
span.set_status(trace.Status(trace.StatusCode.ERROR, str(exc)))
obs.TASK_ATTEMPTS_FAILED.labels(kind=task.kind.value).inc()
await _report(
deps.tasks.fail(task.id, str(exc), worker_id, deps.settings.max_attempts),
task.id,
"fail",
)
else:
# Only provisions go in the provision histogram. The buckets run 10s..1800s
# because they were sized for helm installs; a sub-second `verify` dropped
# into the same series drags the p95 down and quietly stops
# SvcforgeProvisionSlow from ever firing.
if task.kind is TaskKind.PROVISION:
obs.PROVISION_TIME.observe(time.monotonic() - started)
await _report(deps.tasks.complete(task.id, worker_id), task.id, "complete")
finally:
sem.release()
async def run_worker(deps: WorkerDeps, stop: asyncio.Event) -> None:
"""Claim and run until told to stop, then drain what is in flight.
Draining is what makes a rolling deploy invisible. Exiting the `async with` block
awaits every in-flight handler, so a pod that is being replaced finishes the provision
it already started instead of abandoning it half-done for the lease to clean up
five minutes later.
"""
sem = asyncio.Semaphore(deps.settings.worker_concurrency)
worker_id = deps.settings.worker_id
async with asyncio.TaskGroup() as tg:
while not stop.is_set():
await sem.acquire()
if stop.is_set():
sem.release()
break
try:
task = await deps.tasks.claim(worker_id)
except Exception:
# A DB blip must not kill the worker; back off and try again.
log.exception("claim failed")
sem.release()
await sleep_or_stop(stop, deps.settings.poll_interval_s)
continue
if task is None:
sem.release()
await sleep_or_stop(stop, deps.settings.poll_interval_s)
continue
tg.create_task(_run_one(deps, task, sem))
# TaskGroup.__aexit__ awaited the in-flight handlers. Now it is safe to exit 0.
async def _amain() -> None:
settings: Settings = load_settings()
# Before anything else: nothing logged above this line is structured, and the metrics
# the SvcforgeTaskFailed / SvcforgeProvisionSlow alerts query do not exist until the
# registry is up.
obs.setup("svcforge-worker", settings)
settings.check_production()
obs.start_metrics_server(settings.metrics_port)
pool = make_pool(settings.runtime_dsn, settings.pool_min_size, settings.pool_max_size)
await pool.open(wait=True)
deps = WorkerDeps(
pool=pool,
instances=InstanceRepo(pool),
tasks=TaskRepo(pool),
provisioner=HelmProvisioner(helm_bin=settings.helm_bin, timeout_s=int(settings.helm_timeout_s)),
notifier=LogNotifier(),
clock=SystemClock(),
catalog=load_catalog(settings.catalog_path),
settings=settings,
)
stop = asyncio.Event()
install_stop_signals(stop)
try:
await run_worker(deps, stop)
finally:
await pool.close()
def main() -> None:
"""One asyncio.run, at the top, never nested."""
asyncio.run(_amain())
if __name__ == "__main__":
main()