"""TaskRepo: enqueue, complete, fail, backoff, leases.""" from __future__ import annotations from datetime import UTC, datetime, timedelta from svcforge_core.domain.models import TaskKind, TaskState 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 tests.integration.helpers import make_instance async def _task_row(pool: DictPool, task_id: int) -> dict[str, object]: async with pool.connection() as conn, conn.cursor() as cur: await cur.execute("select * from tasks where id = %s", (task_id,)) row = await cur.fetchone() assert row is not None return dict(row) async def test_enqueue_defaults_to_runnable_now(pool: DictPool) -> None: iid = await make_instance(pool) repo = TaskRepo(pool) tid = await repo.enqueue_standalone(iid, TaskKind.PROVISION) row = await _task_row(pool, tid) assert row["state"] == "queued" assert row["attempts"] == 0 async def test_complete_marks_done_and_releases_lock(pool: DictPool) -> None: iid = await make_instance(pool) repo = TaskRepo(pool) tid = await repo.enqueue_standalone(iid, TaskKind.PROVISION) claimed = await repo.claim("w1") assert claimed is not None await repo.complete(claimed.id) row = await _task_row(pool, tid) assert row["state"] == "done" assert row["locked_by"] is None async def test_fail_under_max_attempts_requeues_with_future_run_after(pool: DictPool) -> None: iid = await make_instance(pool) repo = TaskRepo(pool) tid = await repo.enqueue_standalone(iid, TaskKind.PROVISION) claimed = await repo.claim("w1") assert claimed is not None assert claimed.attempts == 1 # attempts increments at CLAIM time, not on failure await repo.fail(tid, "helm exploded", max_attempts=5) row = await _task_row(pool, tid) assert row["state"] == "queued" assert row["last_error"] == "helm exploded" assert row["locked_by"] is None # Backoff pushed it out; it must not be immediately runnable again. run_after = row["run_after"] assert isinstance(run_after, datetime) assert run_after >= datetime.now(UTC) - timedelta(seconds=1) async def test_fail_at_max_attempts_dead_letters_and_marks_instance(pool: DictPool) -> None: """A dead-letter state, not an infinite retry.""" iid = await make_instance(pool) tasks, instances = TaskRepo(pool), InstanceRepo(pool) tid = await tasks.enqueue_standalone(iid, TaskKind.PROVISION) async with pool.connection() as conn, conn.cursor() as cur: await cur.execute("update tasks set attempts = 5 where id = %s", (tid,)) await tasks.fail(tid, "chart not found", max_attempts=5) row = await _task_row(pool, tid) assert row["state"] == "failed" inst = await instances.get(iid, team="platform") assert inst is not None assert inst.state is InstanceState.FAILED assert inst.error == "chart not found" async def test_fail_does_not_resurrect_a_deleted_instance(pool: DictPool) -> None: """Raw SQL must obey the same state machine domain.transition() enforces. LEGAL[DELETED] is empty — deleted is terminal. A deprovision task that exhausts its retries after the instance is already gone must record nothing on it, not drag it back to 'failed'. """ iid = await make_instance(pool, state=InstanceState.DELETED) tasks, instances = TaskRepo(pool), InstanceRepo(pool) tid = await tasks.enqueue_standalone(iid, TaskKind.DEPROVISION) async with pool.connection() as conn, conn.cursor() as cur: await cur.execute("update tasks set attempts = 5 where id = %s", (tid,)) await tasks.fail(tid, "helm uninstall kept failing", max_attempts=5) # The task still dead-letters — that part is unconditional. row = await _task_row(pool, tid) assert row["state"] == "failed" # But the instance stays deleted. inst = await instances.get(iid, team="platform") assert inst is not None assert inst.state is InstanceState.DELETED, "raw SQL bypassed the state machine" assert inst.error is None async def test_fail_truncates_error_to_2kb(pool: DictPool) -> None: iid = await make_instance(pool) repo = TaskRepo(pool) tid = await repo.enqueue_standalone(iid, TaskKind.PROVISION) await repo.claim("w1") await repo.fail(tid, "x" * 9000, max_attempts=5) row = await _task_row(pool, tid) assert isinstance(row["last_error"], str) assert len(row["last_error"]) == 2000 async def test_run_after_in_the_future_is_not_claimable(pool: DictPool) -> None: iid = await make_instance(pool) repo = TaskRepo(pool) future = datetime.now(UTC) + timedelta(hours=1) await repo.enqueue_standalone(iid, TaskKind.PROVISION, run_after=future) assert await repo.claim("w1") is None async def test_reset_expired_leases_recovers_a_dead_workers_task(pool: DictPool) -> None: """No distributed lock survives a power cut. Only the lease recovers this row.""" iid = await make_instance(pool) repo = TaskRepo(pool) tid = await repo.enqueue_standalone(iid, TaskKind.PROVISION) claimed = await repo.claim("worker-that-will-die") assert claimed is not None # Simulate: the worker was SIGKILLed 10 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 '10 minutes' where id = %s", (tid,)) n = await repo.reset_expired_leases(lease_seconds=300) assert n == 1 row = await _task_row(pool, tid) assert row["state"] == TaskState.QUEUED.value assert row["locked_by"] is None # And it is claimable again. assert await repo.claim("w2") is not None async def test_fresh_lease_is_not_reset(pool: DictPool) -> None: iid = await make_instance(pool) repo = TaskRepo(pool) await repo.enqueue_standalone(iid, TaskKind.PROVISION) await repo.claim("w1") assert await repo.reset_expired_leases(lease_seconds=300) == 0