Files
svcforge/services/worker/handlers.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

185 lines
7.7 KiB
Python

"""Task handlers.
Every handler obeys one rule: running it twice must equal running it once.
A worker can be SIGKILLed after helm installed the release but before the DB row says so;
the lease expires, another worker claims the same task, and this function runs again. A
handler that is not idempotent gives the tenant two Elasticsearches and you a bill.
Idempotency is bought in two places: a deterministic `release_name`, and adapters that
state desired state (`helm upgrade --install`) instead of issuing imperative commands.
"""
from __future__ import annotations
from collections.abc import Awaitable, Callable
from typing import Any
from services.worker.deps import WorkerDeps
from svcforge_core.domain.models import CatalogEntry, Instance, Task, TaskKind
from svcforge_core.domain.states import InstanceState
from svcforge_core.errors import SvcforgeError
from svcforge_core.obs import get_logger
from svcforge_core.repo.instances import INSTANCE_COLUMNS
log = get_logger("svcforge.worker")
class HandlerError(SvcforgeError, RuntimeError):
"""A task failed in a way worth retrying. The message lands in tasks.last_error."""
async def _load_instance(task: Task, deps: WorkerDeps) -> Instance:
async with deps.pool.connection() as conn, conn.cursor() as cur:
await cur.execute(
f"select {INSTANCE_COLUMNS} from instances where id = %s", # noqa: S608 - module constant
(task.instance_id,),
)
row = await cur.fetchone()
if row is None:
raise HandlerError(f"instance {task.instance_id} vanished")
return Instance.model_validate(row)
def _values_for(inst: Instance, entry: CatalogEntry) -> dict[str, Any]:
"""Turn a catalog size into helm values."""
size = entry.sizes.get(inst.size)
if size is None:
raise HandlerError(f"size {inst.size!r} not in catalog for {inst.service_type!r}")
return {"replicaCount": size.replicas, "resources": size.resources}
async def handle_provision(task: Task, deps: WorkerDeps) -> None:
"""Install the release and mark the instance ready. Idempotent."""
inst = await _load_instance(task, deps)
if inst.state is InstanceState.READY:
# A previous attempt already finished; the crash was after the work, before the
# bookkeeping. Nothing to do — and re-installing would be the bug.
return
entry = deps.catalog.get(inst.service_type)
if entry is None:
raise HandlerError(f"unknown service_type {inst.service_type!r}")
# Best-effort CAS. It returning False means someone else moved the row; the helm call
# below is idempotent either way, so this is bookkeeping, not a lock.
await deps.instances.update_state(inst.id, InstanceState.REQUESTED, InstanceState.PROVISIONING)
await deps.provisioner.install(
release=inst.release_name,
ns=inst.namespace,
entry=entry,
values=_values_for(inst, entry),
)
endpoint = f"http://{inst.release_name}.{inst.namespace}.svc.cluster.local"
ok = await deps.instances.update_state(
inst.id, InstanceState.PROVISIONING, InstanceState.READY, endpoint=endpoint
)
if ok:
try:
await deps.notifier.send(
"instance.ready",
f"instance {inst.id} is ready at {endpoint}",
{"instance_id": str(inst.id), "team": inst.team, "service_type": inst.service_type},
)
except Exception:
# The provision succeeded and the row is READY; the notification is a courtesy.
# Propagating a webhook timeout would fail the task, and the retry would hit the
# READY early-return and drop the notification anyway — so a flaky notifier
# would turn every provision into a "failed" task.
log.exception("notify.failed", instance_id=str(inst.id))
async def handle_deprovision(task: Task, deps: WorkerDeps) -> None:
"""Remove the release and mark the instance deleted. Idempotent."""
inst = await _load_instance(task, deps)
if inst.state is InstanceState.DELETED:
return
# `helm uninstall` of an already-gone release is not an error to us: the adapter
# swallows not-found, because the desired state — no release — is already true.
await deps.provisioner.uninstall(release=inst.release_name, ns=inst.namespace)
# Raise rather than ignore the CAS result. Swallowing it leaves the release gone, the
# row on `state=ready` with a dangling endpoint, the task marked done — and 60 seconds
# later the drift check re-provisions the thing the tenant asked to delete.
if not await deps.instances.update_state(inst.id, InstanceState.DELETING, InstanceState.DELETED):
raise HandlerError(
f"instance {inst.id} was {inst.state.value}, expected {InstanceState.DELETING.value}"
)
async def handle_upgrade(task: Task, deps: WorkerDeps) -> None:
"""Upgrade the release to the catalog's pinned version, then record it.
`instances.chart_version` is written only AFTER helm reports success. That column is
what the day-2 work-list query compares against, so writing it optimistically would
make the fleet look upgraded while it isn't.
"""
inst = await _load_instance(task, deps)
entry = deps.catalog.get(inst.service_type)
if entry is None:
raise HandlerError(f"unknown service_type {inst.service_type!r}")
if inst.chart_version == entry.chart_version:
return # already there
await deps.provisioner.install(
release=inst.release_name,
ns=inst.namespace,
entry=entry,
values=_values_for(inst, entry),
)
async with deps.pool.connection() as conn, conn.cursor() as cur:
await cur.execute(
"update instances set chart_version = %s, updated_at = now() where id = %s",
(entry.chart_version, inst.id),
)
async def handle_verify(task: Task, deps: WorkerDeps) -> None:
"""Post-upgrade health probe. On failure, halt the whole rollout for this service type.
The work-list query returns nothing while `rollout_state='halted'`, so a bad chart stops
after the first tenant instead of all of them. Clearing it is a deliberate SQL statement:
an automatic un-halt would resume breaking things.
"""
inst = await _load_instance(task, deps)
releases = {r.name for r in await deps.provisioner.list_releases()}
if inst.release_name in releases:
return
# `returning` plus a `where` on the update half says whether THIS call halted the
# rollout. The halt is idempotent; the page is not. Without the distinction, a verify
# that burns its full retry budget sends five identical notifications for one incident.
async with deps.pool.connection() as conn, conn.cursor() as cur:
await cur.execute(
"""insert into catalog_versions (service_type, rollout_state)
values (%s, 'halted')
on conflict (service_type) do update set rollout_state = 'halted'
where catalog_versions.rollout_state <> 'halted'
returning service_type""",
(inst.service_type,),
)
newly_halted = await cur.fetchone() is not None
if newly_halted:
await deps.notifier.send(
"rollout.halted",
f"rollout halted for {inst.service_type}: {inst.release_name} failed verify",
{"instance_id": str(inst.id), "team": inst.team, "service_type": inst.service_type},
)
raise HandlerError(f"verify failed for {inst.release_name}; rollout halted")
HANDLERS: dict[TaskKind, Callable[[Task, WorkerDeps], Awaitable[None]]] = {
TaskKind.PROVISION: handle_provision,
TaskKind.DEPROVISION: handle_deprovision,
TaskKind.UPGRADE: handle_upgrade,
TaskKind.VERIFY: handle_verify,
}