From 3965e09bd2d1040903626899543e4c0634a008da Mon Sep 17 00:00:00 2001 From: Anuraj Jit Saikia Date: Thu, 30 Jul 2026 20:33:46 +0530 Subject: [PATCH] =?UTF-8?q?feat(chat):=20wire=20contract=20=E2=80=94=20ack?= =?UTF-8?q?ed=20sends,=2016k=20reject,=20server-minted=20identity?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Chat sends now use a Socket.io ack union (SendChatMessageAck): the server mints id/kind/sentAt, derives the room from the socket session, rejects empty and >16,000-char messages outright (never slices), and rate-limits by size — token cost = ceil(utf8 bytes of the whole serialized payload / 1000), clamped to bucket capacity so any legal message stays sendable, with a 64 KiB serialized backstop and maxHttpBufferSize at 128 KiB. The client gates the composer while a send is pending, times the ack out at 8s, keeps the draft on rejection, and only clears it when the ack matches the sent text. Includes a character counter past 14,000 and inline error surfacing. Covered by handler unit tests (malformed args, backstop, production guard wiring via call-through spy + behavioral bucket drain), rate-limit unit tests, ChatPanel tests, and a Playwright e2e for over-limit rejection, exact-16k delivery, and rate-limit rejection/recovery with an ordered non-delivery proof. Closes #192 --- client/src/components/ChatPanel.test.tsx | 112 ++++++++++++-- client/src/components/ChatPanel.tsx | 35 ++++- client/src/pages/RoomPage.tsx | 38 ++++- e2e/tests/chat-reactions.spec.js | 176 ++++++++++++++++++++++ server/src/socket/handlers.ts | 55 ++++++- server/src/socket/io.ts | 3 + server/src/socket/rate-limit.ts | 52 +++++-- server/test/handlers.test.js | 177 ++++++++++++++++++++--- server/test/socket-rate-limit.test.js | 76 +++++++++- shared/src/events.ts | 15 +- 10 files changed, 676 insertions(+), 63 deletions(-) diff --git a/client/src/components/ChatPanel.test.tsx b/client/src/components/ChatPanel.test.tsx index af7bda3..57ca6a7 100644 --- a/client/src/components/ChatPanel.test.tsx +++ b/client/src/components/ChatPanel.test.tsx @@ -25,25 +25,45 @@ beforeAll(() => { })); }); -// ChatPanel is fully controlled (input/setInput/onSend come from the parent, which -// also clears the input on send). This harness mirrors RoomPage's real wiring so -// tests exercise the true contract: onSend reads the current input, then clears it. -interface HarnessProps { onSendSpy?: (input: string) => void; messages?: ChatMessage[]; currentUserId?: string; onClose?: () => void } +type ChatAck = { ok: true } | { ok: false; code: string; message: string }; +type AckCallback = (error: Error | null, response?: ChatAck) => void; + +// ChatPanel is fully controlled. This harness mirrors RoomPage's acked send +// contract: the send remains pending until the callback settles, then only a +// successful acknowledgement may clear the original draft. +interface HarnessProps { onSendSpy?: (input: string, ack: AckCallback) => void; messages?: ChatMessage[]; currentUserId?: string; onClose?: () => void } function Harness({ onSendSpy = vi.fn(), messages = [], currentUserId = 'me', onClose = vi.fn() }: HarnessProps) { const [input, setInput] = useState(''); + const [sendError, setSendError] = useState(''); + const [sending, setSending] = useState(false); const handleSend = (e: FormEvent) => { e.preventDefault(); - if (!input.trim()) return; - onSendSpy(input); - setInput(''); + if (!input.trim() || sending) return; + const sentText = input; + setSending(true); + onSendSpy(sentText, (error, response) => { + setSending(false); + if (error) { + setSendError("Couldn't send — try again"); + return; + } + if (response?.ok) { + setInput((current) => current === sentText ? '' : current); + setSendError(''); + return; + } + setSendError(response?.message ?? "Couldn't send — try again"); + }); }; return ( { setInput(value); setSendError(''); }} onSend={handleSend} + sendError={sendError} + sending={sending} currentUserId={currentUserId} onClose={onClose} /> @@ -122,8 +142,8 @@ describe('ChatPanel', () => { }); describe('sending a message', () => { - it('invokes the send handler with the typed text and clears the input', () => { - const onSendSpy = vi.fn(); + it('clears the composer after an ok acknowledgement', () => { + const onSendSpy = vi.fn((_input, ack: AckCallback) => ack(null, { ok: true })); render(); fireEvent.change(composer(), { target: { value: 'Hello team' } }); @@ -132,10 +152,78 @@ describe('ChatPanel', () => { fireEvent.click(sendButton()); expect(onSendSpy).toHaveBeenCalledTimes(1); - expect(onSendSpy).toHaveBeenCalledWith('Hello team'); - // Parent clears the controlled input after send. + expect(onSendSpy.mock.calls[0][0]).toBe('Hello team'); expect(composer()).toHaveValue(''); }); + + it('keeps an over-limit draft, blocks send, and shows its error and counter', () => { + const onSendSpy = vi.fn(); + render(); + const overLimit = 'x'.repeat(16_001); + + fireEvent.change(composer(), { target: { value: overLimit } }); + + expect(composer()).toHaveValue(overLimit); + expect(screen.getByText('16001 / 16000')).toBeInTheDocument(); + expect(screen.getByText(/Messages can be at most 16000 characters/i)).toBeInTheDocument(); + expect(sendButton()).toBeDisabled(); + fireEvent.click(sendButton()); + expect(onSendSpy).not.toHaveBeenCalled(); + }); + + it('keeps the draft and shows the server message when the acknowledgement rejects', () => { + const onSendSpy = vi.fn((_input, ack: AckCallback) => ack(null, { + ok: false, + code: 'MESSAGE_TOO_LONG', + message: 'Messages can be at most 16000 characters.', + })); + render(); + + fireEvent.change(composer(), { target: { value: 'A draft to revise' } }); + fireEvent.click(sendButton()); + + expect(composer()).toHaveValue('A draft to revise'); + expect(screen.getByText('Messages can be at most 16000 characters.')).toBeInTheDocument(); + }); + + it('keeps the draft and offers a retry when the acknowledgement times out', () => { + const onSendSpy = vi.fn((_input, ack: AckCallback) => ack(new Error('operation has timed out'))); + render(); + + fireEvent.change(composer(), { target: { value: 'A draft to retry' } }); + fireEvent.click(sendButton()); + + expect(composer()).toHaveValue('A draft to retry'); + expect(screen.getByText("Couldn't send — try again")).toBeInTheDocument(); + expect(sendButton()).toBeEnabled(); + }); + + it('disables repeat sends while awaiting an acknowledgement and preserves a newer draft', () => { + let acknowledge: AckCallback | undefined; + const onSendSpy = vi.fn((_input, ack: AckCallback) => { acknowledge = ack; }); + render(); + + fireEvent.change(composer(), { target: { value: 'first' } }); + fireEvent.click(sendButton()); + + expect(sendButton()).toBeDisabled(); + fireEvent.change(composer(), { target: { value: 'second' } }); + fireEvent.click(sendButton()); + expect(onSendSpy).toHaveBeenCalledTimes(1); + + acknowledge?.(null, { ok: true }); + expect(composer()).toHaveValue('second'); + }); + + it('shows the character counter only after 14,000 characters', () => { + render(); + + fireEvent.change(composer(), { target: { value: 'x'.repeat(14_000) } }); + expect(screen.queryByText('14000 / 16000')).not.toBeInTheDocument(); + + fireEvent.change(composer(), { target: { value: 'x'.repeat(14_001) } }); + expect(screen.getByText('14001 / 16000')).toBeInTheDocument(); + }); }); describe('empty / whitespace guard', () => { diff --git a/client/src/components/ChatPanel.tsx b/client/src/components/ChatPanel.tsx index 784027f..f6f0703 100644 --- a/client/src/components/ChatPanel.tsx +++ b/client/src/components/ChatPanel.tsx @@ -7,12 +7,16 @@ import { Close as CloseIcon, Send as SendIcon } from '@mui/icons-material'; import { usePanelDialog } from '../hooks/usePanelDialog'; interface ChatSender { id: string; name?: string; avatar?: string } -export interface ChatMessage { id?: string; type?: 'event' | 'chat'; text: string; ts: string | number | Date; sender?: ChatSender } +export interface ChatMessage { id?: string; kind?: 'text'; type?: 'event' | 'chat'; text: string; ts: string | number | Date; sender?: ChatSender } +const CHAT_MESSAGE_LIMIT = 16_000; +const CHAT_COUNTER_THRESHOLD = 14_000; interface ChatPanelProps { messages: ChatMessage[]; input: string; setInput: (value: string) => void; onSend: (event: FormEvent) => void; + sendError?: string; + sending: boolean; currentUserId?: string; onClose: () => void; } @@ -41,10 +45,22 @@ function messageKey(msg: ChatMessage): string { // In-call chat. Desktop: a 372px wide in-flow side column. // Mobile: a bottom sheet (62vh, slides up over the video, with backdrop). -export default function ChatPanel({ messages, input, setInput, onSend, currentUserId, onClose }: ChatPanelProps) { +export default function ChatPanel({ messages, input, setInput, onSend, sendError, sending, currentUserId, onClose }: ChatPanelProps) { const bottomRef = useRef(null); const isMobile = useMediaQuery((theme) => theme.breakpoints.down('sm')); const { initialFocusRef, panelRef, onKeyDown } = usePanelDialog(onClose); + const tooLong = input.length > CHAT_MESSAGE_LIMIT; + const composerError = tooLong + ? `Messages can be at most ${CHAT_MESSAGE_LIMIT} characters.` + : sendError; + + function handleSubmit(event: FormEvent) { + if (tooLong) { + event.preventDefault(); + return; + } + onSend(event); + } useEffect(() => { bottomRef.current?.scrollIntoView({ behavior: 'smooth' }); @@ -189,7 +205,7 @@ export default function ChatPanel({ messages, input, setInput, onSend, currentUs {/* Composer */} - + setInput(e.target.value)} autoComplete="off" + error={Boolean(composerError)} + helperText={composerError} slotProps={{ input: { sx: { borderRadius: 999, bgcolor: 'rgba(255,255,255,0.04)' }, endAdornment: ( - + @@ -210,6 +228,15 @@ export default function ChatPanel({ messages, input, setInput, onSend, currentUs }, }} /> + {input.length > CHAT_COUNTER_THRESHOLD && ( + + {input.length} / {CHAT_MESSAGE_LIMIT} + + )} ); diff --git a/client/src/pages/RoomPage.tsx b/client/src/pages/RoomPage.tsx index ccbe64f..1e2fe09 100644 --- a/client/src/pages/RoomPage.tsx +++ b/client/src/pages/RoomPage.tsx @@ -103,6 +103,8 @@ export default function RoomPage() { const [messages, setMessages] = useState([]); const [users, setUsers] = useState([]); const [input, setInput] = useState(''); + const [chatError, setChatError] = useState(''); + const [chatSendPending, setChatSendPending] = useState(false); // Single right rail: only one of Chat / People / Transcript is open at a time (Meet-style). const [activePanel, setActivePanel] = useState(null); // 'chat' | 'people' | 'transcript' | null const showChat = activePanel === 'chat'; @@ -361,7 +363,14 @@ export default function RoomPage() { )); }); socket.on('chat-message', (msg) => { - setMessages((prev) => [...prev, { type: 'chat', ...msg }]); + setMessages((prev) => [...prev, { + type: 'chat', + id: msg.id, + kind: msg.kind, + sender: msg.sender, + text: msg.text, + ts: msg.sentAt, + }]); const fromOther = msg.sender?.id !== userIdRef.current; if (fromOther) playSound('message'); // When the chat is closed, surface a Meet-style preview + unread badge. @@ -497,10 +506,22 @@ export default function RoomPage() { function sendMessage(e: FormEvent) { e.preventDefault(); - const text = input.trim(); - if (!text) return; - socket.emit('chat-message', { roomId, text }); - setInput(''); + if (!input.trim() || chatSendPending) return; + const sentText = input; + setChatSendPending(true); + socket.timeout(8000).emit('chat-message', { text: sentText }, (err, response) => { + setChatSendPending(false); + if (err || !response) { + setChatError("Couldn't send — try again"); + return; + } + if (response.ok) { + setInput((current) => current === sentText ? '' : current); + setChatError(''); + return; + } + setChatError(response.message); + }); } @@ -1446,8 +1467,13 @@ export default function RoomPage() { { + setInput(value); + setChatError(''); + }} onSend={sendMessage} + sendError={chatError} + sending={chatSendPending} currentUserId={user?.id} onClose={() => setActivePanel(null)} /> diff --git a/e2e/tests/chat-reactions.spec.js b/e2e/tests/chat-reactions.spec.js index 43e31f9..27c5678 100644 --- a/e2e/tests/chat-reactions.spec.js +++ b/e2e/tests/chat-reactions.spec.js @@ -98,6 +98,182 @@ test.describe('two-peer chat relay', () => { }); }); +test.describe('chat send contract', () => { + test('keeps an over-limit draft, shows its error, and delivers nothing to the other peer', async ({ browser }) => { + const { pageA, pageB, close } = await joinedPair(browser); + const chatPanelA = pageA.getByTestId('chat-panel'); + const chatPanelB = pageB.getByTestId('chat-panel'); + const composerA = pageA.getByPlaceholder('Send a message to everyone'); + const overLimit = 'x'.repeat(16_001); + + await openChat(pageA); + await openChat(pageB); + + // fill() assigns the complete draft in one operation; typing 16k keys would + // make this contract test slow and unlike a pasted draft. + await composerA.fill(overLimit); + + await expect(composerA).toHaveValue(overLimit); + await expect(chatPanelA.getByText(/Messages can be at most 16000 characters/i)).toBeVisible(); + await expect(pageA.getByRole('button', { name: 'Send message' })).toBeDisabled(); + + // The send button is disabled, so exercise the OTHER submit path too — + // Enter in the composer — and confirm the draft is still sitting there + // afterwards instead of having been consumed by an attempted send. + await composerA.press('Enter'); + await expect(composerA).toHaveValue(overLimit); + + // Non-delivery only means something once a LATER message has arrived: an + // empty panel satisfies toHaveCount(0) instantly. Send a normal sentinel + // from the same peer and wait for B to render it — the over-limit draft was + // submitted first, had the whole round trip to surface, and still has not. + const sentinel = `sentinel-after-over-limit-${Date.now()}`; + await composerA.fill(sentinel); + await pageA.getByRole('button', { name: 'Send message' }).click(); + await expect(chatPanelB.getByText(sentinel, { exact: true })).toBeVisible({ timeout: 15_000 }); + await expect(chatPanelB.getByText(overLimit, { exact: true })).toHaveCount(0); + + await close(); + }); + + test('delivers an exactly-16,000-character message intact and clears the sender composer', async ({ browser }) => { + const { pageA, pageB, close } = await joinedPair(browser); + const chatPanelB = pageB.getByTestId('chat-panel'); + const composerA = pageA.getByPlaceholder('Send a message to everyone'); + const exactLimit = 'x'.repeat(16_000); + + await openChat(pageA); + await openChat(pageB); + await composerA.fill(exactLimit); + await pageA.getByRole('button', { name: 'Send message' }).click(); + + const receivedMessage = chatPanelB.getByText(exactLimit, { exact: true }); + await expect(receivedMessage).toBeVisible({ timeout: 15_000 }); + await expect(receivedMessage).toHaveText(exactLimit); + await expect(composerA).toHaveValue('', { timeout: 15_000 }); + + await close(); + }); + + // The server's own rejection path, end to end. The chat bucket (capacity 20, + // refill 5/s) charges each send by its serialized payload size, so a burst of + // ~4 KB messages costs ~5 tokens apiece and drains it within a handful of + // sends. What is under test is what RoomPage does with the resulting + // `{ ok:false, code:'RATE_LIMITED', … }` ack: surface its message, KEEP the + // draft, and stay usable once tokens refill. + test('surfaces the server rate-limit rejection, keeps the draft, and recovers', async ({ browser }) => { + // Two-peer setup + a send burst + waiting out a real token refill does not + // fit the 30 s default; the budget is generous rather than a fixed wait. + test.setTimeout(120_000); + + const { pageA, pageB, close } = await joinedPair(browser); + const chatPanelA = pageA.getByTestId('chat-panel'); + const chatPanelB = pageB.getByTestId('chat-panel'); + const composerA = pageA.getByPlaceholder('Send a message to everyone'); + const sendA = pageA.getByRole('button', { name: 'Send message' }); + const rateLimitError = chatPanelA.getByText(/rate limit exceeded/i); + + await openChat(pageA); + await openChat(pageB); + + // Send `text` from A, retrying ONLY while the server keeps refusing it. + // + // Retrying is safe after a RATE_LIMITED ack: the guard dropped the event, so + // nothing was broadcast. It is NOT safe after RoomPage's 8 s ack timeout — + // that timeout does not cancel the already-emitted Socket.IO event, so the + // server may have accepted and broadcast the text while only the ack was + // late. Re-sending there would manufacture a second LEGITIMATE delivery and + // corrupt the exact-count assertion below, so an ambiguous settle records a + // failure and emits nothing further. + async function sendUntilAccepted(text) { + let ambiguousSettle = null; + await expect(async () => { + await composerA.fill(text); + await sendA.click(); + + // Wait for THIS attempt's ack callback to run: RoomPage clears + // chatSendPending inside it, so either the draft cleared (ok ack) or the + // send button came back enabled with the draft intact (a failure path). + // The window exceeds RoomPage's own 8 s ack timeout, so the callback has + // certainly fired one way or the other by the time this resolves. + await expect + .poll( + async () => (await composerA.inputValue()) === '' || (await sendA.isEnabled()), + { timeout: 15_000 }, + ) + .toBe(true); + + if ((await composerA.inputValue()) === '') return; // ok ack → accepted. + + // The draft survived, so this attempt took a failure path — and that + // same callback wrote chatError in the same React batch that cleared + // chatSendPending. The copy on screen right now is therefore THIS + // attempt's: no earlier attempt could have re-enabled the button without + // also overwriting the very error being read here, and an accepted send + // clears it outright. RoomPage's two failure copies are distinct, and + // only the rate-limit one is safe to retry. + if (!(await rateLimitError.isVisible())) { + ambiguousSettle = "the send failed without a RATE_LIMITED error — RoomPage showed its ack-timeout copy, and that timeout does not cancel the emitted event, so the server may already have broadcast this text"; + return; // Ends the retry loop without re-sending; asserted below. + } + throw new Error('rate limited — retry once the chat bucket refills'); + }).toPass({ timeout: 60_000 }); + + expect( + ambiguousSettle, + `chat send must settle on a definitive rejection before it is retried; instead ${ambiguousSettle}`, + ).toBeNull(); + } + + // Burst until the server denies one. Each iteration settles on exactly one + // of two visible outcomes — accepted (the ok ack clears the composer) or + // denied (the error surfaces) — so the loop polls real state, never sleeps. + let blockedDraft = null; + for (let attempt = 0; attempt < 60 && blockedDraft === null; attempt += 1) { + const text = `burst-${attempt}-${'x'.repeat(4_000)}`; + await composerA.fill(text); + await sendA.click(); + await expect + .poll( + async () => (await rateLimitError.isVisible()) || (await composerA.inputValue()) === '', + { timeout: 15_000 }, + ) + .toBe(true); + if (await rateLimitError.isVisible()) blockedDraft = text; + } + + expect(blockedDraft, 'the chat bucket should deny a send within the burst').not.toBeNull(); + + // The rejection is visible on A and the draft survived it. Whether the + // denied send leaked to B is NOT asserted here: the denial ack reaches A + // over a channel independent of any broadcast reaching B, so a count check + // at this instant could simply be outrunning the leak. That is established + // below, ordered behind a later arrival. + await expect(rateLimitError).toBeVisible(); + await expect(composerA).toHaveValue(blockedDraft); + + // …and that very same draft goes through once the bucket refills, so the + // limit is a transient back-off rather than a wedged composer. + await sendUntilAccepted(blockedDraft); + await expect(rateLimitError).toBeHidden(); + + // Now order the non-delivery proof: send a sentinel AFTER the recovery and + // wait for B to render it. Everything A sent earlier has had its full round + // trip by then, so the blocked draft must appear on B EXACTLY once — that + // one copy is the recovery send; a second would be the rate-limited send + // having leaked through despite its RATE_LIMITED ack. + const sentinel = `after-recovery-sentinel-${Date.now()}`; + await sendUntilAccepted(sentinel); + await expect(chatPanelB.getByText(sentinel, { exact: true })).toBeVisible({ timeout: 15_000 }); + + const deliveredDraft = chatPanelB.getByText(blockedDraft, { exact: true }); + await expect(deliveredDraft).toHaveCount(1); + await expect(deliveredDraft).toBeVisible(); + + await close(); + }); +}); + test.describe('chat unread badge', () => { test('a message arriving while the panel is closed badges the toggle, clearing on open', async ({ browser }) => { const { pageA, pageB, close } = await joinedPair(browser); diff --git a/server/src/socket/handlers.ts b/server/src/socket/handlers.ts index 2c61dba..ccbc70f 100644 --- a/server/src/socket/handlers.ts +++ b/server/src/socket/handlers.ts @@ -18,8 +18,18 @@ import { transcriptionConfigured, } from '../transcription/meeting-transcription.js'; import { logger } from '../config/logger.js'; -import { socketRateLimiter } from './rate-limit.js'; +import { chatMessageTokenCost, chatPayloadByteLength, socketRateLimiter } from './rate-limit.js'; import type { Server } from 'socket.io'; +import { randomUUID } from 'node:crypto'; + +const MAX_CHAT_MESSAGE_LENGTH = 16_000; +const MAX_CHAT_MESSAGE_BYTES = 64 * 1024; + +function chatTextFromPayload(payload: unknown): string | undefined { + if (!payload || typeof payload !== 'object') return undefined; + const { text } = payload as { text?: unknown }; + return typeof text === 'string' ? text : undefined; +} // A dropped connection (network blip, reload, server restart) makes Socket.IO // reconnect with a brand-new socket: the old socket fires `disconnect` (→ @@ -104,16 +114,47 @@ export function registerHandlers(io: Server) { if (getRoomUsers(roomId).length === 0) scheduleTranscriptExpiry(roomId); }); - socket.on('chat-message', socketRateLimiter.guard(socket, 'chat', 'chat-message', ({ roomId, text }) => { - if (!roomId || !text || typeof text !== 'string') return; - const trimmed = text.trim().slice(0, 1000); - if (!trimmed) return; + socket.on('chat-message', socketRateLimiter.guard(socket, 'chat', 'chat-message', (...args: unknown[]) => { + const payload = args[0]; + const text = chatTextFromPayload(payload); + const maybeCallback = args.at(-1); + const callback = typeof maybeCallback === 'function' + ? maybeCallback as (response: unknown) => void + : undefined; + const roomId = getUserRoom(socket.id); + if (!roomId) { + return callback?.({ ok: false, code: 'NOT_IN_ROOM', message: 'Join a room before sending a message.' }); + } + if (typeof text !== 'string' || !text.trim()) { + return callback?.({ ok: false, code: 'EMPTY_MESSAGE', message: 'Enter a message before sending.' }); + } + const payloadBytes = chatPayloadByteLength(payload); + if (text.length > MAX_CHAT_MESSAGE_LENGTH || payloadBytes === null || payloadBytes > MAX_CHAT_MESSAGE_BYTES) { + return callback?.({ + ok: false, + code: 'MESSAGE_TOO_LONG', + message: `Messages can be at most ${MAX_CHAT_MESSAGE_LENGTH.toLocaleString()} characters.`, + maxLength: MAX_CHAT_MESSAGE_LENGTH, + }); + } + const messageId = randomUUID(); io.to(roomId).emit('chat-message', { + id: messageId, + kind: 'text', sender: socket.user, - text: trimmed, - ts: Date.now(), + text, + sentAt: Date.now(), }); + return callback?.({ ok: true, messageId }); + }, { + cost: (payload: unknown) => chatMessageTokenCost(payload), + rateLimitAck: (retryAfterMs) => ({ + ok: false, + code: 'RATE_LIMITED', + message: 'Rate limit exceeded — slow down and try again.', + retryAfterMs, + }), })); // Shared transcription is host-controlled and server-authoritative. Clients diff --git a/server/src/socket/io.ts b/server/src/socket/io.ts index 839cc25..060aa1f 100644 --- a/server/src/socket/io.ts +++ b/server/src/socket/io.ts @@ -7,6 +7,9 @@ let io: Server | undefined; export function initSocket(httpServer: HttpServer) { io = new Server(httpServer, { cors: { origin: env.clientUrl, credentials: true }, + // Chat rejects serialized payloads above 64 KiB; PCM audio chunks and SFU + // signaling packets stay far below this 128 KiB transport ingress bound. + maxHttpBufferSize: 128 * 1024, }); return io; } diff --git a/server/src/socket/rate-limit.ts b/server/src/socket/rate-limit.ts index 6895663..8ce98e1 100644 --- a/server/src/socket/rate-limit.ts +++ b/server/src/socket/rate-limit.ts @@ -9,7 +9,7 @@ // resumes the drained bucket instead of resetting it. // // A bucket holds up to `capacity` tokens and refills at `refillPerSec`; every -// guarded event spends one token. When a bucket runs dry the event is DROPPED — +// guarded event spends one token unless it supplies a size-based cost. When a bucket runs dry the event is DROPPED — // the handler never runs — and, if the event carried an ack callback, that // callback receives a structured `{ error, retryAfterMs }` (the same shape the // HTTP 429 body uses) so the client can back off. We never disconnect on an @@ -44,29 +44,48 @@ export interface BucketState { export interface ConsumeResult { allowed: boolean; - /** When denied, ms until enough tokens refill to allow one event. */ + /** When denied, ms until enough tokens refill to cover the requested cost. */ retryAfterMs: number; } /** - * Refill by elapsed time, then try to spend one token. Mutates `state`. + * Refill by elapsed time, then try to spend the requested number of tokens. Mutates `state`. * Pure aside from the mutation + the injected `now`, so it unit-tests cleanly. */ -export function consumeToken(state: BucketState, config: BucketConfig, now: number): ConsumeResult { +export function consumeToken(state: BucketState, config: BucketConfig, now: number, cost = 1): ConsumeResult { const elapsedSec = Math.max(0, (now - state.last) / 1000); state.tokens = Math.min(config.capacity, state.tokens + elapsedSec * config.refillPerSec); state.last = now; - if (state.tokens >= 1) { - state.tokens -= 1; + if (state.tokens >= cost) { + state.tokens -= cost; return { allowed: true, retryAfterMs: 0 }; } - const deficit = 1 - state.tokens; + const deficit = cost - state.tokens; const retryAfterMs = Math.ceil((deficit / config.refillPerSec) * 1000); return { allowed: false, retryAfterMs }; } +/** Serialized byte size of a chat event's first argument, or null if it cannot be encoded. */ +export function chatPayloadByteLength(payload: unknown): number | null { + try { + const serialized = JSON.stringify(payload); + return typeof serialized === 'string' ? Buffer.byteLength(serialized, 'utf8') : null; + } catch { + return null; + } +} + +/** Cost of a complete chat payload in UTF-8 kilobyte tokens. */ +export function chatMessageTokenCost(payload: unknown): number { + const bytes = chatPayloadByteLength(payload); + // A value Socket.IO cannot serialize must never become a cheap rate-limit bypass. + return bytes === null + ? env.rateLimit.socket.chat.capacity + : Math.max(1, Math.ceil(bytes / 1000)); +} + /** * Stable identity a bucket is keyed on: the authenticated user id when present * (socketAuth runs before any handler, so it normally is), else the handshake @@ -119,6 +138,10 @@ export interface SocketRateLimiter { bucket: string, event: string, handler: (...args: A) => void, + options?: { + cost?: (...args: A) => number; + rateLimitAck?: (retryAfterMs: number) => unknown; + }, ): (...args: A) => void; /** Exposed for tests: the actor's bucket state (created on demand). */ getState(socket: Socket, bucket: string): BucketState; @@ -185,6 +208,10 @@ export function createSocketRateLimiter(opts: SocketRateLimiterOptions): SocketR bucket: string, event: string, handler: (...args: A) => void, + options?: { + cost?: (...args: A) => number; + rateLimitAck?: (retryAfterMs: number) => unknown; + }, ): (...args: A) => void { const config = opts.buckets[bucket]; if (!config) throw new Error(`Unknown rate-limit bucket: ${bucket}`); @@ -194,7 +221,14 @@ export function createSocketRateLimiter(opts: SocketRateLimiterOptions): SocketR return (...args: A): void => { const state = getState(socket, bucket); - const { allowed, retryAfterMs } = consumeToken(state, config, now()); + const requestedCost = options?.cost?.(...args) ?? 1; + // Clamp deliberately: a legal message may cost more than capacity (for + // example 16,000 CJK characters). It is admitted once and drains this + // bucket instead of being permanently impossible to send. + const cost = Number.isFinite(requestedCost) && requestedCost > 0 + ? Math.min(Math.ceil(requestedCost), config.capacity) + : 1; + const { allowed, retryAfterMs } = consumeToken(state, config, now(), cost); if (allowed) { state.violations = 0; @@ -229,7 +263,7 @@ export function createSocketRateLimiter(opts: SocketRateLimiterOptions): SocketR // event is silently dropped (no ack channel to answer on). const maybeCallback = args[args.length - 1]; if (typeof maybeCallback === 'function') { - (maybeCallback as (ack: { error: string; retryAfterMs: number }) => void)({ + (maybeCallback as (ack: unknown) => void)(options?.rateLimitAck?.(retryAfterMs) ?? { error: 'Rate limit exceeded — slow down and try again.', retryAfterMs, }); diff --git a/server/test/handlers.test.js b/server/test/handlers.test.js index bd6feca..6ed2700 100644 --- a/server/test/handlers.test.js +++ b/server/test/handlers.test.js @@ -38,6 +38,11 @@ vi.mock('../src/config/logger.js', () => ({ import { registerHandlers } from '../src/socket/handlers.js'; import { registerWebrtcHandlers } from '../src/socket/webrtc.js'; import { registerSfuHandlers } from '../src/socket/sfu-handlers.js'; +// The REAL limiter — handlers.test.js deliberately does not mock it, so the +// chat guard these tests exercise is the production one, wired with production +// options and the production env bucket config. +import { chatMessageTokenCost, socketRateLimiter } from '../src/socket/rate-limit.js'; +import { env } from '../src/config/env.js'; import { addUser, removeUser, getRoomUsers, isUserInRoom, getUserRoom, } from '../src/socket/room-manager.js'; @@ -304,48 +309,176 @@ describe('disconnect grace window', () => { // chat-message // --------------------------------------------------------------------------- describe('chat-message', () => { - it('broadcasts chat-message to the room with sender + text + ts', () => { + it('rejects malformed event arguments without throwing and uses a trailing acknowledgement', () => { + getUserRoom.mockReturnValue(ROOM); + + const noPayload = setup({ ...USER, id: 'user-no-payload' }); + expect(() => noPayload.handlers['chat-message']()).not.toThrow(); + + const nullPayload = setup({ ...USER, id: 'user-null-payload' }); + const nullCallback = vi.fn(); + expect(() => nullPayload.handlers['chat-message'](null, nullCallback)).not.toThrow(); + expect(nullCallback).toHaveBeenCalledWith({ + ok: false, + code: 'EMPTY_MESSAGE', + message: expect.any(String), + }); + + const extraArgument = setup({ ...USER, id: 'user-extra-argument' }); + const callback = vi.fn(); + expect(() => extraArgument.handlers['chat-message']({ text: 'x' }, 'extra', callback)).not.toThrow(); + expect(callback).toHaveBeenCalledWith({ ok: true, messageId: expect.any(String) }); + }); + + it('delivers an exactly 16,000-character message intact with server-minted metadata', () => { const { handlers, ioEmits } = setup(); + const callback = vi.fn(); + const text = 'x'.repeat(16_000); + getUserRoom.mockReturnValue(ROOM); - handlers['chat-message']({ roomId: ROOM, text: 'Hello' }); + handlers['chat-message']({ text }, callback); const msg = ioEmits.find((e) => e.event === 'chat-message'); expect(msg).toBeDefined(); + expect(msg.target).toBe(ROOM); expect(msg.payload.sender).toEqual(USER); - expect(msg.payload.text).toBe('Hello'); - expect(typeof msg.payload.ts).toBe('number'); + expect(msg.payload.text).toBe(text); + expect(msg.payload.id).toMatch(/^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i); + expect(msg.payload.kind).toBe('text'); + expect(typeof msg.payload.sentAt).toBe('number'); + expect(callback).toHaveBeenCalledWith({ ok: true, messageId: msg.payload.id }); }); - it('ignores empty text', () => { + it('rejects empty text without broadcasting it', () => { const { handlers, ioEmits } = setup(); - handlers['chat-message']({ roomId: ROOM, text: '' }); + const callback = vi.fn(); + getUserRoom.mockReturnValue(ROOM); + handlers['chat-message']({ text: '' }, callback); + + expect(callback).toHaveBeenCalledWith({ ok: false, code: 'EMPTY_MESSAGE', message: expect.any(String) }); expect(ioEmits.some((e) => e.event === 'chat-message')).toBe(false); }); - it('ignores whitespace-only text', () => { + it('rejects a sender that is not in a room', () => { const { handlers, ioEmits } = setup(); - handlers['chat-message']({ roomId: ROOM, text: ' ' }); + const callback = vi.fn(); + getUserRoom.mockReturnValue(null); + handlers['chat-message']({ text: 'Hi' }, callback); + + expect(callback).toHaveBeenCalledWith({ ok: false, code: 'NOT_IN_ROOM', message: expect.any(String) }); expect(ioEmits.some((e) => e.event === 'chat-message')).toBe(false); }); - it('ignores non-string text', () => { - const { handlers, ioEmits } = setup(); - handlers['chat-message']({ roomId: ROOM, text: 42 }); - expect(ioEmits.some((e) => e.event === 'chat-message')).toBe(false); + it('rejects a 16,001-character message without broadcasting it', () => { + const { handlers, ioEmits } = setup({ ...USER, id: 'user-oversized-message' }); + const callback = vi.fn(); + getUserRoom.mockReturnValue(ROOM); + + handlers['chat-message']({ text: 'x'.repeat(16_001) }, callback); + + expect(callback).toHaveBeenCalledWith({ + ok: false, + code: 'MESSAGE_TOO_LONG', + message: expect.any(String), + maxLength: 16_000, + }); + expect(ioEmits.some((event) => event.event === 'chat-message')).toBe(false); }); - it('ignores missing roomId', () => { - const { handlers, ioEmits } = setup(); - handlers['chat-message']({ text: 'Hi' }); - expect(ioEmits.some((e) => e.event === 'chat-message')).toBe(false); + it('rejects a payload whose serialized size exceeds the chat byte backstop', () => { + const { handlers, ioEmits } = setup({ ...USER, id: 'user-padded-payload' }); + const callback = vi.fn(); + getUserRoom.mockReturnValue(ROOM); + + handlers['chat-message']({ text: 'ok', padding: 'x'.repeat(64 * 1024) }, callback); + + expect(callback).toHaveBeenCalledWith({ + ok: false, + code: 'MESSAGE_TOO_LONG', + message: expect.any(String), + maxLength: 16_000, + }); + expect(ioEmits.some((event) => event.event === 'chat-message')).toBe(false); }); +}); - it('trims and caps text to 1000 characters', () => { - const { handlers, ioEmits } = setup(); - const long = 'x'.repeat(1500); - handlers['chat-message']({ roomId: ROOM, text: ' ' + long + ' ' }); - const msg = ioEmits.find((e) => e.event === 'chat-message'); - expect(msg.payload.text.length).toBe(1000); +// --------------------------------------------------------------------------- +// chat-message rate limiting — asserts the PRODUCTION guard options, not a +// reconstruction of them. socket/rate-limit.js is not mocked in this file, so +// `handlers['chat-message']` IS the guarded wrapper handlers.js registered, +// spending tokens from the real env-configured chat bucket. +// --------------------------------------------------------------------------- +describe('chat-message rate limiting (production wiring)', () => { + it('guards chat-message on the chat bucket with the production cost + rateLimitAck options', () => { + const guardSpy = vi.spyOn(socketRateLimiter, 'guard'); + try { + setup({ ...USER, id: 'user-guard-options' }, 'sock-guard-options'); + + const chatGuard = guardSpy.mock.calls + .find(([, bucket, event]) => bucket === 'chat' && event === 'chat-message'); + expect(chatGuard, 'handlers.js must guard chat-message on the chat bucket').toBeDefined(); + + const options = chatGuard[4]; + expect(options, 'handlers.js must pass cost + rateLimitAck options').toBeDefined(); + + // cost() is the size-weighted cost of the WHOLE payload — padding beyond + // the text field is charged, so it cannot be a flat per-event 1. + const small = { text: 'hello' }; + const padded = { text: 'ok', padding: 'x'.repeat(15_000) }; + expect(options.cost(small)).toBe(chatMessageTokenCost(small)); + expect(options.cost(padded)).toBe(chatMessageTokenCost(padded)); + expect(options.cost(padded)).toBeGreaterThan(options.cost(small)); + + // rateLimitAck() mints exactly the locked wire shape, echoing retryAfterMs. + const ack = options.rateLimitAck(1_234); + expect(ack).toEqual({ + ok: false, + code: 'RATE_LIMITED', + message: expect.any(String), + retryAfterMs: 1_234, + }); + expect(ack.message.length).toBeGreaterThan(0); + } finally { + guardSpy.mockRestore(); + } + }); + + it('drops an over-budget send through the production guard and acks the locked RATE_LIMITED shape', () => { + const { handlers, ioEmits } = setup( + { ...USER, id: 'user-production-rate-limit' }, + 'sock-production-rate-limit', + ); + getUserRoom.mockReturnValue(ROOM); + + // Each send is deliberately big enough to cost several tokens, so the + // bucket drains in a handful of calls no matter how fast the loop runs + // (the limiter reads the real clock; refill over a few ms is negligible). + const text = 'x'.repeat(4_000); + const costPerSend = chatMessageTokenCost({ text }); + const { capacity } = env.rateLimit.socket.chat; + const sends = Math.ceil(capacity / costPerSend) + 3; + + const acks = []; + for (let i = 0; i < sends; i += 1) { + handlers['chat-message']({ text: `${i}-${text}` }, (ack) => acks.push(ack)); + } + + const accepted = acks.filter((ack) => ack.ok); + const denied = acks.filter((ack) => !ack.ok); + expect(accepted.length).toBeGreaterThan(0); + expect(denied.length).toBeGreaterThan(0); + + expect(denied[0]).toEqual({ + ok: false, + code: 'RATE_LIMITED', + message: expect.any(String), + retryAfterMs: expect.any(Number), + }); + expect(denied[0].message.length).toBeGreaterThan(0); + expect(denied[0].retryAfterMs).toBeGreaterThan(0); + + // A denied send never reaches the room: exactly one broadcast per ok ack. + expect(ioEmits.filter((e) => e.event === 'chat-message')).toHaveLength(accepted.length); }); }); diff --git a/server/test/socket-rate-limit.test.js b/server/test/socket-rate-limit.test.js index 64f6d70..4d73ea6 100644 --- a/server/test/socket-rate-limit.test.js +++ b/server/test/socket-rate-limit.test.js @@ -12,7 +12,8 @@ vi.mock('../src/config/logger.js', () => ({ logger: { info: vi.fn(), warn: vi.fn(), debug: vi.fn(), error: vi.fn() }, })); -import { actorKeyFor, consumeToken, createSocketRateLimiter } from '../src/socket/rate-limit.js'; +import { actorKeyFor, chatMessageTokenCost, consumeToken, createSocketRateLimiter } from '../src/socket/rate-limit.js'; +import { env } from '../src/config/env.js'; // Fake socket supporting multiple 'disconnect' listeners (like real socket.io). function makeSocket({ id = 'sock-1', userId = 'user-1', headers = {}, address = '10.0.0.9' } = {}) { @@ -114,6 +115,40 @@ describe('createSocketRateLimiter — guard', () => { expect(ack.retryAfterMs).toBeGreaterThan(0); }); + it('uses the chat contract rate-limit acknowledgement when a chat send is over budget', () => { + const limiter = createSocketRateLimiter({ + ...opts(), + buckets: { chat: { capacity: 1, refillPerSec: 1 }, signaling: { capacity: 5, refillPerSec: 5 } }, + }); + const socket = makeSocket(); + const handler = vi.fn(); + const callback = vi.fn(); + const wrapped = limiter.guard(socket, 'chat', 'chat-message', handler, { + cost: (payload) => chatMessageTokenCost(payload), + rateLimitAck: (retryAfterMs) => ({ + ok: false, + code: 'RATE_LIMITED', + message: 'Rate limit exceeded — slow down and try again.', + retryAfterMs, + }), + }); + + wrapped({ text: 'first' }, callback); + wrapped({ text: 'second' }, callback); + + expect(handler).toHaveBeenCalledTimes(1); + expect(callback).toHaveBeenCalledTimes(1); + expect(callback).toHaveBeenCalledWith({ + ok: false, + code: 'RATE_LIMITED', + message: expect.any(String), + retryAfterMs: expect.any(Number), + }); + const ack = callback.mock.calls[0][0]; + expect(ack.message).not.toHaveLength(0); + expect(ack.retryAfterMs).toBeGreaterThan(0); + }); + it('drops silently (no throw) when the over-limit event has no ack callback', () => { const limiter = createSocketRateLimiter(opts()); const socket = makeSocket(); @@ -142,6 +177,45 @@ describe('createSocketRateLimiter — guard', () => { expect(handler).toHaveBeenCalledTimes(3); }); + it('admits a legal 16,000-character multibyte chat message and drains the bucket', () => { + const limiter = createSocketRateLimiter({ + ...opts(), + buckets: { chat: { capacity: 20, refillPerSec: 1 }, signaling: { capacity: 5, refillPerSec: 5 } }, + }); + const socket = makeSocket(); + const handler = vi.fn(); + const wrapped = limiter.guard(socket, 'chat', 'chat-message', handler, { + cost: (payload) => chatMessageTokenCost(payload), + }); + + wrapped({ text: '漢'.repeat(16_000) }); + + expect(handler).toHaveBeenCalledTimes(1); + expect(limiter.getState(socket, 'chat').tokens).toBe(0); + }); + + it('charges an unserializable chat payload at the full chat bucket capacity', () => { + expect(chatMessageTokenCost({ text: 'hi', invalid: BigInt(1) })) + .toBe(env.rateLimit.socket.chat.capacity); + }); + + it('charges padding in the complete chat payload instead of its tiny text field', () => { + const limiter = createSocketRateLimiter({ + ...opts(), + buckets: { chat: { capacity: 20, refillPerSec: 1 }, signaling: { capacity: 5, refillPerSec: 5 } }, + }); + const socket = makeSocket({ id: 'padded-socket' }); + const handler = vi.fn(); + const wrapped = limiter.guard(socket, 'chat', 'chat-message', handler, { + cost: (payload) => chatMessageTokenCost(payload), + }); + + wrapped({ text: 'ok', padding: 'x'.repeat(15_000) }); + + expect(handler).toHaveBeenCalledTimes(1); + expect(limiter.getState(socket, 'chat').tokens).toBe(4); + }); + it('disconnects the socket only after sustained flooding past the threshold', () => { const limiter = createSocketRateLimiter(opts()); const socket = makeSocket(); diff --git a/shared/src/events.ts b/shared/src/events.ts index 84e2d95..d6eb6b9 100644 --- a/shared/src/events.ts +++ b/shared/src/events.ts @@ -8,10 +8,21 @@ import type { SocketAck } from './sfu'; export interface RoomUser extends AuthUserDto { socketId?: string } export interface ChatMessagePayload { + id: string; + kind: 'text'; sender: AuthUserDto; text: string; - ts: number; + sentAt: number; } +export type SendChatMessageAck = + | { ok: true; messageId: string } + | { + ok: false; + code: 'EMPTY_MESSAGE' | 'MESSAGE_TOO_LONG' | 'RATE_LIMITED' | 'NOT_IN_ROOM'; + message: string; + maxLength?: number; + retryAfterMs?: number; + }; export interface TranscriptState { active: boolean; startedAt: number | null; @@ -111,7 +122,7 @@ export interface ServerToClientEvents extends WebRtcServerToClientEvents { export interface RoomClientToServerEvents { 'join-room': (roomId: string) => void; 'leave-room': (roomId?: string) => void; - 'chat-message': (payload: { roomId: string; text: string }) => void; + 'chat-message': (payload: { text: string }, callback: (response: SendChatMessageAck) => void) => void; 'transcript-start': (payload: Record, callback: (response: TranscriptControlAck) => void) => void; 'transcript-stop': (payload: Record, callback: (response: TranscriptControlAck) => void) => void; 'transcript-contributor-start': (payload: Record, callback: (response: SocketAck<{ ok: true }>) => void) => void;