Systems Hub
// Real-Time Engineering at Scale

HTTP দিয়ে Realtime
হয় না।
WebSocket লাগে — কিন্তু কীভাবে?

৫০,০০০ student একসাথে class দেখছে। কে কোন class এ আছে? কতজন live দেখছে? Teacher chat করলে সবাই পাচ্ছে?
HTTP polling দিয়ে এটা করলে server মরে যাবে। WebSocket + Socket.IO + Redis দিয়ে কীভাবে এটা handle করা হয়েছে — সেটাই আজকের বিষয়।

Socket.IO Redis Adapter Presence System Live Chat Viewer Count Rate Limiting
WebSocket
02
// The Core Problem

কেন HTTP যথেষ্ট না?

Realtime feature এ HTTP polling vs WebSocket এর তুলনা।

❌ HTTP Polling (ভুল approach)
// Client প্রতি 2 সেকেন্ডে জিজ্ঞেস করে: setInterval(async () => { const res = await fetch( '/api/viewer-count?classId=xyz' ); updateUI(await res.json()); }, 2000); // 50,000 students × 30/min = 1.5M req/min // Server dead 💀
কী সমস্যা?
প্রতিটা request এ নতুন connection। Header overhead বিশাল। Server push করতে পারে না — client ই চাইতে হয়। 50K user মানে 1.5M req/min।
👤
GET /count
🖥️
← আবার 2s পরে, আবার 2s পরে... →
✅ WebSocket (এই system)
// একবার connect — সারাদিন open থাকে const socket = io(serverUrl, { transports: ['websocket'] ← no polling fallback }); // Server push করে যখন দরকার: socket.on('watch:online', (data) => { updateViewerCount(data.onlineWatchers); }); // 50K connections: server push করে // শুধু count বদলালে — efficient! ✅
WebSocket এর advantage
একবার handshake → persistent connection। Server যখন খুশি push করতে পারে। Header overhead নেই। 50K connections maintain করা সম্ভব।
👤
🖥️
← Persistent bidirectional connection →
03
// Architecture

Socket System — Full Architecture

Client connect হলে কী হয় — initialization থেকে event handling পর্যন্ত।

// setupSocket(httpServer) const io = new Server(httpServer, { transports: ['websocket'], ← WS only pingInterval: 25000, ← heartbeat pingTimeout: 20000, maxHttpBufferSize: 64 * 1024, ← 64KB max perMessageDeflate: false, ← CPU save }); // Multi-server support via Redis Adapter io.adapter( createShardedAdapter(pubClient, subClient) ); // Auth middleware → all sockets must pass io.use(socketAuth); // Feature modules register करो startMetricsPublisher(io); startBroadcaster(io); ← viewer count startHistogram(io); ← watch progress registerSocket(io); ← event handlers
4টা Feature Module
🎥
Watch Module
Recorded class viewer count — কে কোন video দেখছে
watch/
📡
Class Module
Live class presence — কতজন live এ আছে
class/
💬
Chat Module
Realtime class chat — batched, rate-limited
chat/
📊
Progress Module
Video progress tracking — histogram
progress/
👨‍💼
Admin Module
Admin room — teachers monitor করতে পারে
admins/
Feature Flag
ENABLE_RECORDED_CLASS_WATCH=false হলে chat + progress module disable। Core presence সবসময় on।
04
// Socket Authentication

Connection এর আগে 4 Layer Security

Anonymous কেউ socket connect করতে পারে না। প্রতিটা connection JWT verify করে।

// socketAuth middleware (io.use()) export default async function socketAuth(socket, next) { // 1️⃣ IP Rate Limit const ip = socket.handshake.headers['x-forwarded-for'] ?.split(',')[0]?.trim(); const allowed = await allowConnectionFromIp(ip); if (!allowed) return next(new Error('Too many connections')); // 2️⃣ Cookie থেকে token পড়া const cookieHeader = socket.handshake.headers.cookie; if (!cookieHeader) return next(new Error('Unauthorized')); // 3️⃣ Host দিয়ে correct cookie key বের করা const host = origin?.replace(/^https?:\/\//, ''); const cookieKey = CookieHelper.refreshCookieName(host); const decoded = verifyRefreshTokenWithSignature(refreshToken); if (!decoded?.id) return next(new Error('Invalid token')); // 4️⃣ User data fetch (Redis cache → DB) const user = await getSocketUser(decoded.role, decoded.id); if (!user) return next(new Error('User not found')); socket.user = user; ← সব handler এ available next(); }
Security Layers
01
IP Rate Limit
প্রতি IP থেকে 10s এ max 60 connections। Redis INCR + EXPIRE। DDoS protection।
02
Cookie Auth
HTTP Cookie থেকে refreshToken। Platform এর host দিয়ে correct cookie key বের করা।
03
JWT Verify
HS256 signature verify, token type check (refresh), decoded.id + role বের করা।
04
User Cache
cache:socket:user:v2:{role}:{id} → Redis 2min TTL → Prisma fallback।
socket.user সব জায়গায়
Auth pass হলে socket.user = { id, role, name, avatar } set হয়। সব event handler এ এটা থেকে user জানা যায় — আর JWT verify করতে হয় না।
05
// Presence System

