merge: add native Discord threads

This commit is contained in:
xmanrui 2026-08-24 03:43:54 +08:00
commit 37992b148a
8 changed files with 1369 additions and 168 deletions

View file

@ -52,7 +52,7 @@ Connect IM bots to DeepSeek Harness by scanning a QR code, using an App Manifest
| QQ | Create a bot with mobile QQ QR scanning, or bind one with AppID + AppSecret | WebSocket connection; private chats show typing and receive one Markdown reply, while mentioned group chats receive only the final answer |
| Slack | Create an app from the bundled App Manifest, then enter a Bot Token (`xoxb-`) and App Token (`xapp-`) | Socket Mode connection; direct DM replies, mention-only channel replies, and preferred native streaming API |
| Telegram | Enter a Bot Token generated by @BotFather | Bot API long polling; DMs work by default and groups respond to mentions or replies, while each bot can optionally enable a private-DM allowlist; streaming uses message edits |
| Discord | Enter a Bot Token generated in the Developer Portal | Gateway v10 connection; direct DM replies, mention-only server replies, and streaming through message edits |
| Discord | Enter a Bot Token generated in the Developer Portal | Gateway v10 connection; direct DM replies; the first mention in a server text or announcement channel creates a native Thread, where follow-up messages no longer need to mention the bot; replies stream through message edits |
| WhatsApp | Scan a QR code with mobile WhatsApp to link a device | WhatsApp Web connection; self-chat only by default, with optional selected-contact and open-response modes; read receipt and typing indicator followed by the final answer |
Other IM platforms can be added through the same channel-adapter structure.
@ -74,7 +74,7 @@ After the model calls the file-return tool, the plugin hands the specified file
| QQ | The bot needs file-message capability and remains subject to QQ's daily upload quota; the bot reports when the quota is exhausted. |
| Slack | The Bot Token needs `files:write`; the Workspace's current policy determines the actual size limit. After changing scopes, re-authorize/reinstall the App and reconnect the bot. |
| Telegram | The bot must be allowed to send documents in the current chat; the Bot API response determines the actual range. |
| Discord | The bot needs **Send Messages**, **Attach Files**, and **Read Message History**. The current account and server capability determine the actual attachment allowance. |
| Discord | Enable **Message Content Intent** in the Developer Portal. The bot needs **Send Messages**, **Create Public Threads**, **Send Messages in Threads**, and **Read Message History**; result-file delivery also requires **Attach Files**. The current account and server capability determine the actual attachment allowance. |
| WhatsApp | The linked session must support Document Messages; the WhatsApp/Baileys response determines the actual range. |
## AI Office Connector

View file

