Systems Hub
// Active Chat System — Full Deep Dive

Class চলছে, সবাই
একসাথে chat করছে।
Backend কী করছে?

৫,০০০ student একটা class এ বসে chat করছে। প্রতি সেকেন্ডে ২০০+ message।
নাইভ approach: সরাসরি DB তে save → DB overwhelmed, UI laggy।

কীভাবে করা হয়েছে? Socket → Token Bucket → Batcher → Redis Stream → Worker → DB — প্রতিটা layer এর কারণ কী?

Socket.IO Rate Limit Ban System Batcher Redis Stream Worker Process Dead Letter
💬
active chat
02
// The Wrong Approach

নাইভ Approach — কেন এটা Fail করে?

বেশিরভাগ developer প্রথমে এটা করে — এবং production এ এসে পস্তায়।

❌ সবাই যা করে (Naive Approach)
// socket handler এ সরাসরি DB call socket.on('chat:send', async (payload) => { // কোনো rate limit নেই // কোনো validation নেই // spam করলেও কিছু হয় না // সরাসরি DB write! await db.query( `INSERT INTO chats VALUES (?)`, [payload.message] ); // সবাইকে emit io.emit('chat:new', payload); });
❌ 5000 user × 10 msg/sec = 50,000 DB write/sec
PostgreSQL সাধারণত 5,000–10,000 write/sec handle করতে পারে। DB dead।
❌ Spam কোনো check নেই
এক user প্রতি মিলিসেকেন্ডে message পাঠাতে পারে। DB flood।
❌ DB fail হলে message হারায়
কোনো retry নেই, কোনো buffer নেই। Error → message gone।
✅ এই System এ Solution
🛡️
Token Bucket Rate Limit
Socket memory এ — Redis call ছাড়াই। 5 burst, 1/sec refill।
🚫
Ban Check (Cached)
15s memo cache — banned student immediately blocked। Redis → Prisma fallback।
📦
In-Memory Batcher
120ms buffer → 50 messages একসাথে emit। Network calls ৯৫%+ কমে।
🌊
Redis Stream Queue
Async DB write। Socket thread block হয় না। Guaranteed delivery।
⚙️
Separate Worker Process
Stream থেকে batch read → Prisma createMany। DB কে controlled pace এ।
03
// The Full Pipeline

Message যাত্রা — ৭ টা Stage

Student "enter" চাপা থেকে DB তে save হওয়া পর্যন্ত — প্রতিটা step কী করে।

👤 Student
sends message
🛡️ Token Bucket
rate check
🚫 Ban Check
restriction memo
✅ Validate
room + payload
3টা জিনিস
একসাথে হয়
📦 Batcher Queue
in-memory buffer
📡 Emit to Room
students see it
👨‍💼 Admin Room
superAdmin sees all
🌊 Redis Stream
async DB queue
⚙️ Worker Process
xreadgroup batch
🐘 PostgreSQL
createMany batch
☠️ Dead Letter
failed msg store
Socket Thread Never Blocks
Redis Stream write async — socket handler millisecond এ return করে। User কোনো lag দেখে না।
3 Parallel Actions
Message accept হলে: batcher queue, admin emit, Redis Stream — তিনটা একসাথে। কোনোটা block করে না।
04
// Validation & Room Check

Message Accept হওয়ার আগে ৫ টা Gate

Invalid data, spam, banned user — কেউ DB পর্যন্ত পৌঁছাতে পারে না।

