Compare commits
1 commit
8fed430afc
...
24bd34f886
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
24bd34f886 |
4 changed files with 305 additions and 48 deletions
101
src/router.ts
101
src/router.ts
|
|
@ -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,61 +80,63 @@ export async function routeMessage(
|
||||||
|
|
||||||
const channel = message.channel as TextChannel;
|
const channel = message.channel as TextChannel;
|
||||||
|
|
||||||
// Show typing indicator while the agent works
|
await channelQueue.enqueue(channelId, async () => {
|
||||||
const typingInterval = setInterval(() => {
|
// Show typing indicator while the agent works
|
||||||
channel.sendTyping().catch(() => {});
|
const typingInterval = setInterval(() => {
|
||||||
}, 5000);
|
channel.sendTyping().catch(() => {});
|
||||||
// Send immediately as well
|
}, 5000);
|
||||||
await channel.sendTyping().catch(() => {});
|
// Send immediately as well
|
||||||
|
await channel.sendTyping().catch(() => {});
|
||||||
|
|
||||||
try {
|
try {
|
||||||
// Load recent conversation history for context
|
// Load recent conversation history for context
|
||||||
const history = db.getConversationHistory(workspace.id, 30);
|
const history = db.getConversationHistory(workspace.id, 30);
|
||||||
|
|
||||||
// Run the agent
|
// Run the agent
|
||||||
const response = await runAgent({
|
const response = await runAgent({
|
||||||
workspacePath: workspace.workspace_path,
|
workspacePath: workspace.workspace_path,
|
||||||
channelName: channel.name,
|
channelName: channel.name,
|
||||||
userMessage: message.content,
|
userMessage: message.content,
|
||||||
conversationHistory: history,
|
conversationHistory: history,
|
||||||
});
|
});
|
||||||
|
|
||||||
// Store user message
|
// Store user message
|
||||||
db.addConversation(
|
|
||||||
workspace.id,
|
|
||||||
message.id,
|
|
||||||
"user",
|
|
||||||
message.content,
|
|
||||||
message.author.username
|
|
||||||
);
|
|
||||||
|
|
||||||
// Split and send response
|
|
||||||
const chunks = splitMessage(response);
|
|
||||||
let firstMsgId: string | null = null;
|
|
||||||
|
|
||||||
for (const chunk of chunks) {
|
|
||||||
const sent = await channel.send(chunk);
|
|
||||||
if (!firstMsgId) {
|
|
||||||
firstMsgId = sent.id;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Store assistant response (use first message id for dedup)
|
|
||||||
if (firstMsgId) {
|
|
||||||
db.addConversation(
|
db.addConversation(
|
||||||
workspace.id,
|
workspace.id,
|
||||||
firstMsgId,
|
message.id,
|
||||||
"assistant",
|
"user",
|
||||||
response,
|
message.content,
|
||||||
null
|
message.author.username
|
||||||
);
|
);
|
||||||
|
|
||||||
|
// Split and send response
|
||||||
|
const chunks = splitMessage(response);
|
||||||
|
let firstMsgId: string | null = null;
|
||||||
|
|
||||||
|
for (const chunk of chunks) {
|
||||||
|
const sent = await channel.send(chunk);
|
||||||
|
if (!firstMsgId) {
|
||||||
|
firstMsgId = sent.id;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Store assistant response (use first message id for dedup)
|
||||||
|
if (firstMsgId) {
|
||||||
|
db.addConversation(
|
||||||
|
workspace.id,
|
||||||
|
firstMsgId,
|
||||||
|
"assistant",
|
||||||
|
response,
|
||||||
|
null
|
||||||
|
);
|
||||||
|
}
|
||||||
|
} catch (error) {
|
||||||
|
const msg = error instanceof Error ? error.message : "Unbekannter Fehler";
|
||||||
|
await channel.send(`Fehler bei der Verarbeitung: ${msg}`).catch(() => {});
|
||||||
|
} finally {
|
||||||
|
clearInterval(typingInterval);
|
||||||
}
|
}
|
||||||
} catch (error) {
|
});
|
||||||
const msg = error instanceof Error ? error.message : "Unbekannter Fehler";
|
|
||||||
await channel.send(`Fehler bei der Verarbeitung: ${msg}`).catch(() => {});
|
|
||||||
} finally {
|
|
||||||
clearInterval(typingInterval);
|
|
||||||
}
|
|
||||||
|
|
||||||
return true;
|
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