from pathlib import Path import pytest from agentci.adapters.storage import Storage from agentci.domain.models import ( Job, JobKind, JobStatus, Workflow, WorkflowKind, WorkflowStatus, ) @pytest.fixture async def storage(tmp_path: Path) -> Storage: migrations = Path(__file__).parents[1] / "src" / "agentci" / "migrations" value = Storage(tmp_path / "state.sqlite3", migrations) await value.initialize() return value def make_job(job_id: str = "job-1") -> Job: return Job( id=job_id, kind=JobKind.PLAN, target_key="alice/repo:issue:3", repo_owner="alice", repo_name="repo", issue_number=3, pr_number=None, requester="alice", message="", comment_id=10, ) async def test_enqueue_is_idempotent_and_claims_fifo(storage: Storage) -> None: assert await storage.enqueue("delivery-1", make_job()) assert not await storage.enqueue("delivery-1", make_job("job-2")) claimed = await storage.claim_next() assert claimed is not None assert claimed.id == "job-1" assert claimed.status is JobStatus.RUNNING assert await storage.claim_next() is None async def test_recovers_running_job_as_failed(storage: Storage) -> None: await storage.enqueue("delivery-1", make_job()) assert await storage.claim_next() is not None recovered = await storage.recover_running() assert [job.id for job in recovered] == ["job-1"] assert await storage.claim_next() is None async def test_persists_and_finds_workflows(storage: Storage, tmp_path: Path) -> None: workflow = Workflow( id="workflow-1", kind=WorkflowKind.PLAN, repo_owner="alice", repo_name="repo", issue_number=3, workspace_path=tmp_path / "repo", base_sha="abc", artifact="# Plan", status=WorkflowStatus.COMPLETED, ) await storage.create_workflow(workflow) loaded = await storage.latest_workflow("alice", "repo", 3, WorkflowKind.PLAN) assert loaded is not None assert loaded.artifact == "# Plan" assert loaded.workspace_path == tmp_path / "repo" async def test_tracks_operational_comments(storage: Storage) -> None: await storage.enqueue("delivery-1", make_job()) await storage.set_job_comment("job-1", "accepted_comment_id", 21) await storage.set_job_comment("job-1", "started_comment_id", 22) assert await storage.operational_comment_ids("alice", "repo", 3) == {21, 22} async def test_failed_followup_does_not_invalidate_completed_workflow( storage: Storage, tmp_path: Path ) -> None: workflow = Workflow( id="workflow-1", kind=WorkflowKind.PLAN, repo_owner="alice", repo_name="repo", issue_number=3, workspace_path=tmp_path / "repo", base_sha="abc", status=WorkflowStatus.COMPLETED, ) await storage.create_workflow(workflow) job = make_job() job.workflow_id = workflow.id await storage.enqueue("delivery-1", job) await storage.fail_job_workflow(job.id) loaded = await storage.latest_workflow("alice", "repo", 3, WorkflowKind.PLAN) assert loaded is not None assert loaded.status is WorkflowStatus.COMPLETED