abbe7a18fc
- P0: hooks.py 3-tuple fix, trigger_dispatcher Contract, contacts/plugin unregister_actions_by_owner - P0: 5 test files — check_permission mocks removed, hardcoded DB credential → env var - P1: attachment_service DmsFile via Contract helper, restore_registry/history_hooks dedup - P1: mail/plugin restore unregister, mcp_client datetime.now(UTC), saved_views/filters patterns - P1: ProtectedRoute fail-closed, 13 test assertion fixes (bcrypt, DB-URLs, SECRET_KEYs) - P2: deprecated notifications → post_system_message (3 files), forgejo Base, report_generator lazy import - P2: webhooks permissions, deps.py/roles.py plugin perms removed, import_export default - P2: address/tags/entity_links patterns removed, worker.py Contract-Umgehungen fixed - P2: 28 frontend TODOs (hardcoded constants, deprecated notification API) - P3: dead code, duplicates, deprecated imports, private attr, __import__ inline - P3: 8 frontend TODOs (LucideIcons, inline styles, XSS, i18n) - ruff: 838 → 0 (612 auto-fix + 246 manual + 27 F821 regression fix) - F821: 30 → 0 (AutomationDefinition, DmsFile, user_id, Path, Any, String) - Contract-Umgehungen: 2 neue gefunden (worker.py:169, worker.py:280) und gefixt
66 lines
1.9 KiB
Python
66 lines
1.9 KiB
Python
"""Redis Pub/Sub helpers for WebSocket multi-worker fanout."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
import uuid
|
|
from collections.abc import Awaitable, Callable
|
|
|
|
from app.core.auth import get_redis
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
async def publish_to_channel(channel: str, message: dict) -> None:
|
|
"""Publish a JSON message to a Redis channel."""
|
|
redis = get_redis()
|
|
await redis.publish(channel, json.dumps(message, default=str))
|
|
|
|
|
|
async def subscribe_to_channel(
|
|
channel: str,
|
|
handler: Callable[[dict], Awaitable[None]],
|
|
) -> asyncio.Task:
|
|
"""Subscribe to a Redis channel and call *handler* for every message.
|
|
|
|
Returns the :class:`asyncio.Task` so the caller can cancel it on disconnect.
|
|
"""
|
|
|
|
async def _subscriber() -> None:
|
|
redis = get_redis()
|
|
pubsub = redis.pubsub()
|
|
await pubsub.subscribe(channel)
|
|
try:
|
|
async for raw in pubsub.listen():
|
|
if raw["type"] == "message":
|
|
try:
|
|
msg = json.loads(raw["data"])
|
|
await handler(msg)
|
|
except Exception:
|
|
logger.exception("Error in pubsub handler for channel %s", channel)
|
|
finally:
|
|
try:
|
|
await pubsub.unsubscribe(channel)
|
|
await pubsub.aclose()
|
|
except Exception:
|
|
logger.debug("PubSub cleanup error for channel %s", channel)
|
|
|
|
return asyncio.create_task(_subscriber())
|
|
|
|
|
|
def get_tenant_channel(tenant_id: uuid.UUID, topic: str) -> str:
|
|
"""Return the Redis channel name for a tenant + topic."""
|
|
return f"ws:{tenant_id}:{topic}"
|
|
|
|
|
|
async def broadcast_to_tenants(
|
|
tenant_id: uuid.UUID,
|
|
topic: str,
|
|
message: dict,
|
|
) -> None:
|
|
"""Publish a message to a tenant-specific Redis channel."""
|
|
channel = get_tenant_channel(tenant_id, topic)
|
|
await publish_to_channel(channel, message)
|