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 <noreply@anthropic.com>
This commit is contained in:
parent
6039e992bd
commit
24bd34f886
4 changed files with 305 additions and 48 deletions
|
|
@ -2,6 +2,9 @@ import { Message, TextChannel } from "discord.js";
|
|||
import { DisclawDatabase } from "./db/database";
|
||||
import { DisclawConfig } from "./config/loader";
|
||||
import { runAgent } from "./agent/runner";
|
||||
import { ChannelQueue } from "./runtime/channel-queue";
|
||||
|
||||
const channelQueue = new ChannelQueue();
|
||||
|
||||
const DISCORD_MAX_LENGTH = 2000;
|
||||
|
||||
|
|
@ -77,6 +80,7 @@ export async function routeMessage(
|
|||
|
||||
const channel = message.channel as TextChannel;
|
||||
|
||||
await channelQueue.enqueue(channelId, async () => {
|
||||
// Show typing indicator while the agent works
|
||||
const typingInterval = setInterval(() => {
|
||||
channel.sendTyping().catch(() => {});
|
||||
|
|
@ -132,6 +136,7 @@ export async function routeMessage(
|
|||
} finally {
|
||||
clearInterval(typingInterval);
|
||||
}
|
||||
});
|
||||
|
||||
return true;
|
||||
}
|
||||
|
|
|
|||
37
src/runtime/channel-queue.ts
Normal file
37
src/runtime/channel-queue.ts
Normal 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;
|
||||
}
|
||||
}
|
||||
115
tests/integration/channel-queue-router.test.ts
Normal file
115
tests/integration/channel-queue-router.test.ts
Normal 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);
|
||||
});
|
||||
});
|
||||
100
tests/unit/channel-queue.test.ts
Normal file
100
tests/unit/channel-queue.test.ts
Normal 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"]);
|
||||
});
|
||||
});
|
||||
Loading…
Reference in a new issue