4a426dbe50
ci / lint (push) Successful in 25s
ci / unit (push) Successful in 59s
ci / types (push) Successful in 1m8s
ci / dockerfile (push) Successful in 13s
ci / chart (push) Successful in 8s
ci / security (push) Successful in 1m2s
ci / integration (push) Successful in 54s
ci / image (api) (push) Successful in 1m6s
ci / image (reconciler) (push) Successful in 3m6s
ci / image (worker) (push) Successful in 2m28s
ci / bump (push) Has been cancelled
trivy took the worker image 39 -> 18 -> 5 findings across two version bumps, and the last 5 (4x golang.org/x/net, 1x Go stdlib) live in kubectl v1.36.2 — the newest kubectl that exists. No release clears them; upstream has not rebuilt against the patched Go yet. Chasing the version further has no end. kubectl was in that image for exactly one call: ensure_namespace before helm. `helm upgrade --install --create-namespace` does the same thing, idempotently, as part of the install it already runs. So the binary goes, and its vendored CVEs go with it. read_secret had no callers. Tradeoff recorded: the namespace no longer gets an svcforge.io/team label, since --create-namespace makes a bare one. Nothing reads that label today.
215 lines
8.8 KiB
Python
215 lines
8.8 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 contextlib
|
|
import signal
|
|
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.settings import Settings, load_settings
|
|
|
|
log = obs.get_logger("svcforge.worker")
|
|
|
|
|
|
async def _sleep_or_stop(stop: asyncio.Event, seconds: float) -> None:
|
|
"""Sleep, but wake immediately on shutdown.
|
|
|
|
`await asyncio.sleep(5)` would make every SIGTERM cost up to five seconds of
|
|
Kubernetes waiting on terminationGracePeriod for no reason.
|
|
"""
|
|
with contextlib.suppress(TimeoutError):
|
|
await asyncio.wait_for(stop.wait(), timeout=seconds)
|
|
|
|
|
|
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.pg_dsn.unicode_string(), 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()
|
|
loop = asyncio.get_running_loop()
|
|
for sig in (signal.SIGTERM, signal.SIGINT):
|
|
# add_signal_handler, NOT signal.signal. signal.signal runs the handler at an
|
|
# arbitrary bytecode boundary on whatever thread the C-level handler lands on,
|
|
# and the event loop will not notice until its next timer fires. This one is
|
|
# loop-safe: the callback runs as a normal loop callback.
|
|
loop.add_signal_handler(sig, stop.set)
|
|
|
|
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()
|