Files
vidconf/workers/tasks/invitations.py
Max Ronzhin 74de51dbd8 feat(email): контактный адрес инстанса, Reply-To в письмах и тестовая отправка
Новая настройка instance_settings.contact_email (включён/адрес, с
валидацией формата) — подставляется в заголовок Reply-To писем
подтверждения регистрации, приглашений и саммари. Админ-эндпоинт
POST /admin/settings/test-email отправляет проверочное письмо синхронно
и возвращает внятный результат (успех либо текст ошибки транспорта),
не раскрывая логин/пароль SMTP.
2026-07-27 00:35:47 +03:00

294 lines
14 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Celery-задача рассылки .ics-приглашений на конференцию.
Ставится из `services.conferences.ConferenceService.create` (при создании
плановой/закреплённой конференции) и `.update` (при правке расписания —
`scheduled_at`/`duration_minutes`/`recurrence`/`title`, см. инкремент
`ics_sequence` там же, а также при изменении СОСТАВА участников без
изменения расписания — без инкремента `ics_sequence`)
через `services.invitations_producer.enqueue_invitations` — backend не
импортирует пакет `workers` напрямую (тот же приём, что
`services.pipeline_producer.enqueue_pipeline`). Получатели по умолчанию
(`emails=None`) — владелец конференции + приглашённые (`conference_invitees`,
ADR-003: зарегистрированные — по email пользователя, внешние — по указанному
email) + для закреплённых дополнительно уникальные участники её прошлых
сеансов (зарегистрированные пользователи и гости с указанным email).
Явный список `emails` — ручная рассылка администратором
(`POST /admin/conferences/{id}/invitations`).
Идемпотентность: `email_deliveries(kind='invitation')` — журнал БЕЗ unique
(переслать обновлённое приглашение при изменении расписания обязано снова
дойти до всех адресатов — задокументированный дубль лучше пропавшего
уведомления об изменении времени встречи).
"""
import logging
import uuid
from datetime import UTC, datetime, timedelta
from typing import Any, Protocol
from zoneinfo import ZoneInfo
from celery.exceptions import MaxRetriesExceededError
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from core.config import get_settings
from core.plugins.config import InstanceConfig
from models.conference import Conference
from models.email_delivery import EmailDelivery
from models.guest import GuestAccess
from models.invitee import ConferenceInvitee
from models.participant import ConferenceParticipant
from models.session import ConferenceSession
from models.user import User
from services.email import EmailAttachment, EmailSendError, create_email_backend
from services.ics import ConferenceHasNoScheduleError, build_invite
from services.instance_settings import load_effective_config
from services.recurrence import RecurrenceRule, expand_occurrences
from workers.celery_app import app
from workers.db import open_session, run_async
# Горизонт поиска первого вхождения повторяющейся серии для текста письма —
# совпадает с `services.ics._OCCURRENCE_SEARCH_HORIZON` (не импортируется,
# приватный модульный атрибут, но значение то же).
_OCCURRENCE_SEARCH_HORIZON = timedelta(days=400)
logger = logging.getLogger(__name__)
RETRY_COUNTDOWN_BASE_S = 60
"""Базовая пауза (сек) перед повтором при временном сбое SMTP; фактический
countdown — `RETRY_COUNTDOWN_BASE_S * (attempt + 1)` (тот же паттерн, что в
`workers.tasks.notify`/`summarize`)."""
_ICS_ATTACHMENT_FILENAME = "invite.ics"
_ICS_MIME_TYPE = "text/calendar; method=REQUEST"
class RetryableTask(Protocol):
"""Минимальный протокол bound-задачи Celery (см. одноимённые протоколы
в `workers.tasks.summarize`/`notify`)."""
request: Any
def retry(self, countdown: int | None = None) -> None: ...
@app.task(
name="workers.tasks.invitations.send_invitations",
bind=True,
max_retries=5,
acks_late=True,
)
def send_invitations(
self: RetryableTask, conference_id: str, emails: list[str] | None = None
) -> None:
"""Точка входа Celery — синхронная обёртка над асинхронной логикой рассылки."""
run_async(lambda: send_invitations_async(self, uuid.UUID(conference_id), emails=emails))
async def send_invitations_async(
task: RetryableTask,
conference_id: uuid.UUID,
*,
emails: list[str] | None = None,
plugins_config: InstanceConfig | None = None,
) -> None:
"""Разослать .ics-приглашение на конференцию `conference_id` получателям.
Guard'ы:
1. Конференция не найдена — выход.
2. Нет ни `scheduled_at`, ни `recurrence` (расписания нет — строить
приглашение не из чего, например, конференция уже разоткреплена и
переведена в мгновенную) — выход без ошибки.
3. Получателей нет (ни явного списка, ни владельца/участников) — выход.
Постоянный отказ конкретного получателя не роняет рассылку остальным
(пропускается с предупреждением); временный сбой транспорта — `task.retry`
с нарастающим countdown, исчерпание попыток — только лог (в отличие от
`notify_session`, здесь нет `pipeline_status`, который можно провалить).
"""
async with open_session() as session:
conference = await session.get(Conference, conference_id)
if conference is None:
logger.warning("send_invitations: конференция %s не найдена", conference_id)
return
if conference.recurrence is None and conference.scheduled_at is None:
logger.info(
"send_invitations: у конференции %s нет расписания — no-op", conference_id
)
return
recipients = (
[email.lower() for email in emails]
if emails is not None
else await _resolve_default_recipients(session, conference)
)
if not recipients:
logger.info("send_invitations: у конференции %s нет получателей", conference_id)
return
cfg = plugins_config or await load_effective_config(session)
settings = get_settings()
join_url = f"{settings.frontend_url}/j/{conference.slug}"
organizer_email = await _resolve_owner_email(session, conference)
try:
ics_bytes = build_invite(
conference,
organizer_email=organizer_email,
join_url=join_url,
display_timezone=cfg.display_timezone,
)
except ConferenceHasNoScheduleError:
logger.warning(
"send_invitations: не удалось построить .ics для конференции %s", conference_id
)
return
attachment = EmailAttachment(
filename=_ICS_ATTACHMENT_FILENAME, content=ics_bytes, mime_type=_ICS_MIME_TYPE
)
title = conference.title or f"{conference.number}"
subject = f"Приглашение на конференцию: {title}"
when = _format_when(conference, display_timezone=cfg.display_timezone)
body_lines = [f"Вас пригласили на конференцию «{title}»."]
if when is not None:
body_lines.append(f"Дата и время: {when}")
body_lines.append(f"Ссылка для входа: {join_url}")
body_lines.append(f"Номер: {conference.number}")
# Состав участников — переиспользуем уже
# посчитанный список получателей, чтобы не делать лишний запрос.
body_lines.append(f"Участники: {', '.join(recipients)}")
body = "\n".join(body_lines)
backend = create_email_backend(settings)
reply_to = cfg.contact_email if cfg.contact_email_enabled else None
for email in recipients:
try:
await backend.send(
to=email,
subject=subject,
body=body,
attachments=[attachment],
reply_to=reply_to,
)
except EmailSendError as exc:
if not exc.retryable:
logger.warning(
"send_invitations: получатель %s конференции %s отклонён сервером — "
"пропущен без повторных попыток: %s",
email,
conference_id,
exc,
)
continue
attempt = getattr(task.request, "retries", 0)
countdown = RETRY_COUNTDOWN_BASE_S * (attempt + 1)
try:
task.retry(countdown=countdown)
except MaxRetriesExceededError:
logger.warning(
"send_invitations: исчерпаны попытки рассылки приглашения "
"конференции %s (SMTP недоступен: %s)",
conference_id,
exc,
)
# Реальный `Task.retry()` сам бросает исключение Retry (не
# возвращает управление) — до сюда доходим только с
# моком/заглушкой `task.retry` в тестах.
return
else:
session.add(
EmailDelivery(
conference_id=conference.id, recipient_email=email, kind="invitation"
)
)
await session.commit()
logger.info(
"send_invitations: конференция %s — приглашение отправлено %d получателям",
conference_id,
len(recipients),
)
def _format_when(conference: Conference, *, display_timezone: str) -> str | None:
"""Дата/время конференции в таймзоне отображения инстанса — для тела письма.
Разовая — `scheduled_at`; закреплённая с повторением — первое вхождение
от `anchor_date` (тот же горизонт поиска, что `services/ics.py`). `None`,
если расписания нет (проверяется раньше в вызывающем коде — сюда такая
конференция не доходит).
"""
dt: datetime
if conference.recurrence is not None:
rule = RecurrenceRule.model_validate(conference.recurrence)
t_from = rule.local_datetime(rule.anchor_date).astimezone(UTC)
occurrences = expand_occurrences(rule, t_from, t_from + _OCCURRENCE_SEARCH_HORIZON)
if not occurrences:
return None
dt = occurrences[0]
elif conference.scheduled_at is not None:
dt = conference.scheduled_at
else:
return None
localized = dt.astimezone(ZoneInfo(display_timezone))
return localized.strftime("%d.%m.%Y %H:%M %Z")
async def _resolve_owner_email(session: AsyncSession, conference: Conference) -> str | None:
"""Email владельца конференции (используется как `ORGANIZER` .ics), либо `None`."""
if conference.owner_id is None:
return None
owner = await session.get(User, conference.owner_id)
return owner.email if owner is not None else None
async def _resolve_default_recipients(session: AsyncSession, conference: Conference) -> list[str]:
"""Получатели по умолчанию: владелец + приглашённые + (для закреплённых) участники прошлых сеансов.
Дедуп по `lower(email)`. Приглашённые
(`conference_invitees`, ADR-003) добавляются независимо от `is_pinned` —
состав задаётся на уровне конференции, а не сеанса. Для незакреплённой
(разовой плановой) конференции без приглашённых сеансов ещё нет —
получатель только один (владелец).
"""
recipients: dict[str, None] = {}
owner_email = await _resolve_owner_email(session, conference)
if owner_email:
recipients.setdefault(owner_email.lower(), None)
invitee_rows = await session.execute(
select(ConferenceInvitee.email, User.email)
.outerjoin(User, ConferenceInvitee.user_id == User.id)
.where(ConferenceInvitee.conference_id == conference.id)
)
for invitee_email, invitee_user_email in invitee_rows.all():
email = invitee_email or invitee_user_email
if email:
recipients.setdefault(email.lower(), None)
if conference.is_pinned:
user_rows = await session.execute(
select(User.email)
.join(ConferenceParticipant, ConferenceParticipant.user_id == User.id)
.join(ConferenceSession, ConferenceParticipant.session_id == ConferenceSession.id)
.where(ConferenceSession.conference_id == conference.id)
.distinct()
)
guest_rows = await session.execute(
select(GuestAccess.email)
.join(ConferenceParticipant, ConferenceParticipant.guest_id == GuestAccess.id)
.join(ConferenceSession, ConferenceParticipant.session_id == ConferenceSession.id)
.where(
ConferenceSession.conference_id == conference.id,
GuestAccess.email.isnot(None),
)
.distinct()
)
for email in [*user_rows.scalars().all(), *guest_rows.scalars().all()]:
if email:
recipients.setdefault(email.lower(), None)
return list(recipients.keys())