CUST-WP-0074-T02: task wait qualifiers in the hub read model
Adds tasks.decision_id and workplan_dependencies.from_task_id (migration
e8f9a0b1c2d3). A from_task_id column is needed rather than
from_workplan_id plus description: the partial unique indexes key on
(from_workplan_id, target, relationship_type), so two tasks in one
workplan waiting on the same target could not both be indexed, and the
derived wait_kind needs per-task attribution. The two existing unique
indexes are narrowed to frontmatter edges (from_task_id IS NULL) and two
task-origin counterparts are added; downgrade deletes task-origin rows
before restoring the old indexes.
POST /workplans/{id}/dependencies/ accepts from_task_id (must belong to
the from workplan; a task cannot depend on itself). TaskCreate/Update/Read
carry decision_id.
Derived read-model fields, no new write routes (rule 5):
- TaskRead.wait_kind: external | human | both | unqualified, null unless
status is wait (dependency rows from this task / needs_human).
- WorkplanRead.blocked_kind: human | external | none, null unless status
is blocked. Human wins over external.
Both come from api/services/wait_kind.py, applied on GET /tasks/,
GET /tasks/{id}, GET /workplans/ and GET /workplans/{id}.
The identifier-migration reference-count proof now includes from_task_id
(FK count guard 22 -> 23).
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Assistant: claude-code
Assistant-Model: sonnet
Assistant-Process: 237582@bnt-lap001
Assistant-Session: f2b3d9f1-8fb9-4b9c-bc2b-837ec5dfc826
This commit is contained in:
parent
e92471df3b
commit
b9997d6fc7
13 changed files with 400 additions and 6 deletions
|
|
@ -56,6 +56,10 @@ class Task(Base, TimestampMixin):
|
||||||
blocking_reason: Mapped[str | None] = mapped_column(Text, nullable=True)
|
blocking_reason: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||||
needs_human: Mapped[bool] = mapped_column(Boolean, default=False, nullable=False, index=True)
|
needs_human: Mapped[bool] = mapped_column(Boolean, default=False, nullable=False, index=True)
|
||||||
intervention_note: Mapped[str | None] = mapped_column(Text, nullable=True)
|
intervention_note: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||||
|
# Record id of the tracked decision a human-gated wait hangs on, e.g.
|
||||||
|
# "CCR-2026-0004" (CUST-WP-0074 Case B). Free text: the decision may live
|
||||||
|
# in another repo and may not be indexed in the hub yet.
|
||||||
|
decision_id: Mapped[str | None] = mapped_column(String(120), nullable=True)
|
||||||
parent_task_id: Mapped[uuid.UUID | None] = mapped_column(
|
parent_task_id: Mapped[uuid.UUID | None] = mapped_column(
|
||||||
UUID(as_uuid=True),
|
UUID(as_uuid=True),
|
||||||
ForeignKey("tasks.id", ondelete="SET NULL", onupdate="CASCADE"),
|
ForeignKey("tasks.id", ondelete="SET NULL", onupdate="CASCADE"),
|
||||||
|
|
|
||||||
|
|
@ -28,7 +28,7 @@ class WorkplanDependency(Base, TimestampMixin):
|
||||||
"to_workplan_id",
|
"to_workplan_id",
|
||||||
"relationship_type",
|
"relationship_type",
|
||||||
unique=True,
|
unique=True,
|
||||||
postgresql_where=text("to_workplan_id IS NOT NULL"),
|
postgresql_where=text("to_workplan_id IS NOT NULL AND from_task_id IS NULL"),
|
||||||
),
|
),
|
||||||
Index(
|
Index(
|
||||||
"uq_wp_dep_task_target",
|
"uq_wp_dep_task_target",
|
||||||
|
|
@ -36,7 +36,25 @@ class WorkplanDependency(Base, TimestampMixin):
|
||||||
"to_task_id",
|
"to_task_id",
|
||||||
"relationship_type",
|
"relationship_type",
|
||||||
unique=True,
|
unique=True,
|
||||||
postgresql_where=text("to_task_id IS NOT NULL"),
|
postgresql_where=text("to_task_id IS NOT NULL AND from_task_id IS NULL"),
|
||||||
|
),
|
||||||
|
# Task-block `depends_on` edges (CUST-WP-0074): the waiting task is the
|
||||||
|
# from side, so two tasks in one workplan may wait on the same target.
|
||||||
|
Index(
|
||||||
|
"uq_wp_dep_from_task_workplan_target",
|
||||||
|
"from_task_id",
|
||||||
|
"to_workplan_id",
|
||||||
|
"relationship_type",
|
||||||
|
unique=True,
|
||||||
|
postgresql_where=text("to_workplan_id IS NOT NULL AND from_task_id IS NOT NULL"),
|
||||||
|
),
|
||||||
|
Index(
|
||||||
|
"uq_wp_dep_from_task_task_target",
|
||||||
|
"from_task_id",
|
||||||
|
"to_task_id",
|
||||||
|
"relationship_type",
|
||||||
|
unique=True,
|
||||||
|
postgresql_where=text("to_task_id IS NOT NULL AND from_task_id IS NOT NULL"),
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -49,6 +67,14 @@ class WorkplanDependency(Base, TimestampMixin):
|
||||||
nullable=False,
|
nullable=False,
|
||||||
index=True,
|
index=True,
|
||||||
)
|
)
|
||||||
|
# Set when the edge originates from a task block rather than workplan
|
||||||
|
# frontmatter; always belongs to `from_workplan_id`.
|
||||||
|
from_task_id: Mapped[uuid.UUID | None] = mapped_column(
|
||||||
|
UUID(as_uuid=True),
|
||||||
|
ForeignKey("tasks.id", ondelete="CASCADE", onupdate="CASCADE"),
|
||||||
|
nullable=True,
|
||||||
|
index=True,
|
||||||
|
)
|
||||||
to_workplan_id: Mapped[uuid.UUID | None] = mapped_column(
|
to_workplan_id: Mapped[uuid.UUID | None] = mapped_column(
|
||||||
UUID(as_uuid=True),
|
UUID(as_uuid=True),
|
||||||
ForeignKey("workplans.id", ondelete="CASCADE", onupdate="CASCADE"),
|
ForeignKey("workplans.id", ondelete="CASCADE", onupdate="CASCADE"),
|
||||||
|
|
@ -72,4 +98,5 @@ class WorkplanDependency(Base, TimestampMixin):
|
||||||
to_workplan: Mapped["Workplan | None"] = relationship( # noqa: F821
|
to_workplan: Mapped["Workplan | None"] = relationship( # noqa: F821
|
||||||
"Workplan", foreign_keys=[to_workplan_id]
|
"Workplan", foreign_keys=[to_workplan_id]
|
||||||
)
|
)
|
||||||
|
from_task: Mapped["Task | None"] = relationship("Task", foreign_keys=[from_task_id]) # noqa: F821
|
||||||
to_task: Mapped["Task | None"] = relationship("Task", foreign_keys=[to_task_id]) # noqa: F821
|
to_task: Mapped["Task | None"] = relationship("Task", foreign_keys=[to_task_id]) # noqa: F821
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,7 @@ from api.schemas.task import (
|
||||||
)
|
)
|
||||||
from api.services.legacy_compat import meter_legacy_body_from_model, meter_legacy_query_param
|
from api.services.legacy_compat import meter_legacy_body_from_model, meter_legacy_query_param
|
||||||
from api.services.lifecycle import status_value, transition_task_status
|
from api.services.lifecycle import status_value, transition_task_status
|
||||||
|
from api.services.wait_kind import annotate_tasks
|
||||||
from api.task_status import normalize_task_status
|
from api.task_status import normalize_task_status
|
||||||
|
|
||||||
router = APIRouter(prefix="/tasks", tags=["tasks"])
|
router = APIRouter(prefix="/tasks", tags=["tasks"])
|
||||||
|
|
@ -69,7 +70,7 @@ async def list_tasks(
|
||||||
if limit is not None:
|
if limit is not None:
|
||||||
q = q.limit(limit)
|
q = q.limit(limit)
|
||||||
result = await session.execute(q)
|
result = await session.execute(q)
|
||||||
return list(result.scalars().all())
|
return list(await annotate_tasks(session, list(result.scalars().all())))
|
||||||
|
|
||||||
|
|
||||||
@router.get("/counts", response_model=list[TaskCountRead])
|
@router.get("/counts", response_model=list[TaskCountRead])
|
||||||
|
|
@ -248,6 +249,7 @@ async def get_task(
|
||||||
task = await session.get(Task, task_id)
|
task = await session.get(Task, task_id)
|
||||||
if task is None:
|
if task is None:
|
||||||
raise HTTPException(status_code=404, detail="Task not found")
|
raise HTTPException(status_code=404, detail="Task not found")
|
||||||
|
await annotate_tasks(session, [task])
|
||||||
return task
|
return task
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -36,9 +36,16 @@ async def _create_dependency(
|
||||||
raise HTTPException(status_code=404, detail="target task not found")
|
raise HTTPException(status_code=404, detail="target task not found")
|
||||||
if workplan_id == body.to_workplan_id:
|
if workplan_id == body.to_workplan_id:
|
||||||
raise HTTPException(status_code=422, detail="a workplan cannot depend on itself")
|
raise HTTPException(status_code=422, detail="a workplan cannot depend on itself")
|
||||||
|
if body.from_task_id is not None:
|
||||||
|
from_task = await session.get(Task, body.from_task_id)
|
||||||
|
if from_task is None or from_task.workplan_id != workplan_id:
|
||||||
|
raise HTTPException(status_code=404, detail="from task not found in this workplan")
|
||||||
|
if body.from_task_id == body.to_task_id:
|
||||||
|
raise HTTPException(status_code=422, detail="a task cannot depend on itself")
|
||||||
|
|
||||||
dep = WorkplanDependency(
|
dep = WorkplanDependency(
|
||||||
from_workplan_id=workplan_id,
|
from_workplan_id=workplan_id,
|
||||||
|
from_task_id=body.from_task_id,
|
||||||
to_workplan_id=body.to_workplan_id,
|
to_workplan_id=body.to_workplan_id,
|
||||||
to_task_id=body.to_task_id,
|
to_task_id=body.to_task_id,
|
||||||
relationship_type=body.relationship_type,
|
relationship_type=body.relationship_type,
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,7 @@ from api.schemas.workplan import (
|
||||||
)
|
)
|
||||||
from api.services.lifecycle import transition_workplan_status
|
from api.services.lifecycle import transition_workplan_status
|
||||||
from api.services.legacy_compat import retire_legacy_route
|
from api.services.legacy_compat import retire_legacy_route
|
||||||
|
from api.services.wait_kind import annotate_workplans
|
||||||
from api.work_record_flavor import (
|
from api.work_record_flavor import (
|
||||||
RESIDUAL_FLAVOR,
|
RESIDUAL_FLAVOR,
|
||||||
WORK_RECORD_FLAVORS,
|
WORK_RECORD_FLAVORS,
|
||||||
|
|
@ -389,7 +390,7 @@ async def list_workplans(
|
||||||
include_residuals: bool = Query(True),
|
include_residuals: bool = Query(True),
|
||||||
session: AsyncSession = Depends(get_session),
|
session: AsyncSession = Depends(get_session),
|
||||||
) -> list[Workplan]:
|
) -> list[Workplan]:
|
||||||
return await _list_workplans(
|
workplans = await _list_workplans(
|
||||||
topic_id=topic_id,
|
topic_id=topic_id,
|
||||||
repo_id=repo_id,
|
repo_id=repo_id,
|
||||||
repo_goal_id=repo_goal_id,
|
repo_goal_id=repo_goal_id,
|
||||||
|
|
@ -400,6 +401,7 @@ async def list_workplans(
|
||||||
include_residuals=include_residuals,
|
include_residuals=include_residuals,
|
||||||
session=session,
|
session=session,
|
||||||
)
|
)
|
||||||
|
return list(await annotate_workplans(session, workplans))
|
||||||
|
|
||||||
|
|
||||||
@router.get("/workplan-index", status_code=status.HTTP_410_GONE)
|
@router.get("/workplan-index", status_code=status.HTTP_410_GONE)
|
||||||
|
|
@ -499,7 +501,9 @@ async def get_workplan(
|
||||||
workplan_id: uuid.UUID,
|
workplan_id: uuid.UUID,
|
||||||
session: AsyncSession = Depends(get_session),
|
session: AsyncSession = Depends(get_session),
|
||||||
) -> Workplan:
|
) -> Workplan:
|
||||||
return await _get_workplan(workplan_id=workplan_id, session=session)
|
wp = await _get_workplan(workplan_id=workplan_id, session=session)
|
||||||
|
await annotate_workplans(session, [wp])
|
||||||
|
return wp
|
||||||
|
|
||||||
|
|
||||||
@router.patch("/{workstream_id}", status_code=status.HTTP_410_GONE)
|
@router.patch("/{workstream_id}", status_code=status.HTTP_410_GONE)
|
||||||
|
|
|
||||||
|
|
@ -45,6 +45,7 @@ class TaskCreate(TaskStatusMixin, TaskFlavorMixin, WorkplanIdCreateMixin):
|
||||||
blocking_reason: str | None = None
|
blocking_reason: str | None = None
|
||||||
needs_human: bool = False
|
needs_human: bool = False
|
||||||
intervention_note: str | None = None
|
intervention_note: str | None = None
|
||||||
|
decision_id: str | None = None
|
||||||
parent_task_id: uuid.UUID | None = None
|
parent_task_id: uuid.UUID | None = None
|
||||||
|
|
||||||
@model_validator(mode="after")
|
@model_validator(mode="after")
|
||||||
|
|
@ -65,6 +66,7 @@ class TaskUpdate(TaskStatusMixin, TaskFlavorMixin):
|
||||||
blocking_reason: str | None = None
|
blocking_reason: str | None = None
|
||||||
needs_human: bool | None = None
|
needs_human: bool | None = None
|
||||||
intervention_note: str | None = None
|
intervention_note: str | None = None
|
||||||
|
decision_id: str | None = None
|
||||||
parent_task_id: uuid.UUID | None = None
|
parent_task_id: uuid.UUID | None = None
|
||||||
# Token passthrough — three tiers (highest precision wins):
|
# Token passthrough — three tiers (highest precision wins):
|
||||||
# 1. tokens_in + tokens_out → exact counts; note defaults to "measured"
|
# 1. tokens_in + tokens_out → exact counts; note defaults to "measured"
|
||||||
|
|
@ -128,7 +130,12 @@ class TaskRead(TaskStatusMixin, WorkplanIdCompatMixin):
|
||||||
blocking_reason: str | None = None
|
blocking_reason: str | None = None
|
||||||
needs_human: bool
|
needs_human: bool
|
||||||
intervention_note: str | None = None
|
intervention_note: str | None = None
|
||||||
|
decision_id: str | None = None
|
||||||
parent_task_id: uuid.UUID | None = None
|
parent_task_id: uuid.UUID | None = None
|
||||||
|
# Derived (CUST-WP-0074 rule 5): "external" | "human" | "both" |
|
||||||
|
# "unqualified" while status is wait, else None. Set by
|
||||||
|
# api.services.wait_kind.annotate_tasks on read routes.
|
||||||
|
wait_kind: str | None = None
|
||||||
created_at: datetime
|
created_at: datetime
|
||||||
updated_at: datetime
|
updated_at: datetime
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -155,6 +155,10 @@ class WorkplanRead(WorkplanStatusMixin):
|
||||||
backing_relative_path: str | None = None
|
backing_relative_path: str | None = None
|
||||||
backing_archived: bool | None = None
|
backing_archived: bool | None = None
|
||||||
backing_synced_at: datetime | None = None
|
backing_synced_at: datetime | None = None
|
||||||
|
# Derived (CUST-WP-0074 rule 5): "human" | "external" | "none" while
|
||||||
|
# status is blocked, else None. Set by
|
||||||
|
# api.services.wait_kind.annotate_workplans on read routes.
|
||||||
|
blocked_kind: str | None = None
|
||||||
created_at: datetime
|
created_at: datetime
|
||||||
updated_at: datetime
|
updated_at: datetime
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -10,6 +10,7 @@ class WorkplanDependencyCreate(BaseModel):
|
||||||
validation_alias=AliasChoices("to_workplan_id", "to_workstream_id"),
|
validation_alias=AliasChoices("to_workplan_id", "to_workstream_id"),
|
||||||
)
|
)
|
||||||
to_task_id: uuid.UUID | None = None
|
to_task_id: uuid.UUID | None = None
|
||||||
|
from_task_id: uuid.UUID | None = None
|
||||||
relationship_type: str = "blocks"
|
relationship_type: str = "blocks"
|
||||||
description: str | None = None
|
description: str | None = None
|
||||||
|
|
||||||
|
|
@ -18,6 +19,7 @@ class WorkplanDependencyRead(BaseModel):
|
||||||
model_config = ConfigDict(from_attributes=True)
|
model_config = ConfigDict(from_attributes=True)
|
||||||
id: uuid.UUID
|
id: uuid.UUID
|
||||||
from_workplan_id: uuid.UUID
|
from_workplan_id: uuid.UUID
|
||||||
|
from_task_id: uuid.UUID | None = None
|
||||||
to_workplan_id: uuid.UUID | None = None
|
to_workplan_id: uuid.UUID | None = None
|
||||||
to_task_id: uuid.UUID | None = None
|
to_task_id: uuid.UUID | None = None
|
||||||
relationship_type: str
|
relationship_type: str
|
||||||
|
|
|
||||||
107
api/services/wait_kind.py
Normal file
107
api/services/wait_kind.py
Normal file
|
|
@ -0,0 +1,107 @@
|
||||||
|
"""Derived wait qualifiers (CUST-WP-0074 rule 5) — read model only.
|
||||||
|
|
||||||
|
A ``wait`` task is *external* when it carries at least one dependency edge
|
||||||
|
(``workplan_dependencies.from_task_id``), *human* when ``needs_human`` is set,
|
||||||
|
*both* when both hold and *unqualified* otherwise. A ``blocked`` workplan is
|
||||||
|
*human* when any of its wait tasks is human or both, else *external* when any
|
||||||
|
is external, else *none*. Nothing here is stored; the routers call the
|
||||||
|
``annotate_*`` helpers and Pydantic picks the attributes up on serialisation.
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import uuid
|
||||||
|
from collections.abc import Iterable, Sequence
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
from sqlalchemy import select
|
||||||
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
|
|
||||||
|
from api.models.task import Task, TaskStatus
|
||||||
|
from api.models.workplan import Workplan
|
||||||
|
from api.models.workplan_dependency import WorkplanDependency
|
||||||
|
from api.task_status import normalize_task_status
|
||||||
|
from api.workplan_status import normalize_workplan_status
|
||||||
|
|
||||||
|
WAIT_KIND_EXTERNAL = "external"
|
||||||
|
WAIT_KIND_HUMAN = "human"
|
||||||
|
WAIT_KIND_BOTH = "both"
|
||||||
|
WAIT_KIND_UNQUALIFIED = "unqualified"
|
||||||
|
|
||||||
|
BLOCKED_KIND_EXTERNAL = "external"
|
||||||
|
BLOCKED_KIND_HUMAN = "human"
|
||||||
|
BLOCKED_KIND_NONE = "none"
|
||||||
|
|
||||||
|
|
||||||
|
def derive_wait_kind(status: Any, *, needs_human: bool, has_dependency: bool) -> str | None:
|
||||||
|
"""Rule 1 / rule 5 for one task. ``None`` unless the task is waiting."""
|
||||||
|
if normalize_task_status(status, default="todo") != "wait":
|
||||||
|
return None
|
||||||
|
if has_dependency and needs_human:
|
||||||
|
return WAIT_KIND_BOTH
|
||||||
|
if needs_human:
|
||||||
|
return WAIT_KIND_HUMAN
|
||||||
|
if has_dependency:
|
||||||
|
return WAIT_KIND_EXTERNAL
|
||||||
|
return WAIT_KIND_UNQUALIFIED
|
||||||
|
|
||||||
|
|
||||||
|
def derive_blocked_kind(status: Any, wait_kinds: Iterable[str | None]) -> str | None:
|
||||||
|
"""Rule 5 for one workplan. ``None`` unless the workplan is blocked."""
|
||||||
|
if normalize_workplan_status(status) != "blocked":
|
||||||
|
return None
|
||||||
|
kinds = {kind for kind in wait_kinds if kind}
|
||||||
|
if WAIT_KIND_HUMAN in kinds or WAIT_KIND_BOTH in kinds:
|
||||||
|
return BLOCKED_KIND_HUMAN
|
||||||
|
if WAIT_KIND_EXTERNAL in kinds:
|
||||||
|
return BLOCKED_KIND_EXTERNAL
|
||||||
|
return BLOCKED_KIND_NONE
|
||||||
|
|
||||||
|
|
||||||
|
async def task_ids_with_dependencies(
|
||||||
|
session: AsyncSession, task_ids: Sequence[uuid.UUID]
|
||||||
|
) -> set[uuid.UUID]:
|
||||||
|
if not task_ids:
|
||||||
|
return set()
|
||||||
|
rows = await session.execute(
|
||||||
|
select(WorkplanDependency.from_task_id).where(
|
||||||
|
WorkplanDependency.from_task_id.in_(list(task_ids))
|
||||||
|
)
|
||||||
|
)
|
||||||
|
return {row[0] for row in rows if row[0] is not None}
|
||||||
|
|
||||||
|
|
||||||
|
async def annotate_tasks(session: AsyncSession, tasks: Sequence[Task]) -> Sequence[Task]:
|
||||||
|
"""Set ``task.wait_kind`` on each ORM task (serialised by TaskRead)."""
|
||||||
|
waiting = [t for t in tasks if normalize_task_status(t.status, default="todo") == "wait"]
|
||||||
|
with_deps = await task_ids_with_dependencies(session, [t.id for t in waiting])
|
||||||
|
for task in tasks:
|
||||||
|
task.wait_kind = derive_wait_kind(
|
||||||
|
task.status,
|
||||||
|
needs_human=bool(task.needs_human),
|
||||||
|
has_dependency=task.id in with_deps,
|
||||||
|
)
|
||||||
|
return tasks
|
||||||
|
|
||||||
|
|
||||||
|
async def annotate_workplans(
|
||||||
|
session: AsyncSession, workplans: Sequence[Workplan]
|
||||||
|
) -> Sequence[Workplan]:
|
||||||
|
"""Set ``workplan.blocked_kind`` on each ORM workplan (serialised by WorkplanRead)."""
|
||||||
|
blocked = [
|
||||||
|
wp for wp in workplans if normalize_workplan_status(wp.status) == "blocked"
|
||||||
|
]
|
||||||
|
kinds_by_workplan: dict[uuid.UUID, list[str | None]] = {wp.id: [] for wp in blocked}
|
||||||
|
if blocked:
|
||||||
|
rows = await session.execute(
|
||||||
|
select(Task).where(
|
||||||
|
Task.workplan_id.in_(list(kinds_by_workplan)),
|
||||||
|
Task.status == TaskStatus.wait,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
wait_tasks = list(rows.scalars().all())
|
||||||
|
await annotate_tasks(session, wait_tasks)
|
||||||
|
for task in wait_tasks:
|
||||||
|
kinds_by_workplan[task.workplan_id].append(task.wait_kind)
|
||||||
|
for wp in workplans:
|
||||||
|
wp.blocked_kind = derive_blocked_kind(wp.status, kinds_by_workplan.get(wp.id, []))
|
||||||
|
return workplans
|
||||||
|
|
@ -386,6 +386,7 @@ async def _task_reference_count(session: AsyncSession, task_id: uuid.UUID) -> in
|
||||||
"(SELECT count(*) FROM progress_events WHERE task_id = :task_id) + "
|
"(SELECT count(*) FROM progress_events WHERE task_id = :task_id) + "
|
||||||
"(SELECT count(*) FROM capability_requests WHERE blocking_task_id = :task_id) + "
|
"(SELECT count(*) FROM capability_requests WHERE blocking_task_id = :task_id) + "
|
||||||
"(SELECT count(*) FROM workplan_dependencies WHERE to_task_id = :task_id) + "
|
"(SELECT count(*) FROM workplan_dependencies WHERE to_task_id = :task_id) + "
|
||||||
|
"(SELECT count(*) FROM workplan_dependencies WHERE from_task_id = :task_id) + "
|
||||||
"(SELECT count(*) FROM suggestions WHERE promoted_task_id = :task_id)"
|
"(SELECT count(*) FROM suggestions WHERE promoted_task_id = :task_id)"
|
||||||
),
|
),
|
||||||
{"task_id": task_id},
|
{"task_id": task_id},
|
||||||
|
|
|
||||||
103
migrations/versions/e8f9a0b1c2d3_task_wait_qualifiers.py
Normal file
103
migrations/versions/e8f9a0b1c2d3_task_wait_qualifiers.py
Normal file
|
|
@ -0,0 +1,103 @@
|
||||||
|
"""task wait qualifiers (CUST-WP-0074-T02)
|
||||||
|
|
||||||
|
Adds ``tasks.decision_id`` (record id of the decision a human-gated wait hangs
|
||||||
|
on) and ``workplan_dependencies.from_task_id`` so that task-block
|
||||||
|
``depends_on`` edges carry the waiting task as the from side. The two existing
|
||||||
|
partial unique indexes are narrowed to frontmatter edges (``from_task_id IS
|
||||||
|
NULL``) and two task-origin counterparts are added, so two tasks in one
|
||||||
|
workplan may wait on the same target.
|
||||||
|
|
||||||
|
Revision ID: e8f9a0b1c2d3
|
||||||
|
Revises: d7e8f9a0b1c2
|
||||||
|
"""
|
||||||
|
import sqlalchemy as sa
|
||||||
|
from alembic import op
|
||||||
|
from sqlalchemy.dialects import postgresql
|
||||||
|
|
||||||
|
revision = "e8f9a0b1c2d3"
|
||||||
|
down_revision = "d7e8f9a0b1c2"
|
||||||
|
branch_labels = None
|
||||||
|
depends_on = None
|
||||||
|
|
||||||
|
|
||||||
|
def upgrade() -> None:
|
||||||
|
op.add_column("tasks", sa.Column("decision_id", sa.String(length=120), nullable=True))
|
||||||
|
|
||||||
|
op.add_column(
|
||||||
|
"workplan_dependencies",
|
||||||
|
sa.Column("from_task_id", postgresql.UUID(as_uuid=True), nullable=True),
|
||||||
|
)
|
||||||
|
op.create_foreign_key(
|
||||||
|
"fk_workplan_dependencies_from_task_id",
|
||||||
|
"workplan_dependencies",
|
||||||
|
"tasks",
|
||||||
|
["from_task_id"],
|
||||||
|
["id"],
|
||||||
|
ondelete="CASCADE",
|
||||||
|
onupdate="CASCADE",
|
||||||
|
)
|
||||||
|
op.create_index(
|
||||||
|
"ix_workplan_dependencies_from_task_id",
|
||||||
|
"workplan_dependencies",
|
||||||
|
["from_task_id"],
|
||||||
|
)
|
||||||
|
|
||||||
|
op.drop_index("uq_wp_dep_workplan_target", table_name="workplan_dependencies")
|
||||||
|
op.drop_index("uq_wp_dep_task_target", table_name="workplan_dependencies")
|
||||||
|
op.create_index(
|
||||||
|
"uq_wp_dep_workplan_target",
|
||||||
|
"workplan_dependencies",
|
||||||
|
["from_workplan_id", "to_workplan_id", "relationship_type"],
|
||||||
|
unique=True,
|
||||||
|
postgresql_where=sa.text("to_workplan_id IS NOT NULL AND from_task_id IS NULL"),
|
||||||
|
)
|
||||||
|
op.create_index(
|
||||||
|
"uq_wp_dep_task_target",
|
||||||
|
"workplan_dependencies",
|
||||||
|
["from_workplan_id", "to_task_id", "relationship_type"],
|
||||||
|
unique=True,
|
||||||
|
postgresql_where=sa.text("to_task_id IS NOT NULL AND from_task_id IS NULL"),
|
||||||
|
)
|
||||||
|
op.create_index(
|
||||||
|
"uq_wp_dep_from_task_workplan_target",
|
||||||
|
"workplan_dependencies",
|
||||||
|
["from_task_id", "to_workplan_id", "relationship_type"],
|
||||||
|
unique=True,
|
||||||
|
postgresql_where=sa.text("to_workplan_id IS NOT NULL AND from_task_id IS NOT NULL"),
|
||||||
|
)
|
||||||
|
op.create_index(
|
||||||
|
"uq_wp_dep_from_task_task_target",
|
||||||
|
"workplan_dependencies",
|
||||||
|
["from_task_id", "to_task_id", "relationship_type"],
|
||||||
|
unique=True,
|
||||||
|
postgresql_where=sa.text("to_task_id IS NOT NULL AND from_task_id IS NOT NULL"),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def downgrade() -> None:
|
||||||
|
op.drop_index("uq_wp_dep_from_task_task_target", table_name="workplan_dependencies")
|
||||||
|
op.drop_index("uq_wp_dep_from_task_workplan_target", table_name="workplan_dependencies")
|
||||||
|
op.drop_index("uq_wp_dep_task_target", table_name="workplan_dependencies")
|
||||||
|
op.drop_index("uq_wp_dep_workplan_target", table_name="workplan_dependencies")
|
||||||
|
# Task-origin rows would collide under the frontmatter-only indexes.
|
||||||
|
op.execute("DELETE FROM workplan_dependencies WHERE from_task_id IS NOT NULL")
|
||||||
|
op.create_index(
|
||||||
|
"uq_wp_dep_workplan_target",
|
||||||
|
"workplan_dependencies",
|
||||||
|
["from_workplan_id", "to_workplan_id", "relationship_type"],
|
||||||
|
unique=True,
|
||||||
|
postgresql_where=sa.text("to_workplan_id IS NOT NULL"),
|
||||||
|
)
|
||||||
|
op.create_index(
|
||||||
|
"uq_wp_dep_task_target",
|
||||||
|
"workplan_dependencies",
|
||||||
|
["from_workplan_id", "to_task_id", "relationship_type"],
|
||||||
|
unique=True,
|
||||||
|
postgresql_where=sa.text("to_task_id IS NOT NULL"),
|
||||||
|
)
|
||||||
|
op.drop_index("ix_workplan_dependencies_from_task_id", table_name="workplan_dependencies")
|
||||||
|
op.drop_constraint(
|
||||||
|
"fk_workplan_dependencies_from_task_id", "workplan_dependencies", type_="foreignkey"
|
||||||
|
)
|
||||||
|
op.drop_column("workplan_dependencies", "from_task_id")
|
||||||
|
op.drop_column("tasks", "decision_id")
|
||||||
126
tests/test_wait_kind.py
Normal file
126
tests/test_wait_kind.py
Normal file
|
|
@ -0,0 +1,126 @@
|
||||||
|
"""CUST-WP-0074-T02: derived wait_kind / blocked_kind read-model fields and the
|
||||||
|
from_task_id dependency edge. Real PostgreSQL via the shared `client` fixture."""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from api.services.wait_kind import derive_blocked_kind, derive_wait_kind
|
||||||
|
from tests.conftest import create_test_domain, create_test_repo, create_test_workplan
|
||||||
|
|
||||||
|
|
||||||
|
class TestDeriveFunctions:
|
||||||
|
def test_wait_kind_by_qualifiers(self):
|
||||||
|
assert derive_wait_kind("todo", needs_human=True, has_dependency=True) is None
|
||||||
|
assert derive_wait_kind("wait", needs_human=False, has_dependency=False) == "unqualified"
|
||||||
|
assert derive_wait_kind("wait", needs_human=False, has_dependency=True) == "external"
|
||||||
|
assert derive_wait_kind("wait", needs_human=True, has_dependency=False) == "human"
|
||||||
|
assert derive_wait_kind("wait", needs_human=True, has_dependency=True) == "both"
|
||||||
|
|
||||||
|
def test_blocked_kind_prefers_human(self):
|
||||||
|
assert derive_blocked_kind("active", ["human"]) is None
|
||||||
|
assert derive_blocked_kind("blocked", []) == "none"
|
||||||
|
assert derive_blocked_kind("blocked", ["unqualified", None]) == "none"
|
||||||
|
assert derive_blocked_kind("blocked", ["external", "unqualified"]) == "external"
|
||||||
|
assert derive_blocked_kind("blocked", ["external", "human"]) == "human"
|
||||||
|
assert derive_blocked_kind("blocked", ["both"]) == "human"
|
||||||
|
|
||||||
|
|
||||||
|
async def _task(client, workplan_id, title, **extra):
|
||||||
|
r = await client.post("/tasks/", json={"workplan_id": workplan_id, "title": title, **extra})
|
||||||
|
assert r.status_code == 201, r.text
|
||||||
|
return r.json()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_wait_kind_and_blocked_kind_on_read_routes(client):
|
||||||
|
domain = await create_test_domain(client)
|
||||||
|
repo = await create_test_repo(client, domain_slug=domain["slug"])
|
||||||
|
upstream = await create_test_workplan(client, repo["id"], slug="upstream", status="active")
|
||||||
|
blocked = await create_test_workplan(client, repo["id"], slug="blocked", status="blocked")
|
||||||
|
|
||||||
|
human = await _task(
|
||||||
|
client, blocked["id"], "human gate", status="wait",
|
||||||
|
needs_human=True, blocking_reason="operator must provision creds",
|
||||||
|
intervention_note="operator must provision creds", decision_id="CCR-2026-0004",
|
||||||
|
)
|
||||||
|
external = await _task(
|
||||||
|
client, blocked["id"], "external", status="wait", blocking_reason="upstream first",
|
||||||
|
)
|
||||||
|
unqualified = await _task(client, blocked["id"], "unqualified", status="wait", blocking_reason="?")
|
||||||
|
todo = await _task(client, blocked["id"], "plain todo")
|
||||||
|
|
||||||
|
r = await client.post(
|
||||||
|
f"/workplans/{blocked['id']}/dependencies/",
|
||||||
|
json={"from_task_id": external["id"], "to_workplan_id": upstream["id"]},
|
||||||
|
)
|
||||||
|
assert r.status_code == 201, r.text
|
||||||
|
assert r.json()["from_task_id"] == external["id"]
|
||||||
|
|
||||||
|
r = await client.get("/tasks/", params={"workplan_id": blocked["id"]})
|
||||||
|
assert r.status_code == 200
|
||||||
|
kinds = {row["title"]: row["wait_kind"] for row in r.json()}
|
||||||
|
assert kinds == {
|
||||||
|
"human gate": "human", "external": "external", "unqualified": "unqualified", "plain todo": None,
|
||||||
|
}
|
||||||
|
r = await client.get(f"/tasks/{human['id']}")
|
||||||
|
assert r.json()["wait_kind"] == "human"
|
||||||
|
assert r.json()["decision_id"] == "CCR-2026-0004"
|
||||||
|
|
||||||
|
r = await client.get(f"/workplans/{blocked['id']}")
|
||||||
|
assert r.json()["blocked_kind"] == "human"
|
||||||
|
r = await client.get(f"/workplans/{upstream['id']}")
|
||||||
|
assert r.json()["blocked_kind"] is None
|
||||||
|
|
||||||
|
# Human gate lifted → the workplan is blocked on the external wait only.
|
||||||
|
r = await client.patch(f"/tasks/{human['id']}", json={"status": "done"})
|
||||||
|
assert r.status_code == 200, r.text
|
||||||
|
r = await client.get("/workplans/", params={"repo_id": repo["id"]})
|
||||||
|
kinds = {row["slug"]: row["blocked_kind"] for row in r.json()}
|
||||||
|
assert kinds == {"upstream": None, "blocked": "external"}
|
||||||
|
|
||||||
|
# Dependency and human flag together → both.
|
||||||
|
r = await client.patch(
|
||||||
|
f"/tasks/{external['id']}",
|
||||||
|
json={"needs_human": True, "intervention_note": "confirm with owner"},
|
||||||
|
)
|
||||||
|
assert r.status_code == 200, r.text
|
||||||
|
r = await client.get(f"/tasks/{external['id']}")
|
||||||
|
assert r.json()["wait_kind"] == "both"
|
||||||
|
del unqualified, todo
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_from_task_must_belong_to_the_from_workplan(client):
|
||||||
|
domain = await create_test_domain(client)
|
||||||
|
repo = await create_test_repo(client, domain_slug=domain["slug"])
|
||||||
|
a = await create_test_workplan(client, repo["id"], slug="a")
|
||||||
|
b = await create_test_workplan(client, repo["id"], slug="b")
|
||||||
|
task_in_b = await _task(client, b["id"], "in b")
|
||||||
|
r = await client.post(
|
||||||
|
f"/workplans/{a['id']}/dependencies/",
|
||||||
|
json={"from_task_id": task_in_b["id"], "to_workplan_id": b["id"]},
|
||||||
|
)
|
||||||
|
assert r.status_code == 404
|
||||||
|
r = await client.post(
|
||||||
|
f"/workplans/{b['id']}/dependencies/",
|
||||||
|
json={"from_task_id": task_in_b["id"], "to_task_id": task_in_b["id"]},
|
||||||
|
)
|
||||||
|
assert r.status_code == 422
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_two_tasks_may_wait_on_the_same_target(client):
|
||||||
|
domain = await create_test_domain(client)
|
||||||
|
repo = await create_test_repo(client, domain_slug=domain["slug"])
|
||||||
|
upstream = await create_test_workplan(client, repo["id"], slug="upstream")
|
||||||
|
wp = await create_test_workplan(client, repo["id"], slug="waiting")
|
||||||
|
t1 = await _task(client, wp["id"], "t1")
|
||||||
|
t2 = await _task(client, wp["id"], "t2")
|
||||||
|
for task in (t1, t2):
|
||||||
|
r = await client.post(
|
||||||
|
f"/workplans/{wp['id']}/dependencies/",
|
||||||
|
json={"from_task_id": task["id"], "to_workplan_id": upstream["id"]},
|
||||||
|
)
|
||||||
|
assert r.status_code == 201, r.text
|
||||||
|
r = await client.get(f"/workplans/{wp['id']}/dependencies/")
|
||||||
|
assert sorted(row["from_task_id"] for row in r.json()) == sorted([t1["id"], t2["id"]])
|
||||||
|
|
@ -188,7 +188,7 @@ def test_every_work_record_foreign_key_cascades_on_update():
|
||||||
for foreign_key in table.foreign_keys
|
for foreign_key in table.foreign_keys
|
||||||
if foreign_key.target_fullname in {"workplans.id", "tasks.id"}
|
if foreign_key.target_fullname in {"workplans.id", "tasks.id"}
|
||||||
]
|
]
|
||||||
assert len(references) == 22
|
assert len(references) == 23
|
||||||
assert all(foreign_key.onupdate == "CASCADE" for foreign_key in references)
|
assert all(foreign_key.onupdate == "CASCADE" for foreign_key in references)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue