Files
Nguyen Minh Phuc 08a529fa63
ci / lint (push) Successful in 24s
ci / types (push) Successful in 34s
ci / unit (push) Successful in 26s
ci / security (push) Successful in 37s
ci / dockerfile (push) Successful in 6s
ci / chart (push) Successful in 7s
ci / integration (push) Successful in 41s
ci / image (api) (push) Successful in 2m9s
ci / image (reconciler) (push) Successful in 2m1s
ci / image (worker) (push) Successful in 2m7s
ci / bump (push) Successful in 13s
Get off Docker Hub, add a catalog values passthrough, fix the dind prune
Docker Hub rate-limits anonymous pulls per source IP and every node here
shares one NAT address, so a busy afternoon fails an unrelated build with
`toomanyrequests`. Nothing in this repo needs to be there.

Every base image now comes from mirror.gcr.io (python, alpine/helm,
postgres) or ghcr.io (uv, trivy). Verified digest-for-digest against
Docker Hub before switching, including the superseded postgres digest
this repo still pins, so every existing pin stays valid — same bytes,
different transport.

catalog.yaml: the three bitnami entries named `bitnamilegacy/<chart>`, a
repo alias nothing in the worker image configures, so they could never
resolve at provision time. All five entries are now `oci://` refs, which
need no `helm repo add`, and all are on latest stable:

  elasticsearch 21.3.15 -> 22.1.6     redis    20.6.2 -> 27.0.15
  postgresql    16.4.5  -> 18.8.0     podinfo  6.7.1  -> 6.14.0

Moving the chart pull is only half of it, though: a bitnami chart
defaults its own images to registry-1.docker.io. CatalogEntry gains a
`values:` dict, merged under the size's replicas and resources, so an
entry can set `global.imageRegistry` and move the image pull too. Size
wins on conflict — otherwise an entry setting replicaCount would make
every size deploy the same shape. Deep merge, because a shallow one
drops sibling keys of a shared nested map.

Bitnami charts reject a substituted registry unless
`global.security.allowInsecureImages` is set. That check is about
provenance, and the mirror serves byte-identical manifests, so it is set
deliberately and only for entries whose digests were verified.

The dind prune had `--filter until=168h` on both prunes, and it got both
cases exactly backwards. `until` reads an image's CREATED time, so it
deleted trivy every leg (a released tool image is always older than any
window) while protecting the dangling build layers it existed to remove.
Measured on node0: 21 dangling images / 5.96GB, and exactly 1 of them
older than 168h. Trivy is protected by a tag now, so the image prune
drops the filter; buildx keeps it, where age genuinely matters.

Tests: +10 unit (deep merge, precedence, no-mutation, and a guard that
fails if any catalog entry points at Docker Hub). Both new guards were
control-tested by breaking the code and watching them fail. The API test
that hardcoded `21.3.15` now reads the catalog — its subject is where
the value comes from, not what it is.
2026-07-21 15:32:12 +00:00

208 lines
8.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 _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,
}