Compare commits

...

1 commit

Author SHA1 Message Date
Nick Tabeling
24bd34f886 feat(DIS-103): ChannelQueue serializes messages per channel (FIFO)
Some checks failed
CI / build-and-test (ubuntu-latest) (pull_request) Has been cancelled
CI / build-and-test (windows-latest) (pull_request) Has been cancelled
CI / lint (pull_request) Has been cancelled
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 <noreply@anthropic.com>
2026-04-09 17:49:02 +02:00
4 changed files with 305 additions and 48 deletions

View file

@ -2,6 +2,9 @@ import { Message, TextChannel } from "discord.js";
import { DisclawDatabase } from "./db/database"; import { DisclawDatabase } from "./db/database";
import { DisclawConfig } from "./config/loader"; import { DisclawConfig } from "./config/loader";
import { runAgent } from "./agent/runner"; import { runAgent } from "./agent/runner";
import { ChannelQueue } from "./runtime/channel-queue";
const channelQueue = new ChannelQueue();
const DISCORD_MAX_LENGTH = 2000; const DISCORD_MAX_LENGTH = 2000;
@ -77,6 +80,7 @@ export async function routeMessage(
const channel = message.channel as TextChannel; const channel = message.channel as TextChannel;
await channelQueue.enqueue(channelId, async () => {
// Show typing indicator while the agent works // Show typing indicator while the agent works
const typingInterval = setInterval(() => { const typingInterval = setInterval(() => {
channel.sendTyping().catch(() => {}); channel.sendTyping().catch(() => {});
@ -132,6 +136,7 @@ export async function routeMessage(
} finally { } finally {
clearInterval(typingInterval); clearInterval(typingInterval);
} }
});
return true; return true;
} }

View file

@ -0,0 +1,37 @@
type Task = () => Promise<void>;
/**
* 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<string, Promise<void>>();
enqueue(channelId: string, task: Task): Promise<void> {
// 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;
}
}

View file

@ -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<void>((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);
});
});

View file

@ -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<void> {
return () =>
new Promise<void>((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<void> = () =>
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<string, unknown> }).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<void>((resolve) => {
setTimeout(() => {
completionOrder.push("slow");
resolve();
}, 50);
})
);
const pFast = queue.enqueue(
"ch-fast",
() =>
new Promise<void>((resolve) => {
setTimeout(() => {
completionOrder.push("fast");
resolve();
}, 10);
})
);
await Promise.all([pSlow, pFast]);
// fast must finish before slow
expect(completionOrder).toEqual(["fast", "slow"]);
});
});