Newer
Older
gnexus-tasks / backend / app / services / push.py
"""Отправка системных уведомлений на устройства пользователя (ТЗ 3.21).

Модуль синхронный и **зовётся только вне event loop**: `pywebpush` ходит в
push-сервис браузера через `requests` и на время запроса блокирует поток. Из REST
это `BackgroundTasks` (FastAPI выполняет их в threadpool), из MCP-тулов —
`Thread(daemon=True)`, у планировщика — свой `asyncio.to_thread`. Каждый вызов
открывает свою сессию: сбой отправки не должен ронять запрос, а чужая сессия —
переживать коммит вызывающего.

Мёртвые подписки убираются лениво, при отправке: ответ 404/410 означает, что
браузер отозвал подписку, а серия прочих ошибок — что устройство потеряно.
Отдельной уборки по расписанию нет намеренно: она бы чистила то, что и так
чистится само, но требовала бы ещё одного тика.
"""

import json
import logging

import requests
from pywebpush import WebPushException, webpush
from sqlalchemy import select
from sqlalchemy.orm import Session

from app.config import get_settings
from app.db import get_session_factory
from app.models import PushSubscription, Task, utcnow
from app.services.push_texts import render, user_language

logger = logging.getLogger(__name__)

# Push-сервис принимает payload до ~4 КБ вместе с шифрованием и заголовками;
# 3500 — запас, чтобы aes128gcm-обвязка не уткнулась в лимит на длинном теле
MAX_PAYLOAD = 3500
# Уведомление, не доставленное за час (телефон офлайн), теряет смысл — оно про
# «сейчас»; push-сервис просто выбросит его, если окно истекло
DEFAULT_TTL = 3600
REQUEST_TIMEOUT = 10
# Сколько подряд неудач терпим, прежде чем счесть подписку мёртвой
MAX_FAILURES = 5
# Коды, которыми push-сервис сообщает «подписки больше нет»
GONE_CODES = (404, 410)


def build_payload(*, kind: str, title: str, body: str, url: str, tag: str) -> str:
    """JSON уведомления для service worker: то, что он покажет, и куда вести по клику.

    Схема минимальна: заголовок, тело, ссылка внутри приложения и `tag` для
    схлопывания — повторное уведомление о той же задаче заменяет предыдущее, а не
    копится в шторке.
    """
    payload = {"kind": kind, "title": title, "body": body, "url": url, "tag": tag}
    data = json.dumps(payload, ensure_ascii=False)
    if len(data.encode()) <= MAX_PAYLOAD:
        return data
    # Усекаем тело, а не выбрасываем уведомление: заголовок и ссылка важнее.
    # Считаем в байтах, а не в символах — кириллица в UTF-8 занимает по два, и
    # «усечь до 3500 символов» на русском теле дало бы вдвое больший payload.
    empty = json.dumps({**payload, "body": ""}, ensure_ascii=False)
    budget = MAX_PAYLOAD - len(empty.encode()) - len("…".encode())
    if budget <= 0:
        return empty
    payload["body"] = body.encode()[:budget].decode("utf-8", "ignore").rstrip() + "…"
    return json.dumps(payload, ensure_ascii=False)


def _drop(db: Session, sub: PushSubscription, reason: str) -> None:
    logger.info("push: подписка %s удалена (%s)", sub.id, reason)
    db.delete(sub)


