Newer
Older
gnexus-tasks / backend / app / api / tasks.py
"""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})
    # force=True: переспрос — LLM анализирует задачу целиком, включая поля,
    # заполненные пользователем (при первичном анализе они не трогаются)
    background.add_task(detail_task, task.id, True)
    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}