mirror of
https://github.com/rennf93/roboco.git
synced 2026-08-03 07:23:24 +02:00
Three coupled claim-path bugs from the 2026-07-24 live incident, fixed at the shared root: - The edge-agnostic sequence bar phantom-held a task behind an unrelated, never-connected same-parent sibling that coincidentally shared a lower raw sequence (stamp_wave_sequence stamps from a partial per-task view). _claim_blocked_by_sequence now branches on is_batch_root_subtask: a MegaTask root-subtask (globally-computed Kahn wave, a deliberate staged-release barrier) keeps the strict rule unchanged; every other same-parent context routes through the pure sequence_blocker_id, which only blocks on a real transitive predecessor via dependency_ids UNIONED with completed_dependency_ids. A task with no same-parent dependency edge at all falls back to the raw bar unchanged (#452 preserved). - The hold surfaced as claim()'s bare None and was misdiagnosed by the verb runner as a concurrent-transition invalid_state. New sequence_hold_reason + a proactive _sequencing_claim_guard return a dedicated Envelope.sequence_held naming the blocker, on both the PENDING and NEEDS_REVISION reclaim paths. - give_me_work offered tasks the claim gate then rejected: both offer paths (list_pending_for_agent, _drop_dependency_held) now consult the bar via the exact claim predicate (is_pending_claim_blocked, extended to NEEDS_REVISION). Co-authored-by: Renn F <rennf93@users.noreply.github.com>
2818 lines
98 KiB
Python
2818 lines
98 KiB
Python
"""TaskService coverage — create + read + list query helpers.
|
|
|
|
The service has 60+ methods covering the full task lifecycle. This file
|
|
covers the read path and the simpler create/list helpers; lifecycle
|
|
transitions (claim, submit_for_qa, complete, ...) are exercised by the
|
|
existing v1-flow integration tests.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from datetime import UTC, datetime, timedelta
|
|
from typing import TYPE_CHECKING, Any, cast
|
|
from uuid import UUID, uuid4
|
|
|
|
import pytest
|
|
import pytest_asyncio
|
|
from roboco.db.tables import AgentTable, AuditLogTable, ProjectTable, TaskTable
|
|
from roboco.models import AgentRole, AgentStatus, Team
|
|
from roboco.models.base import (
|
|
Complexity,
|
|
TaskNature,
|
|
TaskStatus,
|
|
TaskType,
|
|
)
|
|
from roboco.models.task import TaskCreateRequest
|
|
from roboco.services.base import ConflictError
|
|
from roboco.services.task import SoftBlockInfo, TaskService, get_task_service
|
|
from sqlalchemy import select, text
|
|
|
|
if TYPE_CHECKING:
|
|
from collections.abc import AsyncIterator
|
|
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
|
|
@pytest_asyncio.fixture
|
|
async def task_setup(
|
|
db_session: AsyncSession,
|
|
) -> AsyncIterator[dict]:
|
|
agent = AgentTable(
|
|
id=uuid4(),
|
|
name="Dev",
|
|
slug=f"be-dev-{uuid4().hex[:8]}",
|
|
role=AgentRole.DEVELOPER,
|
|
team=Team.BACKEND,
|
|
status=AgentStatus.ACTIVE,
|
|
model_config={},
|
|
system_prompt="dev",
|
|
capabilities=[],
|
|
permissions={},
|
|
metrics={},
|
|
)
|
|
db_session.add(agent)
|
|
await db_session.flush()
|
|
project = ProjectTable(
|
|
id=uuid4(),
|
|
name="T-Proj",
|
|
slug=f"t-proj-{uuid4().hex[:8]}",
|
|
git_url="https://example.com/r.git",
|
|
assigned_cell=Team.BACKEND,
|
|
created_by=agent.id,
|
|
)
|
|
db_session.add(project)
|
|
await db_session.flush()
|
|
yield {
|
|
"svc": TaskService(db_session),
|
|
"agent_id": agent.id,
|
|
"project_id": project.id,
|
|
"db": db_session,
|
|
}
|
|
|
|
|
|
def _req(setup: dict[str, Any], **overrides: Any) -> TaskCreateRequest:
|
|
return TaskCreateRequest(
|
|
title=overrides.pop("title", "t"),
|
|
description=overrides.pop("description", "d"),
|
|
acceptance_criteria=overrides.pop("acceptance_criteria", ["ac"]),
|
|
team=overrides.pop("team", Team.BACKEND),
|
|
created_by=setup["agent_id"],
|
|
project_id=setup["project_id"],
|
|
task_type=overrides.pop("task_type", TaskType.CODE),
|
|
nature=overrides.pop("nature", TaskNature.TECHNICAL),
|
|
estimated_complexity=overrides.pop("estimated_complexity", Complexity.MEDIUM),
|
|
**overrides,
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Create / Get / Update / Delete
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_create_task(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
assert task.id is not None
|
|
assert task.status == TaskStatus.PENDING
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_create_task_with_explicit_backlog_status(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup, status=TaskStatus.BACKLOG))
|
|
assert task.status == TaskStatus.BACKLOG
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_returns_task(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
fetched = await svc.get(task.id)
|
|
assert fetched is not None
|
|
assert fetched.id == task.id
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.get(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_delete_returns_true_on_success(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
assert await svc.delete(task.id) is True
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_delete_returns_false_when_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.delete(uuid4()) is False
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# List queries
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_all(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
a = await svc.create(_req(task_setup, title="a"))
|
|
b = await svc.create(_req(task_setup, title="b"))
|
|
rows = await svc.list_all()
|
|
ids = {t.id for t in rows}
|
|
assert a.id in ids
|
|
assert b.id in ids
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_by_team(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
backend = await svc.create(_req(task_setup, team=Team.BACKEND))
|
|
rows = await svc.list_by_team(Team.BACKEND)
|
|
assert backend.id in {t.id for t in rows}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_by_assignee(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
aid = task_setup["agent_id"]
|
|
task = await svc.create(_req(task_setup, assigned_to=aid))
|
|
rows = await svc.list_by_assignee(aid)
|
|
assert task.id in {t.id for t in rows}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_by_status(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
pending = await svc.create(_req(task_setup))
|
|
rows = await svc.list_by_status(TaskStatus.PENDING)
|
|
assert pending.id in {t.id for t in rows}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_pending(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
pending = await svc.create(_req(task_setup))
|
|
rows = await svc.list_pending()
|
|
assert pending.id in {t.id for t in rows}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_blocked_empty(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_blocked(team=Team.BACKEND)
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_awaiting_qa(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_awaiting_qa(team=Team.BACKEND)
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_awaiting_docs(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_awaiting_docs(team=Team.BACKEND)
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_awaiting_pm_review(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_awaiting_pm_review(team=Team.BACKEND)
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_awaiting_ceo_approval(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_awaiting_ceo_approval()
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_count_by_status(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
await svc.create(_req(task_setup))
|
|
counts = await svc.count_by_status(team=Team.BACKEND)
|
|
assert isinstance(counts, dict)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_count_by_team(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
counts = await svc.count_by_team()
|
|
assert isinstance(counts, dict)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_active_count_for_agent(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
count = await svc.get_active_count(task_setup["agent_id"])
|
|
assert isinstance(count, int)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# list_recent_for_project — the prompter's history digest source
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_recent_for_project_orders_most_recent_activity_first(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
db = task_setup["db"]
|
|
now = datetime.now(UTC)
|
|
|
|
oldest = await svc.create(_req(task_setup, title="oldest"))
|
|
middle = await svc.create(_req(task_setup, title="middle"))
|
|
newest = await svc.create(_req(task_setup, title="newest"))
|
|
|
|
# Distinct activity dates: oldest only has created_at far in the past;
|
|
# middle was touched (updated_at) more recently; newest actually completed
|
|
# (completed_at wins over updated_at/created_at).
|
|
oldest.created_at = now - timedelta(days=10)
|
|
oldest.updated_at = None
|
|
middle.created_at = now - timedelta(days=9)
|
|
middle.updated_at = now - timedelta(days=5)
|
|
newest.created_at = now - timedelta(days=8)
|
|
newest.updated_at = now - timedelta(days=7)
|
|
newest.completed_at = now - timedelta(days=1)
|
|
await db.flush()
|
|
|
|
rows = await svc.list_recent_for_project(task_setup["project_id"])
|
|
ids = [t.id for t in rows]
|
|
assert ids.index(newest.id) < ids.index(middle.id) < ids.index(oldest.id)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_recent_for_project_respects_limit(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
db = task_setup["db"]
|
|
now = datetime.now(UTC)
|
|
|
|
tasks = [await svc.create(_req(task_setup, title=f"t{i}")) for i in range(3)]
|
|
for i, t in enumerate(tasks):
|
|
t.created_at = now - timedelta(days=10 - i) # t0 oldest, t2 newest
|
|
t.updated_at = None
|
|
await db.flush()
|
|
|
|
query_limit = 2
|
|
rows = await svc.list_recent_for_project(
|
|
task_setup["project_id"], limit=query_limit
|
|
)
|
|
assert len(rows) == query_limit
|
|
ids = [t.id for t in rows]
|
|
assert ids == [tasks[2].id, tasks[1].id]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_recent_for_project_scoped_to_project(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
in_scope = await svc.create(_req(task_setup, title="in-scope"))
|
|
|
|
other_project = ProjectTable(
|
|
id=uuid4(),
|
|
name="Other-Proj",
|
|
slug=f"other-proj-{uuid4().hex[:8]}",
|
|
git_url="https://example.com/other.git",
|
|
assigned_cell=Team.BACKEND,
|
|
created_by=task_setup["agent_id"],
|
|
)
|
|
db_session.add(other_project)
|
|
await db_session.flush()
|
|
other_req = _req(task_setup, title="other-project-task")
|
|
other_req.project_id = cast("UUID", other_project.id)
|
|
out_of_scope = await svc.create(other_req)
|
|
|
|
rows = await svc.list_recent_for_project(task_setup["project_id"])
|
|
ids = {t.id for t in rows}
|
|
assert in_scope.id in ids
|
|
assert out_of_scope.id not in ids
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_recent_for_project_excludes_cancelled(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
db = task_setup["db"]
|
|
live = await svc.create(_req(task_setup, title="live"))
|
|
cancelled = await svc.create(_req(task_setup, title="cancelled"))
|
|
cancelled.status = TaskStatus.CANCELLED
|
|
await db.flush()
|
|
|
|
rows = await svc.list_recent_for_project(task_setup["project_id"])
|
|
ids = {t.id for t in rows}
|
|
assert live.id in ids
|
|
assert cancelled.id not in ids
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Subtask hierarchy
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_subtasks_empty(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup))
|
|
subs = await svc.get_subtasks(parent.id)
|
|
assert subs == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_subtasks_returns_children(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup))
|
|
child = await svc.create(_req(task_setup, parent_task_id=parent.id))
|
|
subs = await svc.get_subtasks(parent.id)
|
|
assert child.id in {s.id for s in subs}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_all_descendants(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup))
|
|
child = await svc.create(_req(task_setup, parent_task_id=parent.id))
|
|
grandchild = await svc.create(_req(task_setup, parent_task_id=child.id))
|
|
descendants = await svc.get_all_descendants(parent.id)
|
|
desc_ids = {d.id for d in descendants}
|
|
assert child.id in desc_ids
|
|
assert grandchild.id in desc_ids
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_all_subtasks_terminal_when_no_subtasks(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup))
|
|
assert await svc.all_subtasks_terminal(parent.id) is True
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Agent lookups (gateway helpers)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_resolve_agent_id_for_uuid(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
aid = task_setup["agent_id"]
|
|
resolved = await svc.resolve_agent_id(str(aid))
|
|
assert resolved == aid
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_active_task_for_agent_returns_none(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.get_active_task_for_agent(task_setup["agent_id"]) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_paused_for_agent_empty(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_paused_for_agent(task_setup["agent_id"])
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_assigned_for_agent(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
aid = task_setup["agent_id"]
|
|
task = await svc.create(_req(task_setup, assigned_to=aid))
|
|
rows = await svc.list_assigned_for_agent(aid)
|
|
assert task.id in {t.id for t in rows}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_in_progress_or_claimed(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_in_progress_or_claimed()
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_strategic_for_board(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_strategic_for_board()
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_long_running_blocked(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_long_running_blocked()
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_awaiting_main_pm_all(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_awaiting_main_pm_all()
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_main_pm_agent_returns_optional(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
# Either None (no main_pm seeded) or an AgentTable (committed by a prior
|
|
# test that's leaked through rollback isolation).
|
|
result = await svc.main_pm_agent()
|
|
assert result is None or hasattr(result, "id")
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_agent_for_team_returns_none_when_unseeded(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
# Either None (no qa seeded) or an AgentTable (committed by a prior test
|
|
# that's leaked through rollback isolation — e.g. test_audit_real_query
|
|
# commits a CELL_PM/DEVELOPER for backend, and other heartbeat tests do
|
|
# similar). The contract here is just "the lookup runs and returns
|
|
# something compatible".
|
|
result = await svc.qa_agent_for_team(Team.BACKEND)
|
|
assert result is None or hasattr(result, "id")
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_documenter_for_team_returns_none_when_unseeded(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
result = await svc.documenter_for_team(Team.BACKEND)
|
|
assert result is None or hasattr(result, "id")
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_cell_pm_for_team_returns_none_when_unseeded(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
result = await svc.cell_pm_for_team(Team.BACKEND)
|
|
assert result is None or hasattr(result, "id")
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_agent_for_returns_view_for_known_agent(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
view = await svc.agent_for(task_setup["agent_id"])
|
|
assert view is not None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_agent_for_returns_none_for_unknown(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.agent_for(uuid4()) is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Update + progress + commits
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_update_modifies_fields(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
updated = await svc.update(task.id, title="renamed")
|
|
assert updated is not None
|
|
assert updated.title == "renamed"
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_update_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.update(uuid4(), title="x") is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_add_progress(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
updated = await svc.add_progress(
|
|
task.id, task_setup["agent_id"], "Working on it", percentage=50
|
|
)
|
|
_PCT = 50
|
|
assert updated is not None
|
|
assert len(updated.progress_updates) == 1
|
|
assert updated.progress_updates[0]["percentage"] == _PCT
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_add_progress_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.add_progress(uuid4(), task_setup["agent_id"], "msg") is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_add_checkpoint(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
updated = await svc.add_checkpoint(
|
|
task.id,
|
|
task_setup["agent_id"],
|
|
state_summary="halfway",
|
|
remaining_work=["finish API", "tests"],
|
|
)
|
|
assert updated is not None
|
|
assert len(updated.checkpoints) == 1
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_add_checkpoint_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert (
|
|
await svc.add_checkpoint(
|
|
uuid4(),
|
|
task_setup["agent_id"],
|
|
state_summary="x",
|
|
remaining_work=[],
|
|
)
|
|
is None
|
|
)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_add_commit(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
updated = await svc.add_commit(
|
|
task.id, hash="abc1234", message="Fix bug", agent_id=task_setup["agent_id"]
|
|
)
|
|
assert updated is not None
|
|
assert len(updated.commits) == 1
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_add_commit_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.add_commit(uuid4(), hash="x", message="y") is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Set plan + heartbeat + idle marking
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_set_plan(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
updated = await svc.set_plan(task.id, "step 1\nstep 2")
|
|
assert updated is not None
|
|
assert updated.plan == {"text": "step 1\nstep 2"}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_set_plan_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.set_plan(uuid4(), "x") is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_set_plan_accepts_dict(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
updated = await svc.set_plan(task.id, {"steps": ["a", "b"]})
|
|
assert updated is not None
|
|
assert updated.plan == {"steps": ["a", "b"]}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_evidence_inspected(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
await svc.mark_evidence_inspected(task.id)
|
|
refreshed = await svc.get(task.id)
|
|
assert refreshed is not None
|
|
assert refreshed.qa_evidence_inspected is True
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_agent_idle(task_setup: dict) -> None:
|
|
"""Idle marking just clears current_task_id; smoke test for completion."""
|
|
svc = task_setup["svc"]
|
|
await svc.mark_agent_idle(task_setup["agent_id"])
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# delete cascades to descendants
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_delete_cascades_to_descendants(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup))
|
|
child = await svc.create(_req(task_setup, parent_task_id=parent.id))
|
|
grandchild = await svc.create(_req(task_setup, parent_task_id=child.id))
|
|
await svc.delete(parent.id)
|
|
assert await svc.get(child.id) is None
|
|
assert await svc.get(grandchild.id) is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Build_substitute_update
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_build_substitute_update_for_pm_review(task_setup: dict) -> None:
|
|
"""Build the substitute update payload for AWAITING_PM_REVIEW transition."""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
update_data, pm_slug = await svc.build_substitute_update(
|
|
agent_id=task_setup["agent_id"],
|
|
task=task,
|
|
new_status=TaskStatus.AWAITING_PM_REVIEW,
|
|
reason="too complex",
|
|
details="needs PM input",
|
|
)
|
|
assert update_data["status"] == TaskStatus.AWAITING_PM_REVIEW.value
|
|
assert "[SUBSTITUTE]" in update_data["dev_notes"]
|
|
assert isinstance(pm_slug, str | type(None))
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Lifecycle: pause / resume
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pause_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.pause(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pause_returns_none_when_not_in_progress(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup)) # PENDING — can't pause.
|
|
assert await svc.pause(task.id) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pause_then_resume(task_setup: dict, db_session: AsyncSession) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
# Manually set IN_PROGRESS for testing pause→resume.
|
|
task.status = TaskStatus.IN_PROGRESS
|
|
await db_session.flush()
|
|
paused = await svc.pause(task.id)
|
|
assert paused is not None
|
|
assert paused.status == TaskStatus.PAUSED
|
|
resumed = await svc.resume(task.id)
|
|
assert resumed is not None
|
|
assert resumed.status == TaskStatus.IN_PROGRESS
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_resume_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.resume(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_resume_returns_none_when_not_paused(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup)) # PENDING — can't resume.
|
|
assert await svc.resume(task.id) is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Lifecycle: heartbeat
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_heartbeat_updates_timestamp(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.CLAIMED
|
|
await db_session.flush()
|
|
# No raise — just records the heartbeat.
|
|
await svc.heartbeat(task.id)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_heartbeat_for_missing_is_noop(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
# No raise even when the task doesn't exist.
|
|
await svc.heartbeat(uuid4())
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Lifecycle: cancel
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_cancel_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.cancel(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_cancel_pending_task(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
cancelled = await svc.cancel(task.id)
|
|
assert cancelled is not None
|
|
assert cancelled.status == TaskStatus.CANCELLED
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Activate (PM activation of backlog tasks)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_activate_succeeds_without_session(task_setup: dict) -> None:
|
|
"""Backlog activation no longer requires a linked discussion session —
|
|
the session subsystem is retired; coordination rides task state."""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup, status=TaskStatus.BACKLOG))
|
|
activated = await svc.activate(task.id, agent_role="cell_pm")
|
|
assert activated.status == TaskStatus.PENDING
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# List queries with team filter
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_blocked_for_team(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_blocked_for_team(Team.BACKEND)
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_blocked_all_teams(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_blocked_all_teams()
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_awaiting_pm_review_for_team(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_awaiting_pm_review_for_team(Team.BACKEND)
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_by_team_or_assignee(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_by_team_or_assignee(
|
|
team=Team.BACKEND, agent_id=task_setup["agent_id"]
|
|
)
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Status transitions returning None when status doesn't match
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_submit_for_verification_returns_none_for_missing(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.submit_for_verification(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_submit_for_verification_returns_none_when_not_in_progress(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup)) # PENDING.
|
|
assert await svc.submit_for_verification(task.id) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_submit_for_qa_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.submit_for_qa(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_submit_for_qa_returns_none_when_not_verifying(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
assert await svc.submit_for_qa(task.id) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pass_qa_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.pass_qa(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pass_qa_returns_none_when_invalid_status(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup)) # PENDING — can't pass QA.
|
|
assert await svc.pass_qa(task.id, agent_role="qa") is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_fail_qa_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.fail_qa(uuid4(), notes="x") is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_docs_complete_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.docs_complete(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_block_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.block(uuid4(), blocker_task_id=uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unblock_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.unblock(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_complete_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.complete(uuid4(), agent_id=task_setup["agent_id"]) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_submit_for_pm_review_returns_none_for_missing(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.submit_for_pm_review(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_pr_created_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.mark_pr_created(uuid4(), pr_number=1, pr_url="u") is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unclaim_for_agent_returns_none_for_missing(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.unclaim_for_agent(uuid4(), agent_id=task_setup["agent_id"]) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_resume_for_agent_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.resume_for_agent(uuid4(), agent_id=task_setup["agent_id"]) is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Lifecycle path: in-progress → verifying via submit_for_verification
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_submit_for_verification_flips_to_verifying(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.IN_PROGRESS
|
|
await db_session.flush()
|
|
result = await svc.submit_for_verification(task.id)
|
|
assert result is not None
|
|
assert result.status == TaskStatus.VERIFYING
|
|
assert result.self_verified is True
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# soft_block
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_soft_block_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
result = await svc.soft_block(
|
|
uuid4(),
|
|
SoftBlockInfo(reason="x", blocker_type="dep", what_needed="d"),
|
|
)
|
|
assert result is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_soft_block_in_progress_task(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.IN_PROGRESS
|
|
task.assigned_to = task_setup["agent_id"]
|
|
await db_session.flush()
|
|
blocked = await svc.soft_block(
|
|
task.id,
|
|
SoftBlockInfo(
|
|
reason="waiting on creds",
|
|
blocker_type="external",
|
|
what_needed="API key",
|
|
),
|
|
)
|
|
assert blocked is not None
|
|
assert blocked.status == TaskStatus.BLOCKED
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unblock_restores_to_in_progress(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.IN_PROGRESS
|
|
task.assigned_to = task_setup["agent_id"]
|
|
# A claimed, actively-worked task has a branch; with one, unblock resumes
|
|
# in_progress (a never-claimed/no-branch task returns to pending instead).
|
|
task.branch_name = "feature/backend/abc12345"
|
|
await db_session.flush()
|
|
await svc.soft_block(
|
|
task.id,
|
|
SoftBlockInfo(reason="x", blocker_type="ext", what_needed="y"),
|
|
)
|
|
unblocked = await svc.unblock(task.id)
|
|
assert unblocked is not None
|
|
assert unblocked.status == TaskStatus.IN_PROGRESS
|
|
# Owner restored into both fields so the dev dispatcher respawns it.
|
|
assert unblocked.assigned_to == task_setup["agent_id"]
|
|
assert unblocked.claimed_by == task_setup["agent_id"]
|
|
# A real resume with no fresh claim() call — unblock must flip the
|
|
# owner's fleet marker itself (mirrors _finalize_claim/_qa_or_doc_claim).
|
|
owner_row = await db_session.get(AgentTable, task_setup["agent_id"])
|
|
assert owner_row is not None
|
|
assert owner_row.status == AgentStatus.ACTIVE
|
|
assert owner_row.current_task_id == task.id
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unblock_keeps_owner_for_give_me_work_claim(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""A task claimed via give_me_work (no assigned_to) keeps its owner.
|
|
|
|
block() only stashes blocker_raised_by from assigned_to, so a give_me_work
|
|
claim (claimed_by set, assigned_to null) would otherwise unblock into a
|
|
split-owner state — assigned_to null but claimed_by set — which both the
|
|
dev dispatcher and the PM pool-router try to grab. unblock must put the
|
|
owner back into both fields.
|
|
"""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.IN_PROGRESS
|
|
task.assigned_to = None
|
|
task.claimed_by = task_setup["agent_id"]
|
|
task.branch_name = "feature/backend/abc12345"
|
|
await db_session.flush()
|
|
await svc.soft_block(
|
|
task.id,
|
|
SoftBlockInfo(reason="x", blocker_type="ext", what_needed="y"),
|
|
)
|
|
unblocked = await svc.unblock(task.id)
|
|
assert unblocked is not None
|
|
assert unblocked.status == TaskStatus.IN_PROGRESS
|
|
assert unblocked.assigned_to == task_setup["agent_id"]
|
|
assert unblocked.claimed_by == task_setup["agent_id"]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_admin_set_status_into_review_queue_releases_agent_marker(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""A non-blocked admin override into a review/queue state clears the
|
|
stale claimant's active_claimant_id (M19) — it must release that
|
|
claimant's fleet marker too, or a dead escalation claim keeps reporting
|
|
an agent as active forever."""
|
|
svc = task_setup["svc"]
|
|
dev_id = task_setup["agent_id"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.IN_PROGRESS
|
|
task.assigned_to = dev_id
|
|
task.claimed_by = dev_id
|
|
task.active_claimant_id = dev_id
|
|
await db_session.flush()
|
|
dev_agent = await db_session.get(AgentTable, dev_id)
|
|
assert dev_agent is not None
|
|
dev_agent.status = AgentStatus.ACTIVE
|
|
dev_agent.current_task_id = task.id
|
|
await db_session.flush()
|
|
|
|
out = await svc.admin_set_status(task.id, TaskStatus.AWAITING_QA)
|
|
assert out is not None
|
|
assert out.active_claimant_id is None
|
|
|
|
dev_agent = await db_session.get(AgentTable, dev_id)
|
|
assert dev_agent is not None
|
|
assert dev_agent.current_task_id is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_divert_owned_task_to_pool_idles_prior_owner(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""_divert_owned_task_to_pool clears ownership entirely — the refused
|
|
owner isn't engaged with this task at all anymore, so its fleet marker
|
|
must be released (mirrors _force_unclaim_to_pending's reaper release)."""
|
|
svc = task_setup["svc"]
|
|
owner_id = task_setup["agent_id"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.IN_PROGRESS
|
|
task.assigned_to = owner_id
|
|
task.claimed_by = owner_id
|
|
await db_session.flush()
|
|
owner_agent = await db_session.get(AgentTable, owner_id)
|
|
assert owner_agent is not None
|
|
owner_agent.status = AgentStatus.ACTIVE
|
|
owner_agent.current_task_id = task.id
|
|
await db_session.flush()
|
|
|
|
await svc._divert_owned_task_to_pool(task, note="test diversion")
|
|
|
|
assert task.status == TaskStatus.PENDING
|
|
assert task.assigned_to is None
|
|
owner_agent = await db_session.get(AgentTable, owner_id)
|
|
assert owner_agent is not None
|
|
assert owner_agent.status == AgentStatus.IDLE
|
|
assert owner_agent.current_task_id is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# QA + completion happy paths
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pass_qa_advances_to_awaiting_documentation(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.AWAITING_QA
|
|
task.pr_number = 42
|
|
task.pr_url = "https://github.com/x/y/pull/42"
|
|
await db_session.flush()
|
|
passed = await svc.pass_qa(task.id, notes="LGTM", agent_role="qa")
|
|
assert passed is not None
|
|
assert passed.status == TaskStatus.AWAITING_DOCUMENTATION
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_fail_qa_advances_to_needs_revision(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.AWAITING_QA
|
|
await db_session.flush()
|
|
failed = await svc.fail_qa(task.id, notes="please fix X")
|
|
assert failed is not None
|
|
assert failed.status == TaskStatus.NEEDS_REVISION
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pass_qa_clears_active_claimant_for_doc_claim(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
# M19 follow-on: pass_qa must clear the QA's active_claimant_id so the
|
|
# documenter's doc_claim isn't rejected by the competing-claimant guard.
|
|
# The direct REST route POST /pass-qa calls pass_qa directly, bypassing
|
|
# the qa_pass wrapper — the clear must live in the transition itself.
|
|
svc = task_setup["svc"]
|
|
qa_id = task_setup["agent_id"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.AWAITING_QA
|
|
task.pr_number = 42
|
|
task.pr_url = "https://github.com/x/y/pull/42"
|
|
await db_session.flush()
|
|
# qa_claim (via _qa_or_doc_claim) flips the QA agent's fleet marker —
|
|
# verify it, then verify pass_qa releases it symmetrically.
|
|
qa_claimed = await svc.qa_claim(qa_id, task.id)
|
|
assert qa_claimed is not None
|
|
qa_agent = await db_session.get(AgentTable, qa_id)
|
|
assert qa_agent is not None
|
|
assert qa_agent.status == AgentStatus.ACTIVE
|
|
assert qa_agent.current_task_id == task.id
|
|
passed = await svc.pass_qa(task.id, notes="LGTM", agent_role="qa")
|
|
assert passed is not None
|
|
assert passed.active_claimant_id is None
|
|
qa_agent = await db_session.get(AgentTable, qa_id)
|
|
assert qa_agent is not None
|
|
assert qa_agent.current_task_id is None
|
|
doc = AgentTable(
|
|
id=uuid4(),
|
|
name="Doc",
|
|
slug=f"be-doc-{uuid4().hex[:8]}",
|
|
role=AgentRole.DOCUMENTER,
|
|
team=Team.BACKEND,
|
|
status=AgentStatus.ACTIVE,
|
|
model_config={},
|
|
system_prompt="doc",
|
|
capabilities=[],
|
|
permissions={},
|
|
metrics={},
|
|
)
|
|
db_session.add(doc)
|
|
await db_session.flush()
|
|
claimed = await svc.doc_claim(doc.id, task.id)
|
|
assert claimed is not None
|
|
assert to_uuid(claimed.active_claimant_id) == doc.id
|
|
doc_agent = await db_session.get(AgentTable, doc.id)
|
|
assert doc_agent is not None
|
|
assert doc_agent.status == AgentStatus.ACTIVE
|
|
assert doc_agent.current_task_id == task.id
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_fail_qa_clears_active_claimant(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
# M19 follow-on: fail_qa must clear the QA's active_claimant_id. The
|
|
# direct REST route POST /fail-qa bypasses the qa_fail wrapper, so the
|
|
# clear must live in the transition itself.
|
|
svc = task_setup["svc"]
|
|
qa_id = task_setup["agent_id"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.AWAITING_QA
|
|
await db_session.flush()
|
|
assert await svc.qa_claim(qa_id, task.id) is not None
|
|
failed = await svc.fail_qa(task.id, notes="please fix X")
|
|
assert failed is not None
|
|
assert failed.active_claimant_id is None
|
|
# fail_qa releases the QA agent's fleet marker too.
|
|
qa_agent = await db_session.get(AgentTable, qa_id)
|
|
assert qa_agent is not None
|
|
assert qa_agent.current_task_id is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# pass_qa returns None when not in valid status
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_fail_qa_returns_none_when_not_awaiting_qa(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup)) # PENDING
|
|
assert await svc.fail_qa(task.id, notes="x") is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Reassign
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_reassign_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.reassign(uuid4(), task_setup["agent_id"]) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_reassign_updates_assigned_to(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
new_aid = task_setup["agent_id"]
|
|
reassigned = await svc.reassign(task.id, new_aid)
|
|
assert reassigned is not None
|
|
assert reassigned.assigned_to == new_aid
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# docs_complete + submit_for_pm_review happy paths
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_docs_complete_advances_status(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.AWAITING_DOCUMENTATION
|
|
task.assigned_to = task_setup["agent_id"]
|
|
task.pr_number = 42
|
|
task.pr_url = "https://github.com/x/y/pull/42"
|
|
await db_session.flush()
|
|
completed = await svc.docs_complete(task.id, doc_notes="Wrote docs")
|
|
assert completed is not None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Status transitions: claim path
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
result = await svc.claim(uuid4(), task_setup["agent_id"])
|
|
assert result is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# unclaim_for_reaper
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unclaim_for_reaper_no_op_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
# Should not raise.
|
|
await svc.unclaim_for_reaper(uuid4())
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# escalate_to_ceo
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_escalate_to_ceo_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
result = await svc.escalate_to_ceo(uuid4(), agent_role="cell_pm", notes="x")
|
|
assert result is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# ceo_approve / ceo_reject
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_ceo_approve_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.ceo_approve(uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_ceo_reject_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
result = await svc.ceo_reject(uuid4(), reason="not good enough")
|
|
assert result is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# mark_pr_created happy path + rejection paths
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_pr_created_advances_when_docs_complete(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.AWAITING_DOCUMENTATION
|
|
task.assigned_to = task_setup["agent_id"]
|
|
task.docs_complete = True
|
|
await db_session.flush()
|
|
result = await svc.mark_pr_created(
|
|
task.id, pr_number=42, pr_url="https://github.com/x/y/pull/42"
|
|
)
|
|
assert result is not None
|
|
assert result.pr_created is True
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_pr_created_rejects_terminal_status(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.COMPLETED
|
|
await db_session.flush()
|
|
assert await svc.mark_pr_created(task.id, pr_number=1, pr_url="u") is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# add_progress + add_checkpoint + add_commit happy paths
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_add_progress_appends_update(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
aid = task_setup["agent_id"]
|
|
updated = await svc.add_progress(task.id, aid, "step done", percentage=10)
|
|
assert updated is not None
|
|
assert len(updated.progress_updates) == 1
|
|
again = await svc.add_progress(task.id, aid, "step 2", percentage=20)
|
|
_AFTER = 2
|
|
assert again is not None
|
|
assert len(again.progress_updates) == _AFTER
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# resolve_agent_id paths
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_resolve_agent_id_for_slug(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
# task_setup created an agent with a slug; resolve via slug.
|
|
agent_row = (
|
|
await db_session.execute(
|
|
select(AgentTable).where(AgentTable.id == task_setup["agent_id"])
|
|
)
|
|
).scalar_one()
|
|
resolved = await svc.resolve_agent_id(agent_row.slug)
|
|
assert resolved == task_setup["agent_id"]
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# all_subtasks_terminal — one in flight returns False
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_all_subtasks_terminal_false_when_child_pending(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup))
|
|
await svc.create(_req(task_setup, parent_task_id=parent.id))
|
|
assert await svc.all_subtasks_terminal(parent.id) is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_all_subtasks_terminal_true_when_all_completed(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup))
|
|
child = await svc.create(_req(task_setup, parent_task_id=parent.id))
|
|
child.status = TaskStatus.COMPLETED
|
|
await db_session.flush()
|
|
assert await svc.all_subtasks_terminal(parent.id) is True
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# claim() happy path — skips git op via pre-set branch_name
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_pending_task_with_existing_branch(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""If branch_name is already set, claim skips _ensure_branch_for_task."""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
# Pre-set branch_name so claim doesn't try to do git ops.
|
|
task.branch_name = "feature/backend/abcd1234"
|
|
await db_session.flush()
|
|
claimed = await svc.claim(task.id, task_setup["agent_id"])
|
|
assert claimed is not None
|
|
assert claimed.status == TaskStatus.CLAIMED
|
|
assert claimed.assigned_to == task_setup["agent_id"]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_pending_with_unmet_dependency_returns_none(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""Sequence guardrail: a PENDING task with a non-terminal dependency cannot
|
|
be claimed — even via the raw claim path the orchestrator dispatcher uses.
|
|
Regression for the MegaTask "Main PM claimed every wave at once" bug."""
|
|
svc = task_setup["svc"]
|
|
dep = await svc.create(_req(task_setup, title="wave-0 dependency"))
|
|
task = await svc.create(_req(task_setup, title="wave-1 dependent"))
|
|
task.branch_name = "feature/backend/abcd1234"
|
|
await db_session.flush()
|
|
await svc.add_dependency(task.id, dep.id) # task depends_on dep (still PENDING)
|
|
|
|
assert await svc.claim(task.id, task_setup["agent_id"]) is None
|
|
|
|
# Once the dependency reaches a terminal state, the claim goes through.
|
|
dep.status = TaskStatus.COMPLETED
|
|
await db_session.flush()
|
|
claimed = await svc.claim(task.id, task_setup["agent_id"])
|
|
assert claimed is not None
|
|
assert claimed.status == TaskStatus.CLAIMED
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_sets_agent_active_and_current_task(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""_finalize_claim is the one production chokepoint every claim verb
|
|
routes through — before this fix, nothing ever wrote agent.status=ACTIVE
|
|
or current_task_id, so the fleet/Today-brief breakdown could never show
|
|
a real "active" agent or populate "working[]"."""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.branch_name = "feature/backend/abcd1234"
|
|
await db_session.flush()
|
|
|
|
claimed = await svc.claim(task.id, task_setup["agent_id"])
|
|
assert claimed is not None
|
|
|
|
agent = await db_session.get(AgentTable, task_setup["agent_id"])
|
|
assert agent is not None
|
|
assert agent.status == AgentStatus.ACTIVE
|
|
assert agent.current_task_id == task.id
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_agent_idle_clears_current_task(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""Idling an agent must clear current_task_id — otherwise it keeps
|
|
reporting the last-claimed task as "currently working" forever, since
|
|
nothing else ever clears the column."""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.branch_name = "feature/backend/abcd1234"
|
|
await db_session.flush()
|
|
await svc.claim(task.id, task_setup["agent_id"])
|
|
|
|
await svc.mark_agent_idle(task_setup["agent_id"])
|
|
|
|
agent = await db_session.get(AgentTable, task_setup["agent_id"])
|
|
assert agent is not None
|
|
assert agent.status == AgentStatus.IDLE
|
|
assert agent.current_task_id is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unclaim_for_agent_clears_current_task(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""A voluntary unclaim releases the claim marker on the AGENT too, not
|
|
just the task — otherwise the fleet keeps showing the agent as working
|
|
on a task it just gave up."""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.branch_name = "feature/backend/abcd1234"
|
|
await db_session.flush()
|
|
await svc.claim(task.id, task_setup["agent_id"])
|
|
|
|
result = await svc.unclaim_for_agent(task.id, task_setup["agent_id"])
|
|
assert result is not None
|
|
|
|
agent = await db_session.get(AgentTable, task_setup["agent_id"])
|
|
assert agent is not None
|
|
assert agent.current_task_id is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Sequence claim guardrail (CEO directive: sequence is the bar, independent
|
|
# of dependency_ids — see _claim_blocked_by_sequence).
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_blocked_by_lower_sequence_sibling_no_dependency_edge(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""A same-parent, lower-sequence, non-terminal sibling blocks a claim even
|
|
with NO dependency edge wired — the live 4-revision-subtask bug: a PM
|
|
delegated siblings with sequence 0..3 and no depends_on edges between
|
|
them, so a later sequence claimed alongside an earlier one."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
seq0 = await svc.create(
|
|
_req(task_setup, title="seq-0 sibling", parent_task_id=parent.id, sequence=0)
|
|
)
|
|
seq2 = await svc.create(
|
|
_req(task_setup, title="seq-2 sibling", parent_task_id=parent.id, sequence=2)
|
|
)
|
|
seq0.status = TaskStatus.IN_PROGRESS
|
|
await db_session.flush()
|
|
|
|
assert await svc.claim(seq2.id, task_setup["agent_id"]) is None
|
|
|
|
seq0.status = TaskStatus.COMPLETED
|
|
await db_session.flush()
|
|
seq2.branch_name = "feature/backend/abcd1234"
|
|
await db_session.flush()
|
|
claimed = await svc.claim(seq2.id, task_setup["agent_id"])
|
|
assert claimed is not None
|
|
assert claimed.status == TaskStatus.CLAIMED
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_allows_equal_sequence_siblings_in_parallel(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""Ties run in parallel: two siblings at the same sequence both claim."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
a = await svc.create(
|
|
_req(task_setup, title="tie-a", parent_task_id=parent.id, sequence=1)
|
|
)
|
|
b = await svc.create(
|
|
_req(task_setup, title="tie-b", parent_task_id=parent.id, sequence=1)
|
|
)
|
|
a.branch_name = "feature/backend/aaaa1111"
|
|
b.branch_name = "feature/backend/bbbb2222"
|
|
await db_session.flush()
|
|
|
|
claimed_a = await svc.claim(a.id, task_setup["agent_id"])
|
|
claimed_b = await svc.claim(b.id, task_setup["agent_id"])
|
|
assert claimed_a is not None
|
|
assert claimed_b is not None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_not_blocked_by_cancelled_lower_sequence_sibling(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""A cancelled (terminal) lower-sequence sibling never holds the claim."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
seq0 = await svc.create(
|
|
_req(task_setup, title="seq-0 cancelled", parent_task_id=parent.id, sequence=0)
|
|
)
|
|
seq1 = await svc.create(
|
|
_req(task_setup, title="seq-1", parent_task_id=parent.id, sequence=1)
|
|
)
|
|
seq0.status = TaskStatus.CANCELLED
|
|
seq1.branch_name = "feature/backend/cccc3333"
|
|
await db_session.flush()
|
|
|
|
claimed = await svc.claim(seq1.id, task_setup["agent_id"])
|
|
assert claimed is not None
|
|
assert claimed.status == TaskStatus.CLAIMED
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_blocked_by_null_sequence_sibling_coalesced_to_zero(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""COALESCE(sequence, 0): a NULL-sequence sibling is treated as sequence 0
|
|
and blocks a seq-1 claim, same as an explicit 0. `sequence` is NOT NULL in
|
|
the live schema — this defends the query anyway, matching the `sib_seq or
|
|
0` idiom `has_earlier_incomplete_code_sibling` already uses; the
|
|
constraint is relaxed transiently (rolled back at teardown) to exercise
|
|
it for real."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
null_seq = await svc.create(
|
|
_req(task_setup, title="null-seq sibling", parent_task_id=parent.id, sequence=0)
|
|
)
|
|
seq1 = await svc.create(
|
|
_req(task_setup, title="seq-1", parent_task_id=parent.id, sequence=1)
|
|
)
|
|
await db_session.execute(
|
|
text("ALTER TABLE tasks ALTER COLUMN sequence DROP NOT NULL")
|
|
)
|
|
await db_session.execute(
|
|
text("UPDATE tasks SET sequence = NULL WHERE id = :id"), {"id": null_seq.id}
|
|
)
|
|
await db_session.flush()
|
|
|
|
assert await svc.claim(seq1.id, task_setup["agent_id"]) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_sequence_zero_or_parentless_unaffected(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""Effective sequence 0 and a parentless task are never sequence-gated,
|
|
regardless of other non-terminal same-parent siblings."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
seq0 = await svc.create(
|
|
_req(task_setup, title="seq-0", parent_task_id=parent.id, sequence=0)
|
|
)
|
|
await svc.create(
|
|
_req(task_setup, title="seq-1 noise", parent_task_id=parent.id, sequence=1)
|
|
)
|
|
lone = await svc.create(_req(task_setup, title="parentless", sequence=5))
|
|
seq0.branch_name = "feature/backend/dddd4444"
|
|
lone.branch_name = "feature/backend/eeee5555"
|
|
await db_session.flush()
|
|
|
|
assert (await svc.claim(seq0.id, task_setup["agent_id"])) is not None
|
|
assert (await svc.claim(lone.id, task_setup["agent_id"])) is not None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_blocked_by_sequence_names_distinct_reason(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""`_claim_blocked_by_sequence` names the blocking sibling — a distinct
|
|
reason from `_claim_blocked_by_dependencies` (audit/logs tell them apart
|
|
since this task carries no dependency edge at all)."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
seq0 = await svc.create(
|
|
_req(task_setup, title="seq-0 blocker", parent_task_id=parent.id, sequence=0)
|
|
)
|
|
seq1 = await svc.create(
|
|
_req(task_setup, title="seq-1", parent_task_id=parent.id, sequence=1)
|
|
)
|
|
seq0.status = TaskStatus.IN_PROGRESS
|
|
await db_session.flush()
|
|
|
|
assert await svc._claim_blocked_by_dependencies(seq1) is False
|
|
reason = await svc._claim_blocked_by_sequence(seq1)
|
|
assert reason is not None
|
|
assert "seq-0 blocker" in reason
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# is_pending_claim_blocked — the public dispatch-time probe the orchestrator
|
|
# fetch filter uses to skip a doomed claim attempt (churn reduction).
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_is_pending_claim_blocked_true_for_lower_sequence_sibling(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
seq0 = await svc.create(
|
|
_req(task_setup, title="seq-0 blocker", parent_task_id=parent.id, sequence=0)
|
|
)
|
|
seq1 = await svc.create(
|
|
_req(task_setup, title="seq-1", parent_task_id=parent.id, sequence=1)
|
|
)
|
|
seq0.status = TaskStatus.IN_PROGRESS
|
|
await db_session.flush()
|
|
|
|
assert await svc.is_pending_claim_blocked(seq1.id) is True
|
|
|
|
seq0.status = TaskStatus.COMPLETED
|
|
await db_session.flush()
|
|
assert await svc.is_pending_claim_blocked(seq1.id) is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_is_pending_claim_blocked_true_for_unmet_dependency(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
dep = await svc.create(_req(task_setup, title="dependency"))
|
|
task = await svc.create(_req(task_setup, title="dependent"))
|
|
await db_session.flush()
|
|
await svc.add_dependency(task.id, dep.id)
|
|
|
|
assert await svc.is_pending_claim_blocked(task.id) is True
|
|
|
|
dep.status = TaskStatus.COMPLETED
|
|
await db_session.flush()
|
|
assert await svc.is_pending_claim_blocked(task.id) is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_is_pending_claim_blocked_false_for_clear_task(
|
|
task_setup: dict,
|
|
) -> None:
|
|
"""A ready task (no dependency edge, sequence 0, no parent) is never held."""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup, title="ready"))
|
|
assert await svc.is_pending_claim_blocked(task.id) is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_is_pending_claim_blocked_false_for_missing_task(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.is_pending_claim_blocked(uuid4()) is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_batch_wave_blocked_by_all_wave0_siblings_no_edges(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""STRICTER than dependency edges where both exist: a MegaTask wave-1
|
|
root-subtask waits for EVERY wave-0 root-subtask, not just the one edge
|
|
target the collision analyzer happened to wire — no edges-exist
|
|
exemption. Scoped to `batch_id`-bearing root-subtasks specifically
|
|
(`_build_confirm_batch` stamps `sequence` as a one-shot, globally
|
|
computed Kahn wave index — a deliberate staged-release barrier); a
|
|
plain (non-batch) same-parent sibling instead needs a real dependency
|
|
path — see `test_claim_not_blocked_by_unconnected_sibling_in_different_stream`."""
|
|
svc = task_setup["svc"]
|
|
umbrella = await svc.create(_req(task_setup, title="umbrella"))
|
|
batch_id = uuid4()
|
|
wave0_a = await svc.create(
|
|
_req(task_setup, title="wave0-a", parent_task_id=umbrella.id, sequence=0)
|
|
)
|
|
wave0_b = await svc.create(
|
|
_req(task_setup, title="wave0-b", parent_task_id=umbrella.id, sequence=0)
|
|
)
|
|
wave1 = await svc.create(
|
|
_req(task_setup, title="wave1", parent_task_id=umbrella.id, sequence=1)
|
|
)
|
|
# Only wave0_a gets an edge (the file-overlap-conditioned edge the
|
|
# analyzer wires) — wave0_b shares no surface with wave1, so production
|
|
# never wires an edge to it.
|
|
await svc.add_dependency(wave1.id, wave0_a.id)
|
|
wave0_a.status = TaskStatus.COMPLETED
|
|
# Direct ORM stamp (mirrors `_build_confirm_batch`'s BatchPlacement,
|
|
# bypassing create-time batch-shape validation the same way the rest of
|
|
# this test pokes `.status` directly).
|
|
wave0_a.batch_id = batch_id
|
|
wave0_b.batch_id = batch_id
|
|
wave1.batch_id = batch_id
|
|
await db_session.flush()
|
|
|
|
# The edge is satisfied, but wave0_b is still open with a lower sequence.
|
|
assert await svc._claim_blocked_by_dependencies(wave1) is False
|
|
assert await svc.claim(wave1.id, task_setup["agent_id"]) is None
|
|
|
|
wave0_b.status = TaskStatus.CANCELLED
|
|
wave1.branch_name = "feature/backend/ffff6666"
|
|
await db_session.flush()
|
|
claimed = await svc.claim(wave1.id, task_setup["agent_id"])
|
|
assert claimed is not None
|
|
assert claimed.status == TaskStatus.CLAIMED
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_not_blocked_by_unconnected_sibling_in_different_stream(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""The 2026-07-24 live incident: a frontend cell with 4 independent
|
|
dev-task streams — stream1-b (stamped wave 1, a real dependency on
|
|
stream1-a) must not be phantom-held by stream4-b (wave 0, in_progress)
|
|
just because they share a same parent and a lower raw sequence number.
|
|
Unlike the MegaTask-batch case above (no `batch_id` here), a plain
|
|
same-parent sibling only blocks via a real dependency path."""
|
|
svc = task_setup["svc"]
|
|
cell = await svc.create(_req(task_setup, title="cell"))
|
|
stream1_a = await svc.create(
|
|
_req(task_setup, title="stream1-a", parent_task_id=cell.id)
|
|
)
|
|
stream1_b = await svc.create(
|
|
_req(task_setup, title="stream1-b", parent_task_id=cell.id)
|
|
)
|
|
stream4_b = await svc.create(
|
|
_req(task_setup, title="stream4-b", parent_task_id=cell.id)
|
|
)
|
|
await svc.add_dependency(stream1_b.id, stream1_a.id)
|
|
await svc.stamp_wave_sequence(stream1_a.id)
|
|
await svc.stamp_wave_sequence(stream1_b.id)
|
|
await svc.stamp_wave_sequence(stream4_b.id)
|
|
assert (stream1_a.sequence, stream1_b.sequence, stream4_b.sequence) == (0, 1, 0)
|
|
|
|
# stream1-a (the REAL predecessor) is done; stream4-b (an unconnected
|
|
# sibling in a different stream) is still open with a lower sequence.
|
|
stream1_a.status = TaskStatus.COMPLETED
|
|
stream4_b.status = TaskStatus.IN_PROGRESS
|
|
stream1_b.branch_name = "feature/frontend/aaaa9999"
|
|
await db_session.flush()
|
|
|
|
claimed = await svc.claim(stream1_b.id, task_setup["agent_id"])
|
|
assert claimed is not None, (
|
|
"an unconnected sibling in a different stream must not phantom-hold"
|
|
)
|
|
assert claimed.status == TaskStatus.CLAIMED
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_not_blocked_after_real_dependency_pruned_on_completion(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""The precise drift mechanic: `_unblock_dependents` strips a completed
|
|
dependency from the live `dependency_ids` (moving it to
|
|
`completed_dependency_ids`) the moment it completes — almost always
|
|
BEFORE the dependent is ever claimed. The sequence claim-gate must
|
|
still recognize the pruned edge as real graph info (via
|
|
`completed_dependency_ids`) rather than treating the now-edge-less task
|
|
as a manually-sequenced, edge-less chain and reviving the raw
|
|
strictly-lower-sequence bar against an unrelated sibling."""
|
|
svc = task_setup["svc"]
|
|
cell = await svc.create(_req(task_setup, title="cell"))
|
|
real_predecessor = await svc.create(
|
|
_req(task_setup, title="real predecessor", parent_task_id=cell.id)
|
|
)
|
|
dependent = await svc.create(
|
|
_req(task_setup, title="dependent", parent_task_id=cell.id)
|
|
)
|
|
unrelated = await svc.create(
|
|
_req(task_setup, title="unrelated stream", parent_task_id=cell.id)
|
|
)
|
|
await svc.add_dependency(dependent.id, real_predecessor.id)
|
|
await svc.stamp_wave_sequence(real_predecessor.id)
|
|
await svc.stamp_wave_sequence(dependent.id)
|
|
await svc.stamp_wave_sequence(unrelated.id)
|
|
assert dependent.sequence == 1
|
|
|
|
# Simulate `_unblock_dependents`'s exact effect: the completed
|
|
# predecessor's edge is pruned from the live column and moved to the
|
|
# completed ledger.
|
|
real_predecessor.status = TaskStatus.COMPLETED
|
|
dependent.dependency_ids = []
|
|
dependent.completed_dependency_ids = [real_predecessor.id]
|
|
unrelated.status = TaskStatus.IN_PROGRESS
|
|
dependent.branch_name = "feature/frontend/bbbb8888"
|
|
await db_session.flush()
|
|
|
|
claimed = await svc.claim(dependent.id, task_setup["agent_id"])
|
|
assert claimed is not None, (
|
|
"a pruned-but-once-real edge must still count as graph info, not "
|
|
"revert to the raw edge-less sequence bar"
|
|
)
|
|
assert claimed.status == TaskStatus.CLAIMED
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_needs_revision_reclaim_blocked_by_lower_sequence_sibling(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""A needs_revision reclaim is sequence-gated too: a lower-sequence
|
|
sibling delegated AFTER the first claim must hold the reclaim (the gap a
|
|
PENDING-only scope left open). The dep guard's identical needs_revision
|
|
gap is pre-existing and deliberately untouched."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
low = await svc.create(
|
|
_req(task_setup, title="late-delegated seq-0", parent_task_id=parent.id)
|
|
)
|
|
high = await svc.create(
|
|
_req(
|
|
task_setup, title="seq-1 in revision", parent_task_id=parent.id, sequence=1
|
|
)
|
|
)
|
|
high.status = TaskStatus.NEEDS_REVISION
|
|
high.branch_name = "feature/backend/9999aaaa"
|
|
low.status = TaskStatus.IN_PROGRESS
|
|
await db_session.flush()
|
|
|
|
assert await svc.claim(high.id, task_setup["agent_id"]) is None
|
|
|
|
low.status = TaskStatus.COMPLETED
|
|
await db_session.flush()
|
|
assert await svc.claim(high.id, task_setup["agent_id"]) is not None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Wave stamping at delegation (stamp_wave_sequence) — sequences derive from
|
|
# the wired collision DAG, not a per-sibling ordinal.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_stamp_wave_sequence_independent_siblings_tie_and_claim_parallel(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""Independent siblings (no edges) share wave 0 — both claimable in
|
|
parallel. The raw ordinal gave them 0/1 and the claim gate then
|
|
serialized ALL delegated work fleet-wide."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
a = await svc.create(_req(task_setup, title="indep-a", parent_task_id=parent.id))
|
|
b = await svc.create(_req(task_setup, title="indep-b", parent_task_id=parent.id))
|
|
await svc.stamp_wave_sequence(a.id)
|
|
await svc.stamp_wave_sequence(b.id)
|
|
assert (a.sequence, b.sequence) == (0, 0)
|
|
|
|
a.branch_name = "feature/backend/aaaa0001"
|
|
b.branch_name = "feature/backend/bbbb0002"
|
|
await db_session.flush()
|
|
assert await svc.claim(a.id, task_setup["agent_id"]) is not None
|
|
assert await svc.claim(b.id, task_setup["agent_id"]) is not None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_stamp_wave_sequence_ascends_for_dependent_siblings(
|
|
task_setup: dict,
|
|
) -> None:
|
|
"""Colliding / ordered siblings ascend: each dependent stamps one wave
|
|
above the max of its same-parent dependency targets, and the claim gate
|
|
blocks the later wave while the earlier is open."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
t0 = await svc.create(_req(task_setup, title="wave-0", parent_task_id=parent.id))
|
|
t1 = await svc.create(_req(task_setup, title="wave-1", parent_task_id=parent.id))
|
|
t2 = await svc.create(_req(task_setup, title="wave-2", parent_task_id=parent.id))
|
|
await svc.stamp_wave_sequence(t0.id)
|
|
await svc.add_dependency(t1.id, t0.id)
|
|
await svc.stamp_wave_sequence(t1.id)
|
|
await svc.add_dependency(t2.id, t1.id)
|
|
await svc.stamp_wave_sequence(t2.id)
|
|
assert (t0.sequence, t1.sequence, t2.sequence) == (0, 1, 2)
|
|
|
|
assert await svc._claim_blocked_by_sequence(t1) is not None
|
|
assert await svc.claim(t1.id, task_setup["agent_id"]) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_stamp_wave_sequence_preserves_stamps_and_ignores_nonsibling_deps(
|
|
task_setup: dict,
|
|
) -> None:
|
|
"""Only the NEW task is stamped — an explicitly authored sibling sequence
|
|
(API create / prompter batch wave) is never rewritten — and a dependency
|
|
on a task under a DIFFERENT parent contributes no wave."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
authored = await svc.create(
|
|
_req(task_setup, title="authored seq-3", parent_task_id=parent.id, sequence=3)
|
|
)
|
|
outside = await svc.create(_req(task_setup, title="other-parent task"))
|
|
new = await svc.create(
|
|
_req(
|
|
task_setup,
|
|
title="new sibling",
|
|
parent_task_id=parent.id,
|
|
dependency_ids=[outside.id],
|
|
)
|
|
)
|
|
await svc.stamp_wave_sequence(new.id)
|
|
authored_stamp = 3
|
|
assert new.sequence == 0 # non-sibling dep excluded from the wave
|
|
assert authored.sequence == authored_stamp # PM-authored stamp untouched
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_stamp_wave_sequence_after_collision_wiring_matches_delegate_flow(
|
|
task_setup: dict,
|
|
) -> None:
|
|
"""The delegate-flow composition end to end: create → wire the collision
|
|
DAG → stamp. Disjoint-surface siblings tie at wave 0; a file-overlap
|
|
sibling lands one wave above the sibling it collides with."""
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup, title="parent"))
|
|
|
|
async def _delegate(title: str, surface: list[str]) -> Any:
|
|
t = await svc.create(
|
|
_req(
|
|
task_setup,
|
|
title=title,
|
|
parent_task_id=parent.id,
|
|
intends_to_touch=surface,
|
|
)
|
|
)
|
|
await svc.wire_sibling_collision_dag(parent.id)
|
|
await svc.stamp_wave_sequence(t.id)
|
|
return t
|
|
|
|
a = await _delegate("touches a.py", ["roboco/api/a.py"])
|
|
b = await _delegate("touches b.py", ["roboco/api/b.py"])
|
|
c = await _delegate("touches a.py too", ["roboco/api/a.py"])
|
|
|
|
assert (a.sequence, b.sequence) == (0, 0) # disjoint → parallel
|
|
assert c.sequence == 1 # collides with a → next wave
|
|
assert await svc._claim_blocked_by_sequence(c) is not None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_already_claimed_by_other_returns_none(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
other = AgentTable(
|
|
id=uuid4(),
|
|
name="Other",
|
|
slug=f"other-{uuid4().hex[:8]}",
|
|
role=AgentRole.DEVELOPER,
|
|
team=Team.BACKEND,
|
|
status=AgentStatus.ACTIVE,
|
|
model_config={},
|
|
system_prompt="x",
|
|
capabilities=[],
|
|
permissions={},
|
|
metrics={},
|
|
)
|
|
db_session.add(other)
|
|
await db_session.flush()
|
|
task = await svc.create(_req(task_setup))
|
|
task.branch_name = "feature/backend/abcd1234"
|
|
task.assigned_to = other.id
|
|
await db_session.flush()
|
|
assert await svc.claim(task.id, task_setup["agent_id"]) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_claim_with_allow_reassign_attempts(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
other = AgentTable(
|
|
id=uuid4(),
|
|
name="Other2",
|
|
slug=f"other2-{uuid4().hex[:8]}",
|
|
role=AgentRole.DEVELOPER,
|
|
team=Team.BACKEND,
|
|
status=AgentStatus.ACTIVE,
|
|
model_config={},
|
|
system_prompt="x",
|
|
capabilities=[],
|
|
permissions={},
|
|
metrics={},
|
|
)
|
|
db_session.add(other)
|
|
await db_session.flush()
|
|
task = await svc.create(_req(task_setup))
|
|
task.branch_name = "feature/backend/abcd1234"
|
|
task.assigned_to = other.id
|
|
await db_session.flush()
|
|
# With allow_reassign=True, the assignment-collision gate is bypassed.
|
|
result = await svc.claim(task.id, task_setup["agent_id"], allow_reassign=True)
|
|
# Either succeeds or fails for other reason — just verify it runs.
|
|
assert result is None or result is not None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Lifecycle: list queries with status filter
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_by_team_with_status(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_by_team(Team.BACKEND, status=TaskStatus.PENDING, limit=50)
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_by_assignee_with_status(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
rows = await svc.list_by_assignee(task_setup["agent_id"], status=TaskStatus.PENDING)
|
|
assert isinstance(rows, list)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# unclaim_for_agent + unclaim_for_reaper happy paths
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unclaim_for_agent_releases_claim(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.CLAIMED
|
|
task.assigned_to = task_setup["agent_id"]
|
|
task.claimed_by = task_setup["agent_id"]
|
|
await db_session.flush()
|
|
result = await svc.unclaim_for_agent(task.id, agent_id=task_setup["agent_id"])
|
|
assert result is not None
|
|
assert result.status == TaskStatus.PENDING
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unclaim_for_agent_releases_blocked_to_pool(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""An agent trapped on a `blocked` task can release it back to the pool
|
|
(returns to pending, assignment cleared) instead of churning with no move."""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.BLOCKED
|
|
task.assigned_to = task_setup["agent_id"]
|
|
task.claimed_by = task_setup["agent_id"]
|
|
await db_session.flush()
|
|
result = await svc.unclaim_for_agent(task.id, agent_id=task_setup["agent_id"])
|
|
assert result is not None
|
|
assert result.status == TaskStatus.PENDING
|
|
assert result.assigned_to is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unclaim_for_reaper_resets(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.CLAIMED
|
|
task.assigned_to = task_setup["agent_id"]
|
|
task.claimed_by = task_setup["agent_id"]
|
|
task.active_claimant_id = task_setup["agent_id"]
|
|
await db_session.flush()
|
|
await svc.unclaim_for_reaper(task.id)
|
|
refreshed = await svc.get(task.id)
|
|
assert refreshed is not None
|
|
assert refreshed.status == TaskStatus.PENDING
|
|
# Claim released but ownership preserved (no ownerless pending limbo).
|
|
assert refreshed.active_claimant_id is None
|
|
assert refreshed.assigned_to == task_setup["agent_id"]
|
|
assert refreshed.claimed_by == task_setup["agent_id"]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unclaim_for_reaper_idles_the_provably_dead_holder(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""The reaper's holder is provably dead (heartbeat past TTL) — it must
|
|
stop reporting ACTIVE with a stale current_task_id, or the fleet keeps
|
|
showing a dead agent as "working" until the eventual re-claim."""
|
|
svc = task_setup["svc"]
|
|
agent = await db_session.get(AgentTable, task_setup["agent_id"])
|
|
assert agent is not None
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.CLAIMED
|
|
task.assigned_to = task_setup["agent_id"]
|
|
task.claimed_by = task_setup["agent_id"]
|
|
task.active_claimant_id = task_setup["agent_id"]
|
|
agent.status = AgentStatus.ACTIVE
|
|
agent.current_task_id = task.id
|
|
await db_session.flush()
|
|
|
|
await svc.unclaim_for_reaper(task.id)
|
|
|
|
refreshed_agent = await db_session.get(AgentTable, task_setup["agent_id"])
|
|
assert refreshed_agent is not None
|
|
assert refreshed_agent.status == AgentStatus.IDLE
|
|
assert refreshed_agent.current_task_id is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# resume_for_agent
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_resume_for_agent_paused_task(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.PAUSED
|
|
task.assigned_to = task_setup["agent_id"]
|
|
await db_session.flush()
|
|
result = await svc.resume_for_agent(task.id, agent_id=task_setup["agent_id"])
|
|
# May or may not succeed depending on validation chain.
|
|
assert result is None or result is not None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Heartbeat tracking
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_heartbeat_updates_last_heartbeat(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.CLAIMED
|
|
task.assigned_to = task_setup["agent_id"]
|
|
await db_session.flush()
|
|
await svc.heartbeat(task.id)
|
|
refreshed = await svc.get(task.id)
|
|
assert refreshed is not None
|
|
assert refreshed.last_heartbeat_at is not None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# cancel cascades
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_cancel_with_note(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
cancelled = await svc.cancel(task.id, cancellation_note="not needed")
|
|
assert cancelled is not None
|
|
assert cancelled.status == TaskStatus.CANCELLED
|
|
assert "not needed" in (cancelled.dev_notes or "")
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_cancel_cascades_to_descendants(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
parent = await svc.create(_req(task_setup))
|
|
child = await svc.create(_req(task_setup, parent_task_id=parent.id))
|
|
cancelled = await svc.cancel(parent.id)
|
|
assert cancelled is not None
|
|
refreshed_child = await svc.get(child.id)
|
|
assert refreshed_child is not None
|
|
assert refreshed_child.status == TaskStatus.CANCELLED
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# soft_block + unblock with restore
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unblock_with_restore_returns_none_for_missing(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
result = await svc.unblock_with_restore(
|
|
pm_agent_id=task_setup["agent_id"], task_id=uuid4(), restore=True
|
|
)
|
|
assert result is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# qa_claim/doc_claim
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_claim_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.qa_claim(qa_agent_id=uuid4(), task_id=uuid4()) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_doc_claim_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
assert await svc.doc_claim(doc_agent_id=uuid4(), task_id=uuid4()) is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# pr_gate_claim — single-claimant guard (F114)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _reviewer(name: str) -> AgentTable:
|
|
return AgentTable(
|
|
id=uuid4(),
|
|
name=name,
|
|
slug=f"pr-reviewer-{uuid4().hex[:8]}",
|
|
role=AgentRole.PR_REVIEWER,
|
|
team=Team.BACKEND,
|
|
status=AgentStatus.ACTIVE,
|
|
model_config={},
|
|
system_prompt="reviewer",
|
|
capabilities=[],
|
|
permissions={},
|
|
metrics={},
|
|
)
|
|
|
|
|
|
def _pm(name: str) -> AgentTable:
|
|
return AgentTable(
|
|
id=uuid4(),
|
|
name=name,
|
|
slug=f"be-pm-{uuid4().hex[:8]}",
|
|
role=AgentRole.CELL_PM,
|
|
team=Team.BACKEND,
|
|
status=AgentStatus.ACTIVE,
|
|
model_config={},
|
|
system_prompt="pm",
|
|
capabilities=[],
|
|
permissions={},
|
|
metrics={},
|
|
)
|
|
|
|
|
|
async def _gate_task(task_setup: dict, db_session: AsyncSession) -> Any:
|
|
"""A task sitting in awaiting_pr_review, owned by ``owner``."""
|
|
svc = task_setup["svc"]
|
|
task = await svc.create(_req(task_setup))
|
|
task.status = TaskStatus.AWAITING_PR_REVIEW
|
|
await db_session.flush()
|
|
return task
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pr_gate_claim_rejects_second_reviewer_race(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""A second PR-reviewer race-claiming a gate task already claimed by another
|
|
reviewer is refused, so the first reviewer's claim and subsequent
|
|
pr_pass/pr_fail actor-checks are not overwritten."""
|
|
svc = task_setup["svc"]
|
|
reviewer1 = _reviewer("R1")
|
|
reviewer2 = _reviewer("R2")
|
|
db_session.add_all([reviewer1, reviewer2])
|
|
await db_session.flush()
|
|
task = await _gate_task(task_setup, db_session)
|
|
# reviewer1 already claimed the gate task.
|
|
task.active_claimant_id = reviewer1.id
|
|
task.claimed_by = reviewer1.id
|
|
task.assigned_to = reviewer1.id
|
|
await db_session.flush()
|
|
|
|
rejected = await svc.pr_gate_claim(reviewer2.id, task.id)
|
|
# Refused, not silently stolen.
|
|
assert rejected is None
|
|
await db_session.refresh(task)
|
|
# The first reviewer's claim is intact (NOT overwritten by reviewer2).
|
|
assert task.active_claimant_id == reviewer1.id
|
|
assert task.claimed_by == reviewer1.id
|
|
assert task.assigned_to == reviewer1.id
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pr_gate_claim_allows_first_reviewer_when_pm_owns_root(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""The first reviewer can still claim a gate task owned by the PM at entry
|
|
(submit_for_review does not clear ownership); the guard only rejects a
|
|
competing REVIEWER claim, not the PM owner."""
|
|
svc = task_setup["svc"]
|
|
pm = _pm("PM")
|
|
reviewer = _reviewer("R")
|
|
db_session.add_all([pm, reviewer])
|
|
await db_session.flush()
|
|
task = await _gate_task(task_setup, db_session)
|
|
# The PM owns the coordination root at gate entry.
|
|
task.active_claimant_id = pm.id
|
|
task.claimed_by = pm.id
|
|
task.assigned_to = pm.id
|
|
await db_session.flush()
|
|
|
|
claimed = await svc.pr_gate_claim(reviewer.id, task.id)
|
|
assert claimed is not None
|
|
await db_session.refresh(task)
|
|
# The reviewer took over the gate (the legitimate overclaim).
|
|
assert task.active_claimant_id == reviewer.id
|
|
assert task.claimed_by == reviewer.id
|
|
assert task.assigned_to == reviewer.id
|
|
# pr_gate_claim (via _qa_or_doc_claim) flips the reviewer's fleet marker.
|
|
reviewer_row = await db_session.get(AgentTable, reviewer.id)
|
|
assert reviewer_row is not None
|
|
assert reviewer_row.status == AgentStatus.ACTIVE
|
|
assert reviewer_row.current_task_id == task.id
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pr_gate_claim_idempotent_for_same_reviewer(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
"""A reviewer re-claiming its own gate claim is idempotent (allowed); the
|
|
guard only refuses a different reviewer."""
|
|
svc = task_setup["svc"]
|
|
reviewer = _reviewer("R")
|
|
db_session.add(reviewer)
|
|
await db_session.flush()
|
|
task = await _gate_task(task_setup, db_session)
|
|
task.active_claimant_id = reviewer.id
|
|
task.claimed_by = reviewer.id
|
|
task.assigned_to = reviewer.id
|
|
await db_session.flush()
|
|
|
|
claimed = await svc.pr_gate_claim(reviewer.id, task.id)
|
|
assert claimed is not None
|
|
await db_session.refresh(task)
|
|
assert task.active_claimant_id == reviewer.id
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# qa_pass / qa_fail / cell_pm_complete (404 paths)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_pass_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
result = await svc.qa_pass(
|
|
qa_agent_id=task_setup["agent_id"],
|
|
task_id=uuid4(),
|
|
notes="LGTM, comprehensive review",
|
|
)
|
|
assert result is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_fail_returns_none_for_missing(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
result = await svc.qa_fail(
|
|
qa_agent_id=task_setup["agent_id"],
|
|
task_id=uuid4(),
|
|
notes="needs revision",
|
|
issues=["bug 1"],
|
|
)
|
|
assert result is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# list_pending with dependency filtering
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_pending_filters_tasks_with_unmet_deps(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
blocker = await svc.create(_req(task_setup))
|
|
blocked = await svc.create(_req(task_setup))
|
|
blocked.dependency_ids = [blocker.id]
|
|
await db_session.flush()
|
|
pending = await svc.list_pending(team=Team.BACKEND)
|
|
pending_ids = {t.id for t in pending}
|
|
# blocker has no deps and is pending — should be included
|
|
assert blocker.id in pending_ids
|
|
# blocked depends on a non-terminal task — should be excluded
|
|
assert blocked.id not in pending_ids
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_pending_disabled_dep_filter(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
blocker = await svc.create(_req(task_setup))
|
|
blocked = await svc.create(_req(task_setup))
|
|
blocked.dependency_ids = [blocker.id]
|
|
await db_session.flush()
|
|
pending = await svc.list_pending(team=Team.BACKEND, filter_by_dependencies=False)
|
|
pending_ids = {t.id for t in pending}
|
|
assert blocker.id in pending_ids
|
|
assert blocked.id in pending_ids
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_pending_includes_tasks_when_deps_completed(
|
|
task_setup: dict, db_session: AsyncSession
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
blocker = await svc.create(_req(task_setup))
|
|
blocker.status = TaskStatus.COMPLETED
|
|
blocked = await svc.create(_req(task_setup))
|
|
blocked.dependency_ids = [blocker.id]
|
|
await db_session.flush()
|
|
pending = await svc.list_pending(team=Team.BACKEND)
|
|
pending_ids = {t.id for t in pending}
|
|
assert blocked.id in pending_ids
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Subtree query helpers
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_all_descendants_empty_for_leaf(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
leaf = await svc.create(_req(task_setup))
|
|
descendants = await svc.get_all_descendants(leaf.id)
|
|
assert descendants == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_all_descendants_traverses_three_levels(
|
|
task_setup: dict,
|
|
) -> None:
|
|
svc = task_setup["svc"]
|
|
grand = await svc.create(_req(task_setup))
|
|
parent = await svc.create(_req(task_setup, parent_task_id=grand.id))
|
|
child = await svc.create(_req(task_setup, parent_task_id=parent.id))
|
|
descendants = await svc.get_all_descendants(grand.id)
|
|
desc_ids = {d.id for d in descendants}
|
|
assert parent.id in desc_ids
|
|
assert child.id in desc_ids
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_count_by_status_with_data(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
await svc.create(_req(task_setup))
|
|
counts = await svc.count_by_status(team=Team.BACKEND)
|
|
assert isinstance(counts, dict)
|
|
assert "pending" in counts
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_active_count_zero_for_unknown(task_setup: dict) -> None:
|
|
svc = task_setup["svc"]
|
|
count = await svc.get_active_count(uuid4())
|
|
assert count == 0
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# add_dependency cycle / self-reference (M18)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
async def _seed_minimal_task(db_session: AsyncSession, tid: UUID) -> UUID:
|
|
agent = AgentTable(
|
|
id=uuid4(),
|
|
name="Dev",
|
|
slug=f"be-dev-{uuid4().hex[:8]}",
|
|
role=AgentRole.DEVELOPER,
|
|
team=Team.BACKEND,
|
|
status=AgentStatus.ACTIVE,
|
|
model_config={},
|
|
system_prompt="dev",
|
|
capabilities=[],
|
|
permissions={},
|
|
metrics={},
|
|
)
|
|
db_session.add(agent)
|
|
await db_session.flush()
|
|
project = ProjectTable(
|
|
id=uuid4(),
|
|
name="M-Proj",
|
|
slug=f"m-proj-{uuid4().hex[:8]}",
|
|
git_url="https://example.com/r.git",
|
|
assigned_cell=Team.BACKEND,
|
|
created_by=agent.id,
|
|
)
|
|
db_session.add(project)
|
|
await db_session.flush()
|
|
|
|
db_session.add(
|
|
TaskTable(
|
|
id=tid,
|
|
title="t",
|
|
description="d",
|
|
acceptance_criteria=["ac"],
|
|
status=TaskStatus.PENDING,
|
|
priority=2,
|
|
task_type=TaskType.CODE,
|
|
nature=TaskNature.TECHNICAL,
|
|
project_id=project.id,
|
|
created_by=agent.id,
|
|
team=Team.BACKEND,
|
|
)
|
|
)
|
|
await db_session.flush()
|
|
return cast("UUID", project.id)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_add_dependency_rejects_self_reference(
|
|
db_session: AsyncSession,
|
|
) -> None:
|
|
tid = uuid4()
|
|
await _seed_minimal_task(db_session, tid)
|
|
svc = get_task_service(db_session)
|
|
with pytest.raises(ConflictError, match="self"):
|
|
await svc.add_dependency(tid, tid)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_add_dependency_rejects_cycle(db_session: AsyncSession) -> None:
|
|
a, b, c = uuid4(), uuid4(), uuid4()
|
|
await _seed_minimal_task(db_session, a)
|
|
await _seed_minimal_task(db_session, b)
|
|
await _seed_minimal_task(db_session, c)
|
|
svc = get_task_service(db_session)
|
|
await svc.add_dependency(a, b)
|
|
await svc.add_dependency(b, c)
|
|
with pytest.raises(ConflictError, match="cycle"):
|
|
await svc.add_dependency(c, a)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# M19: _qa_or_doc_claim must lock the task row FOR UPDATE
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_claim_serializes_concurrent_claimants(
|
|
db_session: AsyncSession,
|
|
) -> None:
|
|
"""Two QA agents race-claiming the same awaiting_qa task: without FOR UPDATE
|
|
both pass the status check and overwrite claimed_by (last-write-wins); the
|
|
first claimant's later pass_review actor-mismatches. With FOR UPDATE the
|
|
second claim sees the first's committed claim and returns None."""
|
|
|
|
tid = uuid4()
|
|
await _seed_minimal_task(db_session, tid)
|
|
# Bump the seeded task into awaiting_qa with a PR open.
|
|
|
|
row = (
|
|
await db_session.execute(select(TaskTable).where(TaskTable.id == tid))
|
|
).scalar_one()
|
|
row.status = TaskStatus.AWAITING_QA
|
|
row.pr_number = 1
|
|
row.pr_created = True
|
|
await db_session.flush()
|
|
|
|
qa_a, qa_b = uuid4(), uuid4()
|
|
for aid in (qa_a, qa_b):
|
|
db_session.add(
|
|
AgentTable(
|
|
id=aid,
|
|
name=f"QA-{aid.hex[:4]}",
|
|
slug=f"qa-{aid.hex[:8]}",
|
|
role=AgentRole.QA,
|
|
team=Team.BACKEND,
|
|
status=AgentStatus.ACTIVE,
|
|
model_config={},
|
|
system_prompt="qa",
|
|
capabilities=[],
|
|
permissions={},
|
|
metrics={},
|
|
)
|
|
)
|
|
await db_session.flush()
|
|
|
|
svc = get_task_service(db_session)
|
|
first = await svc.qa_claim(qa_a, tid)
|
|
assert first is not None
|
|
assert to_uuid(first.active_claimant_id) == qa_a
|
|
second = await svc.qa_claim(qa_b, tid)
|
|
assert second is None, (
|
|
"second QA claimant overwrote the first's claimed_by — FOR UPDATE missing"
|
|
)
|
|
|
|
|
|
def to_uuid(v: object) -> UUID | None:
|
|
if v is None:
|
|
return None
|
|
if isinstance(v, UUID):
|
|
return v
|
|
return UUID(str(v))
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# H4: docs_complete + mark_pr_created lock the task row FOR UPDATE
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_parallel_completion_race_single_transition(
|
|
db_session: AsyncSession,
|
|
) -> None:
|
|
"""docs_complete and mark_pr_created race: both see ready_for_pm=True,
|
|
both transition to awaiting_pm_review, both emit audit rows. With FOR
|
|
UPDATE the second caller serializes behind the first and sees the
|
|
already-advanced status — it must NOT transition again."""
|
|
|
|
tid = uuid4()
|
|
await _seed_minimal_task(db_session, tid)
|
|
|
|
row = (
|
|
await db_session.execute(select(TaskTable).where(TaskTable.id == tid))
|
|
).scalar_one()
|
|
row.status = TaskStatus.AWAITING_DOCUMENTATION
|
|
row.docs_complete = True
|
|
row.pr_created = False
|
|
row.branch_name = "feature/x"
|
|
await db_session.flush()
|
|
|
|
svc = get_task_service(db_session)
|
|
t1 = await svc.mark_pr_created(tid, pr_number=42, pr_url="https://x")
|
|
assert t1 is not None and t1.status == TaskStatus.AWAITING_PM_REVIEW
|
|
await svc.docs_complete(tid, doc_notes="docs done")
|
|
rows = (
|
|
(
|
|
await db_session.execute(
|
|
select(AuditLogTable).where(AuditLogTable.target_id == tid)
|
|
)
|
|
)
|
|
.scalars()
|
|
.all()
|
|
)
|
|
transitions = [r for r in rows if r.event_type == "task.awaiting_pm_review"]
|
|
assert len(transitions) <= 1, (
|
|
f"parallel completion emitted {len(transitions)} awaiting_pm_review "
|
|
"audit rows — FOR UPDATE missing"
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# L29: pass_qa / fail_qa accept AWAITING_QA only
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pass_qa_refuses_in_progress(db_session: AsyncSession) -> None:
|
|
"""The gateway enforces AWAITING_QA; the service must match — an
|
|
in_progress QA claim must not pass_qa directly (the audit journey
|
|
would skip the task.awaiting_qa row, breaking QA dwell metrics)."""
|
|
|
|
tid = uuid4()
|
|
await _seed_minimal_task(db_session, tid)
|
|
|
|
row = (
|
|
await db_session.execute(select(TaskTable).where(TaskTable.id == tid))
|
|
).scalar_one()
|
|
row.status = TaskStatus.IN_PROGRESS
|
|
row.pr_number = 1
|
|
row.pr_created = True
|
|
await db_session.flush()
|
|
|
|
svc = get_task_service(db_session)
|
|
result = await svc.pass_qa(tid, notes="ok")
|
|
assert result is None, "pass_qa accepted IN_PROGRESS — no awaiting_qa hop (L29)"
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pass_qa_accepts_awaiting_qa(db_session: AsyncSession) -> None:
|
|
tid = uuid4()
|
|
await _seed_minimal_task(db_session, tid)
|
|
|
|
row = (
|
|
await db_session.execute(select(TaskTable).where(TaskTable.id == tid))
|
|
).scalar_one()
|
|
row.status = TaskStatus.AWAITING_QA
|
|
row.pr_number = 1
|
|
row.pr_created = True
|
|
await db_session.flush()
|
|
|
|
svc = get_task_service(db_session)
|
|
result = await svc.pass_qa(tid, notes="ok")
|
|
assert result is not None and result.status == TaskStatus.AWAITING_DOCUMENTATION
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# L30: mark_pr_created passes audit_agent_id
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_pr_created_audit_attributed_to_developer(
|
|
db_session: AsyncSession,
|
|
) -> None:
|
|
"""The awaiting_pm_review audit row from mark_pr_created must carry the
|
|
developer's agent_id, not NULL. The dev's claimed_by is already clear
|
|
by the time mark_pr_created runs (submit_for_qa cleared it), so the
|
|
caller must pass audit_agent_id explicitly."""
|
|
|
|
tid = uuid4()
|
|
await _seed_minimal_task(db_session, tid)
|
|
|
|
row = (
|
|
await db_session.execute(select(TaskTable).where(TaskTable.id == tid))
|
|
).scalar_one()
|
|
row.status = TaskStatus.AWAITING_DOCUMENTATION
|
|
row.docs_complete = True
|
|
row.pr_created = False
|
|
row.branch_name = "feature/x"
|
|
row.claimed_by = None # already cleared by submit_for_qa
|
|
await db_session.flush()
|
|
|
|
dev_id = uuid4()
|
|
db_session.add(
|
|
AgentTable(
|
|
id=dev_id,
|
|
name="Dev",
|
|
slug=f"dev-{dev_id.hex[:8]}",
|
|
role=AgentRole.DEVELOPER,
|
|
team=Team.BACKEND,
|
|
status=AgentStatus.ACTIVE,
|
|
model_config={},
|
|
system_prompt="dev",
|
|
capabilities=[],
|
|
permissions={},
|
|
metrics={},
|
|
)
|
|
)
|
|
await db_session.flush()
|
|
svc = get_task_service(db_session)
|
|
await svc.mark_pr_created(
|
|
tid, pr_number=7, pr_url="https://x", audit_agent_id=dev_id
|
|
)
|
|
rows = (
|
|
(
|
|
await db_session.execute(
|
|
select(AuditLogTable).where(AuditLogTable.target_id == tid)
|
|
)
|
|
)
|
|
.scalars()
|
|
.all()
|
|
)
|
|
transition = [r for r in rows if r.event_type == "task.awaiting_pm_review"]
|
|
assert transition, "no awaiting_pm_review audit row emitted"
|
|
assert transition[0].agent_id == dev_id, (
|
|
f"audit agent_id={transition[0].agent_id} != dev {dev_id} (L30)"
|
|
)
|