104 lines
3.2 KiB
Python
104 lines
3.2 KiB
Python
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
|