// chat.js — send() function export async function send(io, socket, payload) { // Gate 1: Payload structure check if (!validateWatchChatPayload(payload)) return; // classType ∈ {CLASS_CONTENT, CYCLE_CONTENT} // classId: string, 8–80 chars // message: string, 1–1000 chars // Gate 2: Rate limit (token bucket) if (!consumeChatToken(socket)) { socket.emit('watch:chat:throttled', { message: "একটু ধীরে পাঠাও" }); return; } // Gate 3: Room membership check const room = watchRoom(payload.classType, payload.classId); if (!socket.rooms.has(room)) return; // user actually joined এই class এ? না হলে block // Gate 4: Ban check (cached) if (await ensureChatAllowed(socket)) return; // student banned? → emit blocked message → return // Gate 5: Message sanitize const msg = payload.message.trim().slice(0, 1000); // trim whitespace + hard cut at 1000 chars // ✅ All gates passed → process message }
৫ Gate Details
📋
Gate 1 — Payload Validate
classType allowlist check, classId length 8–80, message non-empty string
~0ms
🪣
Gate 2 — Token Bucket
socket.data.chatBucket — pure JS, no Redis. 5 burst, 1/sec refill
~0ms
🏠
Gate 3 — Room Membership
socket.rooms.has(room) — Socket.IO internal set check. Cross-room injection impossible
~0ms
🚫
Gate 4 — Ban Check
Admin only. 15s memo → Redis cache → Prisma. Banned message বাংলায় পায়
~1ms
✂️
Gate 5 — Sanitize
trim() + slice(0, 1000). XSS নেই, DB overflow নেই
~0ms
Total overhead
সব gate মিলিয়ে ~1-2ms। Socket thread এ কোনো DB call নেই — ultra fast।
05
// Token Bucket Rate Limiter

Token Bucket — কীভাবে Spam Block হয়?

Redis call ছাড়াই, socket এর নিজের memory তে — ultra fast rate limiter।

// rateLimit.js — Token Bucket Algorithm export function consumeChatToken(socket) { const now = Date.now(); // Initial bucket: 5 tokens const bucket = socket.data.chatBucket || { tokens: 5, at: now }; // Time passed → tokens refill const elapsed = (now - bucket.at) / 1000; ← seconds const tokens = Math.min( 5, ← max cap bucket.tokens + elapsed * 1 ← 1 token/sec ); if (tokens < 1) { // No token available → rate limited socket.data.chatBucket = { tokens, at: now }; return false; ← blocked! } // Consume 1 token socket.data.chatBucket = { tokens: tokens - 1, at: now }; return true; ← allowed! }
Live Simulation
Token Bucket State (5 max, 1/sec refill)
5/5 tokens — এখন পাঠাতে পারবে
কেন Redis এ করা হয়নি?
Redis call = ~1ms network latency প্রতিটা message এ। 5000 user × 10msg/s = 50,000 Redis call/s। Socket memory এ = 0ms।
📊
Burst কেন 5?
Real user একসাথে কয়েকটা message দিতে পারে। 5 burst দিলে smooth feel। 1/sec refill মানে sustained 60/min।
🚨
Throttled হলে কী?
watch:chat:throttled event emit + বাংলায় message। Counter track হয় metrics এ।
🔄
Typing Throttle আলাদা
Typing event: min 3s interval। Room > 300 users হলে typing event বন্ধ — server save।
06
// Ban & Restriction System

Student Ban — 3 Layer Cache System

Admin যখন student কে ban করে — পরবর্তী message এ instantly blocked। কিন্তু DB call কম রেখে।