def _deliver(db: Session, user_id: str, payload: str, ttl: int = DEFAULT_TTL) -> int:
    """Разослать payload на все устройства пользователя; вернуть число успехов.

    Коммитит сама: подписки удаляются и обновляются по ходу, и незакоммиченные
    изменения потерялись бы вместе с закрытием сессии.
    """
    settings = get_settings()
    if not settings.push_enabled:
        return 0

    subscriptions = db.scalars(
        select(PushSubscription).where(PushSubscription.user_id == user_id)
    ).all()
    if not subscriptions:
        return 0

    claims = {"sub": settings.vapid_subject or "mailto:noreply@example.com"}
    sent = 0
    for sub in subscriptions:
        try:
            webpush(
                {
                    "endpoint": sub.endpoint,
                    "keys": {"p256dh": sub.p256dh, "auth": sub.auth},
                },
                data=payload,
                vapid_private_key=settings.vapid_private_key,
                vapid_claims=dict(claims),
                content_encoding="aes128gcm",
                timeout=REQUEST_TIMEOUT,
                ttl=ttl,
            )
        except WebPushException as exc:
            if exc.status_code in GONE_CODES:
                _drop(db, sub, f"push-сервис ответил {exc.status_code}")
                continue
            sub.failure_count += 1
            if sub.failure_count >= MAX_FAILURES:
                _drop(db, sub, f"{sub.failure_count} неудач подряд")
            else:
                logger.warning("push: подписка %s, ошибка отправки: %s", sub.id, exc)
        except requests.RequestException as exc:
            # Сеть недоступна — это не приговор подписке, но и не успех
            sub.failure_count += 1
            logger.warning("push: подписка %s, сеть недоступна: %s", sub.id, exc)
        else:
            sub.failure_count = 0
            sub.last_success_at = utcnow()
            sent += 1
    db.commit()
    return sent


def send_to_user(user_id: str, payload: str, ttl: int = DEFAULT_TTL) -> int:
    """Отправить готовый payload всем устройствам пользователя (своя сессия)."""
    db = get_session_factory()()
    try:
        return _deliver(db, user_id, payload, ttl)
    except Exception:  # noqa: BLE001 — фон не имеет права падать наружу
        db.rollback()
        logger.exception("push: не удалось отправить уведомление пользователю %s", user_id)
        return 0
    finally:
        db.close()


def _notify_task(task_id: int, kind: str, **fields: str) -> None:
    """Собрать текст по задаче и отправить его (одна сессия на всё)."""
    db = get_session_factory()()
    try:
        task = db.get(Task, task_id)
        if task is None or not task.user_id:
            return
        rendered = render(kind, user_language(db, task.user_id), title=task.title, **fields)
        if rendered is None:
            return
        title, body = rendered
        payload = build_payload(
            kind=kind,
            title=title,
            body=body,
            url=f"/tasks/{task_id}",
            tag=f"{kind}:{task_id}",
        )
        _deliver(db, task.user_id, payload)
    except Exception:  # noqa: BLE001 — фон не имеет права падать наружу
        db.rollback()
        logger.exception("push: не удалось уведомить о задаче %s", task_id)
    finally:
        db.close()


def notify_agent_claimed(task_id: int, agent_name: str) -> None:
    """Агент взял задачу в работу (task.claimed)."""
    _notify_task(task_id, "task.claimed", name=agent_name)


def notify_agent_review(task_id: int) -> None:
    """Работа закрыта и ждёт приёмки владельца (task.review)."""
    _notify_task(task_id, "task.review")


def save_subscription(
    db: Session,
    user_id: str,
    endpoint: str,
    p256dh: str,
    auth: str,
    user_agent: str | None,
) -> PushSubscription:
    """Upsert подписки по endpoint.

    endpoint уникален глобально, а не в пределах пользователя: тот же браузерный
    профиль может войти другим аккаунтом, и тогда подписка должна переехать к
    нему, а не размножиться. Поэтому чужую строку с тем же endpoint перепривязываем.
    """
    sub = db.scalar(select(PushSubscription).where(PushSubscription.endpoint == endpoint))
    if sub is None:
        sub = PushSubscription(user_id=user_id, endpoint=endpoint, p256dh=p256dh, auth=auth)
        db.add(sub)
    sub.user_id = user_id
    sub.p256dh = p256dh
    sub.auth = auth
    sub.user_agent = (user_agent or "")[:300] or None
    # Переподписка означает живое устройство: обнуляем счётчик неудач, иначе
    # подписка, пережившая пару сбоев, умерла бы сразу после возвращения
    sub.failure_count = 0
    db.flush()
    return sub


def delete_subscription(db: Session, user_id: str, endpoint: str) -> bool:
    """Отписаться идемпотентно: нет подписки — тоже успех."""
    sub = db.scalar(
        select(PushSubscription).where(
            PushSubscription.endpoint == endpoint, PushSubscription.user_id == user_id
        )
    )
    if sub is None:
        return False
    db.delete(sub)
    return True