"""Отправки: s2s webhooks (подпись app/signature.py, docs/05) и
web-push подпискам пользователя (канал push, docs/06).
Секреты: s2s цели — per-target из env (S2S_SECRET_<token_ref>); VAPID — из
app_settings (Настройки → Web Push, write-only). Секреты никогда не
возвращаются через API — воркер читает их значения напрямую.
"""
import json
import os
import time
import httpx
from pywebpush import WebPushException, webpush as pywebpush
from sqlalchemy import select
from app.models import ChannelTarget, Delivery, Event, PushSubscription
from app.settings_store import get_setting_row
from app.signature import make_signature
#: Референс-получателю достаточно этих трёх; остальное — настройка цели.
S2S_TIMEOUT = httpx.Timeout(10.0)
def secret_env_name(token_ref: str) -> str:
"""token_ref «navi-rei» → переменная окружения S2S_SECRET_NAVI_REI."""
cleaned = token_ref.upper().replace("-", "_").replace(".", "_")
return f"S2S_SECRET_{cleaned}"
def resolve_secret(token_ref: str) -> str | None:
return os.environ.get(secret_env_name(token_ref)) or None
def send_s2s(target: ChannelTarget, event: Event, source_name: str) -> tuple[bool, str | None]:
"""POST конверта события в endpoint цели, подписанный по схеме gnexus-auth.
Тело — конверт без изменений + event_id (см. docs/05):
описание источника в него не входит. Успех — любой 2xx; прочее —
провал попытки (воркер поставит ретрай).
"""
cfg = target.config or {}
endpoint = cfg.get("endpoint")
token_ref = cfg.get("token_ref")
if not endpoint:
return False, "s2s: в config цели нет endpoint"
secret = resolve_secret(token_ref) if token_ref else None
if not secret:
return False, (
f"s2s: секрет цели не задан в окружении "
f"({secret_env_name(token_ref) if token_ref else 'нет token_ref в config'})"
)
envelope = {
"event_id": str(event.id),
"source": source_name,
"subject": event.subject,
"action": event.action,
"priority": event.priority,
"payload": dict(event.payload or {}),
}
body = json.dumps(envelope, ensure_ascii=False, separators=(",", ":")).encode()
headers = {
"Content-Type": "application/json",
"X-Gnexus-Event-Id": str(event.id),
"X-Gnexus-Event-Type": f"{source_name}.{event.subject}.{event.action}",
"X-Gnexus-Event-Timestamp": str(int(time.time())),
"X-Gnexus-Signature": make_signature(body, secret),
"X-Synapse-Source": source_name,
"User-Agent": "Synapse/0.1 (gnexus notification hub)",
}
try:
response = httpx.post(
endpoint,
content=body,
headers=headers,
timeout=S2S_TIMEOUT,
verify=cfg.get("verify_tls", True),
follow_redirects=False,
)
except httpx.HTTPError as err:
return False, f"s2s: сеть: {err}"
if 200 <= response.status_code < 300:
return True, None
return False, f"s2s: HTTP {response.status_code} · {response.text[:200]}"
def send_push(
db, delivery: Delivery, event: Event, source_name: str
) -> tuple[bool, str | None]:
"""Web-push подпискам получателя (delivery.recipient_user_id, docs/06).
VAPID — из настроек (app_settings, «Настройки → Web Push»); секретный ключ
write-only, воркеру достаточно факта его наличия. Подписки со статусами
404/410 (браузер отписался/истёк) удаляются; успех — все живые подписки
отправлены. Подписок нет — ошибка «нет активных подписок»: ретраи не
спасут, но аудит объясняет, куда событие не ушло.
"""
public = get_setting_row(db, "vapid_public_key")
private = get_setting_row(db, "vapid_private_key")
subject = get_setting_row(db, "push_subject") or "mailto:synapse@gnexus.space"
if not public or not private:
return False, "push: VAPID-ключи не заданы — Настройки → Web Push"
subscriptions = db.execute(
select(PushSubscription).where(
PushSubscription.user_id == delivery.recipient_user_id
)
).scalars().all()
if not subscriptions:
return False, "push: у пользователя нет активных подписок"
payload = json.dumps(
{
"title": f"Synapse · {source_name}.{event.subject}.{event.action}",
"body": delivery.rendered_message or f"{source_name}.{event.subject}.{event.action}",
"event_id": str(event.id),
"url": "/my-events",
},
ensure_ascii=False,
)
gone, failed = [], []
for sub in subscriptions:
try:
pywebpush(
subscription_info={
"endpoint": sub.endpoint,
"keys": {"p256dh": sub.p256dh, "auth": sub.auth},
},
data=payload,
vapid_private_key=private,
vapid_claims={"sub": subject},
)
except WebPushException as err:
code = getattr(getattr(err, "response", None), "status_code", None)
if code in (404, 410):
gone.append(sub) # браузер отписался/истёк — чистим
continue
failed.append(f"HTTP {code or '?'}")
except Exception as err: # сеть и прочее — попробуем в ретрае
failed.append(str(err)[:120])
for sub in gone:
db.delete(sub)
if failed:
return False, "; ".join(failed)
return True, None