Files
svcforge/tests/integration/test_worker.py
T
Nguyen Minh Phuc c76154aeaa
ci / lint (push) Successful in 34s
ci / unit (push) Successful in 1m41s
ci / types (push) Successful in 1m41s
ci / dockerfile (push) Successful in 18s
ci / security (push) Successful in 1m27s
ci / chart (push) Failing after 1m11s
ci / integration (push) Successful in 1m10s
ci / image (api) (push) Has been skipped
ci / image (reconciler) (push) Has been skipped
ci / image (worker) (push) Has been skipped
ci / bump (push) Has been skipped
review: fix 26 findings from a 4-agent audit
CORRECTNESS
- lost-lease race: complete()/fail() did not check ownership, so a worker whose
  lease expired could mark a task done while another worker was running it, or
  requeue a task someone else owned. Reproduced, fixed with a CAS on
  (state, locked_by), pinned by two regression tests.
- worker died on report failure: _run_one's docstring claimed no exception
  escapes the TaskGroup; fail()/complete() were outside the guarded block, so a
  DB blip cancelled every sibling provision on the pod.
- claim query used an INNER join, which could strand a just-claimed task and
  report 'queue empty'. LEFT join.
- InstanceRepo.set_error bypassed the state machine and had no callers. Deleted.
- handle_deprovision ignored its CAS result, so a wrong-state instance kept a
  dangling endpoint and got re-provisioned by the drift check 60s later.
- handle_verify re-notified on every retry: five pages for one halt.

DEPLOY-BREAKING
- the migration Job could never succeed: no Dockerfile copied migrations/, and
  migrate.py resolved the path relative to the source tree, which only works for
  an editable install. Added COPY + SVCFORGE_MIGRATIONS_DIR.
- ServiceMonitor selector did not match the Service: API metrics never scraped.
- SvcforgeReconcilerStale fired permanently from every pod, because the gauge is
  module-level and every service exports it as 0. Scoped to the reconciler job.
- SvcforgeTaskFailed latched forever on a monotonic counter. Now increase()[15m].
- the digest guard accepted the all-zeros placeholder.
- worker terminationGracePeriodSeconds was 60s against a 600s helm timeout.

DEAD CODE THAT SHOULD NOT HAVE BEEN
- adapters/k8s.py was never called, so tenant namespaces were never created and
  the first provision for a new team would fail. Wired into handle_provision.
- adapters/redis.py was never imported by any service. Rate limiting is now wired
  into the API, failing open.
- Settings.check_production() had no callers. Given an explicit environment and
  called from every entrypoint.

OBSERVABILITY
- the API never called obs.setup(): no JSON logs, no trace correlation, log_json
  silently inert.
- LogNotifier's structured fields were discarded by the stdlib->structlog bridge.
- bind_task_context cleared the 'service' binding for the life of every task.
- split tasks_failed into task_attempts_failed and tasks_dead_lettered.

SECURITY
- trivy correctly blocked the worker/reconciler images: helm 3.16.2 and kubectl
  1.31.2 carry CRITICAL Go stdlib CVEs. Bumped to helm 3.21.3 and kubectl 1.35.3,
  which also closes a four-minor skew against the v1.35.3 cluster.

TESTS THAT COULD NOT FAIL
- the concurrency cap test passed on a fully serial worker.
- the alert/metric cross-check asserted a hardcoded list instead of reading the
  chart, so it could not catch a rename on the chart side.
- fixed OTel tracer-provider pollution between test files.

DOCS
- ARCHITECTURE.md: mermaid diagrams, user stories, and the helm-vs-ArgoCD
  guarantee (verified with --dry-run=server).
- AGENTS.md + CLAUDE.md.
- prose sweep for back-and-forth phrasing across 19 files.
2026-07-18 12:13:49 +00:00

224 lines
8.3 KiB
Python

"""The worker loop: drains on SIGTERM, retries, stays idempotent.
These run in milliseconds against a FakeProvisioner. That is the payoff for putting helm
behind a Protocol in Module 5: the crash-safety properties are testable without a cluster.
"""
from __future__ import annotations
import asyncio
import time
from pathlib import Path
from uuid import UUID
import pytest
from services.worker.deps import WorkerDeps
from services.worker.main import run_worker
from svcforge_core.adapters.clock import SystemClock
from svcforge_core.domain.catalog import load_catalog
from svcforge_core.domain.models import TaskKind
from svcforge_core.domain.states import InstanceState
from svcforge_core.repo.db import DictPool
from svcforge_core.repo.instances import InstanceRepo
from svcforge_core.repo.tasks import TaskRepo
from svcforge_core.settings import Settings
from tests.fakes import FakeNotifier, FakeProvisioner
from tests.integration.helpers import build_instance
CATALOG = load_catalog(Path(__file__).resolve().parents[2] / "catalog.yaml")
def _settings(**over: object) -> Settings:
base: dict[str, object] = {
"pg_dsn": "postgresql://x:x@127.0.0.1:5432/x",
"worker_id": "w-test",
"worker_concurrency": 4,
"poll_interval_s": 0.05,
"max_attempts": 3,
}
base.update(over)
return Settings(**base) # type: ignore[arg-type]
class FakeNamespaceEnsurer:
"""Records the namespaces it was asked to create. Satisfies NamespaceEnsurer."""
def __init__(self) -> None:
self.ensured: list[str] = []
async def ensure_namespace(self, ns: str, labels: dict[str, str] | None = None) -> None:
self.ensured.append(ns)
def _deps(
pool: DictPool,
prov: FakeProvisioner,
notifier: FakeNotifier | None = None,
namespaces: FakeNamespaceEnsurer | None = None,
**over: object,
) -> WorkerDeps:
return WorkerDeps(
pool=pool,
instances=InstanceRepo(pool),
tasks=TaskRepo(pool),
provisioner=prov,
namespaces=namespaces or FakeNamespaceEnsurer(),
notifier=notifier or FakeNotifier(),
clock=SystemClock(),
catalog=CATALOG,
settings=_settings(**over),
)
async def _seed(pool: DictPool, **kw: object) -> tuple[str, int]:
inst = build_instance(**kw) # type: ignore[arg-type]
async with pool.connection() as conn:
await InstanceRepo(pool).create(conn, inst)
tid = await TaskRepo(pool).enqueue_standalone(inst.id, TaskKind.PROVISION)
return str(inst.id), tid
async def _state_of(pool: DictPool, tid: int) -> str:
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute("select state from tasks where id = %s", (tid,))
row = await cur.fetchone()
assert row is not None
return str(row["state"])
async def test_worker_provisions_and_marks_ready(pool: DictPool) -> None:
iid, tid = await _seed(pool)
prov = FakeProvisioner()
notifier = FakeNotifier()
stop = asyncio.Event()
worker = asyncio.create_task(run_worker(_deps(pool, prov, notifier), stop))
await asyncio.sleep(0.5)
stop.set()
await asyncio.wait_for(worker, timeout=5)
assert await _state_of(pool, tid) == "done"
inst = await InstanceRepo(pool).get(UUID(iid), team="platform")
assert inst is not None
assert inst.state is InstanceState.READY
assert inst.endpoint is not None
assert len(prov.installed) == 1
# Assert the notification fired, and that it happened on the FIRST attempt.
#
# Without this the handler could raise after marking the instance ready — the task
# requeues, the retry hits the idempotency early-return, and everything above still
# passes while the worker is quietly crashing on every provision. Idempotency is
# supposed to make crashes survivable, not invisible; asserting attempts==1 is what
# keeps a masked crash from reading as success.
assert notifier.events() == ["instance.ready"]
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute("select attempts from tasks where id = %s", (tid,))
row = await cur.fetchone()
assert row is not None
assert row["attempts"] == 1, "task was retried: the handler raised after doing the work"
async def test_sigterm_drains_in_flight(pool: DictPool) -> None:
"""Stop is requested mid-provision: the worker must FINISH the task, then exit.
Abandoning it would not lose the task — the lease would recover it — but only after
five minutes of a tenant watching 'provisioning'. Draining costs two seconds.
"""
_, tid = await _seed(pool)
prov = FakeProvisioner(delay=2.0)
stop = asyncio.Event()
started = time.monotonic()
worker = asyncio.create_task(run_worker(_deps(pool, prov), stop))
await asyncio.sleep(0.5)
stop.set() # mid-flight: the handler is still inside its 2s install
await asyncio.wait_for(worker, timeout=5)
elapsed = time.monotonic() - started
assert await _state_of(pool, tid) == "done", "worker abandoned an in-flight task"
assert elapsed >= 2.0, "worker returned before the in-flight task finished"
assert elapsed < 5.0
async def test_idle_worker_stops_promptly(pool: DictPool) -> None:
"""Nothing queued: stop must wake the poll sleep, not wait it out."""
stop = asyncio.Event()
worker = asyncio.create_task(run_worker(_deps(pool, FakeProvisioner(), poll_interval_s=5.0), stop))
await asyncio.sleep(0.2)
started = time.monotonic()
stop.set()
await asyncio.wait_for(worker, timeout=2)
assert time.monotonic() - started < 1.0, "stop did not interrupt the poll sleep"
async def test_failed_task_is_requeued_with_backoff(pool: DictPool) -> None:
_, tid = await _seed(pool)
prov = FakeProvisioner(fail_on={"platform-elasticsearch"})
stop = asyncio.Event()
worker = asyncio.create_task(run_worker(_deps(pool, prov), stop))
await asyncio.sleep(0.6)
stop.set()
await asyncio.wait_for(worker, timeout=5)
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute("select state, attempts, last_error from tasks where id = %s", (tid,))
row = await cur.fetchone()
assert row is not None
assert row["state"] == "queued" # requeued, not failed — attempts remain
assert row["attempts"] >= 1
assert row["last_error"]
async def test_provision_twice_installs_once(pool: DictPool) -> None:
"""The idempotency claim, executed.
Simulates the crash window: the release is installed and the instance is READY, but
the task got re-queued (worker died before reporting). Re-running must not re-install.
"""
from services.worker.handlers import handle_provision
inst = build_instance(state=InstanceState.REQUESTED)
async with pool.connection() as conn:
await InstanceRepo(pool).create(conn, inst)
tid = await TaskRepo(pool).enqueue_standalone(inst.id, TaskKind.PROVISION)
task = await TaskRepo(pool).claim("w1")
assert task is not None
prov = FakeProvisioner()
deps = _deps(pool, prov)
await handle_provision(task, deps)
await handle_provision(task, deps) # the redelivery
assert len(prov.installed) == 1, "second run re-installed: handler is not idempotent"
assert tid == task.id
@pytest.mark.parametrize("concurrency", [1, 4])
async def test_concurrency_cap_is_respected(pool: DictPool, concurrency: int) -> None:
"""The semaphore is what stops one worker from starting 200 helm processes."""
for _ in range(6):
await _seed(pool)
prov = FakeProvisioner(delay=0.2)
stop = asyncio.Event()
worker = asyncio.create_task(run_worker(_deps(pool, prov, worker_concurrency=concurrency), stop))
await asyncio.sleep(0.5)
stop.set()
await asyncio.wait_for(worker, timeout=10)
# BOTH bounds. `<= concurrency` alone passes on a worker whose semaphore is broken to 1,
# or that lost concurrency entirely — it only proves the cap is not exceeded, never that
# concurrency exists. With 6 tasks seeded and a 0.2s handler, a working worker reaches
# its cap exactly.
assert prov.max_concurrent <= concurrency, f"ran {prov.max_concurrent} at once, cap was {concurrency}"
assert prov.max_concurrent == concurrency, (
f"only reached {prov.max_concurrent} concurrent with a cap of {concurrency}: "
"the worker is not actually running tasks in parallel"
)