From 45314bb7d92e1d3217ee5044e2adc1338cd1fa0a Mon Sep 17 00:00:00 2001 From: muhsar27 Date: Sun, 30 Aug 2026 08:52:34 +0100 Subject: [PATCH] docs: add architecture and implementation documentation for SSE integration --- backend/docs/SSE_ARCHITECTURE.md | 45 ++++++++++++++++++------------ backend/docs/SSE_IMPLEMENTATION.md | 2 +- 2 files changed, 28 insertions(+), 19 deletions(-) diff --git a/backend/docs/SSE_ARCHITECTURE.md b/backend/docs/SSE_ARCHITECTURE.md index cc45e4f6..3bf547c3 100644 --- a/backend/docs/SSE_ARCHITECTURE.md +++ b/backend/docs/SSE_ARCHITECTURE.md @@ -3,27 +3,27 @@ ## System Flow ``` -┌─────────────────┐ -│ Blockchain │ -│ Indexer │ -│ (Stellar) │ -└────────┬────────┘ - │ Events - ▼ +┌─────────────────────────────────────────────────────────┐ +│ Stellar Blockchain / Soroban │ +│ - On-chain stream contract executions & ledger events │ +└────────────────────────────┬────────────────────────────┘ + │ On-Chain Events (via Soroban RPC poll) + ▼ ┌─────────────────────────────────────────────────────────┐ │ Backend Server │ │ │ │ ┌──────────────────────────────────────────────────┐ │ -│ │ Stream Controller │ │ -│ │ - Creates/updates streams │ │ -│ │ - Calls sseService.broadcast() │ │ +│ │ Soroban Event Worker (Indexer) │ │ +│ │ - Polls Soroban RPC for confirmed events │ │ +│ │ - Persists stream state & events to Database │ │ +│ │ - Calls sseService.broadcastToStream/Admin() │ │ │ └──────────────┬───────────────────────────────────┘ │ -│ │ │ +│ │ Dispatch events (indexer-driven) │ │ ▼ │ │ ┌──────────────────────────────────────────────────┐ │ │ │ SSE Service │ │ │ │ - Manages client connections │ │ -│ │ - Filters by subscription │ │ +│ │ - Filters by subscription (stream/user/admin) │ │ │ │ - Broadcasts to matching clients │ │ │ └──────────────┬───────────────────────────────────┘ │ │ │ │ @@ -40,6 +40,8 @@ └─────────────────────────────────────┘ ``` +> **Note on Indexer-Driven Event Origin**: SSE broadcast events originate asynchronously from the background indexer worker (`SorobanEventWorker`) only after transaction confirmation on the Stellar ledger, not synchronously from HTTP API controllers (`stream.controller.ts`, etc.). When a user submits an action (create, pause, withdraw, top-up, cancel), API controllers do not broadcast SSE events directly; events are dispatched once the Soroban event is polled and confirmed on-chain. Additionally, background workers like `StreamRunwayWorker` may dispatch computed alerts (e.g. `STREAM_LOW_BALANCE`). + ## Connection Flow ``` @@ -72,7 +74,7 @@ Client Server │ [Auto Reconnect - 1s] │ │ │ │ GET /events/subscribe │ - ├──────────────────────────────>│ + │ ├──────────────────────────────>│ │ │ │ 200 OK │ │<──────────────────────────────┤ @@ -122,9 +124,9 @@ Client Server └─────────────────┘ Flow: -1. Backend 1 receives stream creation +1. Backend 1 (SorobanEventWorker) indexes confirmed on-chain event 2. Backend 1 publishes to Redis: "stream-events" -3. All backends (1, 2, 3) receive message +3. All backends (1, 2, 3) receive message via Redis subscriber 4. Each backend broadcasts to its connected clients 5. Total: 33 clients receive the event ``` @@ -132,20 +134,27 @@ Flow: ## Event Broadcasting Logic ```typescript -// Broadcast to specific stream +// Broadcast to specific stream (called by SorobanEventWorker / StreamRunwayWorker) sseService.broadcastToStream("123", "stream.created", data) ↓ Filter clients: subscription includes "123" or "*" ↓ Send to matching clients -// Broadcast to user -sseService.broadcastToUser("GABC...", "stream.created", data) +// Broadcast to user (called by StreamRunwayWorker) +sseService.broadcastToUser("GABC...", "STREAM_LOW_BALANCE", data) ↓ Filter clients: subscription includes "user:GABC..." or "*" ↓ Send to matching clients +// Broadcast to admin (called by SorobanEventWorker for admin/fee events) +sseService.broadcastToAdmin("stream.fee_config_updated", data) + ↓ + Filter clients: admin subscribers ("admin" or "*") + ↓ + Send to matching clients + // Broadcast to all sseService.broadcast("stream.created", data) ↓ diff --git a/backend/docs/SSE_IMPLEMENTATION.md b/backend/docs/SSE_IMPLEMENTATION.md index fe9de92d..1e7174ca 100644 --- a/backend/docs/SSE_IMPLEMENTATION.md +++ b/backend/docs/SSE_IMPLEMENTATION.md @@ -193,7 +193,7 @@ import Redis from 'ioredis'; const redis = new Redis(process.env.REDIS_URL); const subscriber = new Redis(process.env.REDIS_URL); -// Publisher (in stream controller) +// Publisher (in Soroban event worker / indexer) redis.publish('stream-events', JSON.stringify({ event: 'stream.created', data: mockStream,