"""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 _deep_merge(base: dict[str, Any], override: dict[str, Any]) -> dict[str, Any]: """`override` wins, except where both sides hold a dict — then merge those too. Shallow `base | override` would be wrong the moment two layers touch different keys of the same nested map: `{"global": {"imageRegistry": ...}}` overridden by `{"global": {"storageClass": ...}}` silently drops the registry, and the pod pulls from somewhere nobody chose. """ out = dict(base) for key, value in override.items(): current = out.get(key) if isinstance(current, dict) and isinstance(value, dict): out[key] = _deep_merge(current, value) else: out[key] = value return out def _values_for(inst: Instance, entry: CatalogEntry) -> dict[str, Any]: """Catalog values, with the requested size's replicas and resources on top. Size last, deliberately. An entry that sets `replicaCount` in its own `values:` would otherwise beat the size the tenant actually asked for, and every size would deploy the same shape. """ 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 _deep_merge(entry.values, {"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, }