// studentBannedCheck.js async function resolveRestriction(socket) { const now = Date.now(); // Layer 1: Socket memory memo (15s) const memo = socket.data.restrictionMemo; if (memo && now - memo.at < 15_000) { return memo.value; ← instant, no I/O } // Layer 2: Redis cache (30s TTL) const restriction = await getCachedRestriction( socket.user.id, RestrictionType.MEDIA_COMMENT ); // getOrLoadStrictCache → Redis hit // or Prisma.studentRestriction.findFirst() // Layer 3: Build view + cache in memo const value = buildRestrictionView(restriction); socket.data.restrictionMemo = { value, at: now }; return value; } // banned হলে client এ পাঠানো হয়: socket.emit('watch:chat:blocked', { restricted: true, message: "তোমার মেসেজ দেওয়ার অনুমতি নেই।", "অনুগ্রহ করে ৫ মিনিট পরে চেষ্টা করো।", bannedUntil: restriction.bannedUntil, remainingMs: 300000, ← countdown দেখাতে permanent: false });
3 Cache Layer Hierarchy
Layer 1: Socket Memo (15s)
socket.data.restrictionMemo — RAM এ। 0ms। 15s এর মধ্যে আবার message দিলে এখানেই জবাব।
0ms
🔴
Layer 2: Redis Cache (30s TTL)
getOrLoadStrictCache — jitter 10%, lock TTL 5s। Thundering herd নেই।
~1ms
🐘
Layer 3: Prisma DB (cache miss)
studentRestriction.findFirst — active ban (bannedUntil > now OR null). FULL বা MEDIA_COMMENT type।
~10ms
Restriction Types
FULL — সব কিছু থেকে ban
MEDIA_COMMENT — শুধু chat/comment ban
bannedUntil: null — permanent ban
bannedUntil: Date — temporary, countdown দেখায়
Admin only restriction check
socket.user.role !== 'student' হলে check skip — teacher/admin কে restrict করা হয় না।
07
// Message Batcher

120ms Batcher — একটা Network Call এ ৫০ Message

প্রতিটা message আলাদা emit করলে network overhead বিশাল। Batcher দিয়ে ৯৫%+ calls কমানো হয়েছে।

// batcher.js — in-memory buffer const buffers = new Map(); ← room → msg[] let timer = null; export function queueChatMessage({ classType, classId, message }) { const room = watchRoom(classType, classId); let list = buffers.get(room) || []; buffers.set(room, list); // Max 500 messages in buffer if (list.length >= 500) { metrics.inc('chat_dropped_overflow'); return; ← spam protect } list.push(message); ensureTimer(); ← 120ms interval start } function flushRoom(room, list) { const batch = list.splice(0, 50); if (batch.length === 1) { // Single: direct emit (no overhead) ioRef.to(room).emit('watch:chat:message', batch[0]); } else { // Multiple: batch emit (1 network call!) ioRef.to(room).emit('watch:chat:batch', batch); } }
Without vs With Batcher
❌ Without Batcher (20 messages in 120ms)
emit(msg1) → 5000 sockets
emit(msg2) → 5000 sockets
... × 20 times
= 20 network calls × 5000 = 100,000 sends
✅ With Batcher (same 20 messages)
buffer.push(msg1...msg20)
120ms tick → flushRoom()
emit([msg1,...,msg20]) → 5000 sockets
= 1 network call × 5000 = 5,000 sends
Single vs Batch Event
1 message → watch:chat:message (single object)
2+ messages → watch:chat:batch (array)
Buffer Overflow Protection
Room buffer max 500. Overflow হলে drop + metric count। Spam attack এ server রক্ষা পায়।
08
// Redis Stream — Async DB Write

Redis Stream — কেন DB Directly লিখি না?

Message accept হওয়ার পরে DB write async। Redis Stream দিয়ে guaranteed, ordered, at-least-once delivery।

// chat.queue.js — enqueue to Redis Stream export async function enqueueChat(message) { return redisConnection.xadd( 'stream:watch-chat', ← stream name 'MAXLEN', '~', 100000, ← max 100K entries '*', ← auto-generate ID 'data', JSON.stringify(message) ); } // '~' = approximate trim (efficient) // Old entries auto-evict
// Stream entry looks like: // ID: 1700000000000-0 // data: {"id":"uuid","message":"...", // "classType":"CLASS_CONTENT", // "classId":"uuid", // "sender":{id,name,avatar,role}} // Worker reads with Consumer Group: await redis.xreadgroup( 'GROUP', 'chat-group', consumer, 'COUNT', 200, ← 200 entries/batch 'BLOCK', 5000, ← wait 5s for data 'STREAMS', stream, '>' ); // '>' = undelivered to this group only
Stream Entries Visual
1700001000001-0 {"message":"ভালো class!"} ✓ ACK
1700001000002-0 {"message":"স্যার বুঝলাম না"} ✓ ACK
1700001000003-0 {"message":"আরেকবার বলুন"} ⏳ pending
1700001000004-0 {"message":"thank you sir"} ⏳ pending
কেন Redis Stream? (Queue vs Stream)
Queue (BullMQ/RabbitMQ): job এক consumer নেয়।
Redis Stream: Consumer Group — multiple worker, ordered, replay possible।
MAXLEN ~ 100,000
Stream size বাধা — পুরনো entries auto trim। Memory controlled। '~' মানে approximate — performance trade-off।
Non-blocking enqueue
enqueueChat().catch(err => console.error) — fire and forget। Queue fail হলে socket user affected হয় না।
09
// Chat Worker Process

