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
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.
208 lines
8.7 KiB
Python
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,
|
|
}
|