@ -55,7 +55,7 @@ Connect IM bots to DeepSeek Harness by scanning a QR code, using an App Manifest
| QQ | 使用手机 QQ 扫码创建机器人,或使用 AppID + AppSecret 手动绑定 | WebSocket 长连接;私聊显示“正在输入”并以单条 Markdown 回复,群聊被 @ 后只发送最终答案 |
| Slack | 使用预置 App Manifest 创建应用,再填写 Bot Token(`xoxb-`)和 App Token(`xapp-`) | Socket Mode 长连接;私聊直接回复,频道被 @ 后响应,优先使用官方流式消息 API |
| Telegram | 使用 @BotFather 生成的 Bot Token | Bot API 长轮询;默认私聊直接响应、群聊被提及或回复时响应,也可为每个机器人独立启用私聊白名单安全模式;通过编辑消息流式显示回答 |
| Discord | 使用 Developer Portal 生成的 Bot Token | Gateway v10 长连接;私信直接回复,服务器频道被提及时响应,通过编辑消息流式显示回答 |
| Discord | 使用 Developer Portal 生成的 Bot Token | Gateway v10 长连接;私信直接回复;服务器文字/公告频道首次 @ 后创建原生 Thread,后续在线程中无需重复 @,并通过编辑消息流式显示回答 |
| WhatsApp | 使用手机 WhatsApp 扫码关联设备 | WhatsApp Web 长连接;默认仅响应账号自聊,也可切换到指定联系人或开放响应模式;显示已读和“正在输入”,再发送最终回答 |
其他 IM 平台可继续按同一渠道适配器结构接入。
@ -77,7 +77,7 @@ Connect IM bots to DeepSeek Harness by scanning a QR code, using an App Manifest
| QQ | 机器人需具备文件消息能力,并受 QQ 当日文件上传配额约束;额度耗尽时会明确提示稍后重试。 |
| Slack | Bot Token 需有 `files:write`;实际大小上限由 Workspace 当前策略决定。已有 App 新增或变更 Scope 后,必须重新授权/安装 App 并重新连接机器人。 |
| Telegram | 机器人必须能在当前聊天发送文档,实际可发送范围以 Bot API 返回为准。 |
| Discord | 机器人需有 **Send Messages**、**Attach Files** 和 **Read Message History** 权限;实际附件额度由当前账号与服务器能力决定。 |
| Discord | Developer Portal 的 Bot 设置中需启用 **Message Content Intent**;机器人需有 **Send Messages**、**Create Public Threads**、**Send Messages in Threads** 和 **Read Message History** 权限;发送结果文件还需 **Attach Files**。实际附件额度由当前账号与服务器能力决定。 |
| WhatsApp | 当前绑定会话需支持 Document Message,实际可发送范围以 WhatsApp/Baileys 返回为准。 |
## AI Office Connector

File diff suppressed because one or more lines are too long

View file