Dedicated Worker — Batch Save + Dead Letter

আলাদা Node.js process — stream থেকে batch read → Prisma batch write। Fail হলে dead letter, না হারায়।

// chat.worker.js — main loop async function handleBatch(entries) { const { parsedMessages, streamIds, badIds } = parseEntries(entries); // Bad JSON → dead letter for (const id of badIds) { await deadLetter(id, null, 'invalid json'); } try { // Try bulk save first (fastest) await saveChats(parsedMessages); ← createMany await ack(streamIds); ← xack all return; } catch (error) { // Batch failed → try one by one for (let i = 0; i < parsedMessages.length; i++) { try { await saveChatOne(parsedMessages[i]); await ack([streamIds[i]]); } catch (err) { // Single also failed → dead letter await deadLetter(streamIds[i], parsedMessages[i], err.message); } } } } // Every 20 batches: reclaim stuck msgs await redis.xautoclaim(stream, group, consumer, 30000, ← 30s idle → reclaim '0', 'COUNT', 200);
Worker Flow Decision Tree
📥
xreadgroup (200 entries, 5s block)
Consumer Group এ register। অন্য worker এর message নেবে না।
🔍
parseEntries → valid/bad split
JSON parse। Bad entries → dead letter immediately। Good entries → save।
createMany batch (fastest path)
Prisma createMany + skipDuplicates। 200 rows একটা INSERT। Success → xack all।
🔄
Batch fail → one-by-one retry
কোন row এ problem? One-by-one দিয়ে বের করো। Good ones save, bad ones dead letter।
☠️
Dead Letter Stream
stream:watch-chat:dead — max 10K. reason + original data। Debug করা যায়।
♻️
xautoclaim — 30s stuck → reclaim
Worker crash হলে 30s পরে অন্য worker এর stuck message নেয়। No message lost।
10
// DB Schema + Auto Cleanup

activeChats Table + Automated Cleanup

Messages DB তে কীভাবে store হয়? আর পুরনো data কীভাবে automatically delete হয়?

// buildActiveChatData.js — DB row builder export function buildActiveChatData(data) { return { id: data.id, ← crypto.randomUUID() message: data.message, // Timestamp as BigInt (microsecond precise) messageCreatedAt: BigInt(data.createdAt), // Content type router classContentId: data.classType === 'CLASS_CONTENT' ? data.classId : null, cycleContentId: data.classType === 'CYCLE_CONTENT' ? data.classId : null, // Sender type router (mutually exclusive) studentId: data.sender.role === 'student' ? data.sender.id : null, adminId: data.sender.role === 'admin' ? data.sender.id : null, superAdminId: data.sender.role === 'superAdmin' ? data.sender.id : null, }; } // service.js: await prisma.activeChat.createMany({ data, skipDuplicates: true, ← idempotent re-delivery });
// activeChat.services.js — Cron Cleanup // Runs daily via cronjob const RETENTION_DAYS = 5; ← 5 দিনের বেশি রাখে না const BATCH_SIZE = 5_000; ← একবারে 5K delete const BATCH_DELAY_MS = 200; ← বিরতি দিয়ে delete const MAX_RUNTIME_MS = 20 * 60 * 1000; ← 20min budget // Redis lock — single instance only const lock = await redis.set( 'lock:cron:activeChat:cleanup', String(process.pid), 'EX', 3600, 'NX' ); if (lock !== 'OK') return; ← অন্য instance skip // Batch delete loop with raw SQL while (true) { const deleted = await prisma.$executeRaw` DELETE FROM "activeChats" WHERE id IN ( SELECT id FROM "activeChats" WHERE "createdAt" < ${cutoff} ORDER BY "createdAt" ASC LIMIT 5000 )`; if (!deleted) break; ← nothing left await sleep(200); ← DB breathe }
5
days retention
5K
rows/batch delete
200ms
batch delay
11
// Admin Monitoring

