Systems Hub
রহিম message পাঠাল —
পর্দার আড়ালে ৭টা ঘটনা ঘটল
একটা simple "send" button-এর পেছনে
Token Bucket, Redis Stream, Batcher, Worker —
এই সব না থাকলে হাজার user একসাথে chat করতে পারত না
৭-Stage Journey ৫ Gate Security Token Bucket Redis Stream 120ms Batcher
Production Chat System বুঝব আজকে
SLIDE 01

Naive Approach কেন Fail করে?

সহজ পদ্ধতিতে বানালে কী হয়?
Naive: রহিম send করলে → সাথে সাথে DB INSERT → সাথে সাথে সবাইকে broadcast
// Naive — সহজ কিন্তু ভুল socket.on('message', async (data) => { // প্রতি message = একটা DB query await prisma.message.create({ data: { content: data.content, ... } }); // প্রতি message = সব client-কে blast io.emit('new-message', data); });
💥
১০০০ জন একসাথে message পাঠাল → ১০০০টা DB query → database slow → সবার chat freeze → server crash
Production approach: Message আসলে সাথে সাথে DB-তে যাবে না। Pipeline-এ যাবে।

Production Pipeline

  • ৫ Gate → message accept/reject
  • Redis-এ push → queue
  • ১২০ms batch → bulk insert
  • Dedicated worker process করে
  • তারপর broadcast
🚀
১০০০ message → ১টা bulk DB insert। Server হাসতে হাসতে সামলায়।
SLIDE 02

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

রহিমের "হ্যালো" কীভাবে করিমের কাছে পৌঁছায়

Client Send — রহিম socket-এ 'send-message' emit করে

5 Gate Check — Auth ✓ Enrollment ✓ Ban ✓ Rate Limit ✓ Content ✓ — সব pass?

Optimistic Broadcast — DB-তে যাওয়ার আগেই room-এ emit (instant feel)

Redis Stream Push — message-কে Redis stream-এ push করো (queue)

120ms Buffer — Worker ১২০ms অপেক্ষা করে, জমা করে

Bulk DB Insert — ১২০ms-এর সব message একটা prisma.createMany()

Acknowledge — DB-তে save হলে client-কে confirm পাঠাও (✓✓ blue tick)

SLIDE 03

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

একটা gate-ও fail করলে message যাবে না

Authentication — JWT valid? User আছে? না থাকলে disconnect।

Enrollment — রহিম কি এই class-এ? Redis cache check → miss হলে DB।

Ban Check — রহিম banned? Redis: banned:userId key আছে? → reject।

Token Bucket Rate Limit — শেষ ১ মিনিটে বেশি message? → slow down।

Content Validation — Empty? ১০০০ char বেশি? Forbidden words? → reject।

