Files
svcforge/services/worker/main.py
T
Nguyen Minh Phuc c53734d2bc
ci / lint (push) Successful in 33s
ci / types (push) Successful in 43s
ci / unit (push) Successful in 32s
ci / security (push) Successful in 57s
ci / dockerfile (push) Successful in 7s
ci / chart (push) Successful in 8s
ci / integration (push) Successful in 55s
ci / image (api) (push) Successful in 3m39s
ci / image (reconciler) (push) Successful in 2m53s
ci / image (worker) (push) Successful in 2m14s
ci / bump (push) Successful in 16s
docs: add USER_GUIDE.md, tighten comments, fix CLI needing a DSN
The comment pass is prose-only: every distinct "why" is kept, the
narration around it is not. Verified by AST-comparing each changed file
against HEAD with docstrings stripped — only the two files below differ
in executable code.

Two real fixes fell out of the read-through:

* The CLI documented itself as never touching the database, then called
  load_settings(), which requires SVCFORGE_PG_DSN. It refused to start
  without a Postgres URL it never opens. It now has its own two-field
  ClientSettings; the orphaned api_url/api_token are dropped from
  Settings, where nothing else read them.
* repo/db.py had the DictRow alias comment and the ERROR_MAX_CHARS
  comment run together above the wrong symbol.

USER_GUIDE.md is the caller-facing guide the README only gestured at:
auth, catalog, every endpoint with curl, the lifecycle, the error table,
rate limiting, the CLI, client generation, an end-to-end poll loop.

It records two facts about the live deployment rather than documenting a
flow nobody can run. SVCFORGE_JWKS_URL points at a realm with no IdP
behind it, so the API logs "JWKS warm-up failed" at startup and every
/v1 request is a 401. And `helm repo list` in the worker returns no
repositories, so the three bitnamilegacy/ catalog entries cannot resolve
at provision time; only the oci:// entries can.

make lint clean, 76 unit + 111 integration tests pass.
2026-07-21 14:51:26 +00:00

191 lines
7.6 KiB
Python

"""The claim loop.
Poll every 5 seconds. Claim while a semaphore slot is free. Run the handler. Report. The
poll is not a placeholder for something better: LISTEN/NOTIFY would shave latency, but it
is fire-and-forget so it can never replace the poll, and it does not exist on pgbouncer's
transaction pooler.
"""
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.
`_run_one` runs inside a TaskGroup, which 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. The task itself is safe either way: it stays `running` and the 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, 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 threaded through every function that might log.
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 new trace, putting the POST that caused the work in a
# different trace from 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. Its buckets run 10s..1800s
# for helm installs, so a sub-second `verify` in 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` awaits every
in-flight handler, so a pod being replaced finishes the provision it started instead of
abandoning it 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 above this line logs structured, and the metrics the
# SvcforgeTaskFailed / SvcforgeProvisionSlow alerts query do not exist until it runs.
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()