API and SaaS Layer

WebSocket Real-Time Events

Technical details of the ASGI WebSocket server — connection lifecycle, JWT authentication, channel authorization via regex patterns, Redis Pub/Sub stream forwarding, and graceful cleanup.

⏱️ 7 min read📊 Level: Intermediate

Real-time feedback is critical for algorithmic trading. DepthSight uses a standalone WebSocket server (api/websocket_server.py, ~365 lines) to stream active positions, logs, and transaction status to users with sub-second latency.


ASGI WebSocket Architecture

The WebSocket service runs as a decoupled service next to the main FastAPI API.

Rendering diagram...

Key Design Decisions:

  • Decoupled Load: By separating WebSockets, high-frequency logging streams do not block REST API requests.
  • Redis Pub/Sub Backend: When a client requests a channel subscription, the server launches an asynchronous task that listens to Redis and forwards events directly to the client's WebSocket.
  • Per-Channel Dedicated Tasks: Each subscribed channel gets its own Redis Pub/Sub subscription and asyncio.Task, providing isolation — if one channel's Redis connection fails, other channels remain unaffected.

Connection Lifecycle

1. WebSocket Handshake (/ws endpoint, line 220)

The connection URL must contain a valid JWT query parameter:

wss://api.depthsight.pro/ws?token=<JWT_TOKEN>

2. Token Authentication (lines 227–258)

On connection:

Sources:

If the token is invalid or the user does not exist, the connection is closed with 4003 Forbidden (WS_1008_POLICY_VIOLATION).

3. Subscription Message Loop (lines 261–333)

After authentication, the server enters a while True loop receiving JSON messages:

Subscribe action:

{"action": "subscribe", "channel": "user_logs:42"}
  1. Checks channel authorization via _is_channel_allowed(channel, user_id).
  2. Checks per-client channel limit (MAX_CHANNELS_PER_CLIENT = 50).
  3. Creates a new Redis client and spawns an asyncio.Task for redis_channel_listener.
  4. Stores (task, redis_client) in active_listeners[websocket][channel].

Unsubscribe action:

{"action": "unsubscribe", "channel": "user_logs:42"}
  1. Cancels the listener task.
  2. Closes the Redis client.
  3. Removes entry from active_listeners[websocket].

Channel Authorization (_is_channel_allowed)

The authorization function (lines 94–145) uses regex pattern matching to enforce access control:

Protected Channel Patterns

These channels require that the user_id embedded in the channel name matches the authenticated user:

Sources:

Legacy Unscoped Channels (Denied)

Older channel names without user_id scope are explicitly blocked to prevent cross-user data leaks:

Sources:

Logic Summary

  1. If the channel is a legacy unscoped channel → deny (returns False).
  2. If the channel matches a protected pattern → extract user_id from channel; allow only if it matches the authenticated user.
  3. All other channels → allow (public/global channels like "global_market_updates").

Redis Pub/Sub Stream Forwarding

The redis_channel_listener coroutine (lines 163–217) is the core of the real-time data pipeline:

Sources:

Key behaviors:

  • JSON Serialization: Messages from Redis are deserialized from JSON and wrapped in {"topic": channel, "payload": ...}.
  • Auto-Reconnect: On Redis connection errors, the inner loop breaks, and the outer loop retries after 2 seconds — the WebSocket connection to the user remains open.
  • Polling Interval: 10ms (await asyncio.sleep(0.01)) — balances CPU usage vs latency.

Cleanup on Disconnect

The finally block (lines 342–361) ensures complete resource cleanup:

Sources:

This handles:

  • Client disconnect (WebSocketDisconnect).
  • Server-side errors.
  • Abnormal closures.

The active_listeners dict (line 49) is the global state tracker:

Sources:

Keyed by WebSocket object, valued by a dict containing "_user_id" and channel_name → (asyncio.Task, aredis.Redis) pairs.


Security Architecture Summary

LayerMechanismProtection Against
TransportWSS (WebSocket Secure)Eavesdropping, MITM
AuthenticationJWT token in query parameterUnauthorized connections
AuthorizationRegex-based _is_channel_allowed()Cross-user data access
Channel LimitMAX_CHANNELS_PER_CLIENT = 50Resource exhaustion
Connection Trackingactive_listeners dictOrphaned resource leaks