async function validateMessage(socket, data) { // Gate 1: auth (middleware-এ done) // Gate 2: enrollment const enrolled = await redis.get( `enroll:${socket.user.id}:${data.classId}` ); if (!enrolled) throw new Error('Not enrolled'); // Gate 3: ban const banned = await redis.get( `banned:${socket.user.id}` ); if (banned) throw new Error('Banned'); // Gate 4: token bucket await checkTokenBucket(socket.user.id); // Gate 5: content if (!data.content?.trim()) throw new Error('Empty message'); if (data.content.length > 1000) throw new Error('Too long'); return true; }
SLIDE 04

Token Bucket — Spam Block করার পদ্ধতি

রহিম যদি ১ সেকেন্ডে ১০০ message পাঠায়?

Token Bucket Logic

  • রহিমের bucket: ১০টা token
  • প্রতি message → ১ token খরচ
  • প্রতি ১ সেকেন্ড → ২ token refill
  • Token = 0 → message block
  • স্বাভাবিক user: কখনো শেষ হবে না
  • Spammer: ৫ সেকেন্ডেই block
💡
উদাহরণ: করিম ১ মিনিটে ৫০টা message পাঠাল (স্বাভাবিক) → কোনো সমস্যা নেই। রহিম ১ সেকেন্ডে ৫০টা → block।
async function checkTokenBucket(userId) { const key = `bucket:chat:${userId}`; // Redis-এ atomic operation const [tokens] = await redis .multi() .get(key) .exec(); const current = tokens ? parseInt(tokens[1]) : 10; if (current <= 0) { throw new Error( 'বেশি দ্রুত message করছ, একটু থামো' ); } // Token কমাও, ৬০ সেকেন্ড TTL await redis.set(key, current - 1, 'EX' , 60); } // Refill — separate cron job // প্রতি সেকেন্ডে ২ token add করো // কিন্তু max 10-এর বেশি নয়
SLIDE 05

Student Ban — ৩ Layer Cache

প্রতি message-এ DB-তে ban check করা যাবে না

Layer 1: Socket Memory

  • Connection-এর সময় ban status load করো
  • socket.isBanned = true/false
  • প্রতি message-এ in-memory check — 0ms

Layer 2: Redis Cache

  • socket.isBanned miss হলে Redis check
  • key: banned:userId
  • TTL: ৫ মিনিট
  • ~0.1ms response time

Layer 3: Database

  • Redis-এও নেই → DB query
  • Result Redis-এ cache করো
  • শুধু প্রতি ৫ মিনিটে একবার
// Connection-এ load (Layer 1) io.use(async (socket, next) => { const banned = await redis.get( `banned:${socket.user.id}` ); socket.isBanned = !!banned; next(); }); // Message-এ check async function checkBan(socket) { // Layer 1: memory if (socket.isBanned) throw err; // Layer 2: Redis const redisBan = await redis.get( `banned:${socket.user.id}` ); if (redisBan) { socket.isBanned = true; throw err; } // Layer 3: DB (rare case) const dbBan = await prisma.ban.findFirst({ where: { userId: socket.user.id, isActive: true } }); if (dbBan) { await redis.setex( `banned:${socket.user.id}`, 300, '1' ); socket.isBanned = true; throw err; } }
SLIDE 06

Redis Stream + Dedicated Worker

Message queue → Bulk DB insert

Redis Stream কী?

  • একটা ordered message queue
  • Producer (socket server) push করে
  • Consumer (worker) pull করে process
  • Message হারায় না — persistent
// Producer — Socket server socket.on('send-message', async (data) => { await validateMessage(socket, data); // Optimistic broadcast (instant) io.to(`class:${data.classId}`) .emit('new-message', { ...data, senderId: socket.user.id, pending: true }); // Redis Stream-এ push await redis.xadd('chat:stream', '*', { content: data.content, senderId: socket.user.id, classId: data.classId, socketId: socket.id }); });
// Consumer — Dedicated Worker Process const buffer = []; async function worker() { while (true) { // Redis Stream থেকে read const msgs = await redis.xread( 'COUNT', 100, 'BLOCK', 120, // 120ms wait 'STREAMS', 'chat:stream', lastId ); if (msgs) buffer.push(...msgs); // 120ms পর bulk insert if (buffer.length > 0) { await prisma.message.createMany({ data: buffer.map(m => ({ content: m.content, senderId: m.senderId, classId: m.classId, })) }); // Acknowledge — blue tick buffer.forEach(m => { io.to(m.socketId).emit('msg-saved'); }); buffer.length = 0; } } } worker();
SLIDE 07

activeChats Table — কে এখন কোথায়?

SuperAdmin দেখবে কোন class-এ কতজন chat করছে
-- activeChats table CREATE TABLE activeChats ( id UUID PRIMARY KEY, userId UUID REFERENCES User(id), classId UUID REFERENCES Class(id), socketId VARCHAR(100), joinedAt TIMESTAMP DEFAULT NOW(), lastActive TIMESTAMP DEFAULT NOW(), UNIQUE(userId, classId) -- duplicate নেই ); -- Join করলে upsert INSERT INTO activeChats (userId, classId, socketId, lastActive) VALUES (...) ON CONFLICT (userId, classId) DO UPDATE SET socketId = EXCLUDED.socketId, lastActive = NOW(); -- Disconnect হলে delete DELETE FROM activeChats WHERE socketId = $1;

SuperAdmin Dashboard

  • সব active class দেখতে পাবে
  • প্রতি class-এ কতজন — realtime
  • কোন user কোন class-এ আছে
  • Idle user detect (lastActive পুরনো)
// SuperAdmin: সব class-এর activity const activeSummary = await prisma .activeChats.groupBy({ by: ['classId'], _count: { userId: true }, orderBy: { _count: { userId: 'desc' } } }); // [{ classId: 'c1', _count: { userId: 450 } }, // { classId: 'c2', _count: { userId: 230 } }]
SuperAdmin একটা screen-এ দেখবে: class-101-এ ৪৫০ জন, class-202-এ ২৩০ জন — realtime।
SLIDE 08

কীভাবে হাজার User Handle করে?

প্রতিটা layer-এর ভূমিকা

Fast Response

  • Optimistic broadcast → instant UI
  • 3-layer ban cache → 0ms check
  • In-memory token bucket

DB Protection

  • 120ms batch → bulk insert
  • Redis Stream buffer
  • Dedicated worker — socket server আলাদা

Spam Control

  • Token bucket per user
  • 5 gate validation
  • Content length limit
🎯
১০,০০০ জন একসাথে message পাঠাল → ৫ gate check (Redis, in-memory) → Optimistic broadcast → Redis Stream-এ push → Worker ১২০ms পর ১টা bulk INSERT। Database একটুও কাঁপল না।
সব একসাথে ধাপে ধাপে