From bb286e54b558c1fbcf4f951144a50962d295bf8c Mon Sep 17 00:00:00 2001 From: Tsiry Sandratraina Date: Sat, 17 Jan 2026 11:41:45 +0300 Subject: [PATCH] Batch events for WebSocket streaming Send events as JSON arrays to improve throughput (batch size 50). Increase PAGE_SIZE to 500 and send queued events in batches. Update client to handle batched messages and reduce ping interval to 30s. Adjust progress/logging to report events every 500. --- tap/scripts/test-client.ts | 21 ++++++++--- tap/src/main.ts | 77 +++++++++++++++++++++----------------- 2 files changed, 58 insertions(+), 40 deletions(-) diff --git a/tap/scripts/test-client.ts b/tap/scripts/test-client.ts index 4e76ff78..737c108e 100644 --- a/tap/scripts/test-client.ts +++ b/tap/scripts/test-client.ts @@ -23,7 +23,7 @@ ws.onopen = () => { console.log("📤 Sending ping..."); ws.send("ping"); - // Send ping every 5 seconds + // Send ping every 30 seconds (less frequent to not interfere with fast streaming) setInterval(() => { if (ws.readyState === WebSocket.OPEN) { const now = Date.now(); @@ -33,7 +33,7 @@ ws.onopen = () => { ); ws.send("ping"); } - }, 5000); + }, 30000); }; ws.onmessage = async (event) => { @@ -43,7 +43,18 @@ ws.onmessage = async (event) => { try { const data = JSON.parse(event.data); - if (data.type === "connected") { + + // Handle batched messages (array of events) + if (Array.isArray(data)) { + messageCount += data.length; + if (messageCount % 500 === 0 || messageCount <= 50) { + console.log( + `📨 [${elapsed}s] Batch received: ${data.length} events (total: ${messageCount})`, + ); + } + } + // Handle single messages + else if (data.type === "connected") { console.log(`📨 [${elapsed}s] Connection confirmed: ${data.message}`); } else if (data.type === "heartbeat") { console.log(`💓 [${elapsed}s] Heartbeat received`); @@ -59,10 +70,10 @@ ws.onmessage = async (event) => { console.log(`📨 [${elapsed}s] Message #${messageCount}: ${event.data}`); } - if (messageCount % 100 === 0) { + if (messageCount % 500 === 0) { const rate = (messageCount / parseFloat(elapsed)).toFixed(2); console.log( - `📊 Progress: ${messageCount} messages received in ${elapsed}s (${rate} msg/s)`, + `📊 Progress: ${messageCount} events received in ${elapsed}s (${rate} events/s)`, ); } diff --git a/tap/src/main.ts b/tap/src/main.ts index 4814205d..55b151f9 100644 --- a/tap/src/main.ts +++ b/tap/src/main.ts @@ -2,14 +2,12 @@ import { ctx } from "./context.ts"; import logger from "./logger.ts"; import schema from "./schema/mod.ts"; import { asc, inArray } from "drizzle-orm"; -import { omit } from "@es-toolkit/es-toolkit/compat"; import type { SelectEvent } from "./schema/event.ts"; import { assureAdminAuth, parseTapEvent } from "@atproto/tap"; import { addToBatch, flushBatch } from "./batch.ts"; -const PAGE_SIZE = 100; -const YIELD_EVERY_N_PAGES = 5; -const YIELD_DELAY_MS = 100; +const PAGE_SIZE = 500; +const BATCH_SEND_SIZE = 50; const ADMIN_PASSWORD = Deno.env.get("TAP_ADMIN_PASSWORD")!; interface ClientState { @@ -43,13 +41,16 @@ function safeSend( return false; } +function formatEvent(evt: SelectEvent): string { + const { createdAt: _createdAt, record, ...rest } = evt; + if (record) { + return JSON.stringify({ ...rest, record: JSON.parse(record) }); + } + return JSON.stringify(rest); +} + export function broadcastEvent(evt: SelectEvent) { - const message = JSON.stringify({ - ...omit(evt, "createdAt", "record"), - ...(evt.record && { - record: JSON.parse(evt.record), - }), - }); + const message = formatEvent(evt); for (const [socket, state] of connectedClients.entries()) { if (socket.readyState === WebSocket.OPEN) { @@ -201,6 +202,8 @@ Deno.serve( logger.info`📄 Fetching page ${page}... (${totalEvents} events sent so far)`; } + // Batch send events for better performance + const batchMessages: string[] = []; for (let i = 0; i < events.length; i++) { const evt = events[i]; @@ -209,22 +212,23 @@ Deno.serve( return; } - const success = safeSend( - socket, - JSON.stringify({ - ...omit(evt, "createdAt", "record"), - ...(evt.record && { - record: JSON.parse(evt.record), - }), - }), - totalEvents, - ); - - if (success) { - totalEvents++; - } else { - logger.error`❌ Failed to send event at index ${totalEvents}, stopping pagination`; - return; + batchMessages.push(formatEvent(evt)); + + // Send batch when full or at end of page + if ( + batchMessages.length >= BATCH_SEND_SIZE || + i === events.length - 1 + ) { + const batchMessage = `[${batchMessages.join(",")}]`; + const success = safeSend(socket, batchMessage, totalEvents); + + if (success) { + totalEvents += batchMessages.length; + batchMessages.length = 0; // Clear batch + } else { + logger.error`❌ Failed to send batch at ${totalEvents}, stopping pagination`; + return; + } } } @@ -247,18 +251,21 @@ Deno.serve( if (queuedCount > 0) { logger.info`📦 Sending ${queuedCount} queued events...`; + // Batch send queued events + const queueMessages: string[] = []; for (const evt of clientState.queue) { if (socket.readyState !== WebSocket.OPEN) break; - safeSend( - socket, - JSON.stringify({ - ...omit(evt, "createdAt", "record"), - ...(evt.record && { - record: JSON.parse(evt.record), - }), - }), - ); + queueMessages.push(formatEvent(evt)); + + if (queueMessages.length >= BATCH_SEND_SIZE) { + safeSend(socket, `[${queueMessages.join(",")}]`); + queueMessages.length = 0; + } + } + + if (queueMessages.length > 0) { + safeSend(socket, `[${queueMessages.join(",")}]`); } clientState.queue = []; -- 2.51.2