SuperAdmin সব Class এর Chat দেখে

Teacher class এ কী বলছে, student কী লিখছে — SuperAdmin একটা dashboard থেকে সব monitor করে।

// admins/register.js — connection এ export default function registerWatchAdmin(io, socket) { // Only superAdmin gets admin room if (socket.user?.role !== 'superAdmin') { return; ← student/admin skip করে } socket.join(adminRoom()); ← 'watch:admin:room' } // admins/emitter.js — message আসলে export function emitAdminChat(io, message) { io.to(adminRoom()) .emit('watch:admin:chat', { ...message, // Extra info for admin: classType: message.classType, classId: message.classId, }); } // chat.js — send() এ, batcher এর পাশাপাশি: emitAdminChat(io, adminMessage); // ← students দেখে না, শুধু superAdmin
classType + classId admin message এ কেন?
Admin room এ সব class এর message আসে। classType+classId দিয়ে admin বুঝতে পারে কোন class থেকে এলো। Student room এ এটা দরকার নেই।
Admin vs Student Room
👤 Student Room
Rahim (Student)
স্যার এইটা বুঝলাম না
Karim (Student)
আমিও একটু confuse
class এর room ID দিয়ে — শুধু ঐ class এর member দেখে
👨‍💼 SuperAdmin Room
CLASS_CONTENT • abc123
Rahim: স্যার এইটা বুঝলাম না
CYCLE_CONTENT • def456
Sakib: ধন্যবাদ স্যার!
'watch:admin:room' — সব class এর message একসাথে
Dual Emit — No Extra DB Call
Message একই object — student room এ simplified, admin room এ classType/classId extra। কোনো additional query নেই।
Batcher vs Admin Room
Student room → batcher queue (120ms delay, batched)। Admin room → emitAdminChat() instant। Admin real-time monitoring নির্ভুল।
12
// Summary — Full Picture

এই Chat System কীভাবে হাজার User Handle করে?

প্রতিটা layer এর design decision — কেন করা হয়েছে।

🪣
Token Bucket
Socket memory তে। Redis call শূন্য। Spam mathematically impossible। 0ms overhead।
🚫
3-Layer Ban
Socket memo 15s → Redis 30s → Prisma। DB call 99% কমে। Banned instant feel।
📦
120ms Batcher
50 msg → 1 emit। Network calls 95%+। Student UI smooth — lag নেই।
🌊
Redis Stream
Async DB write। Socket thread free। At-least-once delivery guaranteed।
⚙️
Batch Worker
createMany 200 rows। Separate process — DB failure socket affect করে না।
☠️
Dead Letter
Failed message হারায় না। Dead stream এ যায়। Debug করা যায়।
🗑️
5-Day Cleanup
Cron job, Redis lock, 5K batch delete + 200ms pause। DB free থাকে।
👨‍💼
Admin Room
SuperAdmin সব class monitor। Student room আলাদা। Instant emit।
মূল কথা
এই system এর প্রতিটা layer এর একটাই উদ্দেশ্য — socket thread কে সবসময় free রাখা। Rate limit → ban check → batcher → stream — সব এই লক্ষ্যে। Socket এ যত কম I/O, তত বেশি users handle করা যায়।