feat(discord): route guild conversations into threads

This commit is contained in:
xmanrui 2026-08-24 02:31:15 +08:00
parent 54496f9ba4
commit 47cf028914
7 changed files with 1102 additions and 164 deletions

View file

@ -15,8 +15,10 @@ import {
validDiscordToken,
} from '../../../src/channels/discord/discord-api.mjs';
import {
DiscordBotClient,
DiscordRuntime,
normalizeDiscordMessage,
resolveDiscordMessageRoute,
} from '../../../src/channels/discord/discord-runtime.mjs';
import {
DISCORD_ENDPOINTS,
@ -116,6 +118,52 @@ test('Discord API retries one rate-limited message request', async () => {
assert.equal(attempts, 2);
});
test('Discord API gets a channel and starts one thread from a source message', async () => {
const requests = [];
const api = new DiscordApi({
token: TOKEN,
fetchImpl: async (url, options) => {
requests.push({ url, options });
if (options.method === 'GET') {
return jsonResponse({ id: '222222222222222223', type: 0 });
}
return jsonResponse({
id: '111111111111111112',
type: 11,
parent_id: '222222222222222223',
});
},
});
const channel = await api.getChannel({ channelId: '222222222222222223' });
const thread = await api.startThreadFromMessage({
channelId: channel.id,
messageId: '111111111111111112',
name: 'one focused task',
});
assert.equal(channel.type, 0);
assert.equal(thread.type, 11);
assert.equal(requests[0].url.pathname, '/api/v10/channels/222222222222222223');
assert.equal(
requests[1].url.pathname,
'/api/v10/channels/222222222222222223/messages/111111111111111112/threads',
);
assert.equal(requests[1].options.method, 'POST');
assert.deepEqual(JSON.parse(requests[1].options.body), { name: 'one focused task' });
const controller = new AbortController();
const reason = new DOMException('stopped before routing', 'AbortError');
controller.abort(reason);
assert.throws(() => api.startThreadFromMessage({
channelId: channel.id,
messageId: '111111111111111113',
name: 'cancelled',
signal: controller.signal,
}), (error) => error === reason);
assert.equal(requests.length, 2);
});
test('Discord API uploads a result file as a native attachment and preserves the reply', async () => {
let request;
const api = new DiscordApi({
@ -404,6 +452,15 @@ test('Discord normalizes DMs and only addressed server messages', () => {
content: '',
}, '1234567890123456789');
assert.equal(unmentionedReply.addressed, false);
assert.equal(normalizeDiscordMessage({
id: '111111111111111114',
channel_id: '222222222222222224',
guild_id: '444444444444444444',
type: 21,
author: { id: '333333333333333334', bot: false },
content: '',
}, '1234567890123456789'), null);
});
test('Discord preserves existing channel and thread addressing before native routing', () => {
@ -465,6 +522,278 @@ test('Discord preserves existing channel and thread addressing before native rou
assert.equal(unmentionedThread.addressed, false);
});
test('Discord routes mentioned text and announcement messages into native threads', async () => {
const botId = '1234567890123456789';
for (const fixture of [
{ parentType: 0, threadType: 11, content: 'build the report' },
{ parentType: 5, threadType: 10, content: 'publish the report' },
]) {
const starts = [];
const cached = [];
const message = {
id: fixture.parentType === 0 ? '111111111111111140' : '111111111111111141',
channel_id: fixture.parentType === 0
? '222222222222222240' : '222222222222222241',
guild_id: '444444444444444444',
author: { id: '333333333333333333', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> ${fixture.content}`,
};
const route = await resolveDiscordMessageRoute(message, botId, {
api: {
async getChannel({ channelId }) {
assert.equal(channelId, message.channel_id);
return { id: channelId, type: fixture.parentType };
},
async startThreadFromMessage(options) {
starts.push(options);
return {
id: message.id,
type: fixture.threadType,
parent_id: message.channel_id,
owner_id: botId,
};
},
},
onChannel: (channel) => cached.push(channel.id),
});
assert.equal(starts.length, 1);
assert.deepEqual({
channelId: starts[0].channelId,
messageId: starts[0].messageId,
name: starts[0].name,
}, {
channelId: message.channel_id,
messageId: message.id,
name: fixture.content,
});
assert.equal(route.conversationId, message.id);
assert.deepEqual(route.replyTarget, { channelId: message.id });
assert.deepEqual(route.conversationRoute, {
peerId: message.channel_id,
threadId: message.id,
managed: true,
});
assert.equal(route.addressed, true);
assert.equal(route.requiresMention, false);
assert.deepEqual(cached, [message.channel_id, message.id]);
}
});
test('Discord keeps direct-message routing stable without a channel lookup', async () => {
const route = await resolveDiscordMessageRoute({
id: '111111111111111142',
channel_id: '222222222222222242',
author: { id: '333333333333333333', bot: false },
content: 'private task',
}, '1234567890123456789', {
api: {
async getChannel() { assert.fail('DM routing must not query a Guild channel'); },
},
});
assert.equal(route.conversationId, '222222222222222242');
assert.deepEqual(route.conversationRoute, { peerId: '222222222222222242' });
assert.deepEqual(route.replyTarget, {
channelId: '222222222222222242',
replyToMessageId: '111111111111111142',
});
});
test('Discord reuses all native thread types and only relaxes mentions for bot-owned threads', async () => {
const botId = '1234567890123456789';
for (const type of [10, 11, 12]) {
for (const managed of [false, true]) {
const channelId = `2222222222222222${type}${managed ? 1 : 0}`;
const message = {
id: `1111111111111112${type}${managed ? 1 : 0}`,
channel_id: channelId,
guild_id: '444444444444444444',
author: { id: '333333333333333333', bot: false },
mentions: [],
content: 'continue',
};
const route = await resolveDiscordMessageRoute(message, botId, {
channel: {
id: channelId,
type,
parent_id: '222222222222222299',
owner_id: managed ? botId : '999999999999999999',
},
api: {
async getChannel() { assert.fail('cached thread must be reused'); },
async startThreadFromMessage() { assert.fail('must not create a nested thread'); },
},
});
assert.equal(route.conversationId, channelId);
assert.equal(route.addressed, managed);
assert.equal(route.requiresMention, !managed);
assert.deepEqual(route.replyTarget, {
channelId,
replyToMessageId: message.id,
});
}
}
});
test('Discord recovers an already-created thread after an uncertain or duplicate create result', async () => {
const botId = '1234567890123456789';
const message = {
id: '111111111111111150',
channel_id: '222222222222222250',
guild_id: '444444444444444444',
author: { id: '333333333333333333', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> recover`,
};
let starts = 0;
const route = await resolveDiscordMessageRoute(message, botId, {
channel: { id: message.channel_id, type: 0 },
api: {
async getChannel({ channelId }) {
assert.equal(channelId, message.id);
return {
id: message.id,
type: 11,
parent_id: message.channel_id,
owner_id: botId,
};
},
async startThreadFromMessage() {
starts += 1;
const error = new Error('request outcome unknown');
error.status = 500;
throw error;
},
},
});
assert.equal(starts, 1);
assert.equal(route.conversationRoute.threadId, message.id);
assert.equal(route.conversationRoute.managed, true);
assert.deepEqual(route.replyTarget, { channelId: message.id });
});
test('Discord falls back before routing when the parent is unsupported or thread creation is rejected', async () => {
const botId = '1234567890123456789';
const base = {
guild_id: '444444444444444444',
author: { id: '333333333333333333', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> fallback`,
};
const unsupported = await resolveDiscordMessageRoute({
...base,
id: '111111111111111160',
channel_id: '222222222222222260',
}, botId, {
channel: { id: '222222222222222260', type: 15 },
api: {
async getChannel() { assert.fail('cached forum channel must be reused'); },
async startThreadFromMessage() { assert.fail('Forum must not use Start Thread from Message'); },
},
});
assert.equal(unsupported.conversationId, '222222222222222260');
assert.equal(unsupported.conversationRoute.fallback, 'unsupported-channel');
assert.equal(
unsupported.replyTarget.notice,
'当前频道不支持自动创建 Thread,已直接在当前频道回复。',
);
const rejected = await resolveDiscordMessageRoute({
...base,
id: '111111111111111161',
channel_id: '222222222222222261',
}, botId, {
channel: { id: '222222222222222261', type: 0 },
api: {
async getChannel({ channelId }) {
assert.equal(channelId, '111111111111111161');
const error = new Error('Unknown Channel');
error.status = 404;
throw error;
},
async startThreadFromMessage() {
const error = new Error('Missing Permissions');
error.status = 403;
throw error;
},
},
});
assert.equal(rejected.conversationId, '222222222222222261');
assert.equal(rejected.conversationRoute.fallback, 'thread-create-failed');
assert.equal(rejected.replyTarget.notice, '无法创建 Thread,已直接在当前频道回复。');
});
test('Discord never runs a parent-channel fallback when thread creation remains uncertain', async () => {
const botId = '1234567890123456789';
const notices = [];
const message = {
id: '111111111111111162',
channel_id: '222222222222222262',
guild_id: '444444444444444444',
author: { id: '333333333333333333', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> uncertain`,
};
await assert.rejects(() => resolveDiscordMessageRoute(message, botId, {
channel: { id: message.channel_id, type: 0 },
api: {
async startThreadFromMessage() { throw new TypeError('socket reset after dispatch'); },
async getChannel({ channelId }) {
assert.equal(channelId, message.id);
const error = new Error('Unknown Channel');
error.status = 404;
throw error;
},
async createMessage(options) {
notices.push(options);
return { id: '777777777777777773' };
},
},
}), (error) => error.code === 'discord-thread-create-uncertain');
assert.equal(notices.length, 1);
assert.equal(notices[0].channelId, message.channel_id);
assert.equal(notices[0].replyToMessageId, message.id);
assert.match(notices[0].content, /结果暂时无法确认/);
});
test('Discord merges a deterministic Thread fallback notice into one delivered answer', async () => {
const creates = [];
const edits = [];
const client = new DiscordBotClient({
api: {
async createMessage(options) {
creates.push(options);
return { id: `88888888888888888${creates.length}` };
},
async editMessage(options) {
edits.push(options);
return { id: options.messageId };
},
},
});
const target = {
channelId: '222222222222222263',
replyToMessageId: '111111111111111163',
notice: '无法创建 Thread,已直接在当前频道回复。',
};
await client.sendText(target, 'final answer');
assert.equal(creates.length, 1);
assert.equal(creates[0].content, `${target.notice}\n\nfinal answer`);
await client.sendText(target, 'follow-up notice');
assert.equal(creates[1].content, 'follow-up notice');
const streamTarget = {
channelId: '222222222222222264',
notice: '当前频道不支持自动创建 Thread,已直接在当前频道回复。',
};
const stream = await client.openStream(streamTarget);
await stream.finish('streamed answer');
assert.equal(creates[2].content, `${streamTarget.notice}\n\n正在处理…`);
assert.equal(edits[0].content, `${streamTarget.notice}\n\nstreamed answer`);
});
test('Discord keeps image attachments in images and exposes ordinary attachments as files', async () => {
const png = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]);
const calls = [];
@ -575,6 +904,329 @@ class FakeSocket {
}
}
test('Discord runtime keeps one isolated Session per managed thread and reuses it without mentions', async () => {
const botId = '1234567890123456789';
const parentChannelId = '222222222222222270';
const unrelatedThreadId = '222222222222222271';
const sessions = new Map();
const seen = new Set();
const starts = [];
const asks = [];
const deliveries = [];
let nextSession = 1;
let nextProviderMessage = 1;
let socket;
const api = {
getCurrentUser: async () => ({ id: botId, bot: true }),
getGatewayBot: async () => ({ url: 'wss://gateway.discord.gg' }),
async getChannel({ channelId }) {
assert.fail(`cached channel expected, got API lookup for ${channelId}`);
},
async startThreadFromMessage({ channelId, messageId, name }) {
starts.push({ channelId, messageId, name });
await new Promise((resolve) => setImmediate(resolve));
return {
id: messageId,
type: 11,
parent_id: channelId,
owner_id: botId,
};
},
async sendTyping({ channelId }) {
deliveries.push({ operation: 'typing', channelId });
},
async createMessage({ channelId, content, replyToMessageId }) {
const id = `8888888888888888${String(nextProviderMessage).padStart(2, '0')}`;
nextProviderMessage += 1;
deliveries.push({ operation: 'create', channelId, content, replyToMessageId, id });
return { id };
},
async editMessage({ channelId, messageId, content }) {
deliveries.push({ operation: 'edit', channelId, messageId, content });
return { id: messageId };
},
};
const harness = {
ensureRunning: async () => true,
async createSession() {
const sessionId = `session-${nextSession}`;
nextSession += 1;
return sessionId;
},
sessionExists: async () => true,
async ask(sessionId, text, options) {
asks.push({ sessionId, text, files: options?.files?.length ?? 0 });
return `answer:${text}`;
},
};
const state = {
sessionFor: (key) => sessions.get(key) ?? null,
async setSession(key, sessionId) { sessions.set(key, sessionId); },
async clearSession(key) { sessions.delete(key); },
hasSeen: (messageId) => seen.has(messageId),
async markSeen(messageId) { seen.add(messageId); },
};
const runtime = new DiscordRuntime({
config: { botId: 'discord_test', platformId: botId, name: 'Harness Discord' },
token: TOKEN,
harness,
state,
createApi: () => api,
createWebSocket: () => {
socket = new FakeSocket();
queueMicrotask(() => socket.emit('message', {
data: JSON.stringify({ op: 10, d: { heartbeat_interval: 45_000 } }),
}));
return socket;
},
random: () => 0.5,
logger: { warn() {}, error(...args) { assert.fail(args.join(' ')); } },
});
await runtime.start();
socket.emit('message', {
data: JSON.stringify({
op: 0,
t: 'GUILD_CREATE',
s: 2,
d: {
id: '444444444444444444',
channels: [{ id: parentChannelId, type: 0 }],
threads: [{
id: unrelatedThreadId,
type: 11,
parent_id: parentChannelId,
owner_id: '999999999999999999',
}],
},
}),
});
const firstMessage = {
id: '111111111111111170',
channel_id: parentChannelId,
guild_id: '444444444444444444',
author: { id: '333333333333333330', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> first task`,
};
const firstPacket = JSON.stringify({ op: 0, t: 'MESSAGE_CREATE', s: 3, d: firstMessage });
socket.emit('message', { data: firstPacket });
socket.emit('message', { data: firstPacket });
await eventually(() => runtime.status.messagesReplied === 1);
assert.equal(starts.length, 1);
assert.equal(sessions.has(`group:${parentChannelId}`), false);
assert.equal(sessions.get(`group:${firstMessage.id}`), 'session-1');
assert.deepEqual(asks, [{ sessionId: 'session-1', text: 'first task', files: 0 }]);
assert.equal(deliveries.every((entry) => entry.channelId === firstMessage.id), true);
assert.equal(deliveries.find((entry) => entry.operation === 'create').replyToMessageId, undefined);
socket.emit('message', {
data: JSON.stringify({
op: 0,
t: 'MESSAGE_CREATE',
s: 4,
d: {
id: '111111111111111171',
channel_id: firstMessage.id,
guild_id: '444444444444444444',
author: { id: '333333333333333330', bot: false },
mentions: [],
content: 'continue task',
attachments: [{
id: '555555555555555571',
filename: 'notes.txt',
content_type: 'text/plain',
size: 4,
url: 'https://cdn.discordapp.com/attachments/222/555/notes.txt',
}],
},
}),
});
await eventually(() => runtime.status.messagesReplied === 2);
assert.deepEqual(asks[1], { sessionId: 'session-1', text: 'continue task', files: 1 });
const statusMessageId = '111111111111111179';
socket.emit('message', {
data: JSON.stringify({
op: 0,
t: 'MESSAGE_CREATE',
s: 5,
d: {
id: statusMessageId,
channel_id: firstMessage.id,
guild_id: '444444444444444444',
author: { id: '333333333333333330', bot: false },
mentions: [],
content: '/status',
},
}),
});
await eventually(() => seen.has(statusMessageId));
assert.equal(asks.length, 2);
assert.equal(deliveries.some((entry) => entry.operation === 'create'
&& entry.channelId === firstMessage.id
&& entry.replyToMessageId === statusMessageId
&& /连接正常/.test(entry.content)), true);
socket.emit('message', {
data: JSON.stringify({
op: 0,
t: 'MESSAGE_CREATE',
s: 6,
d: {
id: '111111111111111172',
channel_id: unrelatedThreadId,
guild_id: '444444444444444444',
author: { id: '333333333333333331', bot: false },
mentions: [],
content: 'must be ignored',
},
}),
});
await eventually(() => seen.has('111111111111111172'));
assert.equal(runtime.status.messagesRejected, 1);
assert.equal(asks.length, 2);
socket.emit('message', {
data: JSON.stringify({
op: 0,
t: 'MESSAGE_CREATE',
s: 7,
d: {
id: '111111111111111173',
channel_id: unrelatedThreadId,
guild_id: '444444444444444444',
author: { id: '333333333333333331', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> explicit thread task`,
},
}),
});
await eventually(() => runtime.status.messagesReplied === 3);
assert.equal(sessions.get(`group:${unrelatedThreadId}`), 'session-2');
assert.deepEqual(asks[2], { sessionId: 'session-2', text: 'explicit thread task', files: 0 });
for (const [index, senderId] of ['333333333333333332', '333333333333333333'].entries()) {
socket.emit('message', {
data: JSON.stringify({
op: 0,
t: 'MESSAGE_CREATE',
s: 8 + index,
d: {
id: `11111111111111118${index}`,
channel_id: parentChannelId,
guild_id: '444444444444444444',
author: { id: senderId, bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> user ${index}`,
},
}),
});
}
await eventually(() => runtime.status.messagesReplied === 5);
assert.equal(starts.length, 3);
assert.notEqual(
sessions.get('group:111111111111111180'),
sessions.get('group:111111111111111181'),
);
assert.equal(sessions.has(`group:${parentChannelId}`), false);
await runtime.stop();
});
test('Discord runtime records one uncertain Thread result and suppresses Gateway replays', async () => {
const botId = '1234567890123456789';
const parentChannelId = '222222222222222290';
const messageId = '111111111111111190';
const seen = new Set();
const notices = [];
const errors = [];
let starts = 0;
let socket;
const runtime = new DiscordRuntime({
config: { botId: 'discord_test', platformId: botId, name: 'Harness Discord' },
token: TOKEN,
harness: {
ensureRunning: async () => true,
createSession: async () => assert.fail('uncertain routing must not create a Session'),
ask: async () => assert.fail('uncertain routing must not run a Prompt'),
},
state: {
hasSeen: (id) => seen.has(id),
async markSeen(id) { seen.add(id); },
},
createApi: () => ({
getCurrentUser: async () => ({ id: botId, bot: true }),
getGatewayBot: async () => ({ url: 'wss://gateway.discord.gg' }),
async startThreadFromMessage() {
starts += 1;
throw new TypeError('connection closed after dispatch');
},
async getChannel({ channelId }) {
assert.equal(channelId, messageId);
const error = new Error('Unknown Channel');
error.status = 404;
throw error;
},
async createMessage(options) {
notices.push(options);
return { id: '777777777777777790' };
},
}),
createWebSocket: () => {
socket = new FakeSocket();
queueMicrotask(() => socket.emit('message', {
data: JSON.stringify({ op: 10, d: { heartbeat_interval: 45_000 } }),
}));
return socket;
},
random: () => 0.5,
logger: { warn() {}, error(...args) { errors.push(args); } },
});
await runtime.start();
socket.emit('message', {
data: JSON.stringify({
op: 0,
t: 'GUILD_CREATE',
s: 2,
d: {
id: '444444444444444444',
channels: [{ id: parentChannelId, type: 0 }],
threads: [],
},
}),
});
const packet = {
op: 0,
t: 'MESSAGE_CREATE',
s: 3,
d: {
id: messageId,
channel_id: parentChannelId,
guild_id: '444444444444444444',
author: { id: '333333333333333333', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> uncertain task`,
},
};
socket.emit('message', { data: JSON.stringify(packet) });
socket.emit('message', { data: JSON.stringify(packet) });
await eventually(() => seen.has(messageId));
assert.equal(starts, 1);
assert.equal(notices.length, 1);
assert.match(notices[0].content, /结果暂时无法确认/);
assert.equal(runtime.status.messagesReceived, 0);
assert.equal(runtime.status.messagesReplied, 0);
socket.emit('message', { data: JSON.stringify({ ...packet, s: 4 }) });
await new Promise((resolve) => setImmediate(resolve));
assert.equal(starts, 1);
assert.equal(notices.length, 1);
assert.ok(errors.length >= 1);
await runtime.stop();
});
test('Discord runtime identifies on Gateway v10 and becomes ready', async () => {
let socket;
const abortMark = deferred();

View file

@ -1414,6 +1414,49 @@ test('a group question only accepts an addressed reply from the initiating actor
assert.equal(bridge.status.messagesRejected, 1);
});
test('a managed native thread can keep group isolation without asking for another mention', async () => {
const fixture = stateFixture({ 'group:managed-thread': 'session-managed' });
const sent = [];
const submitted = deferred();
const bridge = createBridge({
state: fixture.state,
bot: { sendText: async (target, text) => sent.push({ target, text }) },
harness: {
sessionExists: async () => true,
ask: async (sessionId, _text, options) => {
await options.onInteraction(questionInteraction({
id: 'managed-thread-question',
sessionId,
questions: [{ id: 'choice', question: '请选择下一步' }],
respond: async (result) => {
submitted.resolve(result);
return { accepted: true };
},
}));
await submitted.promise;
return '已继续';
},
},
});
const route = {
kind: 'group',
conversationId: 'managed-thread',
addressed: true,
requiresMention: false,
};
const processing = bridge.accept(message('managed-start', '开始', route));
await eventually(() => sent.some(({ text }) => text.includes('请选择下一步')));
assert.doesNotMatch(sent[0].text, /群聊中请 @机器人/);
await bridge.accept(message('managed-answer', '继续', route));
assert.deepEqual((await submitted.promise).value.answer.answers, [{
id: 'choice',
selected: [],
custom: '继续',
}]);
await processing;
});
test('deduplicates replays and safely closes recovered questions and approvals', async () => {
const fixture = stateFixture();
const sent = [];