66eb6cb0ee
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.
198 lines
8.1 KiB
Python
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()
|