@ -120,6 +120,29 @@ export class DiscordApi {
return this.#request('gateway/bot', { ...options, method: 'GET' });
}
getChannel({ channelId, signal } = {}) {
return this.#request(`channels/${snowflake(channelId, 'channel id')}`, {
method: 'GET',
signal,
});
}
startThreadFromMessage({ channelId, messageId, name, signal } = {}) {
const threadName = cleanString(name);
if (!threadName || [...threadName].length > 100) {
throw new TypeError('Discord thread name must contain 1-100 characters');
}
if (signal?.aborted) throw abortReason(signal);
return this.#request(
`channels/${snowflake(channelId, 'channel id')}/messages/${snowflake(messageId, 'message id')}/threads`,
{
method: 'POST',
signal,
body: { name: threadName },
},
);
}
createMessage({ channelId, content, replyToMessageId, signal }) {
return this.#request(`channels/${snowflake(channelId, 'channel id')}/messages`, {
method: 'POST',

View file

@ -4,8 +4,15 @@ import { fetchImageBuffer } from '../shared/image-prompt.mjs';
import { DiscordApi } from './discord-api.mjs';
import { createDiscordBridgeStatus, DiscordHarnessBridge } from './discord-bridge.mjs';
const DISCORD_GATEWAY_INTENTS = (1 << 0) | (1 << 9) | (1 << 12);
const DISCORD_GATEWAY_INTENTS = (1 << 0) | (1 << 9) | (1 << 12) | (1 << 15);
const THREAD_RECOVERY_TIMEOUT_MS = 5_000;
const RECONNECT_DELAYS_MS = Object.freeze([1_000, 3_000, 5_000, 10_000, 30_000]);
const GUILD_TEXT = 0;
const GUILD_ANNOUNCEMENT = 5;
const ANNOUNCEMENT_THREAD = 10;
const PUBLIC_THREAD = 11;
const PRIVATE_THREAD = 12;
const THREAD_TYPES = new Set([ANNOUNCEMENT_THREAD, PUBLIC_THREAD, PRIVATE_THREAD]);
const IMAGE_MEDIA_TYPES = new Set(['image/jpeg', 'image/png', 'image/webp', 'image/gif']);
const IMAGE_FILE_TYPES = new Map([
['.jpg', 'image/jpeg'],
@ -59,6 +66,94 @@ function stripBotMention(text, botId) {
return text.replace(new RegExp(`<@!?${botId}>`, 'g'), '').trim();
}
function cleanThreadName(message, botId) {
const name = stripBotMention(message?.content ?? '', botId).replace(/\s+/g, ' ').trim()
|| 'DeepSeek Harness';
return [...name].slice(0, 100).join('');
}
function isThreadChannel(channel) {
return THREAD_TYPES.has(Number(channel?.type));
}
function isThreadFromMessage(channel, message) {
return isThreadChannel(channel)
&& String(channel?.id ?? '') === String(message?.id ?? '')
&& String(channel?.parent_id ?? '') === String(message?.channel_id ?? '');
}
function rememberChannel(channel, callback) {
if (channel?.id) callback?.(channel);
return channel;
}
function withConversationRoute(normalized, channel, botId, {
fromSourceMessage = false,
fallback = null,
notice = null,
} = {}) {
const channelId = String(channel?.id ?? normalized.conversationId);
const thread = isThreadChannel(channel);
const managed = thread && String(channel?.owner_id ?? '') === String(botId);
const parentId = thread && channel?.parent_id ? String(channel.parent_id) : channelId;
return {
...normalized,
conversationId: channelId,
addressed: normalized.addressed || managed,
requiresMention: thread ? !managed : normalized.kind === 'group',
replyTarget: {
channelId,
...(!fromSourceMessage && normalized.replyTarget?.replyToMessageId
? { replyToMessageId: normalized.replyTarget.replyToMessageId }
: {}),
...(notice ? { notice } : {}),
},
connectionTestTarget: { channelId },
conversationRoute: {
peerId: parentId,
...(thread ? { threadId: channelId, managed } : {}),
...(fallback ? { fallback } : {}),
},
};
}
function uncertainThreadCreate(error) {
const status = Number(error?.status);
return !Number.isInteger(status) || status >= 500;
}
function threadCreateUncertain(cause) {
const error = new Error('Discord Thread creation result is uncertain', { cause });
error.code = 'discord-thread-create-uncertain';
return error;
}
async function findCreatedThread(api, message, { signal, onChannel } = {}) {
try {
const channel = rememberChannel(await api.getChannel({
channelId: String(message.id),
signal,
}), onChannel);
return isThreadFromMessage(channel, message) ? channel : null;
} catch (error) {
if (error?.status === 404) return null;
throw error;
}
}
async function sendThreadUncertainNotice(api, normalized, signal) {
try {
await api.createMessage({
channelId: normalized.replyTarget.channelId,
replyToMessageId: normalized.replyTarget.replyToMessageId,
content: 'Thread 创建结果暂时无法确认。若已创建,请在对应 Thread 中重试;若未创建,请稍后重新 @机器人。',
signal,
});
} catch (error) {
if (signal?.aborted) throw signal.reason ?? error;
}
}
function attachmentMediaType(attachment) {
const value = typeof attachment?.content_type === 'string'
? attachment.content_type.split(';', 1)[0].trim().toLowerCase() : '';
@ -109,7 +204,8 @@ function discordFileSource(attachment, fetchImpl) {
}
export function normalizeDiscordMessage(message, botId, { fetchImpl = fetch } = {}) {
if (!message?.id || !message?.channel_id || !message?.author?.id) return null;
if (!message?.id || !message?.channel_id || !message?.author?.id
|| Number(message.type) === 21) return null;
const direct = !message.guild_id;
const addressed = direct
|| message.mentions?.some((mention) => String(mention?.id) === String(botId));
@ -135,9 +231,84 @@ export function normalizeDiscordMessage(message, botId, { fetchImpl = fetch } =
};
}
export async function resolveDiscordMessageRoute(message, botId, {
api,
channel,
fetchImpl = fetch,
signal,
onChannel,
} = {}) {
const normalized = normalizeDiscordMessage(message, botId, { fetchImpl });
if (!normalized || normalized.senderIsBot) return normalized;
signal?.throwIfAborted();
if (normalized.kind === 'direct') {
return withConversationRoute(normalized, { id: normalized.conversationId, type: 1 }, botId);
}
if (!api || typeof api.getChannel !== 'function') {
throw new TypeError('Discord route resolution requires the Discord API');
}
const sourceChannel = rememberChannel(channel ?? await api.getChannel({
channelId: normalized.conversationId,
signal,
}), onChannel);
if (isThreadChannel(sourceChannel)) {
return withConversationRoute(normalized, sourceChannel, botId);
}
if (!normalized.addressed) return withConversationRoute(normalized, sourceChannel, botId);
const sourceType = Number(sourceChannel?.type);
if (sourceType !== GUILD_TEXT && sourceType !== GUILD_ANNOUNCEMENT) {
return withConversationRoute(normalized, sourceChannel, botId, {
fallback: 'unsupported-channel',
notice: '当前频道不支持自动创建 Thread,已直接在当前频道回复。',
});
}
signal?.throwIfAborted();
let created;
try {
created = rememberChannel(await api.startThreadFromMessage({
channelId: normalized.conversationId,
messageId: normalized.messageId,
name: cleanThreadName(message, botId),
signal,
}), onChannel);
if (!isThreadFromMessage(created, message)) {
throw new Error('Discord returned an invalid thread for the source message');
}
} catch (error) {
let recovered;
try {
recovered = await findCreatedThread(api, message, {
signal: AbortSignal.timeout(THREAD_RECOVERY_TIMEOUT_MS),
onChannel,
});
} catch (recoveryError) {
if (signal?.aborted) throw signal.reason ?? error;
await sendThreadUncertainNotice(api, normalized, signal);
throw threadCreateUncertain(recoveryError);
}
if (recovered) {
return withConversationRoute(normalized, recovered, botId, { fromSourceMessage: true });
}
if (signal?.aborted) throw signal.reason ?? error;
if (uncertainThreadCreate(error)) {
await sendThreadUncertainNotice(api, normalized, signal);
throw threadCreateUncertain(error);
}
return withConversationRoute(normalized, sourceChannel, botId, {
fallback: 'thread-create-failed',
notice: '无法创建 Thread,已直接在当前频道回复。',
});
}
return withConversationRoute(normalized, created, botId, { fromSourceMessage: true });
}
export class DiscordBotClient {
#api;
#signal;
#deliveredNotices = new WeakSet();
constructor({ api, signal }) {
this.#api = api;
@ -145,7 +316,9 @@ export class DiscordBotClient {
}
async sendText(target, text) {
const chunks = splitMessageText(text, 1_900);
const notice = !this.#deliveredNotices.has(target) && target?.notice
? String(target.notice) : null;
const chunks = splitMessageText(notice ? `${notice}\n\n${text}` : text, 1_900);
const providerMessageIds = [];
for (const [index, chunk] of chunks.entries()) {
const result = await this.#api.createMessage({
@ -154,6 +327,7 @@ export class DiscordBotClient {
replyToMessageId: index === 0 ? target.replyToMessageId : undefined,
signal: this.#signal,
});
if (notice && index === 0) this.#deliveredNotices.add(target);
if (typeof result?.id === 'string' && result.id) providerMessageIds.push(result.id);
}
return { providerMessageIds };
@ -173,21 +347,25 @@ export class DiscordBotClient {
}
async openStream(target) {
const notice = !this.#deliveredNotices.has(target) && target?.notice
? String(target.notice) : null;
const decorate = (content) => notice ? `${notice}\n\n${content}` : content;
const stream = createEditableMessageStream({
limit: 1_900,
limit: notice ? 1_800 : 1_900,
create: async (content) => {
const message = await this.#api.createMessage({
channelId: target.channelId,
content,
content: decorate(content),
replyToMessageId: target.replyToMessageId,
signal: this.#signal,
});
if (notice) this.#deliveredNotices.add(target);
return message.id;
},
edit: (messageId, content) => this.#api.editMessage({
channelId: target.channelId,
messageId,
content,
content: decorate(content),
signal: this.#signal,
}),
sendRemainder: (content) => this.#api.createMessage({
@ -241,6 +419,8 @@ export class DiscordRuntime {
#generation = 0;
#stopped = true;
#starting = null;
#channels = new Map();
#routing = new Map();
constructor({
config,
@ -299,6 +479,8 @@ export class DiscordRuntime {
this.#resumeUrl = null;
this.#sequence = null;
this.#reconnectAttempt = 0;
this.#channels.clear();
this.#routing.clear();
this.#status.startedAt = new Date().toISOString();
this.#status.connectionState = 'connecting';
this.#status.lastError = null;
@ -432,11 +614,21 @@ export class DiscordRuntime {
markReady();
} else if (packet.t === 'RESUMED') {
markReady();
} else if (packet.t === 'GUILD_CREATE') {
for (const channel of [...(packet.d?.channels ?? []), ...(packet.d?.threads ?? [])]) {
this.#rememberChannel(channel);
}
} else if (packet.t === 'CHANNEL_CREATE' || packet.t === 'CHANNEL_UPDATE'
|| packet.t === 'THREAD_CREATE' || packet.t === 'THREAD_UPDATE') {
this.#rememberChannel(packet.d);
} else if (packet.t === 'CHANNEL_DELETE' || packet.t === 'THREAD_DELETE') {
if (packet.d?.id) this.#channels.delete(String(packet.d.id));
} else if (packet.t === 'THREAD_LIST_SYNC') {
for (const channel of packet.d?.threads ?? []) this.#rememberChannel(channel);
} else if (packet.t === 'MESSAGE_CREATE') {
const message = normalizeDiscordMessage(packet.d, this.#config.platformId);
const bridge = this.#bridge;
if (message && bridge) {
void bridge.accept(message).catch((error) => {
if (packet.d && bridge) {
void this.#acceptMessage(packet.d, bridge).catch((error) => {
if (generation !== this.#generation || this.#stopped) return;
this.#logger.error?.(
`[dsh-im:discord] bot ${this.#config.botId} message handling failed:`,
@ -475,6 +667,38 @@ export class DiscordRuntime {
});
}
#rememberChannel(channel) {
if (!channel?.id) return;
this.#channels.set(String(channel.id), channel);
}
async #acceptMessage(message, bridge) {
const messageId = String(message?.id ?? '');
if (!messageId || this.#state.hasSeen(messageId)) return;
let route = this.#routing.get(messageId);
if (!route) {
route = resolveDiscordMessageRoute(message, this.#config.platformId, {
api: this.#api,
channel: this.#channels.get(String(message.channel_id)),
signal: this.#abortController?.signal,
onChannel: (resolved) => this.#rememberChannel(resolved),
});
this.#routing.set(messageId, route);
void route.finally(() => {
if (this.#routing.get(messageId) === route) this.#routing.delete(messageId);
}).catch(() => undefined);
}
try {
const normalized = await route;
if (normalized) await bridge.accept(normalized);
} catch (error) {
if (error?.code === 'discord-thread-create-uncertain') {
await this.#state.markSeen(messageId);
}
throw error;
}
}
#sendGateway(socket, payload) {
if (socket.readyState !== 1) return;
socket.send(JSON.stringify(payload));
@ -541,6 +765,7 @@ export class DiscordRuntime {
this.#socket = null;
this.#bridge = null;
this.#api = null;
this.#routing.clear();
try {
if (socket && socket.readyState < 2) socket.close(1000, 'Plugin stopped');
} catch (error) {

View file

@ -482,7 +482,7 @@ export class TextHarnessBridge {
key: conversationKey,
actor: senderId,
target,
requiresMention: message.kind === 'group',
requiresMention: message.kind === 'group' && message.requiresMention !== false,
}),
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
files: message.files,

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,496 @@ 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', () => {
const botId = '1234567890123456789';
const parentChannelId = '222222222222222223';
const fixtures = [
{ label: 'guild text', channelId: parentChannelId, channelType: 0 },
{ label: 'announcement', channelId: '222222222222222224', channelType: 5 },
{ label: 'announcement thread', channelId: '222222222222222225', channelType: 10 },
{ label: 'public thread', channelId: '222222222222222226', channelType: 11 },
{ label: 'private thread', channelId: '222222222222222227', channelType: 12 },
];
for (const [index, fixture] of fixtures.entries()) {
const message = normalizeDiscordMessage({
id: `11111111111111112${index}`,
channel_id: fixture.channelId,
guild_id: '444444444444444444',
channel_type: fixture.channelType,
author: { id: `33333333333333333${index}`, bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> ${fixture.label}`,
}, botId);
assert.equal(message.kind, 'group');
assert.equal(message.addressed, true);
assert.equal(message.conversationId, fixture.channelId);
assert.deepEqual(message.replyTarget, {
channelId: fixture.channelId,
replyToMessageId: `11111111111111112${index}`,
});
}
const firstUser = normalizeDiscordMessage({
id: '111111111111111130',
channel_id: parentChannelId,
guild_id: '444444444444444444',
author: { id: '333333333333333330', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> first`,
}, botId);
const secondUser = normalizeDiscordMessage({
id: '111111111111111131',
channel_id: parentChannelId,
guild_id: '444444444444444444',
author: { id: '333333333333333331', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> second`,
}, botId);
assert.equal(firstUser.conversationId, secondUser.conversationId);
const unmentionedThread = normalizeDiscordMessage({
id: '111111111111111132',
channel_id: '222222222222222226',
guild_id: '444444444444444444',
author: { id: '333333333333333330', bot: false },
mentions: [],
content: 'continue',
}, botId);
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 });
const foreignMessage = {
...message,
id: '111111111111111151',
content: `<@${botId}> recover foreign thread`,
};
const foreign = await resolveDiscordMessageRoute(foreignMessage, botId, {
channel: { id: foreignMessage.channel_id, type: 0 },
api: {
async getChannel({ channelId }) {
assert.equal(channelId, foreignMessage.id);
return {
id: foreignMessage.id,
type: 11,
parent_id: foreignMessage.channel_id,
owner_id: '999999999999999999',
};
},
async startThreadFromMessage() {
const error = new Error('Thread already exists');
error.status = 400;
throw error;
},
},
});
assert.equal(foreign.addressed, true, 'the explicit source mention still starts this turn');
assert.equal(foreign.requiresMention, true);
assert.equal(foreign.conversationRoute.managed, false);
assert.deepEqual(foreign.replyTarget, { channelId: foreignMessage.id });
});
test('Discord converges a post-dispatch abort but sends no create request after a pre-dispatch abort', async () => {
const botId = '1234567890123456789';
const base = {
channel_id: '222222222222222252',
guild_id: '444444444444444444',
author: { id: '333333333333333333', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> abort boundary`,
};
const before = new AbortController();
const beforeReason = new DOMException('cancelled before dispatch', 'AbortError');
before.abort(beforeReason);
let preDispatchCalls = 0;
await assert.rejects(() => resolveDiscordMessageRoute({
...base,
id: '111111111111111152',
}, botId, {
channel: { id: base.channel_id, type: 0 },
signal: before.signal,
api: {
async getChannel() { preDispatchCalls += 1; },
async startThreadFromMessage() { preDispatchCalls += 1; },
},
}), (error) => error === beforeReason);
assert.equal(preDispatchCalls, 0);
const after = new AbortController();
const afterReason = new DOMException('cancelled after dispatch', 'AbortError');
const message = { ...base, id: '111111111111111153' };
let starts = 0;
let lookups = 0;
const recovered = await resolveDiscordMessageRoute(message, botId, {
channel: { id: message.channel_id, type: 0 },
signal: after.signal,
api: {
async startThreadFromMessage() {
starts += 1;
after.abort(afterReason);
throw afterReason;
},
async getChannel({ channelId, signal }) {
lookups += 1;
assert.equal(channelId, message.id);
assert.notEqual(signal, after.signal);
assert.equal(signal.aborted, false);
return {
id: message.id,
type: 11,
parent_id: message.channel_id,
owner_id: botId,
};
},
},
});
assert.equal(starts, 1);
assert.equal(lookups, 1);
assert.equal(recovered.conversationId, message.id);
assert.equal(recovered.conversationRoute.managed, true);
});
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 streamed text and result files on the final created Thread target', async () => {
const botId = '1234567890123456789';
const message = {
id: '111111111111111165',
channel_id: '222222222222222265',
guild_id: '444444444444444444',
author: { id: '333333333333333333', bot: false },
mentions: [{ id: botId }],
content: `<@${botId}> deliver in thread`,
};
const operations = [];
const api = {
async getChannel() {
assert.fail('the cached parent channel must be reused');
},
async startThreadFromMessage() {
return {
id: message.id,
type: 11,
parent_id: message.channel_id,
owner_id: botId,
};
},
async createMessage(options) {
operations.push({ operation: 'create', ...options });
return { id: '888888888888888881' };
},
async editMessage(options) {
operations.push({ operation: 'edit', ...options });
return { id: options.messageId };
},
async createFileMessage(options) {
operations.push({ operation: 'file', ...options });
return { id: '888888888888888882' };
},
async sendTyping(options) {
operations.push({ operation: 'typing', ...options });
},
};
const route = await resolveDiscordMessageRoute(message, botId, {
channel: { id: message.channel_id, type: 0 },
api,
});
const client = new DiscordBotClient({ api });
await client.sendTyping(route.replyTarget);
const stream = await client.openStream(route.replyTarget);
await stream.finish('thread final answer');
await client.sendFile(route.replyTarget, {
fileName: 'result.txt',
bytes: Buffer.from('thread result'),
});
assert.deepEqual(operations.map(({ operation }) => operation), [
'typing', 'create', 'edit', 'file',
]);
assert.equal(operations.every(({ channelId }) => channelId === message.id), true);
assert.equal(operations.some(({ channelId }) => channelId === message.channel_id), false);
assert.equal(operations.find(({ operation }) => operation === 'create').replyToMessageId, undefined);
assert.equal(operations.find(({ operation }) => operation === 'file').replyToMessageId, undefined);
});
test('Discord keeps image attachments in images and exposes ordinary attachments as files', async () => {
@ -516,6 +1054,374 @@ 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 createRuntime = () => 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(' ')); } },
});
let runtime = createRuntime();
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();
runtime = createRuntime();
await runtime.start();
socket.emit('message', {
data: JSON.stringify({
op: 0,
t: 'GUILD_CREATE',
s: 20,
d: {
id: '444444444444444444',
channels: [{ id: parentChannelId, type: 0 }],
threads: [{
id: firstMessage.id,
type: 11,
parent_id: parentChannelId,
owner_id: botId,
}],
},
}),
});
socket.emit('message', {
data: JSON.stringify({
op: 0,
t: 'MESSAGE_CREATE',
s: 21,
d: {
id: '111111111111111174',
channel_id: firstMessage.id,
guild_id: '444444444444444444',
author: { id: '333333333333333330', bot: false },
mentions: [],
content: 'continue after restart',
},
}),
});
await eventually(() => runtime.status.messagesReplied === 1);
assert.deepEqual(asks[5], {
sessionId: 'session-1',
text: 'continue after restart',
files: 0,
});
assert.equal(starts.length, 3, 'restart must reuse the persisted bot-owned Thread');
assert.equal(sessions.get(`group:${firstMessage.id}`), 'session-1');
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();
@ -563,7 +1469,7 @@ test('Discord runtime identifies on Gateway v10 and becomes ready', async () =>
assert.equal(runtime.status.ready, true);
const identify = socket.sent.find((packet) => packet.op === 2);
assert.equal(identify.d.token, TOKEN);
assert.equal(identify.d.intents, 4_609);
assert.equal(identify.d.intents, 37_377);
assert.equal(identify.d.properties.browser, 'dsh-im');
for (const [id, sequence] of [

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 = [];