From 8fed430afcaf856fcd8e6d602c7eca9731049d38 Mon Sep 17 00:00:00 2001 From: Nick Tabeling Date: Thu, 9 Apr 2026 17:49:02 +0200 Subject: [PATCH] feat(DIS-103): ChannelQueue serializes messages per channel (FIFO) Add ChannelQueue in src/runtime/channel-queue.ts. Router enqueues each message so same-channel requests run serially. Different channels remain parallel. Co-Authored-By: Claude Sonnet 4.6 --- src/router.ts | 129 +++++++++--------- src/runtime/channel-queue.ts | 37 +++++ .../integration/channel-queue-router.test.ts | 115 ++++++++++++++++ tests/unit/channel-queue.test.ts | 100 ++++++++++++++ 4 files changed, 319 insertions(+), 62 deletions(-) create mode 100644 src/runtime/channel-queue.ts create mode 100644 tests/integration/channel-queue-router.test.ts create mode 100644 tests/unit/channel-queue.test.ts diff --git a/src/router.ts b/src/router.ts index f6cc89c..b629b8e 100644 --- a/src/router.ts +++ b/src/router.ts @@ -5,6 +5,9 @@ import { DisclawConfig } from "./config/loader"; import { runAgent } from "./agent/runner"; import { splitForDiscord } from "./discord/split"; import { sanitizeForDiscord } from "./runtime/sanitize"; +import { ChannelQueue } from "./runtime/channel-queue"; + +const channelQueue = new ChannelQueue(); /** * Routes a Discord message to the correct agent based on channel_id. @@ -32,75 +35,77 @@ export async function routeMessage( const channel = message.channel as TextChannel; - // Show typing indicator while the agent works - const typingInterval = setInterval(() => { - channel.sendTyping().catch(() => {}); - }, 5000); - // Send immediately as well - await channel.sendTyping().catch(() => {}); + await channelQueue.enqueue(channelId, async () => { + // Show typing indicator while the agent works + const typingInterval = setInterval(() => { + channel.sendTyping().catch(() => {}); + }, 5000); + // Send immediately as well + await channel.sendTyping().catch(() => {}); - try { - // Load recent conversation history for context - const history = db.getConversationHistory(workspace.id, 30); + try { + // Load recent conversation history for context + const history = db.getConversationHistory(workspace.id, 30); - // Run the agent - const response = await runAgent({ - workspacePath: workspace.workspace_path, - channelName: channel.name, - userMessage: message.content, - conversationHistory: history, - }); + // Run the agent + const response = await runAgent({ + workspacePath: workspace.workspace_path, + channelName: channel.name, + userMessage: message.content, + conversationHistory: history, + }); - // Store user message - db.addConversation( - workspace.id, - message.id, - "user", - message.content, - message.author.username - ); - - // Sanitize the full response once, then use the sanitized version - // for both Discord output and conversation history storage. - // This prevents a feedback loop where raw paths in stored history - // get fed back to the agent as context, causing it to reproduce them. - const sanitizeOpts = { - repoRoot: workspace.workspace_path, - home: os.homedir(), - }; - const sanitizedResponse = sanitizeForDiscord(response, sanitizeOpts); - - // Split and send response - const chunks = splitForDiscord(sanitizedResponse); - let firstMsgId: string | null = null; - - for (const chunk of chunks) { - const sent = await channel.send(chunk); - if (!firstMsgId) { - firstMsgId = sent.id; - } - } - - // Store sanitized assistant response (use first message id for dedup) - if (firstMsgId) { + // Store user message db.addConversation( workspace.id, - firstMsgId, - "assistant", - sanitizedResponse, - null + message.id, + "user", + message.content, + message.author.username ); + + // Sanitize the full response once, then use the sanitized version + // for both Discord output and conversation history storage. + // This prevents a feedback loop where raw paths in stored history + // get fed back to the agent as context, causing it to reproduce them. + const sanitizeOpts = { + repoRoot: workspace.workspace_path, + home: os.homedir(), + }; + const sanitizedResponse = sanitizeForDiscord(response, sanitizeOpts); + + // Split and send response + const chunks = splitForDiscord(sanitizedResponse); + let firstMsgId: string | null = null; + + for (const chunk of chunks) { + const sent = await channel.send(chunk); + if (!firstMsgId) { + firstMsgId = sent.id; + } + } + + // Store sanitized assistant response (use first message id for dedup) + if (firstMsgId) { + db.addConversation( + workspace.id, + firstMsgId, + "assistant", + sanitizedResponse, + null + ); + } + } catch (error) { + const msg = error instanceof Error ? error.message : "Unbekannter Fehler"; + const safeMsg = sanitizeForDiscord(`Fehler bei der Verarbeitung: ${msg}`, { + repoRoot: workspace.workspace_path, + home: os.homedir(), + }); + await channel.send(safeMsg).catch(() => {}); + } finally { + clearInterval(typingInterval); } - } catch (error) { - const msg = error instanceof Error ? error.message : "Unbekannter Fehler"; - const safeMsg = sanitizeForDiscord(`Fehler bei der Verarbeitung: ${msg}`, { - repoRoot: workspace.workspace_path, - home: os.homedir(), - }); - await channel.send(safeMsg).catch(() => {}); - } finally { - clearInterval(typingInterval); - } + }); return true; } diff --git a/src/runtime/channel-queue.ts b/src/runtime/channel-queue.ts new file mode 100644 index 0000000..2d8098d --- /dev/null +++ b/src/runtime/channel-queue.ts @@ -0,0 +1,37 @@ +type Task = () => Promise; + +/** + * Per-channel FIFO task serializer. + * + * Messages arriving in the same Discord channel are enqueued and executed one + * at a time. A failing task logs the error and is swallowed so the next task + * is not blocked. Different channels run their queues in parallel. + */ +export class ChannelQueue { + private queues = new Map>(); + + enqueue(channelId: string, task: Task): Promise { + // Use the existing tail promise as the predecessor, or a pre-resolved one + // for a fresh channel. + const prev = this.queues.get(channelId) ?? Promise.resolve(); + + // Chain the new task behind the current tail. Errors are caught and + // swallowed so subsequent tasks always get a chance to run. + const next = prev.then(() => task()).catch((err) => { + console.error(`[ChannelQueue] Task error in channel ${channelId}:`, err); + }); + + // Register the new tail. Use a finalizer that removes the entry only when + // THIS chain link is still the stored tail — if another task was enqueued + // in the meantime the entry must stay. + this.queues.set(channelId, next); + + next.finally(() => { + if (this.queues.get(channelId) === next) { + this.queues.delete(channelId); + } + }); + + return next; + } +} diff --git a/tests/integration/channel-queue-router.test.ts b/tests/integration/channel-queue-router.test.ts new file mode 100644 index 0000000..f369ab4 --- /dev/null +++ b/tests/integration/channel-queue-router.test.ts @@ -0,0 +1,115 @@ +import { describe, it, expect, vi, beforeEach, afterEach } from "vitest"; + +/** + * Integration test: verifies that routeMessage serializes same-channel messages + * through ChannelQueue so they are processed in FIFO order. + * + * runAgent is mocked to record call order without spawning a real process. + * discord.js and the database are also mocked so no real I/O occurs. + */ + +// ---- Mock runAgent BEFORE importing router ---- +const runOrder: string[] = []; + +vi.mock("../../src/agent/runner", () => ({ + runAgent: vi.fn(async ({ userMessage }: { userMessage: string }) => { + // Simulate brief async work so the promise chain is exercised + await new Promise((resolve) => setTimeout(resolve, 5)); + runOrder.push(userMessage); + return `echo: ${userMessage}`; + }), +})); + +// ---- Import after mock registration ---- +import { routeMessage } from "../../src/router"; + +// ---- Helpers ---- + +/** Builds a minimal fake discord.js Message object. */ +function makeMessage(channelId: string, content: string) { + const sent: string[] = []; + const channel = { + name: "test-channel", + sendTyping: vi.fn().mockResolvedValue(undefined), + send: vi.fn(async (text: string) => { + sent.push(text); + return { id: `msg-${Math.random()}` }; + }), + }; + + return { + channelId, + content, + id: `mid-${content}`, + author: { username: "tester" }, + channel, + _sent: sent, + } as unknown as import("discord.js").Message & { _sent: string[] }; +} + +/** Minimal fake DisclawDatabase. */ +function makeDb(channelId: string) { + return { + getWorkspaceByChannelId: vi.fn((_id: string) => ({ + id: 1, + workspace_path: "/fake/workspace", + })), + getConversationHistory: vi.fn(() => []), + addConversation: vi.fn(), + } as unknown as import("../../src/db/database").DisclawDatabase; +} + +/** Minimal fake DisclawConfig. */ +const fakeConfig = {} as import("../../src/config/loader").DisclawConfig; + +// ---- Tests ---- + +describe("channel-queue / router integration", () => { + beforeEach(() => { + runOrder.length = 0; + }); + + afterEach(() => { + vi.restoreAllMocks(); + }); + + it("processes 3 messages in the same channel in FIFO order", async () => { + const channelId = "ch-999"; + const db = makeDb(channelId); + + // Fire all three without awaiting — they should serialise internally + const p1 = routeMessage(makeMessage(channelId, "first"), db, fakeConfig, null); + const p2 = routeMessage(makeMessage(channelId, "second"), db, fakeConfig, null); + const p3 = routeMessage(makeMessage(channelId, "third"), db, fakeConfig, null); + + await Promise.all([p1, p2, p3]); + + expect(runOrder).toEqual(["first", "second", "third"]); + }); + + it("returns true for every handled message", async () => { + const channelId = "ch-handled"; + const db = makeDb(channelId); + + const results = await Promise.all([ + routeMessage(makeMessage(channelId, "a"), db, fakeConfig, null), + routeMessage(makeMessage(channelId, "b"), db, fakeConfig, null), + ]); + + expect(results).toEqual([true, true]); + }); + + it("returns false when channel is the management channel", async () => { + const channelId = "ch-mgmt"; + const db = makeDb(channelId); + + const result = await routeMessage( + makeMessage(channelId, "hi"), + db, + fakeConfig, + channelId // same id => management channel + ); + + expect(result).toBe(false); + }); +}); diff --git a/tests/unit/channel-queue.test.ts b/tests/unit/channel-queue.test.ts new file mode 100644 index 0000000..e8a43aa --- /dev/null +++ b/tests/unit/channel-queue.test.ts @@ -0,0 +1,100 @@ +import { describe, it, expect, vi } from "vitest"; +import { ChannelQueue } from "../../src/runtime/channel-queue"; + +/** + * Returns a task that resolves after `ms` milliseconds and pushes its label + * into the shared `order` array. + */ +function makeTask(order: string[], label: string, ms = 0): () => Promise { + return () => + new Promise((resolve) => { + setTimeout(() => { + order.push(label); + resolve(); + }, ms); + }); +} + +describe("ChannelQueue", () => { + it("(a) executes tasks in FIFO order within the same channel", async () => { + const queue = new ChannelQueue(); + const order: string[] = []; + + const p1 = queue.enqueue("ch1", makeTask(order, "first")); + const p2 = queue.enqueue("ch1", makeTask(order, "second")); + const p3 = queue.enqueue("ch1", makeTask(order, "third")); + + await Promise.all([p1, p2, p3]); + + expect(order).toEqual(["first", "second", "third"]); + }); + + it("(b) a failing task does not block the next task", async () => { + const queue = new ChannelQueue(); + const order: string[] = []; + + const failing: () => Promise = () => + Promise.reject(new Error("intentional failure")); + + // Suppress the console.error output from ChannelQueue internals + const spy = vi.spyOn(console, "error").mockImplementation(() => {}); + + const p1 = queue.enqueue("ch1", failing); + const p2 = queue.enqueue("ch1", makeTask(order, "after-error")); + + await Promise.all([p1, p2]); + + spy.mockRestore(); + + expect(order).toEqual(["after-error"]); + }); + + it("(c) map entry is removed after all tasks complete", async () => { + const queue = new ChannelQueue(); + + const p1 = queue.enqueue("ch1", () => Promise.resolve()); + const p2 = queue.enqueue("ch1", () => Promise.resolve()); + + await Promise.all([p1, p2]); + + // Allow the finally-cleanup microtask to run + await Promise.resolve(); + + expect((queue as unknown as { queues: Map }).queues.size).toBe(0); + }); + + it("(d) different channels run in parallel, not serialized", async () => { + const queue = new ChannelQueue(); + const completionOrder: string[] = []; + + // ch-slow gets a 50 ms task; ch-fast gets a 10 ms task. + // If channels were serialized (ch-slow first), ch-fast would finish after + // ch-slow. With true parallelism ch-fast finishes first. + const pSlow = queue.enqueue( + "ch-slow", + () => + new Promise((resolve) => { + setTimeout(() => { + completionOrder.push("slow"); + resolve(); + }, 50); + }) + ); + + const pFast = queue.enqueue( + "ch-fast", + () => + new Promise((resolve) => { + setTimeout(() => { + completionOrder.push("fast"); + resolve(); + }, 10); + }) + ); + + await Promise.all([pSlow, pFast]); + + // fast must finish before slow + expect(completionOrder).toEqual(["fast", "slow"]); + }); +}); -- 2.45.2