Files
agentci/tests/test_storage.py
T

152 lines
4.9 KiB
Python

import sqlite3
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
async def test_opencode_migration_preserves_and_tags_legacy_session_ids(tmp_path: Path) -> None:
legacy_migrations = tmp_path / "legacy-migrations"
legacy_migrations.mkdir()
migrations = Path(__file__).parents[1] / "src" / "agentci" / "migrations"
(legacy_migrations / "001_initial.sql").write_text(
(migrations / "001_initial.sql").read_text()
)
database = tmp_path / "legacy.sqlite3"
legacy = Storage(database, legacy_migrations)
await legacy.initialize()
with sqlite3.connect(database) as connection:
connection.execute(
"""
INSERT INTO workflows (
id, kind, repo_owner, repo_name, issue_number, base_sha,
workspace_path, primary_session_id, reviewer_session_id,
artifact, status, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
"legacy-workflow",
"plan",
"alice",
"repo",
3,
"abc",
str(tmp_path / "repo"),
"legacy-primary",
"legacy-reviewer",
"# Preserved plan",
"completed",
"2026-07-20T00:00:00+00:00",
"2026-07-20T00:00:00+00:00",
),
)
migrated = Storage(database, migrations)
await migrated.initialize()
loaded = await migrated.latest_workflow("alice", "repo", 3, WorkflowKind.PLAN)
assert loaded is not None
assert loaded.artifact == "# Preserved plan"
assert loaded.primary_session_id == "legacy-primary"
assert loaded.reviewer_session_id == "legacy-reviewer"
assert loaded.runtime == "codex"