diff --git a/alembic/versions/20261003_0005_event_expired.py b/alembic/versions/20261003_0005_event_expired.py new file mode 100644 index 0000000..ce6b0c5 --- /dev/null +++ b/alembic/versions/20261003_0005_event_expired.py @@ -0,0 +1,28 @@ +""" Статус события expired: ttl_seconds реально применяется, ретеншн чистит старое + +Revision ID: a1c4e7d52b98 +Revises: c8e4f2b77d10 +Create Date: 2026-10-03 +""" + +from alembic import op + +revision = "a1c4e7d52b98" +down_revision = "c8e4f2b77d10" +branch_labels = None +depends_on = None + +STATUS_NEW = "status IN ('queued', 'processing', 'done', 'failed', 'expired')" +STATUS_OLD = "status IN ('queued', 'processing', 'done', 'failed')" + + +def upgrade() -> None: + op.drop_constraint("ck_event_status", "events", type_="check") + op.create_check_constraint("ck_event_status", "events", STATUS_NEW) + + +def downgrade() -> None: + # События в expired обратно в queued не возвращаем — просто роняем + # constraint; их нет после первого прогона ретеншна. + op.drop_constraint("ck_event_status", "events", type_="check") + op.create_check_constraint("ck_event_status", "events", STATUS_OLD) \ No newline at end of file diff --git a/app/config.py b/app/config.py index 569c693..9c4a4b1 100644 --- a/app/config.py +++ b/app/config.py @@ -43,6 +43,11 @@ # first — только первое подошедшее по weight routing_match_mode: str = "all" + # Ретеншн (beat-задача synapse.expire_events): события старше N дней + # удаляются вместе с доставками; ttl_seconds удаляет раньше — по своей + # шкале. 0 — окно по возрасту отключено (ttl продолжит работать). + retention_days: int = 30 + # SPA лежит в образе рядом с приложением; в dev используется vite-сервер. spa_dist_dir: Path = BASE_DIR / "spa_static" diff --git a/app/models/events.py b/app/models/events.py index 2584720..c33b5ce 100644 --- a/app/models/events.py +++ b/app/models/events.py @@ -35,7 +35,10 @@ __tablename__ = "events" __table_args__ = ( CheckConstraint("priority IN ('low','normal','high','critical')", name="ck_priority"), - CheckConstraint("status IN ('queued','processing','done','failed')", name="ck_event_status"), + # expired — истёкший ttl_seconds: воркер не маршрутизирует, ретеншн удалит. + CheckConstraint( + "status IN ('queued','processing','done','failed','expired')", name="ck_event_status" + ), Index("ix_events_dedup", "dedup_key", "created_at"), ) @@ -49,7 +52,8 @@ payload: Mapped[dict] = mapped_column(JSONB, default=dict) dedup_key: Mapped[str | None] = mapped_column(String(255)) - # created_at + ttl_seconds; воркер не берёт события с expires_at < now. + # created_at + ttl_seconds (считается на приёме); воркер не маршрутизирует + # истёкшее (status=expired), ретеншн удалит (docs/05, «Ретеншн»). expires_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) # Резерв контракта: в MVP scheduling не реализован. scheduled_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) diff --git a/app/worker/celery_app.py b/app/worker/celery_app.py index 157c4ec..170be77 100644 --- a/app/worker/celery_app.py +++ b/app/worker/celery_app.py @@ -16,13 +16,17 @@ celery_app.conf.update( task_default_queue="synapse", task_track_started=True, - # Ретраи доставок: beat стартует вместе с воркером (`worker -B`), - # раз в 60 с просмотр pending с наступившим next_retry_at. + # Beat стартует вместе с воркером (`worker -B`): + # раз в 60 с — ретраи; ежечасно — ретеншн (ttl_seconds + RETENTION_DAYS). beat_schedule={ "retry-due-deliveries": { "task": "synapse.retry_due", "schedule": 60.0, }, + "expire-old-events": { + "task": "synapse.expire_events", + "schedule": 3600.0, + }, }, task_acks_late=True, worker_prefetch_multiplier=1, diff --git a/app/worker/tasks.py b/app/worker/tasks.py index d35d4df..9678d32 100644 --- a/app/worker/tasks.py +++ b/app/worker/tasks.py @@ -1,4 +1,5 @@ -"""Задачи воркера: ping, ingest (Routing Engine), retry_due (ретраи доставок). +"""Задачи воркера: ping, ingest (Routing Engine), retry_due (ретраи), +expire_events (ретеншн, beat ежечасно). Приём (api) кладёт только конверт в БД; здесь живёт Routing Engine: - подбор правил по условиям (source/subjects/actions/priority_min/payload); @@ -7,6 +8,9 @@ лог пользователя (payload.user_id, см. docs/05), s2s — подписано (senders.py, docs/05); Telegram/Email — заглушки под #29; провалы — ретрай-цикл через beat (synapse.retry_due). +- истёкшие ttl_seconds события не маршрутизируются (status=expired) и + удаляются beat-задачей synapse.expire_events вместе с событиями старше + RETENTION_DAYS (см. «Ретеншн» в docs/05). """ import uuid @@ -14,7 +18,7 @@ from celery import shared_task from jinja2 import Template -from sqlalchemy import select +from sqlalchemy import delete, select from app.config import get_settings from app.database import SessionLocal @@ -111,6 +115,26 @@ return template_text +def _user_recipient(event: Event) -> tuple[str | None, str | None]: + """payload.user_id → (адресат, причина пропуска). + + Смысловое значение поля — sub gnexus-auth: непустая строка (число + допускаем — приводится к строке). Отсутствие поля — не ошибка источника + (None, None); тип-не-строка — ошибка источника, причина уйдёт в аудит + skipped-доставки, иначе доставка ушла бы «в никуда» молча. + """ + raw = (event.payload or {}).get("user_id") + if raw is None: + return None, None + if isinstance(raw, str): + rid = raw.strip()[:64] + return (rid, None) if rid else (None, None) + if isinstance(raw, int) and not isinstance(raw, bool): + rid = str(raw)[:64] + return (rid, None) if rid else (None, None) + return None, f"payload.user_id имеет тип {type(raw).__name__} — не похоже на sub gnexus-auth" + + def _throttled( db, rule: RoutingRule, channel: str, target_id: int | None, now: datetime, recipient: str | None = None, @@ -161,6 +185,14 @@ if event.status != "queued": return f"skip: event {event_id} в статусе {event.status}" + now = datetime.now(UTC) + # ttl_seconds: событие бесполезно после срока (docs/05) — не маршрутизируем, + # удалять будет beat-задача expire_events. + if event.expires_at is not None and event.expires_at < now: + event.status = "expired" + db.commit() + return f"skip: event {event_id} истёк ttl" + source = db.get(Source, event.source_id) event.status = "processing" @@ -189,11 +221,10 @@ ).scalars().all() for act in actions: # Канал user адресует пользователя из payload.user_id - # (конвенция, docs/05); без него доставить некому. - recipient = None + # (конвенция, docs/05); без валидного sub доставить некому. + recipient, user_reason = None, None if act.channel == "user": - rid = str((event.payload or {}).get("user_id") or "").strip() - recipient = rid[:64] or None + recipient, user_reason = _user_recipient(event) # шумодав: срезаем (с аудит-записью), critical проходит if event.priority != "critical" and _throttled( @@ -251,7 +282,7 @@ rendered_message=_render( act.template or rule.template or "", event, source.name ), - last_error="payload.user_id не задан — некому адресовать", + last_error=user_reason or "payload.user_id не задан — некому адресовать", ) ) n_skipped += 1 @@ -351,4 +382,31 @@ seconds=RETRY_DELAYS[min(delivery.attempts - 1, len(RETRY_DELAYS) - 1)] ) db.commit() - return f"retry_due: обработано {processed}" \ No newline at end of file + return f"retry_due: обработано {processed}" + + +@celery_app.task(name="synapse.expire_events") +def expire_events() -> str: + """Ретеншн (beat, ежечасно): удаление событий с истёкшим ttl_seconds + и событий старше RETENTION_DAYS; доставки падают каскадом — в том + числе записи личного лога канала user (/api/v1/me/events читает их). + RETENTION_DAYS <= 0 отключает удаление по возрасту (ttl продолжит). + """ + settings = get_settings() + now = datetime.now(UTC) + ids: list[uuid.UUID] = [] + with SessionLocal() as db: + ids += db.execute( + select(Event.id).where(Event.expires_at.is_not(None), Event.expires_at < now) + ).scalars().all() + if settings.retention_days > 0: + cutoff = now - timedelta(days=settings.retention_days) + ids += db.execute( + select(Event.id).where(Event.created_at < cutoff).limit(5000) + ).scalars().all() + ids = list(dict.fromkeys(ids)) # уникальные, порядок неважен + if not ids: + return "expire_events: удалять нечего" + db.execute(delete(Event).where(Event.id.in_(ids))) # deliveries — каскадом + db.commit() + return f"expire_events: удалено {len(ids)} событий" \ No newline at end of file diff --git a/docs/05-ingestion-api.md b/docs/05-ingestion-api.md index a8c005b..8aacf06 100644 --- a/docs/05-ingestion-api.md +++ b/docs/05-ingestion-api.md @@ -39,7 +39,7 @@ | `priority` | нет | `low \| normal \| high \| critical`, дефолт `normal`. Фиксированная шкала, не число. `critical` в будущем = обход шумодавов, отдельная очередь. | | `payload` | нет | Любой JSON. Опциональная JSON Schema типа enforced'ится воркером (не на приёме). Шаблоны каналов должны быть устойчивы к отсутствию полей. | | `dedup_key` | нет | Ретраи источника: тот же ключ в TTL-окне → `deduplicated: true, id: <первый>`. Дубли в каналах раздражают сильнее всего. | -| `ttl_seconds` | нет | Событие «контейнер упал» бесполезно через час. Дефолт — «вечно», поле обязательно с первого дня — ретрофит дороже. | +| `ttl_seconds` | нет | Событие «контейнер упал» бесполезно через час: воркер не маршрутизирует истёкшее (статус `expired`), ретеншн удалит (см. «Ретеншн»). Дефолт — «вечно», поле обязательно с первого дня — ретрофит дороже. | | `scheduled_at` | резерв | Не реализуется в MVP; поле зарезервировано. | Не существует и не появится в клиентском контракте: `recipients`, `topics`, `channel` — получатели/каналы/подписки живут только внутри Synapse (правила в админке). `tags` из обсуждений **убраны** из v1; вернутся, когда появится реальное правило, требующее тега (решение обсуждалось 2026-10-03). @@ -48,6 +48,8 @@ Если событие связано с конкретным пользователем, источник кладёт его id gnexus-auth (`sub`) в `payload.user_id` — это часть «что случилось», а не получатель: конверт по-прежнему ничего не знает о получателях. Что Synapse делает с этой привязкой — решают правила: событие попадает в **личный лог пользователя** (`/api/v1/me/events`) только если подошло правило маршрутизации с действием канала `user` (цель не задаётся); правила без такого действия — «событие только для админа». Позже тот же механизм направит push-каналы (#29) конкретному пользователю. +Формат значения: непустая строка (число тоже примем — приведётся к строке). Всё прочее (массив, объект, bool) — при маршрутизации записывается как `skipped` с причиной в `/admin/deliveries`: доставка «в никуда» обязана быть видимой, не тихой. Внутри s2s-тела payload уходит как есть — получатель-сервис получает `payload.user_id` без специальных полей. + ## Ответы | Код | Когда | @@ -116,6 +118,15 @@ Провал доставки (сеть, 5xx провайдера) → обратно-экспоненциальная пауза: `30 c → 2 м → 10 м → 30 м`, всего 5 попыток → `failed`. Скан due-доставок — Celery beat раз в минуту (`synapse.retry_due`). Пока канал не реализован (#29/#28), доставки честно проходят этот цикл и падают в `failed` — таблица сразу показывает поведение ретраев. +### Ретеншн + +Удаление — beat-задачей `synapse.expire_events` (ежечасно): + +- событие с истёкшим `ttl_seconds` воркер не маршрутизирует (статус `expired`), ретеншн удаляет его; доставки падают каскадом; +- события старше `RETENTION_DAYS` (дни, дефолт 30) удаляются независимо от ttl; `0` — окно по возрасту отключено (ttl продолжит работать). + +Следствие: `/api/v1/me/events` и разделы админки читают текущее состояние, а не архив — удалили событие, пропали и записи о нём (включая личный лог канала `user`). + ## Доставка s2s: подпись и верификация Когда правило направляет событие в системный webhook (Navi и другие сервисы), Synapse подписывает доставку. Схема — **ровно та же, что у вебхуков gnexus-auth** (`WebhookSignature.php`), так что принимающий код в экосистеме один и тот же. diff --git a/frontend/src/ui.js b/frontend/src/ui.js index c4cd2ff..557948a 100644 --- a/frontend/src/ui.js +++ b/frontend/src/ui.js @@ -14,6 +14,7 @@ queued: "info", processing: "info", done: "success", + expired: "warning", // истёкшший ttl_seconds — не маршрутизировано, ретеншн удалит }; export function badgeFor(status) {