"""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.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.runtime import install_stop_signals, sleep_or_stop 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()