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
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.
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.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()
|