কতজন দেখছে? — Viewer Count Architecture

এক class এ হাজার হাজার viewer। তাদের count করতে হবে exact বা approximate — Redis ZSET + HLL দিয়ে।

Redis Data Structures
watch:presence:{type}:{id} ZSET
member: "student:abc123"
score: timestamp (last heartbeat)
watch:hll:{type}:{id}:{bucket} HyperLogLog
500+ viewer হলে approximate count
Error rate ~0.8% — acceptable
watch:meta:{type}:{id} HASH
field: member, value: JSON user info
Max 500 entries stored (exact list)
// Exact (≤500) vs Approximate (>500) return { count: Number(result?.[0] || 0), isExact: Number(result?.[1]) === 1, }; // isExact=false → "~5,230 watching" // isExact=true → "342 watching"
Join → Heartbeat → Leave Flow
📥
watch:join event
socket.join(room) → markActive() buffer → updateMeta() HASH
💓
watch:heartbeat (30s)
ZSET score update → user still alive signal
📤
watch:leave / disconnect
ZREM from ZSET + HDEL from HASH + markRoomDirty()
📡
Broadcaster (1s tick)
dirty rooms → leader election → count → emit watch:online
Presence Buffer — Batched Write
Heartbeat প্রতিটা Redis এ instantly লিখলে N×1000 write/sec। Buffer এ collect করে, 1s interval এ pipeline flush। Writes ৯০%+ কমে।
06
// Broadcaster Logic

Smart Broadcaster — কখন emit করবে?

প্রতি 1s এ সব room broadcast করলে chaos হবে। Smart conditions দিয়ে শুধু দরকারে emit।

// broadcaster.js — key logic async function processRoom(classType, classId) { // 1. Leader election (Redis NX lock) // Multi-server এ শুধু একজনই emit করবে const leader = await acquireLeader(classType, classId); if (!leader) return; // অন্য server handle করছে // 2. Count fetch const { count, isExact } = await getOnlineCount({...}); // 3. shouldEmit? — Smart condition const emit = await shouldEmit(classType, classId, count); if (!emit) return; // count same হলে skip // 4. Broadcast to room ioRef.to(watchRoom(classType, classId)) .emit('watch:online', { onlineWatchers: count, approximate: !isExact, updatedAt: Date.now(), }); } // shouldEmit logic: const minDelta = Math.max( CHANGE_MIN_DELTA, // min 1 Math.floor(count * 0.02) // or 2% of count ); // 10,000 viewer room: 200 বদলালে emit // 10 viewer room: 1 বদলালেই emit
4 Smart Conditions
🎯
Dirty Room Only
markRoomDirty() call হলেই room process হয়। কেউ join/leave না করলে — broadcast নেই।
🏆
Leader Election
Redis NX lock — 3s TTL। Multi-server এ শুধু একটা server emit করে। Duplicate broadcast হয় না।
📊
Change Threshold
Count এর 2% বা minimum 1 — যেটা বেশি। 10K room: 200 viewer change না হলে emit নেই।
Force Refresh
30s একবার force emit — count same থাকলেও। Client স্টেল data দেখবে না।
07
// Chat System

Realtime Chat — Token Bucket + Batcher

Class চলাকালীন student chat করে। Rate limit না থাকলে spam, delay থাকলে laggy — এই balance কীভাবে করা হয়েছে।

Token Bucket Rate Limiter
// rateLimit.js — per socket, in-memory export function consumeChatToken(socket) { const now = Date.now(); const bucket = socket.data.chatBucket || { tokens: 5, at: now }; ← BURST=5 const elapsed = (now - bucket.at) / 1000; const tokens = Math.min( 5, ← max 5 bucket.tokens + elapsed * 1 ← 1/sec refill ); if (tokens < 1) { socket.emit('watch:chat:throttled'); return false; ← rate limited! } socket.data.chatBucket = { tokens: tokens - 1, at: now }; return true; }
Token Bucket Visual (5 burst, 1/sec refill)
5 messages burst → empty → 1/sec refill
Chat Message Batcher
// batcher.js — 120ms flush interval function flushRoom(room, list) { const batch = list.splice(0, 50); ← max 50 if (batch.length === 1) { // Single message: direct emit ioRef.to(room).emit('watch:chat:message', batch[0]); } else { // Multiple: batch emit (1 network call!) ioRef.to(room).emit('watch:chat:batch', batch); } }
কেন Batch করা?
৫০০ student room এ একসাথে ২০ message আসলে — ২০টা emit না করে একটা array emit। Network call ১, payload compression হয়।
Typing Throttle
typing event min 3s interval। Room size > 300 হলে typing event disabled — server overload থেকে বাঁচায়।
Room Buffer Overflow
Room buffer max 500 messages। বেশি হলে chat_dropped_overflow metric। Data drop হয় — spam protect।
08
// Horizontal Scaling

