"""Планировщик напоминаний (ТЗ 3.21, подраздел «Напоминания»).
Тонкая обвязка вокруг правил из `services/reminders.py`: правила ничего не знают
ни о часах, ни о базе, а здесь они получают «сегодня» по локальному поясу,
отсеиваются по журналу отправленного (`push_deliveries`) и уходят в push.
Задача живёт в lifespan приложения, а не в отдельном процессе: напоминания —
часть приложения, и его перезапуск не должен оставлять сироту (см. `app/main.py`).
Итерация идёт в `asyncio.to_thread`: внутри синхронные SQLAlchemy и `requests`, и
в event loop их пускать нельзя.
Дедуп — вставкой в `push_deliveries` **до** отправки. Тик 15-минутный, а правила
дневные, поэтому без журнала одна и та же просрочка приходила бы сотню раз;
обратный размен (упасть между вставкой и отправкой) даёт пропуск, а не дубль —
для напоминания это правильная сторона ошибки.
"""
import asyncio
import logging
from datetime import UTC, date, datetime
from zoneinfo import ZoneInfo
from sqlalchemy import select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from app.config import get_settings
from app.db import get_session_factory
from app.models import AppSetting, PushDelivery, PushSubscription, Task
from app.services import push, push_texts, reminders
logger = logging.getLogger(__name__)
# Тик планировщика. Правила дневные, но пятнадцать минут дают попадание в окно
# сводки и в тихие часы с запасом, а стоят почти ничего
TICK_SECONDS = 900
# Тот же ключ, что `QUIET_HOURS_KEY` в api/settings.py: импортировать оттуда значило
# бы развернуть зависимость «сервис → API» ради одного строкового литерала
QUIET_HOURS_KEY = "quiet_hours"
async def scheduler_loop() -> None:
"""Крутится, пока живёт приложение: тик, поток на итерацию, пауза.
Падение одной итерации логируем и продолжаем: сбой на конкретной задаче не
повод перестать напоминать обо всех остальных до перезапуска API.
"""
while True:
try:
await asyncio.to_thread(run_iteration)
except asyncio.CancelledError:
raise # остановка приложения — не сбой итерации
except Exception: # noqa: BLE001 — цикл обязан пережить любую итерацию
logger.exception("напоминания: итерация планировщика упала")
await asyncio.sleep(TICK_SECONDS)
def quiet_hours(db: Session, user_id: str) -> str | None:
"""Тихое окно пользователя («HH:MM-HH:MM») или None, если не задано."""
row = db.get(AppSetting, (user_id, QUIET_HOURS_KEY))
return row.value if row is not None else None
def run_iteration(now: datetime | None = None) -> None:
"""Один проход: собрать напоминания всем подписанным устройствам и отправить.
`now` — только для тестов и ручных прогонов; в бою берётся текущий момент.
Пользователи без подписок пропускаются сразу: напоминание некуда доставить,
а его журнал всё равно не нужен — при подписке оно придёт в свой день.
"""
settings = get_settings()
if not settings.push_enabled:
return
moment = (now or datetime.now(UTC)).astimezone(ZoneInfo(settings.reminder_timezone))
today = moment.date()
db = get_session_factory()()
try:
user_ids = db.scalars(select(PushSubscription.user_id).distinct()).all()
for user_id in user_ids:
try:
_run_for_user(db, user_id, moment, today)
except Exception: # noqa: BLE001 — один пользователь не роняет остальных
db.rollback()
logger.exception("напоминания: сбой у пользователя %s", user_id)
finally:
db.close()
def _run_for_user(db: Session, user_id: str, moment: datetime, today: date) -> None:
"""Собрать и отправить напоминания одного пользователя."""
# Тихие часы: пропускаем целиком и **не пишем в журнал** — условие напоминания
# дневное («сегодня день D»), поэтому следующий тик после окна отправит его сам
if reminders.in_quiet_hours(quiet_hours(db, user_id), moment):
return
tasks = db.scalars(select(Task).where(Task.user_id == user_id)).all()
language = push_texts.user_language(db, user_id)
pending: list[tuple[reminders.Reminder, str]] = []
for task in tasks:
for item in reminders.task_reminders(task, today):
payload = _payload(item, language)
if payload is not None:
pending.append((item, payload))
if reminders.summary_due(moment):
counts = reminders.summary_counts(tasks, today)
# Сводка без новостей — молчание, и в журнал она тоже не пишется: иначе
# завтрашняя сводка считалась бы отправленной сегодняшней
if any(counts.values()):
item = reminders.Reminder(
reminders.KIND_SUMMARY,
today.isoformat(),
None,
{"parts": push_texts.summary_body(language, counts)},
)
payload = _payload(item, language, url="/")
if payload is not None:
pending.append((item, payload))
if not pending:
return
fresh = _claim(db, user_id, pending)
db.commit()
for item, payload in fresh:
push.send_to_user(user_id, payload)
logger.info("напоминания: %s → %s (%s)", item.kind, user_id, item.ref_key)
def _payload(item: reminders.Reminder, language: str, url: str | None = None) -> str | None:
"""Собрать payload напоминания; неизвестный вид (нет текста) → None."""
rendered = push_texts.render(item.kind, language, **item.fields)
if rendered is None:
return None
title, body = rendered
target = url if url is not None else f"/tasks/{item.task_id}"
# tag как у событий агента: напоминание об одной задаче заменяет предыдущее,
# а не копится в шторке. У сводки задачи нет — её тег и есть вид
tag = f"{item.kind}:{item.task_id}" if item.task_id is not None else item.kind
return push.build_payload(kind=item.kind, title=title, body=body, url=target, tag=tag)
def _claim(
db: Session, user_id: str, pending: list[tuple[reminders.Reminder, str]]
) -> list[tuple[reminders.Reminder, str]]:
"""Отметить напоминания отправленными; вернуть те, что отправляются впервые.
Отметка идёт до отправки, поэтому повторный проход (или перезапуск API) те же
напоминания уже не отправит. SAVEPOINT на каждую вставку — чтобы нарушение
уникальности не откатило соседние отметки той же транзакции.
"""
fresh: list[tuple[reminders.Reminder, str]] = []
for item, payload in pending:
savepoint = db.begin_nested()
db.add(PushDelivery(user_id=user_id, kind=item.kind, ref_key=item.ref_key))
try:
savepoint.commit()
except IntegrityError:
savepoint.rollback()
continue
fresh.append((item, payload))
return fresh