Files
svcforge/services/worker/main.py
T
Nguyen Minh Phuc e971e04d75
ci / lint (push) Successful in 28s
ci / types (push) Successful in 37s
ci / unit (push) Successful in 27s
ci / security (push) Successful in 38s
ci / dockerfile (push) Successful in 6s
ci / chart (push) Successful in 9s
ci / integration (push) Successful in 49s
ci / image (api) (push) Successful in 2m25s
ci / image (reconciler) (push) Successful in 2m45s
ci / image (worker) (push) Successful in 2m37s
ci / bump (push) Successful in 24s
fix: shared runtime helper must live where every image packages it
services/_runtime.py crashed the worker and reconciler on boot with
`ModuleNotFoundError: No module named 'services._runtime'`, while every unit and
integration gate was green and the API rolled out fine.

The cause is packaging, not code. Each service Dockerfile copies only its own
`services/<svc>/` subdir — `services/` itself is a namespace package with no
__init__.py, so a file added at the `services/` root is never copied into any
image. Tests import from the source tree, where the file exists, so nothing
below the image boundary could catch it. The API survived only because it does
not import the helper.

Moved to svcforge_core.runtime, which `COPY libs/ libs/` packages into every
image, next to adapters/tempyaml.py for the same reason.

Added an import smoke test to the image job: `docker run --entrypoint python
<image> -c "import services.<svc>.main"` loads the whole transitive graph inside
the built image and fails the build before the digest is pushed. This is the
one check the test suite structurally cannot perform — it runs against source,
the image is a different filesystem — and it is exactly the gap this bug fell
through.
2026-07-21 02:40:38 +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.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()