Multiple Server — Redis Adapter

একটা server এ ৫০K connection সম্ভব না। Multiple server দরকার — Redis Adapter দিয়ে সব server sync।

Without Adapter (ভুল)
Server 1
User A, B, C
Server 2
User D, E, F
সমস্যা
Server 1 এ Admin broadcast করলে শুধু User A,B,C পায়। Server 2 এর D,E,F কিছু পায় না।
✅ With Redis Sharded Adapter
Server 1
Redis Pub/Sub
Server 2
কীভাবে কাজ করে?
Server 1 emit করলে Redis Pub/Sub এ publish হয়। Server 2 subscribe করা → message পায় → নিজের connected users কে forward করে।
Scale Config (scale.js)
PRESENCE config
BUCKET_MS: 30,000ms (heartbeat window)
WINDOW_BUCKETS: 3 (90s active window)
EXACT_TRACK_LIMIT: 500 (ZSET vs HLL)
FLUSH_INTERVAL: 1000ms (buffer flush)
BROADCAST config
TICK_MS: 1000ms (check interval)
MIN_INTERVAL: 3000ms (leader lock TTL)
MAX_ROOMS_PER_TICK: 400 (throughput cap)
FORCE_REFRESH: 30s (stale guard)
CHAT config
FLUSH_MS: 120ms (batch interval)
BURST: 5, REFILL: 1/sec (token bucket)
MAX_MESSAGE: 1000 chars
TYPING_MIN_INTERVAL: 3000ms
09
// Production Resilience

Graceful Shutdown & Metrics

Server বন্ধ হলে data হারাবে না। আর সব কিছু monitor করতে metrics।

// Graceful Shutdown (SIGTERM/SIGINT) const shutdown = async () => { // 1. Stop broadcaster timer stopBroadcaster(); // 2. Flush pending chat messages stopChatBatcher(); // 3. Flush presence buffer → Redis await stopPresenceBuffer(); ← data saved! // 4. Flush video progress → DB await stopProgressBuffer(); ← no loss! // 5. Stop histogram await stopHistogram(); // 6. Close socket server io.close(); }; process.once('SIGTERM', shutdown); ← k8s/pm2 process.once('SIGINT', shutdown); ← Ctrl+C
কেন এটা Critical?
Deploy এর সময় server restart হয়। Buffer এ থাকা presence data, chat message, video progress flush না করলে হারিয়ে যায়। Graceful shutdown সব save করে তারপর বন্ধ হয়।
Real-time Metrics
📈
connections / disconnections
Total socket connect/disconnect count
📊
watch_join / heartbeat / user_list
Recorded class viewer engagement metrics
🚀
broadcast_emitted / skipped
কতবার emit হলো, কতবার skip হলো — efficiency track
💬
chat_emitted / dropped_overflow
Chat throughput, overflow drop rate
🗂️
presence_flush_calls / rooms
Buffer flush efficiency — batch size tracking
🎛️
Socket Admin UI
@socket.io/admin-ui — visual monitoring dashboard
10
// Summary

এই Socket System কেন Production-Ready?

প্রতিটা design decision এর কারণ আছে — performance, security, resilience।

🔐
JWT + IP Rate Limit
Connection এর আগেই auth। IP থেকে 10s এ 60+ connection block। Anonymous socket সম্ভব না।
📊
ZSET + HLL Dual Mode
≤500 viewer: exact count। >500: HyperLogLog approximate। Memory efficient, accurate enough।
🎯
Smart Broadcast
Dirty room only + leader election + change threshold + force refresh। Unnecessary emit শূন্য।
🏎️
Presence Buffer
In-memory buffer, 1s pipeline flush। N×1000 Redis write → batch। Throughput ৯০%+ বাড়ে।
🔥
Token Bucket Chat
5 burst, 1/sec refill। Spam impossible। 120ms batch emit। 50 msg একটা network call।
🌐
Redis Sharded Adapter
Multiple server sync। Pub/Sub দিয়ে cross-server broadcast। Infinite horizontal scale।
মূল কথা
Socket system মানে শুধু io.on('connection') না — মানে auth layer, presence buffering, smart broadcasting, rate limiting, graceful shutdown — সব একসাথে। এই architecture দিয়ে লাখ user handle করা সম্ভব।