ZCore LogoZCore
Api reference

StreamManager

API reference for the real-time event streaming engine using Redis PubSub, tunable queue limits, and bounded memory queues.

StreamManager coordinates active listener queues and Redis PubSub operations. It facilitates cluster-wide event routing via subscription channels formatted as stream:user:<user_id>, falling back to local memory queues if Redis is unconfigured or offline.

Class Definition

from zcore.web import StreamManager

class StreamManager:
    def __init__(self) -> None:
        ...

Global Setup Function

init_stream_redis

Initializes the shared Redis connection client for the streaming subsystem.

from typing import Any
from zcore.web import init_stream_redis

def init_stream_redis(client: Any) -> None: ...

Prop

Type


Properties

Prop

Type


Methods

subscribe

Subscribes a user, returning a bounded async listener queue. Queue capacity is dynamically resolved from settings.STREAM_QUEUE_MAXSIZE (Default: 100). Initializes background Redis PubSub listeners if this is the first active subscription for a user on this node.

async def subscribe(self, user_id: Any) -> asyncio.Queue[Any]: ...

Prop

Type

unsubscribe

Unsubscribes a specific listener queue for a user. Shuts down background Redis tasks if no active listeners remain in the local registry.

async def unsubscribe(self, user_id: Any, queue: asyncio.Queue[Any]) -> None: ...

Prop

Type

subscription

Asynchronous context manager safely wrapping active user event streams. Guarantees automatic cleanup and queue unregistration upon block exit or connection drops.

@asynccontextmanager
async def subscription(
    self, 
    user_id: Any
) -> AsyncGenerator[asyncio.Queue[Any], None]: ...

Prop

Type

publish

Publishes an event payload to a target user's stream. Broadcasts cluster-wide across Redis PubSub channel stream:user:{user_id}, falling back to local memory delivery if unconfigured.

async def publish(self, user_id: Any, data: dict[str, Any]) -> None: ...

Prop

Type

start_listening

Explicitly initiates background Redis PubSub pattern-matching subscribers (stream:user:*).

async def start_listening(self) -> None: ...

Usage Example (Server-Sent Events)

from fastapi import APIRouter
from fastapi.responses import StreamingResponse
from zcore.web import StreamManager

router = APIRouter()
stream_manager = StreamManager()

@router.get("/notifications/stream/{user_id}")
async def stream_notifications(user_id: str):
    async def event_generator():
        # Automatically registers and safely tears down queue on client disconnect
        async with stream_manager.subscription(user_id) as queue:
            while True:
                payload = await queue.get()
                yield f"data: {payload}\n\n"

    return StreamingResponse(event_generator(), media_type="text/event-stream")

Slow-Consumer Overflow Protection: If a client consumes messages too slowly and their bounded queue reaches maximum capacity (settings.STREAM_QUEUE_MAXSIZE, default: 100 items), StreamManager catches asyncio.QueueFull, immediately unregisters the stalled queue, and drops it to protect host node memory from buffer bloat.

On this page