"""API задач M1/M2: быстрый захват, стек, детализация (LLM-предложение), CRUD."""
from typing import Any, cast
from fastapi import APIRouter, BackgroundTasks, HTTPException, Query, Response
from sqlalchemy import Select, case, or_, select
from app.actor import Actor
from app.dependencies import ActorDep, DbDep, UserIdDep
from app.models import CoinEvent, Project, Tag, Task, utcnow
from app.realtime import publish
from app.schemas import (
DEADLINE_PERIODS,
RECUR_KINDS,
AcceptIn,
ApproveIn,
SuggestIn,
TaskCreate,
TaskOut,
TaskUpdate,
)
from app.services import claims, tasklog
from app.services.closing import handle_task_closed
from app.services.detailing import apply_proposal, detail_task
from app.services.options import pick_options
from app.services.push import notify_agent_claimed, notify_agent_review
from app.services.xp import (
XP_CREATE_TASK,
coins_for_create,
grant_create_xp,
)
router = APIRouter(prefix="/api/tasks", tags=["tasks"])
VALID_STATUSES = {"to_do", "in_progress", "done", "cancelled", "deferred"}
def _get_task_or_404(db: Any, task_id: int, user_id: str) -> Task:
task = cast(
Task | None,
db.scalar(select(Task).where(Task.id == task_id, Task.user_id == user_id)),
)
if task is None:
raise HTTPException(status_code=404, detail="Task not found")
return task
def _get_task_for_update(db: Any, task_id: int, user_id: str) -> Task:
"""Задача с блокировкой строки (Postgres FOR UPDATE; в SQLite игнорируется).
Параллельные PATCH одной задачи сериализуются: двойное закрытие не гонит
двух наград (уникальный XpEvent.task_id — вторая линия обороны).
"""
task = cast(
Task | None,
# of=Task: запрос тянет project joined-ом — Postgres запрещает FOR UPDATE
# на nullable-стороне outer join, блокируем только строку задачи
db.scalar(
select(Task).where(Task.id == task_id, Task.user_id == user_id).with_for_update(of=Task)
),
)
if task is None:
raise HTTPException(status_code=404, detail="Task not found")
return task
def _mandate_error(exc: Exception) -> HTTPException:
"""Мандат нарушен → 403; занятая другим агентом задача → 409 (конфликт)."""
status = 409 if isinstance(exc, claims.LeaseBusyError) else 403
return HTTPException(status_code=status, detail=str(exc))
def _mandate(actor: Actor, task: Task, *, sets_eligible: bool = False) -> None:
"""Рамка мандата (ТЗ 3.20): агент — 403/409, если вышел за неё.
sets_eligible — запрос меняет флаг доступности: его ставит только владелец.
"""
try:
if sets_eligible:
claims.ensure_owner_sets_eligible(actor)
claims.ensure_agent_can_touch(task, actor)
except claims.MandateError as exc:
raise _mandate_error(exc) from None
def _archived_project_ids(db: Any, user_id: str) -> Select[tuple[int]]:
"""id проектов в архиве — их задачи в рабочих видах не показываются."""
return select(Project.id).where(Project.is_archived.is_(True), Project.user_id == user_id)
def _validate_parent(db: Any, task: Task | None, parent_id: int | None, user_id: str) -> None:
"""Родитель должен существовать (и быть своим); циклы в дереве запрещены."""
if parent_id is None:
return
if task is not None and parent_id == task.id:
raise HTTPException(status_code=400, detail="Task cannot be its own parent")
parent = cast(
Task | None,
db.scalar(select(Task).where(Task.id == parent_id, Task.user_id == user_id)),
)
if parent is None:
raise HTTPException(status_code=400, detail="Unknown parent task")
if task is not None:
ancestor: Task | None = parent
while ancestor is not None:
if ancestor.id == task.id:
raise HTTPException(status_code=400, detail="Cycle in task tree")
ancestor = ancestor.parent
@router.post("")
async def create_task(
schema: TaskCreate,
db: DbDep,
user_id: UserIdDep,
actor: ActorDep,
background: BackgroundTasks,
response: Response,
) -> dict[str, int]:
"""Быстрый захват: достаточно title — задача попадает в стек (raw, to_do).
С parent_task_id — создание подзадачи (M3). Сразу в фоне запускается
автодетализация (LLM-предложение метаданных).
"""
_validate_parent(db, None, schema.parent_task_id, user_id)
task = Task(
user_id=user_id,
title=schema.title,
description=schema.description,
parent_task_id=schema.parent_task_id,
# Доступность агенту (ТЗ 3.20) и автор задачи: «создано агентом» — признак
# самой задачи, иначе список фильтровался бы join'ом с токенами
ai_eligible=schema.ai_eligible,
created_by_kind=actor.kind,
created_by_name=actor.name or None,
)
db.add(task)
db.flush()
tasklog.log_event(db, task, tasklog.KIND_CREATED, actor)
# Микронаграда за создание (ТЗ 3.13): без конфетти, просто тост —
# заголовки X-Created-* фронт буферизует отдельно от награды за закрытие
grant_create_xp(db, "create_task", user_id)
db.add(CoinEvent(user_id=user_id, source="create_task", amount=coins_for_create("create_task")))
response.headers["X-Created-XP"] = str(XP_CREATE_TASK)
response.headers["X-Created-Coins"] = str(coins_for_create("create_task"))
# Коммитим до планирования фоновой работы: BackgroundTasks выполняются ДО
# teardown-коммита зависимости get_db, а detail_task открывает свою сессию.
db.commit()
publish(user_id, "task.changed", {"id": task.id})
publish(user_id, "xp.changed", {"celebrate": False})
background.add_task(detail_task, task.id)
return {"id": task.id}
@router.get("", response_model=list[TaskOut])
async def list_tasks(
db: DbDep,
user_id: UserIdDep,
detail_state: str | None = Query(None),
status: str | None = Query(None),
project_id: int | None = Query(None),
tag_id: int | None = Query(None),
tag: str | None = Query(None, description="Подстрока имени тега (без учёта регистра)"),
priority: str | None = Query(None),
priority_min: int | None = Query(
None, ge=0, le=10, description="Нижняя граница приоритета (шкала 0-10)"
),
parent_id: int | None = Query(None),
# Мандат агента (ТЗ 3.20): «доступно агенту», «создано мной / агентом», приёмка
ai_eligible: bool | None = Query(None),
created_by_kind: str | None = Query(None, pattern="^(user|agent)$"),
accept_state: str | None = Query(None, pattern="^(pending|accepted|rejected)$"),
# Тип задачи (0.91): фильтр списка «Регулярные» (ТЗ 3.6)
task_type: str | None = Query(None, pattern="^(one_time|recurring)$"),
include_archived: bool = Query(False),
) -> list[Task]:
# Завершённые задачи сортируются по дате закрытия, остальные — по добавлению
closed = case((Task.status == "done", Task.done_at), else_=Task.created_at)
stmt = (
select(Task)
.where(Task.user_id == user_id)
.order_by(closed.desc().nullslast())
)
if detail_state:
stmt = stmt.where(Task.detail_state == detail_state)
if status:
stmt = stmt.where(Task.status == status)
if project_id:
stmt = stmt.where(Task.project_id == project_id)
elif not include_archived:
# задачи архивных проектов — только в истории (ТЗ 3.8.1)
stmt = stmt.where(
or_(
Task.project_id.is_(None),
Task.project_id.not_in(_archived_project_ids(db, user_id)),
)
)
if parent_id is not None:
stmt = stmt.where(Task.parent_task_id == parent_id)
if ai_eligible is not None:
stmt = stmt.where(Task.ai_eligible.is_(ai_eligible))
if created_by_kind:
stmt = stmt.where(Task.created_by_kind == created_by_kind)
if accept_state:
stmt = stmt.where(Task.accept_state == accept_state)
if task_type:
stmt = stmt.where(Task.task_type == task_type)
if tag_id:
stmt = stmt.where(Task.tags.any(Tag.id == tag_id))
# Поиск по подстроке имени тега (ТЗ 3.6, 0.70): ввод в списке задач не выбирает
# тег из выпадашки, а ищет по нему. autoescape глушит служебные знаки LIKE —
# иначе «%» во вводе совпал бы с чем угодно, а «_» — с любым символом
query = (tag or "").strip()
if query:
stmt = stmt.where(Task.tags.any(Tag.name.icontains(query, autoescape=True)))
# Фильтр по приоритету (ТЗ 3.6) — по градациям шкалы 0-10 (см. 3.13):
# без приоритета задача считается «very_low»
if priority:
ranges = {
"very_low": (None, 2),
"low": (3, 4),
"medium": (5, 6),
"high": (7, 8),
"urgent": (9, 10),
}
if priority not in ranges:
raise HTTPException(status_code=422, detail=f"Unknown priority grade: {priority}")
low, high = ranges[priority]
if low is None:
stmt = stmt.where(or_(Task.priority.is_(None), Task.priority <= high))
else:
stmt = stmt.where(Task.priority.between(low, high))
# Отсечка по нижней границе шкалы (ТЗ 3.6, 0.74) — ею работает переключатель
# «Приоритетные» в списке задач: «выше среднего» = high+urgent, то есть 7-10.
# Задачи без приоритета (very_low, NULL) не проходят сравнение — так и надо
if priority_min is not None:
stmt = stmt.where(Task.priority >= priority_min)
return list(db.scalars(stmt).all())
@router.get("/{task_id}", response_model=TaskOut)
async def get_task(task_id: int, db: DbDep, user_id: UserIdDep) -> Task:
return _get_task_or_404(db, task_id, user_id)
@router.patch("/{task_id}", response_model=TaskOut)
async def update_task(
task_id: int,
schema: TaskUpdate,
db: DbDep,
user_id: UserIdDep,
actor: ActorDep,
response: Response,
background: BackgroundTasks,
) -> Task:
task = _get_task_for_update(db, task_id, user_id)
data = schema.model_dump(exclude_unset=True)
old_status, old_eligible = task.status, task.ai_eligible
old_accept = task.accept_state
# Мандат (ТЗ 3.20): без флага доступности агент задачу не меняет, а сам флаг
# ставит только владелец — иначе «пометил и закрыл» обходило бы правило
_mandate(actor, task, sets_eligible="ai_eligible" in data)
if "title" in data and not str(data["title"]).strip():
raise HTTPException(status_code=422, detail="title cannot be empty")
if "status" in data and data["status"] not in VALID_STATUSES:
raise HTTPException(status_code=422, detail=f"Unknown status: {data['status']}")
if "deadline_period" in data and data["deadline_period"] not in (None,) + DEADLINE_PERIODS:
raise HTTPException(
status_code=422, detail=f"Unknown deadline period: {data['deadline_period']}"
)
if "recur_kind" in data and data["recur_kind"] not in (None,) + RECUR_KINDS:
raise HTTPException(
status_code=422, detail=f"Unknown recurrence rule: {data['recur_kind']}"
)
if "task_type" in data and data["task_type"] not in ("one_time", "recurring"):
raise HTTPException(status_code=422, detail=f"Unknown task type: {data['task_type']}")
if "project_id" in data and data["project_id"] is not None:
if (
db.scalar(
select(Project).where(Project.id == data["project_id"], Project.user_id == user_id)
)
is None
):
raise HTTPException(status_code=400, detail="Unknown project")
if "parent_task_id" in data:
_validate_parent(db, task, data["parent_task_id"], user_id)
if "tag_ids" in data:
tag_ids = data.pop("tag_ids") or []
tags = db.scalars(select(Tag).where(Tag.id.in_(tag_ids), Tag.user_id == user_id)).all()
if len(tags) != len(set(tag_ids)):
raise HTTPException(status_code=400, detail="Unknown tag id in tag_ids")
task.tags = list(tags)
# Пометка источника закрытия (геймификация): не поле задачи, в setattr не идёт
earned_via = data.pop("earned_via", None)
# Комментарий к закрытию (ТЗ 3.20): уходит в журнал, агенту обязателен
close_comment = data.pop("close_comment", None)
# Агент закрывает задачу с отчётом: без него владельцу нечего принимать.
# Требование — на настоящий переход в done, а не на повторный PATCH
if (
data.get("status") == "done"
and old_status != "done"
and actor.is_agent
and not (close_comment or "").strip()
):
raise HTTPException(status_code=422, detail="Close comment is required for an AI agent")
earned_event = None
for field, value in data.items():
setattr(task, field, value)
if data.get("status") == "done":
if task.done_at is None:
task.done_at = utcnow()
# Регулярная (ТЗ 3.5) + геймификация (ТЗ 3.13) — общий путь с MCP:
# закрытие агентом даёт те же спавн/XP/монеты, что и из UI
outcome = handle_task_closed(db, task, via_options=earned_via == "options")
event = outcome.event
if event is not None:
response.headers["X-Earned-XP"] = str(event.amount)
response.headers["X-Plant-Rarity"] = event.rarity
if event.via_options:
response.headers["X-Options-Close"] = "1"
if task.mentally_hard:
# «Ментально сложная» — фронт украшает праздник («сила воли»)
response.headers["X-Mentally-Hard"] = "1"
response.headers["X-Earned-Coins"] = str(outcome.coins)
earned_event = event
task.done_by_kind = actor.kind
# Приёмка (ТЗ 3.20) — отдельная плоскость: агенту работу ещё примут,
# владелец, закрывший задачу сам, приёмкой не занимается
task.accept_state = "pending" if actor.is_agent else None
elif not actor.is_agent:
# Владелец перезакрыл уже закрытую задачу — его закрытие снимает
# ожидание приёмки: принимать работу он будет не сам у себя
task.accept_state = None
elif data.get("status") not in (None, "done") and task.done_at is not None:
# Выход из done: отметка выполнения сбрасывается. Повторное закрытие
# снова пройдёт через ветку наград, но XP grant_task_xp повторно не
# даст, а спавн удержит spawned_at — дубликатов не будет
task.done_at = None
# Журнал (ТЗ 3.20): в него идёт всё жизненное — закрытие, смена статуса,
# выдача мандата. Записи только о настоящих переменах: PATCH, ничего не
# изменивший, журнал не засоряет
if data.get("status") and task.status != old_status:
if task.status == "done" and old_status != "done":
tasklog.log_event(db, task, tasklog.KIND_COMPLETED, actor, close_comment)
else:
tasklog.log_event(
db,
task,
tasklog.KIND_STATUS,
actor,
tasklog.status_transition(old_status, task.status),
)
if "ai_eligible" in data and task.ai_eligible != old_eligible:
tasklog.log_event(
db, task, tasklog.KIND_ELIGIBLE, actor, "on" if task.ai_eligible else "off"
)
# Аренда (ТЗ 3.20): действие владельца её снимает (задача снова свободна для
# агентов), а выход из работы — тем более: держать аренду на закрытой задаче
# значило бы «оживить» её следующим heartbeat'ом агента
if not actor.is_agent or task.status not in ("to_do", "in_progress"):
claims.clear_claim(task)
db.flush()
db.refresh(task) # перечитать связи (project/tags) после обновления
# События после коммита: подписчик SSE может сразу перечитать данные
db.commit()
publish(user_id, "task.changed", {"id": task.id})
if task.accept_state == "pending" and old_accept != "pending":
# Работа агента ждёт приёмки (ТЗ 3.20): владельцу — пометка и тост
publish(user_id, "task.review", {"id": task.id, "title": task.title})
# …и системное уведомление на устройства (ТЗ 3.21): приёмки ждут ровно
# тогда, когда вкладка закрыта. Отправка — вне ответа: pywebpush ходит
# в чужой push-сервис и блокирует поток
background.add_task(notify_agent_review, task.id)
if earned_event is not None:
publish(
user_id,
"xp.changed",
{"amount": earned_event.amount, "rarity": earned_event.rarity, "celebrate": True},
)
publish(user_id, "garden.changed", {"reason": "task"})
return task
@router.post("/{task_id}/claim", response_model=TaskOut)
async def claim_task(
task_id: int, db: DbDep, user_id: UserIdDep, actor: ActorDep, background: BackgroundTasks
) -> Task:
"""Взять задачу в работу (ТЗ 3.20): аренда на 2 часа, повторный вызов — продление."""
task = _get_task_for_update(db, task_id, user_id)
try:
fresh = claims.claim_task(task, actor)
except claims.MandateError as exc:
raise _mandate_error(exc) from None
if fresh:
tasklog.log_event(db, task, tasklog.KIND_CLAIMED, actor)
db.flush()
db.refresh(task)
db.commit()
# Продление (heartbeat) молчит: иначе агент, работающий часами, засыпал бы
# владельца событиями — а в задаче ничего не изменилось
if fresh:
publish(user_id, "task.changed", {"id": task.id})
publish(user_id, "task.claimed", {"id": task.id, "claimed_by": task.claimed_by})
# Системное уведомление (ТЗ 3.21): агент берётся за работу как раз тогда,
# когда владелец закрыл вкладку — иначе узнал бы об этом через два часа
background.add_task(notify_agent_claimed, task.id, task.claimed_by or "")
return task
@router.post("/{task_id}/release", response_model=TaskOut)
async def release_task(
task_id: int, db: DbDep, user_id: UserIdDep, actor: ActorDep, reason: str = Query("")
) -> Task:
"""Вернуть задачу в свободные: агент понял, что она ему не по силам (ТЗ 3.20)."""
task = _get_task_for_update(db, task_id, user_id)
try:
had_lease = claims.release_claim(task, actor)
except claims.MandateError as exc:
raise _mandate_error(exc) from None
if had_lease:
tasklog.log_event(db, task, tasklog.KIND_RELEASED, actor, reason or None)
db.flush()
db.refresh(task)
db.commit()
if had_lease:
publish(user_id, "task.changed", {"id": task.id})
return task
def _ensure_owner_accepts(actor: Actor) -> None:
"""Работу агента принимает и возвращает только владелец (ТЗ 3.20)."""
if actor.is_agent:
raise HTTPException(status_code=403, detail="Acceptance is the owner's decision")
@router.post("/{task_id}/accept", response_model=TaskOut)
async def accept_task(
task_id: int,
db: DbDep,
user_id: UserIdDep,
actor: ActorDep,
schema: AcceptIn | None = None,
) -> Task:
"""Принять работу агента: pending → accepted. Задача остаётся закрытой."""
task = _get_task_for_update(db, task_id, user_id)
_ensure_owner_accepts(actor)
if task.accept_state != "pending":
raise HTTPException(status_code=409, detail="Task is not waiting for acceptance")
task.accept_state = "accepted"
tasklog.log_event(db, task, tasklog.KIND_ACCEPTED, actor, schema.comment if schema else None)
db.flush()
db.refresh(task)
db.commit()
publish(user_id, "task.changed", {"id": task.id})
return task
@router.post("/{task_id}/reject", response_model=TaskOut)
async def reject_task(
task_id: int,
db: DbDep,
user_id: UserIdDep,
actor: ActorDep,
schema: AcceptIn | None = None,
) -> Task:
"""Вернуть работу агента в работу: задача снова to_do, аренда снята.
XP и монеты не отзываются (ТЗ 3.13, «только позитив»), а spawned_at не
трогаем: повторное закрытие после возврата не даст ни второй награды
(идемпотентность по задаче), ни второго экземпляра регулярной (ТЗ 8.7).
"""
task = _get_task_for_update(db, task_id, user_id)
_ensure_owner_accepts(actor)
if task.accept_state != "pending":
raise HTTPException(status_code=409, detail="Task is not waiting for acceptance")
task.accept_state = "rejected"
task.status = "to_do"
task.done_at = None
claims.clear_claim(task)
tasklog.log_event(db, task, tasklog.KIND_REJECTED, actor, schema.comment if schema else None)
db.flush()
db.refresh(task)
db.commit()
publish(user_id, "task.changed", {"id": task.id})
return task
@router.post("/{task_id}/approve", response_model=TaskOut)
async def approve_task(
task_id: int, db: DbDep, user_id: UserIdDep, schema: ApproveIn | None = None
) -> Task:
"""Утверждение детализации: raw → approved.
schema.apply_proposal=True — принять предложение автодетализации перед
утверждением (кнопка «Да, всё верно»); иначе предложенные метаданные
игнорируются (после ручной правки их уже применил PATCH).
"""
task = _get_task_or_404(db, task_id, user_id)
if schema and schema.apply_proposal and task.ai_proposal:
apply_proposal(db, task, task.ai_proposal, user_id)
task.detail_state = "approved"
task.approved_at = utcnow()
db.flush()
db.refresh(task)
db.commit()
publish(user_id, "task.changed", {"id": task.id})
return task
@router.post("/{task_id}/redetail", response_model=TaskOut)
async def redetail_task(
task_id: int, db: DbDep, user_id: UserIdDep, background: BackgroundTasks
) -> Task:
"""Перезапустить автодетализацию (например, после появления новых тегов/проектов).
Переспрос — это сброс предложения и анализ заново по текущему тексту задачи
(владелец мог дополнить заголовок или описание). Правила анализа те же, что
у первичного: теги дополняются, приоритет и оценка — только в пустые поля,
поэтому переспрос не может перетереть значения пользователя.
"""
task = _get_task_or_404(db, task_id, user_id)
if task.detail_state != "raw":
# Утверждённая задача уже детализирована: повтор — no-op, предложение
# не сбрасываем, фоновую LLM не будим (worker всё равно работает
# только с raw)
return task
task.ai_proposal = None
# Коммит до фоновой работы, чтобы worker не прочитал старое состояние
# и его запись не перезаписалась teardown-коммитом (см. create_task).
db.commit()
publish(user_id, "task.changed", {"id": task.id})
background.add_task(detail_task, task.id)
return task
@router.post("/suggest", response_model=list[TaskOut])
async def suggest_tasks(schema: SuggestIn, db: DbDep, user_id: UserIdDep) -> list[Task]:
"""Режим «3 варианта»: до трёх задач под доступное время."""
return pick_options(db, schema.available_minutes, user_id)
@router.delete("/{task_id}")
async def delete_task(
task_id: int, db: DbDep, user_id: UserIdDep, actor: ActorDep
) -> dict[str, bool]:
task = _get_task_or_404(db, task_id, user_id)
# Журнал переживает задачу: запись делает снапшот заголовка (ТЗ 3.20)
tasklog.log_deleted(db, task, actor)
db.delete(task)
db.commit()
publish(user_id, "task.deleted", {"id": task_id})
return {"ok": True}