Files
svcforge/tests/integration/test_reconciler.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

569 lines
23 KiB
Python

"""The four checks, against a real Postgres and a fake cluster.
Integration, not unit, because there is nothing to unit test: each check is a query and a
transaction. The behaviour worth asserting — that the CAS and the insert commit together,
that the idempotency guard is a real `not exists`, that a second tick does not double-
enqueue — lives entirely in the part a mock would replace.
The cluster is faked; the database is not. FakeProvisioner is a dict of releases, which is
all the drift check needs: drift is "helm says X, the DB says Y", and a dict says X.
"""
from __future__ import annotations
from datetime import UTC, datetime, timedelta
from typing import Any
from uuid import UUID
import pytest
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
from services.reconciler.main import (
ReconcilerDeps,
check_drift,
check_lease_expiry,
check_ttl,
check_version_drift,
tick,
)
from svcforge_core.domain.models import CatalogEntry, SizeSpec, TaskKind, TaskState
from svcforge_core.domain.states import InstanceState
from svcforge_core.obs import RECONCILER_LAST_TICK
from svcforge_core.repo.db import DictPool
from svcforge_core.repo.instances import InstanceRepo
from svcforge_core.repo.reconcile import ReconcileRepo
from svcforge_core.repo.tasks import TaskRepo
from svcforge_core.settings import Settings
from tests.fakes import FakeClock, FakeNotifier, FakeProvisioner
from tests.integration.helpers import build_instance
NOW = datetime(2026, 7, 17, 12, 0, tzinfo=UTC) # a Friday
OLD_VERSION = "21.3.19"
NEW_VERSION = "21.3.20"
# '0 3 * * 0' is 03:00 Sunday. From a Friday noon that is always in the future, which is
# the whole assertion of the window test.
SUNDAY_0300_HCM = "0 3 * * 0|Asia/Ho_Chi_Minh"
def _entry(*, version: str = NEW_VERSION, security: bool = False) -> CatalogEntry:
return CatalogEntry(
service_type="elasticsearch",
chart="bitnamilegacy/elasticsearch",
chart_version=version,
security=security,
sizes={"small": SizeSpec(replicas=1, resources={})},
)
def _settings(**over: object) -> Settings:
return Settings(pg_dsn="postgresql://u:p@localhost:5432/db", **over) # type: ignore[arg-type]
def _deps(
pool: DictPool,
provisioner: FakeProvisioner,
*,
catalog: dict[str, CatalogEntry] | None = None,
max_in_flight: int = 1,
**settings_over: object,
) -> ReconcilerDeps:
return ReconcilerDeps(
pool=pool,
instances=InstanceRepo(pool),
tasks=TaskRepo(pool),
reconcile=ReconcileRepo(pool),
provisioner=provisioner,
notifier=FakeNotifier(),
clock=FakeClock(NOW),
catalog=catalog if catalog is not None else {"elasticsearch": _entry()},
settings=_settings(**settings_over),
own_team="platform",
max_in_flight=max_in_flight,
)
async def _seed(
pool: DictPool,
*,
state: InstanceState = InstanceState.READY,
team: str = "platform",
chart_version: str = OLD_VERSION,
expires_at: datetime | None = None,
maintenance_window: str | None = None,
) -> tuple[UUID, str, str]:
"""Insert one instance. Returns (id, release_name, namespace).
`expires_at` and `maintenance_window` go in with SQL rather than through
`InstanceRepo.create`: create() does not write `maintenance_window` at all (nothing but
the day-2 work list reads it), and that is a fact about the repo, not a gap in it.
"""
inst = build_instance(team=team, state=state, chart_version=chart_version)
repo = InstanceRepo(pool)
async with pool.connection() as conn:
await repo.create(conn, inst)
async with conn.cursor() as cur:
await cur.execute(
"update instances set expires_at = %s, maintenance_window = %s where id = %s",
(expires_at, maintenance_window, inst.id),
)
return inst.id, inst.release_name, inst.namespace
async def _tasks_for(pool: DictPool, instance_id: UUID) -> list[dict[str, Any]]:
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute(
"select id, kind, state, run_after, traceparent from tasks where instance_id = %s order by id",
(instance_id,),
)
return list(await cur.fetchall())
async def _state_of(pool: DictPool, instance_id: UUID) -> str:
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute("select state from instances where id = %s", (instance_id,))
row = await cur.fetchone()
assert row is not None
return str(row["state"])
# --- Check 1: drift ------------------------------------------------------------------------
async def test_drift_reprovisions_a_ready_instance_whose_release_vanished(
pool: DictPool,
) -> None:
"""`helm uninstall` by hand. Nobody sends an event; the next tick notices anyway.
This is the acceptance path from Module 7:
helm uninstall <release> -n <ns> && python -m services.reconciler.main --once
select kind, state from tasks order by id desc limit 1 -> provision | queued
"""
instance_id, _, _ = await _seed(pool)
deps = _deps(pool, FakeProvisioner()) # empty cluster: the release is gone
await check_drift(deps)
tasks = await _tasks_for(pool, instance_id)
assert [(t["kind"], t["state"]) for t in tasks] == [(TaskKind.PROVISION.value, TaskState.QUEUED.value)]
async def test_drift_leaves_the_instance_in_provisioning_not_failed(pool: DictPool) -> None:
"""The two-hop state change, and why it matters.
`handle_provision` returns early on a `ready` row and ends with a
`provisioning -> ready` CAS. Hand it anything else and helm runs but the bookkeeping
lands nowhere. So the reconciler must leave the row in `provisioning` — via `failed`,
because LEGAL has no `ready -> provisioning` edge — before the worker can claim it.
"""
instance_id, _, _ = await _seed(pool)
await check_drift(_deps(pool, FakeProvisioner()))
assert await _state_of(pool, instance_id) == InstanceState.PROVISIONING.value
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute("select error from instances where id = %s", (instance_id,))
row = await cur.fetchone()
assert row is not None
assert "drift" in row["error"] # the tenant gets told why, not just that
async def test_drift_ignores_an_instance_whose_release_is_present(pool: DictPool) -> None:
"""The happy path is the same code. It must enqueue nothing at all."""
instance_id, release, namespace = await _seed(pool)
provisioner = FakeProvisioner()
await provisioner.install(release, namespace, _entry(), {})
await check_drift(_deps(pool, provisioner))
assert await _tasks_for(pool, instance_id) == []
assert await _state_of(pool, instance_id) == InstanceState.READY.value
async def test_drift_is_idempotent_across_ticks(pool: DictPool) -> None:
"""Two ticks, one task. A provision takes minutes; ticks are 60 seconds apart.
Without the guard the second tick sees a `provisioning` row — no longer `ready`, so the
drift branch skips it. The guard is what covers the case where it is `ready` again
before the task is done.
"""
instance_id, _, _ = await _seed(pool)
deps = _deps(pool, FakeProvisioner())
await check_drift(deps)
await check_drift(deps)
assert len(await _tasks_for(pool, instance_id)) == 1
async def test_drift_never_deletes_an_orphan_release(pool: DictPool) -> None:
"""A release the DB has never heard of. Log it, bill nobody, delete nothing.
v1 policy, and it is a policy about evidence: "no row in this table" is not proof the
release is unowned. It might belong to another tool, another team, or a migration that
is half done. An operator deletes it after reading the log.
"""
provisioner = FakeProvisioner()
await provisioner.install("someone-elses-redis", "other-ns", _entry(), {})
await check_drift(_deps(pool, provisioner))
assert "someone-elses-redis" in provisioner.releases
assert provisioner.uninstalled == []
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute("select count(*) as n from tasks")
row = await cur.fetchone()
assert row is not None
assert row["n"] == 0
async def test_drift_does_not_call_an_in_flight_provision_an_orphan(pool: DictPool) -> None:
"""A `requested` instance a worker is installing right now is not an orphan.
`known_releases` covers every row in any state for exactly this reason. Scope it to
`ready` and every provision in progress gets reported as an orphan on every tick, which
trains everyone to ignore the orphan log.
"""
_, release, namespace = await _seed(pool, state=InstanceState.REQUESTED)
provisioner = FakeProvisioner()
await provisioner.install(release, namespace, _entry(), {})
await check_drift(_deps(pool, provisioner))
assert provisioner.uninstalled == []
# --- Check 2: lease expiry -----------------------------------------------------------------
async def test_lease_expiry_returns_a_dead_workers_task_to_the_queue(pool: DictPool) -> None:
"""SIGKILL leaves `running` with `locked_by` set and nobody running it.
No cleanup code in the worker can fix this, because the worker is the part that died.
The lease is the only thing that recovers the row.
"""
instance_id, _, _ = await _seed(pool)
tasks = TaskRepo(pool)
await tasks.enqueue_standalone(instance_id, TaskKind.PROVISION)
claimed = await tasks.claim("worker-that-is-about-to-die")
assert claimed is not None
# The worker died six minutes ago and never reported.
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute(
"update tasks set locked_at = now() - interval '6 minutes' where id = %s",
(claimed.id,),
)
await check_lease_expiry(_deps(pool, FakeProvisioner(), lease_seconds=300))
rows = await _tasks_for(pool, instance_id)
assert rows[0]["state"] == TaskState.QUEUED.value
async def test_lease_expiry_leaves_a_live_worker_alone(pool: DictPool) -> None:
"""A task claimed a second ago is not a dead worker. Reclaiming it would double-provision.
Handlers are idempotent, so a wrongly-freed lease is survivable, though it still costs a
duplicated helm run — which is why lease_seconds sits above helm's own --timeout.
"""
instance_id, _, _ = await _seed(pool)
tasks = TaskRepo(pool)
await tasks.enqueue_standalone(instance_id, TaskKind.PROVISION)
assert await tasks.claim("worker-1") is not None
await check_lease_expiry(_deps(pool, FakeProvisioner(), lease_seconds=300))
rows = await _tasks_for(pool, instance_id)
assert rows[0]["state"] == TaskState.RUNNING.value
# --- Check 3: TTL --------------------------------------------------------------------------
async def test_ttl_expired_instance_goes_to_deleting_with_a_deprovision_task(
pool: DictPool,
) -> None:
"""The check that stops a demo cluster from becoming a permanent line on the bill."""
instance_id, _, _ = await _seed(pool, expires_at=datetime.now(UTC) - timedelta(minutes=1))
await check_ttl(_deps(pool, FakeProvisioner()))
tasks = await _tasks_for(pool, instance_id)
assert [(t["kind"], t["state"]) for t in tasks] == [(TaskKind.DEPROVISION.value, TaskState.QUEUED.value)]
# `deleting` before the worker claims it: handle_deprovision ends with a
# `deleting -> deleted` CAS, and a `ready` row would leave the DB advertising an
# endpoint for a release helm has already removed.
assert await _state_of(pool, instance_id) == InstanceState.DELETING.value
async def test_ttl_ignores_an_instance_that_has_not_expired(pool: DictPool) -> None:
"""And ignores one with no expires_at at all: null means no TTL, not expired."""
live, _, _ = await _seed(pool, expires_at=datetime.now(UTC) + timedelta(hours=1))
forever, _, _ = await _seed(pool, expires_at=None)
await check_ttl(_deps(pool, FakeProvisioner()))
assert await _tasks_for(pool, live) == []
assert await _tasks_for(pool, forever) == []
async def test_ttl_is_idempotent_across_ticks(pool: DictPool) -> None:
"""A deprovision takes minutes and ticks are 60s apart. One task, not four."""
instance_id, _, _ = await _seed(pool, expires_at=datetime.now(UTC) - timedelta(minutes=1))
deps = _deps(pool, FakeProvisioner())
await check_ttl(deps)
await check_ttl(deps)
await check_ttl(deps)
assert len(await _tasks_for(pool, instance_id)) == 1
async def test_ttl_recovers_a_deleting_instance_whose_task_was_never_enqueued(
pool: DictPool,
) -> None:
"""The API's DELETE crashed between the CAS and the enqueue. This is the sweep it relies on.
That statement order is chosen *because* this check exists. The other order leaves a
deprovision task pointing at a `ready` instance, and a worker tears down a live service
nobody asked to delete.
"""
instance_id, _, _ = await _seed(pool, state=InstanceState.DELETING)
await check_ttl(_deps(pool, FakeProvisioner()))
tasks = await _tasks_for(pool, instance_id)
assert [t["kind"] for t in tasks] == [TaskKind.DEPROVISION.value]
assert await _state_of(pool, instance_id) == InstanceState.DELETING.value
async def test_ttl_re_enqueues_after_a_deprovision_exhausted_its_attempts(
pool: DictPool,
) -> None:
"""`done` and `failed` are not outstanding. A transient outage must not strand the row.
The guard asks "is one queued or running", not "has one ever existed" — otherwise a
deprovision that burned its five attempts during a cluster outage would leave the
instance billing forever with nothing left to retry it.
"""
instance_id, _, _ = await _seed(pool, state=InstanceState.DELETING)
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute(
"insert into tasks (instance_id, kind, state) values (%s, %s, %s)",
(instance_id, TaskKind.DEPROVISION.value, TaskState.FAILED.value),
)
await check_ttl(_deps(pool, FakeProvisioner()))
states = [t["state"] for t in await _tasks_for(pool, instance_id)]
assert TaskState.QUEUED.value in states
# --- Check 4: version drift ----------------------------------------------------------------
async def test_version_drift_enqueues_one_upgrade_for_the_own_team_instance_first(
pool: DictPool,
) -> None:
"""max_in_flight=1 across three stale instances, and it picks ours.
Eating your own dog food is an `order by`: we are the tenant who finds out the chart is
broken, and the halt stops the other two before they ever hear about it.
"""
await _seed(pool, team="payments")
await _seed(pool, team="search")
ours, _, _ = await _seed(pool, team="platform")
await check_version_drift(_deps(pool, FakeProvisioner()))
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute("select instance_id, kind from tasks")
rows = list(await cur.fetchall())
assert len(rows) == 1
assert rows[0]["instance_id"] == ours
assert rows[0]["kind"] == TaskKind.UPGRADE.value
async def test_version_drift_enqueues_nothing_while_the_rollout_is_halted(
pool: DictPool,
) -> None:
"""One column stops the fleet. A failed verify writes it; a human clears it with SQL."""
instance_id, _, _ = await _seed(pool)
async with pool.connection() as conn, conn.cursor() as cur:
await cur.execute(
"insert into catalog_versions (service_type, rollout_state) values ('elasticsearch', 'halted')"
)
await check_version_drift(_deps(pool, FakeProvisioner()))
assert await _tasks_for(pool, instance_id) == []
async def test_version_drift_ignores_an_instance_already_on_the_catalog_version(
pool: DictPool,
) -> None:
"""`chart_version` is written only after helm succeeds, which is what makes this the query."""
instance_id, _, _ = await _seed(pool, chart_version=NEW_VERSION)
await check_version_drift(_deps(pool, FakeProvisioner()))
assert await _tasks_for(pool, instance_id) == []
async def test_version_drift_parks_the_upgrade_until_the_maintenance_window(
pool: DictPool,
) -> None:
"""The queue does the waiting, in `where run_after <= now()`.
The waiting is done by the queue rather than by a scheduler or an in-memory timer: a
task parked in Postgres until 03:00 Sunday survives a reconciler restart. That is the
whole reason `run_after` exists.
"""
instance_id, _, _ = await _seed(pool, maintenance_window=SUNDAY_0300_HCM)
await check_version_drift(_deps(pool, FakeProvisioner()))
tasks = await _tasks_for(pool, instance_id)
assert len(tasks) == 1
assert tasks[0]["run_after"] > NOW # a Friday; the next 03:00 Sunday is days away
async def test_version_drift_bypasses_the_window_for_a_security_bump(pool: DictPool) -> None:
"""A CVE with a public exploit does not wait until Sunday. That is what `security:` is for."""
instance_id, _, _ = await _seed(pool, maintenance_window=SUNDAY_0300_HCM)
catalog = {"elasticsearch": _entry(security=True)}
await check_version_drift(_deps(pool, FakeProvisioner(), catalog=catalog))
tasks = await _tasks_for(pool, instance_id)
assert len(tasks) == 1
assert tasks[0]["run_after"] <= NOW
async def test_version_drift_is_idempotent_across_ticks(pool: DictPool) -> None:
"""The guard that makes max_in_flight mean anything.
The instance stays on the work list for the whole duration of its own upgrade —
`chart_version` is only written on success — and for the hours it spends parked waiting
for 03:00. Without the guard, `max_in_flight=1` is sixty tasks an hour against one
release.
"""
instance_id, _, _ = await _seed(pool, maintenance_window=SUNDAY_0300_HCM)
deps = _deps(pool, FakeProvisioner())
for _ in range(3):
await check_version_drift(deps)
assert len(await _tasks_for(pool, instance_id)) == 1
async def test_version_drift_skips_one_bad_window_and_keeps_going(pool: DictPool) -> None:
"""One tenant's typo must not freeze everyone else's security rollout."""
broken, _, _ = await _seed(pool, team="payments", maintenance_window="not a cron|Asia/Ho_Chi_Minh")
deps = _deps(pool, FakeProvisioner(), max_in_flight=5)
await check_version_drift(deps) # must not raise
assert await _tasks_for(pool, broken) == []
# --- The tick ------------------------------------------------------------------------------
class _AngryProvisioner(FakeProvisioner):
"""A cluster that cannot be reached. The drift check's worst day."""
async def list_releases(self) -> list[Any]:
raise RuntimeError("dial tcp: i/o timeout")
async def test_tick_runs_the_other_three_checks_when_one_blows_up(pool: DictPool) -> None:
"""A helm binary that cannot reach the API server must not stop TTLs from expiring.
This is the entire argument for wrapping each check independently, and it is asserted
rather than assumed because the failure mode — a tick that dies on check one — looks
exactly like a tick that found nothing to do.
"""
expired, _, _ = await _seed(pool, expires_at=datetime.now(UTC) - timedelta(minutes=1))
deps = _deps(pool, _AngryProvisioner())
await tick(deps) # must not raise
assert [t["kind"] for t in await _tasks_for(pool, expired)] == [TaskKind.DEPROVISION.value]
async def test_tick_sets_the_gauges_and_the_heartbeat(pool: DictPool) -> None:
"""queue_depth after the checks, not before, and the heartbeat unconditionally.
The heartbeat is what `SvcforgeReconcilerStale` reads. It answers "is the loop running",
not "is everything fine" — the checks have their own alerts, and an alert that means two
things gets muted.
"""
from prometheus_client import REGISTRY
await _seed(pool, expires_at=datetime.now(UTC) - timedelta(minutes=1))
deps = _deps(pool, _AngryProvisioner()) # one check fails; the heartbeat still ticks
await tick(deps)
assert REGISTRY.get_sample_value("svcforge_queue_depth") == 1.0
assert REGISTRY.get_sample_value("svcforge_instances", {"state": "deleting"}) == 1.0
assert RECONCILER_LAST_TICK._value.get() == pytest.approx(NOW.timestamp())
@pytest.fixture(scope="session")
def tracing() -> InMemorySpanExporter:
"""A real tracer provider for the process, collecting spans in memory.
Global because OTEL's is: `trace.set_tracer_provider` takes once per process, and
`obs.tracer()` resolves it at call time. Session-scoped so the second call never
happens from HERE.
The internals reset is load-bearing rather than cosmetic. `set_tracer_provider` is
one-shot: a second call logs "Overriding of current TracerProvider is not allowed" at
WARNING and is otherwise ignored. Any earlier test that builds a FastAPI app calls
`obs.setup()` and burns that one shot, after which this fixture silently installs
nothing, `get_finished_spans()` returns `[]`, and the failure reads as "the reconciler
stopped writing traceparents" rather than "another test got there first". Clearing the
module globals is the only way to take the shot back.
"""
trace._TRACER_PROVIDER = None
trace._TRACER_PROVIDER_SET_ONCE._done = False
exporter = InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(SimpleSpanProcessor(exporter))
trace.set_tracer_provider(provider)
return exporter
async def test_a_task_the_tick_enqueues_carries_the_ticks_traceparent(
pool: DictPool,
tracing: InMemorySpanExporter,
) -> None:
"""Nothing propagates a trace through a table. The column is written at insert or never.
Driven through `tick`, not through `check_drift` with a span wrapped around it by the
test — that version passed while the real entrypoint wrote null on every row, because
the only span in the production path (`helm.list`) had already closed by the time the
insert ran. A test that supplies the context under test proves the propagator works and
nothing about this service.
"""
instance_id, _, _ = await _seed(pool)
tracing.clear()
await tick(_deps(pool, FakeProvisioner()))
tasks = await _tasks_for(pool, instance_id)
traceparent = tasks[0]["traceparent"]
assert traceparent is not None, "the reconciler's own tasks are unjoinable to its tick"
# Same trace as the tick's span, which is the entire point of storing the column.
tick_spans = [s for s in tracing.get_finished_spans() if s.name == "reconciler.tick"]
assert len(tick_spans) == 1
assert traceparent.split("-")[1] == format(tick_spans[0].context.trace_id, "032x")