Caching & Real-Time Streaming
Deep dive into ZCore's distributed Redis cache with resilient in-memory LRU fallback and the cluster-wide PubSub streaming engine.
ZCore provides unified infrastructure for high-speed distributed data access and real-time event streaming, designed to operate seamlessly with Redis in clustered environments while gracefully degrading to in-memory fallbacks when standalone.
1. The BaseCache Dual-Layer Resilience Architecture
BaseCache acts as a transparent, fault-tolerant caching gateway. When configured with a Redis connection URL (init_cache(redis_url)), all operations are directed to distributed Redis keyspaces.
Graceful Network Degradation:
If a network partition or Redis outage occurs mid-request, ZCore catches the connection error, logs the diagnostic warning, and immediately serves/stores the record in the local, thread-safe in-memory TTLLRUCache. When Redis recovers, operations automatically route back to distributed caching without restarting the application.
Dynamic Framework Boundaries
Both BaseCache and TTLLRUCache dynamically bind to centralized framework configuration settings with sensible fallbacks:
- Default Lifespan (
CACHE_DEFAULT_TTL): Fallback TTL (default: 3600 seconds) applied automatically whenttlis omitted incache.set(). - Local In-Memory Capacity (
CACHE_LOCAL_MAXSIZE): Maximum key volume (default: 1000 items) retained in the local fallback store before LRU eviction triggers. - Garbage Collection Sweep (
CACHE_EVICTION_INTERVAL): Interval in seconds (default: 60 seconds) between non-blocking background sweeps purging expired records.
Automatic Pydantic V2 Deserialization
You can supply target_type=TaskResponse to cache.get(). ZCore automatically deserializes the raw JSON string and executes TaskResponse.model_validate(), returning strongly-typed model instances:
# Fetches from Redis/Local cache and validates schema automatically
task: TaskResponse | None = await cache.get("task_123", target_type=TaskResponse)Background Garbage Collection (weakref.WeakSet)
TTLLRUCache instances are registered in a global weakref.WeakSet. This allows the background asynchronous eviction loop (_start_eviction_loop) to sweep expired keys across all active caches at intervals defined by settings.CACHE_EVICTION_INTERVAL without creating circular memory references or preventing clean de-allocation.
2. Real-Time Streaming Engine (StreamManager)
For real-time features like WebSockets, Server-Sent Events (SSE), and live notifications, ZCore provides the StreamManager.
Cluster-Wide Message Propagation
StreamManager connects to Redis PubSub using pattern subscription (stream:user:*). When any worker node in your cluster publishes an event to a user, all nodes receive the frame and route it to the active local queues for that user.
# Stream Manager Subscription Lifecycle (e.g., in an SSE or WebSocket route)
from fastapi.responses import StreamingResponse
from zcore.web import StreamManager
stream_manager = StreamManager()
@router.get("/tasks/stream/{user_id}")
async def stream_user_events(user_id: str): # Accepts polymorphic user IDs (UUID, int, or str)
async def event_generator():
# Automatically subscribes and handles cleanup upon connection closure
async with stream_manager.subscription(user_id) as queue:
while True:
data = await queue.get()
yield f"data: {data}\n\n"
return StreamingResponse(event_generator(), media_type="text/event-stream")Slow-Consumer Overflow Protection
Each connected client listener is allocated a bounded asyncio.Queue whose capacity is resolved dynamically from settings.STREAM_QUEUE_MAXSIZE (default: 100 items):
- If a publisher pushes events faster than a client can consume them and the queue reaches capacity,
StreamManagercatchesasyncio.QueueFull. - It immediately unregisters and drops the overflowing queue, safeguarding node memory against buffer bloat and memory exhaustion.
Standalone Mode: If Redis is uninitialized, StreamManager routes events directly through local memory queues, allowing real-time features to function in single-node development environments out of the box.
Security & Authentication Architecture
Explore ZCore's cryptographic services (Argon2id, JWT), the BaseAuth template pipeline, fail-fast production assertions, and scope permissions.
Testing Infrastructure (ZTestClient)
Explore how ZCore orchestrates IoC sandboxes, savepoint database rollbacks, context mocking, and application lifespans.