feat: add native result-file delivery

This commit is contained in:
xmanrui 2026-08-23 20:02:06 +08:00
parent 2803bbcbab
commit d472dd5a75
65 changed files with 10205 additions and 255 deletions

View file

@ -1,4 +1,5 @@
import { randomUUID } from 'node:crypto';
import { extname } from 'node:path';
import { fetchImageBuffer, ImagePromptError } from '../shared/image-prompt.mjs';
@ -25,8 +26,67 @@ function nonEmptyString(value) {
}
function safeProviderCode(value) {
const code = nonEmptyString(value);
return code && /^[A-Za-z0-9_.:-]{1,160}$/.test(code) ? code : undefined;
const code = value === undefined || value === null ? null : String(value).trim();
return code && /^-?[A-Za-z0-9_.:-]{1,160}$/.test(code) ? code : undefined;
}
function preserveArtifactMetadata(target, source) {
if (Number.isInteger(source?.status)) target.status = source.status;
if (source?.providerCode !== undefined) target.providerCode = source.providerCode;
return target;
}
function dingtalkArtifactError(cause, { fallback = 'artifact-provider-rejected' } = {}) {
if (cause?.code?.startsWith?.('artifact-')) return cause;
const status = Number(cause?.status);
const providerCode = safeProviderCode(cause?.providerCode);
const providerText = providerCode ?? '';
let code = fallback;
let message = 'DingTalk could not prepare the file for delivery.';
if (status === 401 || status === 403 || providerCode === '401' || providerCode === '403'
|| /(?:permission|forbidden|unauthor|access.?denied|\.auth(?:\.|$))/i.test(providerText)) {
code = 'artifact-permission-required';
message = 'DingTalk denied permission to send the file.';
} else if (status === 413 || providerCode === '413'
|| /(?:too.?large|size.?limit)/i.test(providerText)) {
code = 'artifact-too-large';
message = 'The file exceeds DingTalk\'s size limit.';
} else if (status === 429 || providerCode === '429'
|| /(?:rate.?limit|too.?many|throttl)/i.test(providerText)) {
code = 'artifact-rate-limited';
message = 'DingTalk rate-limited file delivery.';
} else if (fallback === 'artifact-provider-rejected') {
message = 'DingTalk rejected the file message.';
}
const error = new Error(message, { cause });
error.code = code;
return preserveArtifactMetadata(error, cause);
}
function uncertainDingtalkDelivery(cause) {
const error = new Error('DingTalk file delivery result is uncertain', { cause });
error.code = 'artifact-delivery-uncertain';
return preserveArtifactMetadata(error, cause);
}
function rejectedProviderResponse(value) {
if (!value || typeof value !== 'object') return null;
for (const field of ['errcode', 'code']) {
if (value[field] !== undefined && value[field] !== 0 && value[field] !== '0') {
return safeProviderCode(value[field]) ?? 'rejected';
}
}
return null;
}
function classifyDingtalkFinalDeliveryError(error, signal) {
if (signal?.aborted) throw abortError(signal);
const status = Number(error?.status);
if (error?.code === 'network-error' || error?.code === 'timeout'
|| error?.code === 'invalid-response' || (status >= 500 && status < 600)) {
return uncertainDingtalkDelivery(error);
}
return dingtalkArtifactError(error);
}
function secureDingtalkDownloadUrl(value) {
@ -174,6 +234,74 @@ async function requestJson(fetchImpl, url, {
}
}
async function requestMultipart(fetchImpl, url, { body, signal, timeoutMs = 60_000 } = {}) {
const controller = new AbortController();
let timedOut = false;
const onAbort = () => controller.abort(signal?.reason);
if (signal?.aborted) throw abortError(signal);
signal?.addEventListener('abort', onAbort, { once: true });
const timer = setTimeout(() => {
timedOut = true;
controller.abort();
}, timeoutMs);
try {
const response = await fetchImpl(url, {
method: 'POST',
body,
signal: controller.signal,
redirect: 'error',
});
let value;
let parseError;
try {
value = await response.json();
} catch (error) {
parseError = error;
}
if (!response.ok) {
throw new DingtalkApiError(
'http-error',
`钉钉服务请求失败(HTTP ${response.status})。`,
{ status: response.status, providerCode: safeProviderCode(value?.code ?? value?.errcode) },
);
}
if (parseError) {
throw new DingtalkApiError(
'invalid-response',
'钉钉服务返回了无法解析的响应。',
{ cause: parseError },
);
}
return value;
} catch (error) {
if (signal?.aborted) throw abortError(signal);
if (timedOut) throw new DingtalkApiError('timeout', '钉钉服务请求超时。', { cause: error });
if (error instanceof DingtalkApiError) throw error;
throw new DingtalkApiError('network-error', '暂时无法完成钉钉文件上传请求。', { cause: error });
} finally {
clearTimeout(timer);
signal?.removeEventListener('abort', onAbort);
}
}
function normalizeFileTarget(target) {
const robotCode = nonEmptyString(target?.robotCode);
if (!robotCode) throw new TypeError('DingTalk robotCode is required');
if (target?.type === 'group') {
const openConversationId = nonEmptyString(target.openConversationId);
if (openConversationId) return { type: 'group', robotCode, openConversationId };
}
if (target?.type === 'user') {
const userId = nonEmptyString(target.userId);
if (userId) return { type: 'user', robotCode, userId };
}
throw new TypeError('DingTalk file target is invalid');
}
function dingtalkFileType(fileName) {
return extname(fileName).slice(1).toLowerCase();
}
function normalizeCardTarget(target) {
if (target?.type === 'user') {
const userId = nonEmptyString(target.userId);
@ -648,6 +776,91 @@ export function createDingtalkApi({
return true;
},
async sendFile({ clientId, clientSecret, target, file, signal }) {
if (!file || typeof file !== 'object'
|| typeof file.fileName !== 'string' || !file.fileName
|| !Buffer.isBuffer(file.bytes)) {
throw new TypeError('A DingTalk file is required');
}
const normalizedTarget = normalizeFileTarget(target);
const fileType = dingtalkFileType(file.fileName);
let token;
try {
token = await accessToken({ clientId, clientSecret, signal });
} catch (error) {
if (signal?.aborted) throw abortError(signal);
const status = Number(error?.status);
const fallback = error?.code === 'http-error' && status >= 400 && status < 500
? 'artifact-provider-rejected'
: 'artifact-provider-failed';
throw dingtalkArtifactError(error, { fallback });
}
const uploadUrl = new URL('media/upload', DINGTALK_REGISTRATION_BASE_URL);
uploadUrl.searchParams.set('access_token', token);
uploadUrl.searchParams.set('type', 'file');
const form = new FormData();
form.append(
'media',
new Blob([file.bytes], { type: file.mediaType ?? 'application/octet-stream' }),
file.fileName,
);
let uploaded;
try {
uploaded = await requestMultipart(fetchImpl, uploadUrl, { body: form, signal });
} catch (error) {
if (signal?.aborted) throw abortError(signal);
const status = Number(error?.status);
const fallback = error?.code === 'http-error' && status >= 400 && status < 500
? 'artifact-provider-rejected'
: 'artifact-provider-failed';
throw dingtalkArtifactError(error, { fallback });
}
const uploadRejection = rejectedProviderResponse(uploaded);
if (uploadRejection || !nonEmptyString(uploaded?.media_id)) {
throw dingtalkArtifactError(new DingtalkApiError(
'upload-rejected',
'钉钉服务拒绝了文件上传。',
{ providerCode: uploadRejection ?? 'missing-media-id' },
));
}
signal?.throwIfAborted();
const messageBody = {
robotCode: normalizedTarget.robotCode,
msgKey: 'sampleFile',
msgParam: JSON.stringify({
mediaId: uploaded.media_id,
fileName: file.fileName,
fileType,
}),
...(normalizedTarget.type === 'group'
? { openConversationId: normalizedTarget.openConversationId }
: { userIds: [normalizedTarget.userId] }),
};
const pathname = normalizedTarget.type === 'group'
? 'v1.0/robot/groupMessages/send'
: 'v1.0/robot/oToMessages/batchSend';
let response;
try {
response = await requestJson(fetchImpl, endpoint(apiBase, pathname), {
body: messageBody,
headers: { 'x-acs-dingtalk-access-token': token },
signal,
action: '文件消息发送',
});
} catch (error) {
throw classifyDingtalkFinalDeliveryError(error, signal);
}
const sendRejection = rejectedProviderResponse(response);
if (sendRejection) {
throw dingtalkArtifactError(new DingtalkApiError(
'send-rejected',
'钉钉服务拒绝了文件消息。',
{ providerCode: sendRejection },
));
}
return response;
},
clearAccessToken(clientId) {
const appKey = nonEmptyString(clientId);
if (appKey) tokenCache.delete(appKey);

View file

@ -30,6 +30,16 @@ import {
promptContentForMessage,
} from '../shared/image-prompt.mjs';
import { rememberConnectionTestTarget } from '../shared/connection-test.mjs';
import {
materializeOutboundArtifact,
releaseOutboundArtifact,
} from '../shared/semantic/artifact.mjs';
import {
createArtifactFailureReceipt,
createDeliveryReceipt,
mergeDeliveryReceipts,
providerMessageIdsFor,
} from '../shared/semantic/delivery.mjs';
const CARD_INITIAL_TEXT = '已连接 DeepSeek Harness,正在思考…';
const CARD_ERROR_TEXT = '消息处理失败,请稍后重试。';
@ -62,6 +72,13 @@ function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
function dingtalkFileProviderIds(result) {
const ids = providerMessageIdsFor(result);
const processQueryKey = nonEmptyString(result?.processQueryKey);
if (processQueryKey && !ids.includes(processQueryKey)) ids.push(processQueryKey);
return ids;
}
function safeErrorDiagnostic(error) {
const chain = [];
const seen = new Set();
@ -195,6 +212,40 @@ function cardTarget(message, sender) {
return { type: 'user', userId: sender };
}
function fileTarget(message, sender, clientId) {
const robotCode = nonEmptyString(message?.robotCode) ?? clientId;
if (String(message?.conversationType) === '2') {
return {
type: 'group',
openConversationId: nonEmptyString(message?.conversationId),
robotCode,
};
}
return { type: 'user', userId: sender, robotCode };
}
function artifactFailureText(fileName, error) {
const name = String(fileName ?? '结果文件').replace(/[\r\n]+/g, ' ').trim() || '结果文件';
switch (error?.code) {
case 'artifact-delivery-uncertain':
return `结果文件「${name}」发送结果未能确认,请先检查聊天内是否已收到,不要立即重试。`;
case 'artifact-permission-required':
return `结果文件「${name}」已生成,但钉钉应用或机器人缺少文件消息权限。请开通应用 qyapi_base 权限,并确认机器人具备文件消息发送能力。`;
case 'artifact-too-large':
return `结果文件「${name}」超过当前钉钉机器人可发送的文件大小,未发送。`;
case 'artifact-rate-limited':
return `结果文件「${name}」暂时被钉钉限流,未能发送,请稍后重试。`;
case 'artifact-provider-rejected':
return `结果文件「${name}」已生成,但钉钉拒绝了该文件消息,请检查文件类型和机器人文件消息配置。`;
case 'artifact-invalid':
case 'artifact-changed':
case 'artifact-unavailable':
return `结果文件「${name}」暂时无法读取或准备发送,请确认文件仍可访问后重试。`;
default:
return `结果文件「${name}」已生成,但暂时未能通过钉钉发送,请稍后重试。`;
}
}
function progressText(update) {
if (update?.type === 'text' && nonEmptyString(update.text)) return update.text;
if (update?.type === 'tool') {
@ -596,7 +647,7 @@ export class DingtalkHarnessBridge {
});
cardStarted = await cardStream.start(CARD_INITIAL_TEXT);
}
const { answer } = await askInWorkspaceSession({
const { answer, artifacts = [] } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
@ -619,11 +670,41 @@ export class DingtalkHarnessBridge {
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
},
});
const streamed = cardStarted && await cardStream.finish(answer);
if (!streamed) await this.#send(sessionWebhook, answer);
const answerText = typeof answer === 'string' && answer.trim()
? answer
: artifacts.length > 0 ? '结果文件已生成。' : answer;
let textDeliveryError = null;
let textReceipt = null;
let streamed = false;
try {
streamed = cardStarted && await cardStream.finish(answerText);
if (streamed) {
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'dingtalk-card',
});
} else {
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'dingtalk-text',
providerMessageIds: await this.#send(sessionWebhook, answerText),
});
}
} catch (error) {
textDeliveryError = error;
}
const delivery = await this.#deliverArtifacts(
fileTarget(message, sender, this.#clientId),
sessionWebhook,
messageId,
artifacts,
textReceipt,
);
if (textDeliveryError && !delivery.userVisible) throw textDeliveryError;
increment(this.#status, 'messagesReplied');
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
return delivery.receipt;
} catch (error) {
if (error?.code === 'turn-stopped') {
if (cardStarted) await cardStream.finish('已停止。').catch(() => undefined);
@ -940,16 +1021,86 @@ export class DingtalkHarnessBridge {
}
async #send(sessionWebhook, text) {
const providerMessageIds = [];
for (const chunk of splitDingtalkText(text, this.#maxMessageChars)) {
this.#signal?.throwIfAborted();
await this.#api.sendText({
const result = await this.#api.sendText({
clientId: this.#clientId,
clientSecret: this.#clientSecret,
sessionWebhook,
text: chunk,
signal: this.#signal,
});
providerMessageIds.push(...providerMessageIdsFor(result));
}
return providerMessageIds;
}
async #deliverArtifacts(target, sessionWebhook, replyTo, artifacts, baseReceipt) {
const receipts = baseReceipt ? [baseReceipt] : [];
let userVisible = Boolean(baseReceipt);
for (const artifact of artifacts) {
this.#signal?.throwIfAborted();
try {
if (typeof this.#api.sendFile !== 'function') {
const unavailable = new Error('DingTalk file delivery is unavailable');
unavailable.code = 'artifact-provider-unavailable';
throw unavailable;
}
const file = await materializeOutboundArtifact(artifact, {
signal: this.#signal,
});
const result = await this.#api.sendFile({
clientId: this.#clientId,
clientSecret: this.#clientSecret,
target,
file,
signal: this.#signal,
});
receipts.push(createDeliveryReceipt({
deliveryId: file.deliveryKey,
presentation: 'dingtalk-file',
providerMessageIds: dingtalkFileProviderIds(result),
artifacts: [{ artifactId: file.artifactId, outcome: 'sent' }],
}));
userVisible = true;
this.#status.artifactsSent = (this.#status.artifactsSent ?? 0) + 1;
} catch (error) {
if (this.#signal?.aborted) throw error;
this.#status.artifactSendErrors = (this.#status.artifactSendErrors ?? 0) + 1;
this.#logger.warn?.(
`[dsh-dingtalk] result file delivery failed (${error?.code ?? 'unknown'})`,
);
let noticeSent = false;
const providerMessageIds = await this.#send(
sessionWebhook,
artifactFailureText(artifact?.fileName, error),
).then((ids) => {
noticeSent = true;
return ids;
}).catch(() => []);
const failureReceipt = createArtifactFailureReceipt({
artifactId: artifact?.artifactId ?? 'unknown',
deliveryId: artifact?.deliveryKey ?? artifact?.artifactId ?? 'unknown',
error,
providerMessageIds,
});
receipts.push(failureReceipt);
if (noticeSent || failureReceipt.artifacts[0]?.outcome === 'unknown') userVisible = true;
} finally {
releaseOutboundArtifact(artifact);
}
}
const receipt = receipts.length === 0
? null
: receipts.length === 1
? receipts[0]
: mergeDeliveryReceipts({
deliveryId: replyTo,
presentation: baseReceipt ? 'dingtalk-text-and-files' : 'dingtalk-files',
receipts,
});
return { receipt, userVisible };
}
}

View file

@ -1,4 +1,9 @@
import { createHash } from 'node:crypto';
const DEFAULT_BASE_URL = 'https://discord.com/api/v10/';
const DEFAULT_FILE_UPLOAD_TIMEOUT_MS = 120_000;
const DISCORD_PERMISSION_ERRORS = new Set([50001, 50013]);
const DISCORD_TOO_LARGE_ERRORS = new Set([40005]);
function cleanString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
@ -9,6 +14,58 @@ function requestSignal(signal, timeoutMs) {
return signal ? AbortSignal.any([signal, timeout]) : timeout;
}
function abortReason(signal) {
return signal?.reason instanceof Error
? signal.reason
: new DOMException('The operation was aborted', 'AbortError');
}
function positiveTimeout(value, name) {
if (!Number.isInteger(value) || value < 1) throw new TypeError(`${name} must be a positive integer`);
return value;
}
function preserveProviderMetadata(target, source) {
if (source?.providerCode !== undefined) target.providerCode = source.providerCode;
if (source?.retry_after !== undefined) {
target.retry_after = source.retry_after;
target.retryAfter = source.retry_after;
}
if (Number.isInteger(source?.status)) target.status = source.status;
return target;
}
function discordArtifactProviderError(cause) {
const providerCode = Number(cause?.providerCode);
const status = Number(cause?.status);
const message = cleanString(cause?.message) ?? '';
let code = 'artifact-provider-rejected';
let summary = 'Discord rejected the attachment.';
if (status === 401 || status === 403 || DISCORD_PERMISSION_ERRORS.has(providerCode)) {
code = 'artifact-permission-required';
summary = 'Discord denied permission to send the attachment.';
} else if (status === 413 || DISCORD_TOO_LARGE_ERRORS.has(providerCode)
|| /(?:request|attachment|file).{0,24}too large/i.test(message)) {
code = 'artifact-too-large';
summary = 'The attachment exceeds Discord\'s size limit.';
} else if (status === 429) {
code = 'artifact-rate-limited';
summary = 'Discord rate-limited attachment delivery.';
} else if (status >= 500) {
code = 'artifact-delivery-uncertain';
summary = 'Discord attachment delivery result is uncertain.';
}
const error = new Error(summary, { cause });
error.code = code;
return preserveProviderMetadata(error, cause);
}
function uncertainDiscordDelivery(cause) {
const error = new Error('Discord attachment delivery result is uncertain', { cause });
error.code = 'artifact-delivery-uncertain';
return preserveProviderMetadata(error, cause);
}
function delay(ms, signal) {
return new Promise((resolve, reject) => {
if (signal?.aborted) {
@ -39,13 +96,20 @@ export class DiscordApi {
#token;
#fetch;
#baseUrl;
#fileUploadTimeoutMs;
constructor({ token, fetchImpl = fetch, baseUrl = DEFAULT_BASE_URL }) {
constructor({
token,
fetchImpl = fetch,
baseUrl = DEFAULT_BASE_URL,
fileUploadTimeoutMs = DEFAULT_FILE_UPLOAD_TIMEOUT_MS,
}) {
if (!validDiscordToken(token)) throw new TypeError('Discord Bot Token is invalid');
if (typeof fetchImpl !== 'function') throw new TypeError('DiscordApi requires fetch');
this.#token = token.trim();
this.#fetch = fetchImpl;
this.#baseUrl = new URL(baseUrl);
this.#fileUploadTimeoutMs = positiveTimeout(fileUploadTimeoutMs, 'fileUploadTimeoutMs');
}
getCurrentUser(options = {}) {
@ -74,6 +138,54 @@ export class DiscordApi {
});
}
async createFileMessage({ channelId, file, replyToMessageId, signal }) {
if (!file || typeof file !== 'object'
|| typeof file.fileName !== 'string' || !file.fileName
|| !Buffer.isBuffer(file.bytes)) {
throw new TypeError('A Discord attachment is required');
}
const deliverySeed = cleanString(file.deliveryKey) ?? cleanString(file.artifactId);
const nonce = deliverySeed
? createHash('sha256').update(deliverySeed).digest('hex').slice(0, 25)
: undefined;
const payload = new FormData();
payload.append('payload_json', JSON.stringify({
allowed_mentions: { parse: [], replied_user: false },
attachments: [{ id: 0, filename: file.fileName }],
...(nonce ? { nonce, enforce_nonce: true } : {}),
...(replyToMessageId ? {
message_reference: {
message_id: snowflake(replyToMessageId, 'message id'),
channel_id: snowflake(channelId, 'channel id'),
fail_if_not_exists: false,
},
} : {}),
}));
payload.append(
'files[0]',
new Blob([file.bytes], { type: file.mediaType ?? 'application/octet-stream' }),
file.fileName,
);
const targetChannelId = snowflake(channelId, 'channel id');
if (signal?.aborted) throw abortReason(signal);
const uploadSignal = requestSignal(signal, this.#fileUploadTimeoutMs);
try {
return await this.#request(`channels/${targetChannelId}/messages`, {
method: 'POST',
signal: uploadSignal,
timeoutMs: this.#fileUploadTimeoutMs,
body: payload,
multipart: true,
});
} catch (error) {
if (signal?.aborted) throw abortReason(signal);
if (error?.code?.startsWith?.('discord-')) {
throw discordArtifactProviderError(error);
}
throw uncertainDiscordDelivery(error);
}
}
editMessage({ channelId, messageId, content, signal }) {
return this.#request(
`channels/${snowflake(channelId, 'channel id')}/messages/${snowflake(messageId, 'message id')}`,
@ -100,6 +212,7 @@ export class DiscordApi {
timeoutMs = 15_000,
expectBody = true,
retry = true,
multipart = false,
}) {
let response;
try {
@ -107,10 +220,10 @@ export class DiscordApi {
method,
headers: {
authorization: `Bot ${this.#token}`,
'content-type': 'application/json',
...(multipart ? {} : { 'content-type': 'application/json' }),
'user-agent': 'DeepSeek-Harness-dsh-im (https://github.com/xmanrui/dsh-im, 1.0.2)',
},
...(body === undefined ? {} : { body: JSON.stringify(body) }),
...(body === undefined ? {} : { body: multipart ? body : JSON.stringify(body) }),
signal: requestSignal(signal, timeoutMs),
redirect: 'error',
});
@ -124,17 +237,32 @@ export class DiscordApi {
try {
parsed = await response.json();
} catch {
if (expectBody) throw new Error(`Discord ${method} returned invalid JSON`);
if (expectBody) {
const error = new Error(`Discord ${method} returned invalid JSON`);
error.status = response?.status;
throw error;
}
}
}
if (response.status === 429 && retry) {
const retryAfterMs = Math.min(10_000, Math.max(50, Number(parsed?.retry_after) * 1_000 || 1_000));
await delay(retryAfterMs, signal);
return this.#request(path, { method, body, signal, timeoutMs, expectBody, retry: false });
return this.#request(path, {
method, body, signal, timeoutMs, expectBody, retry: false, multipart,
});
}
if (!response.ok) {
const error = new Error(cleanString(parsed?.message) ?? `Discord API failed with HTTP ${response.status}`);
error.code = `discord-${response.status}`;
error.status = response.status;
if (Number.isInteger(parsed?.code) || typeof parsed?.code === 'string') {
error.providerCode = parsed.code;
}
const retryAfter = Number(parsed?.retry_after);
if (Number.isFinite(retryAfter) && retryAfter >= 0) {
error.retry_after = retryAfter;
error.retryAfter = retryAfter;
}
throw error;
}
return expectBody ? parsed : null;

View file

@ -112,7 +112,7 @@ export function normalizeDiscordMessage(message, botId, { fetchImpl = fetch } =
};
}
class DiscordBotClient {
export class DiscordBotClient {
#api;
#signal;
@ -123,22 +123,32 @@ class DiscordBotClient {
async sendText(target, text) {
const chunks = splitMessageText(text, 1_900);
let result = null;
const providerMessageIds = [];
for (const [index, chunk] of chunks.entries()) {
result = await this.#api.createMessage({
const result = await this.#api.createMessage({
channelId: target.channelId,
content: chunk,
replyToMessageId: index === 0 ? target.replyToMessageId : undefined,
signal: this.#signal,
});
if (typeof result?.id === 'string' && result.id) providerMessageIds.push(result.id);
}
return result;
return { providerMessageIds };
}
sendTyping(target) {
return this.#api.sendTyping({ channelId: target.channelId, signal: this.#signal });
}
sendFile(target, file) {
return this.#api.createFileMessage({
channelId: target.channelId,
file,
replyToMessageId: target.replyToMessageId,
signal: this.#signal,
});
}
async openStream(target) {
const stream = createEditableMessageStream({
limit: 1_900,
@ -162,6 +172,7 @@ class DiscordBotClient {
content,
signal: this.#signal,
}),
messageIdForResult: (message) => message?.id,
});
return stream.start();
}

View file

@ -34,6 +34,15 @@ import {
} from '../shared/preset-command.mjs';
import { runWorkspaceCommand, resolveSessionListWorkspace, workspacePathSnapshot } from '../shared/workspace-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
import {
materializeOutboundArtifact,
releaseOutboundArtifact,
} from '../shared/semantic/artifact.mjs';
import {
createArtifactFailureReceipt,
createDeliveryReceipt,
mergeDeliveryReceipts,
} from '../shared/semantic/delivery.mjs';
import {
MENU_PAGE_SIZE,
completionCard,
@ -126,6 +135,33 @@ function safeErrorText(error) {
}
}
function artifactFailureText(fileName, error) {
const name = String(fileName ?? '结果文件').replace(/[\r\n]+/g, ' ').trim() || '结果文件';
switch (error?.code) {
case 'artifact-permission-required':
return `结果文件「${name}」已生成,但机器人缺少飞书文件上传权限。请为应用添加 im:resource 并完成必要审批后重试。`;
case 'artifact-too-large':
return `结果文件「${name}」超过飞书 30 MB 上限,未发送。`;
case 'artifact-empty':
return `结果文件「${name}」为空,飞书不允许发送空文件。`;
case 'artifact-changed':
case 'artifact-invalid':
case 'artifact-unavailable':
return `结果文件「${name}」暂时无法读取或准备发送,请确认文件仍可访问后重试。`;
case 'artifact-rate-limited':
return `结果文件「${name}」暂时被飞书限流,未能发送,请稍后重试。`;
case 'artifact-delivery-uncertain':
return `结果文件「${name}」发送结果未能确认,请先检查聊天内是否已收到,不要立即重试。`;
default:
return `结果文件「${name}」已生成,但暂时未能发送,请稍后重试。`;
}
}
function answerTextForDelivery(answer, artifacts) {
if (typeof answer === 'string' && answer.trim()) return answer;
return artifacts.length > 0 ? '结果文件已生成。' : answer;
}
function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
@ -540,7 +576,10 @@ export class FeishuHarnessBridge {
.then(() => this.#handle(event, key, { alreadyRecorded }));
const settled = finalize
? work
.then(() => this.#finishReaction(messageId, processingReaction, 'DONE'))
.then(async (receipt) => {
await this.#finishReaction(messageId, processingReaction, 'DONE');
return receipt;
})
.catch((error) => this.#handleMessageFailure(
event,
messageId,
@ -736,10 +775,11 @@ export class FeishuHarnessBridge {
this.#logger.info?.(`[dsh-feishu] processing ${event.message.chat_type} message ${messageId}`);
try {
await this.#answerWithStream(event, key, message);
const receipt = await this.#answerWithStream(event, key, message);
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
return receipt;
} finally {
await this.#cancelPendingInteraction(key);
await this.#approvals.closeRoute(key);
@ -1578,6 +1618,81 @@ export class FeishuHarnessBridge {
};
}
async #sendAnswerText(chatId, answer, { deliveryId, presentation }) {
const providerMessageIds = [];
for (const chunk of splitText(answer)) {
this.#signal?.throwIfAborted();
const messageId = await this.#send(chatId, chunk);
if (messageId) providerMessageIds.push(messageId);
}
return createDeliveryReceipt({
deliveryId,
presentation,
providerMessageIds,
});
}
async #deliverArtifacts(chatId, replyTo, artifacts = [], baseReceipt) {
const receipts = baseReceipt ? [baseReceipt] : [];
let failureNoticeVisible = false;
for (const artifact of artifacts) {
this.#signal?.throwIfAborted();
try {
if (typeof this.#channel?.sendFile !== 'function') {
const unavailable = new Error('Feishu file delivery is unavailable');
unavailable.code = 'artifact-provider-unavailable';
throw unavailable;
}
const file = await materializeOutboundArtifact(artifact, {
signal: this.#signal,
});
receipts.push(await this.#channel.sendFile(chatId, file, {
replyTo,
signal: this.#signal,
}));
this.#status.artifactsSent = (this.#status.artifactsSent ?? 0) + 1;
} catch (error) {
if (this.#signal?.aborted) throw error;
this.#status.artifactSendErrors = (this.#status.artifactSendErrors ?? 0) + 1;
this.#logger.warn?.(
`[dsh-feishu] result file delivery failed (${error?.code ?? 'unknown'})`,
);
let noticeMessageId = null;
try {
noticeMessageId = await this.#send(chatId, artifactFailureText(artifact?.fileName, error));
failureNoticeVisible = true;
} catch {
this.#logger.warn?.('[dsh-feishu] unable to send the safe result-file failure notice');
}
receipts.push(createArtifactFailureReceipt({
artifactId: artifact?.artifactId ?? 'unknown',
deliveryId: artifact?.deliveryKey ?? artifact?.artifactId ?? 'unknown',
error,
providerMessageIds: noticeMessageId ? [noticeMessageId] : [],
}));
} finally {
releaseOutboundArtifact(artifact);
}
}
if (receipts.length === 0) {
return {
receipt: createDeliveryReceipt({
deliveryId: replyTo,
presentation: 'feishu-files',
}),
failureNoticeVisible,
};
}
const receipt = receipts.length === 1
? receipts[0]
: mergeDeliveryReceipts({
deliveryId: baseReceipt?.deliveryId ?? artifacts[0]?.deliveryKey ?? replyTo,
presentation: baseReceipt ? 'feishu-text-and-files' : 'feishu-files',
receipts,
});
return { receipt, failureNoticeVisible };
}
async #answerWithStream(event, key, message) {
const chatId = event.message.chat_id;
const messageId = event.message.message_id;
@ -1586,7 +1701,7 @@ export class FeishuHarnessBridge {
? await promptContentForMessage(message, { signal: this.#signal })
: undefined;
if (!this.#channel?.stream) {
const { answer } = await askInWorkspaceSession({
const { answer, artifacts = [] } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
@ -1596,15 +1711,41 @@ export class FeishuHarnessBridge {
existsOptions: { signal: this.#signal },
askOptions: this.#interactionAskOptions(event, key),
});
for (const chunk of splitText(answer)) await this.#send(chatId, chunk);
let textReceipt;
let textSendError = null;
try {
textReceipt = await this.#sendAnswerText(
chatId,
answerTextForDelivery(answer, artifacts),
{
deliveryId: messageId,
presentation: 'feishu-text',
},
);
} catch (error) {
textSendError = error;
this.#logger.warn?.(
'[dsh-feishu] final text delivery failed; continuing with result files:',
error,
);
}
const delivery = await this.#deliverArtifacts(chatId, messageId, artifacts, textReceipt);
const artifactDispatched = delivery.receipt.artifacts.some(
({ outcome }) => outcome === 'sent' || outcome === 'unknown',
);
if (textSendError && !artifactDispatched && !delivery.failureNoticeVisible) {
throw textSendError;
}
this.#status.streamFallbacks = (this.#status.streamFallbacks ?? 0) + 1;
return;
return delivery.receipt;
}
let promptStarted = false;
let completedAnswer = '';
let completedArtifacts = [];
let stream;
try {
await this.#channel.stream(chatId, {
stream = await this.#channel.stream(chatId, {
markdown: async (controller) => {
promptStarted = true;
const askOptions = {
@ -1614,7 +1755,7 @@ export class FeishuHarnessBridge {
this.#status.streamUpdates = (this.#status.streamUpdates ?? 0) + 1;
},
};
({ answer: completedAnswer } = await askInWorkspaceSession({
const completed = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
@ -1623,26 +1764,56 @@ export class FeishuHarnessBridge {
createOptions: { signal: this.#signal },
existsOptions: { signal: this.#signal },
askOptions,
}));
await controller.setContent(completedAnswer);
});
completedAnswer = completed.answer;
completedArtifacts = completed.artifacts ?? [];
await controller.setContent(answerTextForDelivery(completedAnswer, completedArtifacts));
},
}, { replyTo: messageId });
this.#status.streamResponses = (this.#status.streamResponses ?? 0) + 1;
} catch (error) {
this.#status.streamErrors = (this.#status.streamErrors ?? 0) + 1;
if (completedAnswer) {
if (completedAnswer || completedArtifacts.length > 0) {
this.#logger.warn?.(
'[dsh-feishu] native stream failed after generation; sending final text:',
error.message,
);
for (const chunk of splitText(completedAnswer)) await this.#send(chatId, chunk);
let textReceipt;
let textSendError = null;
try {
textReceipt = await this.#sendAnswerText(
chatId,
answerTextForDelivery(completedAnswer, completedArtifacts),
{
deliveryId: messageId,
presentation: 'feishu-text-fallback',
},
);
} catch (fallbackError) {
textSendError = fallbackError;
this.#logger.warn?.(
'[dsh-feishu] fallback text delivery failed; continuing with result files:',
fallbackError,
);
}
const delivery = await this.#deliverArtifacts(
chatId,
messageId,
completedArtifacts,
textReceipt,
);
const artifactDispatched = delivery.receipt.artifacts.some(
({ outcome }) => outcome === 'sent' || outcome === 'unknown',
);
if (textSendError && !artifactDispatched && !delivery.failureNoticeVisible) {
throw textSendError;
}
this.#status.streamFallbacks = (this.#status.streamFallbacks ?? 0) + 1;
return;
return delivery.receipt;
}
if (promptStarted) throw error;
this.#logger.warn?.('[dsh-feishu] native stream unavailable; using text fallback:', error.message);
const { answer } = await askInWorkspaceSession({
const { answer, artifacts = [] } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
@ -1652,9 +1823,46 @@ export class FeishuHarnessBridge {
existsOptions: { signal: this.#signal },
askOptions: this.#interactionAskOptions(event, key),
});
for (const chunk of splitText(answer)) await this.#send(chatId, chunk);
let textReceipt;
let textSendError = null;
try {
textReceipt = await this.#sendAnswerText(
chatId,
answerTextForDelivery(answer, artifacts),
{
deliveryId: messageId,
presentation: 'feishu-text-fallback',
},
);
} catch (fallbackError) {
textSendError = fallbackError;
this.#logger.warn?.(
'[dsh-feishu] fallback text delivery failed; continuing with result files:',
fallbackError,
);
}
const delivery = await this.#deliverArtifacts(chatId, messageId, artifacts, textReceipt);
const artifactDispatched = delivery.receipt.artifacts.some(
({ outcome }) => outcome === 'sent' || outcome === 'unknown',
);
if (textSendError && !artifactDispatched && !delivery.failureNoticeVisible) {
throw textSendError;
}
this.#status.streamFallbacks = (this.#status.streamFallbacks ?? 0) + 1;
return delivery.receipt;
}
const delivery = await this.#deliverArtifacts(
chatId,
messageId,
completedArtifacts,
createDeliveryReceipt({
deliveryId: messageId,
presentation: 'feishu-cardkit',
providerMessageIds: stream?.messageId ? [stream.messageId] : [],
}),
);
this.#status.streamResponses = (this.#status.streamResponses ?? 0) + 1;
return delivery.receipt;
}
async #processInteractionReply(event, messageId, key, expected, processingReaction) {

View file

@ -1,6 +1,22 @@
import { createHash } from 'node:crypto';
import { trackOutboundArtifactProviderPromise } from '../shared/semantic/artifact.mjs';
import { createDeliveryReceipt } from '../shared/semantic/delivery.mjs';
const STREAM_ELEMENT_ID = 'stream_md';
const DEFAULT_INITIAL_TEXT = '已连接 DeepSeek Harness,正在思考…';
const MAX_STREAM_CHARS = 28000;
const MAX_FILE_OPERATION_TIMEOUT_MS = 120_000;
const FILE_DELIVERY_ERRORS = new Map([
[99991672, ['artifact-permission-required', 'Feishu file delivery requires the im:resource permission.']],
[234006, ['artifact-too-large', 'The result file exceeds Feishu\'s size limit.']],
[234010, ['artifact-empty', 'Feishu does not accept empty files.']],
[230017, ['artifact-provider-rejected', 'Feishu rejected the uploaded file ownership.']],
[230020, ['artifact-rate-limited', 'Feishu temporarily rate-limited file delivery.']],
[230049, ['artifact-delivery-uncertain', 'Feishu could not confirm the file message result.']],
[230055, ['artifact-provider-rejected', 'Feishu rejected the file message type.']],
]);
function assertApiSuccess(operation, response) {
if (response?.code && response.code !== 0) {
@ -9,6 +25,103 @@ function assertApiSuccess(operation, response) {
return response;
}
function providerErrorCode(cause) {
const pending = [cause];
const seen = new Set();
let fallback;
while (pending.length > 0) {
const value = pending.shift();
if (!value || seen.has(value)) continue;
if (typeof value === 'object') seen.add(value);
if (Array.isArray(value)) {
pending.push(...value);
continue;
}
const code = Number(value?.code);
if (Number.isFinite(code) && code !== 0) {
if (FILE_DELIVERY_ERRORS.has(code)) return code;
fallback ??= code;
}
pending.push(value?.response?.data, value?.data, value?.error, value?.cause);
}
return fallback;
}
function fileDeliveryError(stage, cause, providerCode, { uncertain = false } = {}) {
const explicitCode = providerCode === undefined || providerCode === null
? undefined
: Number(providerCode);
const code = Number.isFinite(explicitCode) && explicitCode !== 0
? explicitCode
: providerErrorCode(cause);
const fallback = Number.isFinite(code)
? ['artifact-provider-rejected', `Feishu rejected file ${stage}.`]
: uncertain
? ['artifact-delivery-uncertain', 'Feishu could not confirm the file message result.']
: ['artifact-provider-failed', `Feishu file ${stage} failed.`];
const [errorCode, message] = FILE_DELIVERY_ERRORS.get(code) ?? fallback;
const error = new Error(message, { cause });
error.code = errorCode;
if (Number.isFinite(code)) error.providerCode = code;
return error;
}
function boundedFileTimeout(value, name) {
if (!Number.isInteger(value) || value < 1 || value > MAX_FILE_OPERATION_TIMEOUT_MS) {
throw new TypeError(`${name} must be an integer between 1 and ${MAX_FILE_OPERATION_TIMEOUT_MS}`);
}
return value;
}
function abortReason(signal) {
return signal?.reason ?? new DOMException('The operation was aborted', 'AbortError');
}
function operationTimeout(stage) {
const error = new Error(`Feishu file ${stage} timed out.`);
error.code = 'provider-timeout';
return error;
}
function waitForFileOperation(operation, { signal, timeoutMs, stage }) {
signal?.throwIfAborted();
const deadline = new AbortController();
const operationSignal = signal
? AbortSignal.any([signal, deadline.signal])
: deadline.signal;
return new Promise((resolve, reject) => {
let settled = false;
const finish = (callback, value) => {
if (settled) return;
settled = true;
clearTimeout(timer);
operationSignal.removeEventListener('abort', onAbort);
callback(value);
};
const onAbort = () => finish(
reject,
signal?.aborted ? abortReason(signal) : operationTimeout(stage),
);
const timer = setTimeout(() => deadline.abort(), timeoutMs);
operationSignal.addEventListener('abort', onAbort, { once: true });
Promise.resolve().then(() => operation(operationSignal)).then(
(value) => finish(resolve, value),
(error) => finish(reject, error),
);
if (operationSignal.aborted) onAbort();
});
}
function deliveryUuid(file, chatId) {
const digest = createHash('sha256')
.update(`${file.deliveryKey}\u0000${chatId}`)
.digest('hex')
.slice(0, 40);
return `dshim_${digest}`;
}
function summaryOf(text) {
const summary = String(text ?? '').replace(/\s+/g, ' ').trim();
return summary.length <= 50 ? summary : `${summary.slice(0, 49)}…`;
@ -39,10 +152,19 @@ function streamingCard(initialText) {
export class VerifiedFeishuChannel {
#client;
#initialText;
#fileUploadTimeoutMs;
#fileMessageTimeoutMs;
constructor({ client, initialText = DEFAULT_INITIAL_TEXT }) {
constructor({
client,
initialText = DEFAULT_INITIAL_TEXT,
fileUploadTimeoutMs = MAX_FILE_OPERATION_TIMEOUT_MS,
fileMessageTimeoutMs = MAX_FILE_OPERATION_TIMEOUT_MS,
}) {
this.#client = client;
this.#initialText = initialText;
this.#fileUploadTimeoutMs = boundedFileTimeout(fileUploadTimeoutMs, 'fileUploadTimeoutMs');
this.#fileMessageTimeoutMs = boundedFileTimeout(fileMessageTimeoutMs, 'fileMessageTimeoutMs');
}
async stream(chatId, input, options = {}) {
@ -107,6 +229,110 @@ export class VerifiedFeishuChannel {
}
}
async sendFile(chatId, file, { replyTo, signal } = {}) {
signal?.throwIfAborted();
if (typeof chatId !== 'string' || !chatId) throw new TypeError('chatId is required');
if (!file || typeof file !== 'object'
|| typeof file.artifactId !== 'string' || !file.artifactId
|| typeof file.deliveryKey !== 'string' || !file.deliveryKey
|| typeof file.fileName !== 'string' || !file.fileName
|| !Buffer.isBuffer(file.bytes)) {
throw new TypeError('A materialized result file is required');
}
let uploaded;
try {
uploaded = await waitForFileOperation((operationSignal) => {
operationSignal.throwIfAborted();
const pending = this.#client.im.v1.file.create({
data: {
file_type: 'stream',
file_name: file.fileName,
file: file.bytes,
},
});
trackOutboundArtifactProviderPromise(file, pending);
return pending;
}, {
signal,
timeoutMs: this.#fileUploadTimeoutMs,
stage: 'upload',
});
} catch (error) {
if (signal?.aborted) throw abortReason(signal);
throw fileDeliveryError('upload', error);
}
signal?.throwIfAborted();
const fileKey = uploaded?.file_key;
if (typeof fileKey !== 'string' || !fileKey) {
throw fileDeliveryError('upload', undefined, uploaded?.code);
}
const uuid = deliveryUuid(file, chatId);
const content = JSON.stringify({ file_key: fileKey });
const request = replyTo
? {
path: { message_id: replyTo },
data: { msg_type: 'file', content, uuid },
}
: {
params: { receive_id_type: 'chat_id' },
data: { receive_id: chatId, msg_type: 'file', content, uuid },
};
const send = () => {
const pending = replyTo
? this.#client.im.v1.message.reply(request)
: this.#client.im.v1.message.create(request);
trackOutboundArtifactProviderPromise(file, pending);
return pending;
};
let response;
try {
response = await waitForFileOperation(async (operationSignal) => {
operationSignal.throwIfAborted();
let result;
try {
result = await send();
} catch (error) {
if (providerErrorCode(error) !== 230049) throw error;
result = { code: 230049 };
}
operationSignal.throwIfAborted();
// Feishu documents 230049 as an uncertain asynchronous send result.
// Reuse the same file_key and UUID once so the provider can deduplicate.
if (Number(result?.code) === 230049) {
result = await send();
operationSignal.throwIfAborted();
}
return result;
}, {
signal,
timeoutMs: this.#fileMessageTimeoutMs,
stage: 'message send',
});
} catch (error) {
if (signal?.aborted) throw abortReason(signal);
throw fileDeliveryError('message send', error, undefined, { uncertain: true });
}
if (Number.isFinite(Number(response?.code)) && Number(response.code) !== 0) {
throw fileDeliveryError('message send', undefined, response.code, { uncertain: true });
}
const messageId = response?.data?.message_id;
if (typeof messageId !== 'string' || !messageId) {
throw fileDeliveryError('message send', undefined, undefined, { uncertain: true });
}
return createDeliveryReceipt({
deliveryId: file.deliveryKey,
presentation: 'feishu-file',
providerMessageIds: [messageId],
artifacts: [{
artifactId: file.artifactId,
outcome: 'sent',
}],
});
}
async #sendCard(chatId, cardId, replyTo) {
const content = JSON.stringify({ type: 'card', data: { card_id: cardId } });
const response = replyTo

View file

@ -9,6 +9,7 @@ export const REQUIRED_TENANT_SCOPES = Object.freeze([
'im:message:send_as_bot',
'im:message.reactions:write_only',
'im:message:recall',
'im:resource',
'cardkit:card:write',
]);

View file

@ -26,8 +26,20 @@ import {
imagePromptUserMessage,
promptContentForMessage,
} from '../shared/image-prompt.mjs';
import {
materializeOutboundArtifact,
releaseOutboundArtifact,
trackOutboundArtifactProviderPromise,
} from '../shared/semantic/artifact.mjs';
import {
createArtifactFailureReceipt,
createDeliveryReceipt,
mergeDeliveryReceipts,
providerMessageIdsFor,
} from '../shared/semantic/delivery.mjs';
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
const DEFAULT_FILE_UPLOAD_TIMEOUT_MS = 120_000;
export const QQ_IMAGE_HOSTS = Object.freeze([
'.myqcloud.com',
@ -120,6 +132,82 @@ function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
function artifactFailureText(fileName, error) {
const name = String(fileName ?? '结果文件').replace(/[\r\n]+/g, ' ').trim() || '结果文件';
if (error?.name === 'UploadDailyLimitExceededError') {
return `结果文件「${name}」已生成,但 QQ 今日文件上传额度已用完,请稍后重试。`;
}
switch (error?.code) {
case 'artifact-delivery-uncertain':
return `结果文件「${name}」的发送结果未能确认,请先检查聊天内是否已收到,不要立即重试。`;
case 'artifact-permission-required':
return `结果文件「${name}」已生成,但当前 QQ 机器人没有文件消息权限。`;
case 'artifact-too-large':
return `结果文件「${name}」超过当前 QQ 机器人可发送的文件大小,未发送。`;
case 'artifact-empty':
return `结果文件「${name}」为空,QQ 不允许发送空文件。`;
case 'artifact-changed':
case 'artifact-invalid':
case 'artifact-unavailable':
return `结果文件「${name}」暂时无法读取或准备发送,请确认文件仍可访问后重试。`;
case 'artifact-rate-limited':
return `结果文件「${name}」暂时被 QQ 限流,未能发送,请稍后重试。`;
case 'artifact-provider-rejected':
return `结果文件「${name}」已生成,但 QQ 拒绝了该文件或文件消息。`;
default:
return `结果文件「${name}」已生成,但暂时未能通过 QQ 发送,请稍后重试。`;
}
}
function answerTextForDelivery(answer, artifacts) {
if (typeof answer === 'string' && answer.trim()) return answer;
return artifacts.length > 0 ? '结果文件已生成。' : answer;
}
function qqArtifactError(error, { dispatched = false } = {}) {
if (error?.code?.startsWith?.('artifact-') || error?.name === 'UploadDailyLimitExceededError') {
return error;
}
const status = Number(error?.httpStatus);
const wrapped = new Error('QQ file delivery failed', { cause: error });
if (status === 401 || status === 403) wrapped.code = 'artifact-permission-required';
else if (status === 413) wrapped.code = 'artifact-too-large';
else if (status === 429) wrapped.code = 'artifact-rate-limited';
else if (status === 400 || status === 404) {
wrapped.code = 'artifact-provider-rejected';
} else {
wrapped.code = dispatched ? 'artifact-delivery-uncertain' : 'artifact-provider-failed';
}
return wrapped;
}
function abortReason(signal) {
return signal?.reason instanceof Error
? signal.reason
: new DOMException('The operation was aborted', 'AbortError');
}
function waitWithSignal(promise, signal) {
if (!signal) return promise;
signal.throwIfAborted();
return new Promise((resolve, reject) => {
let settled = false;
const finish = (callback, value) => {
if (settled) return;
settled = true;
signal.removeEventListener('abort', onAbort);
callback(value);
};
const onAbort = () => finish(reject, abortReason(signal));
signal.addEventListener('abort', onAbort, { once: true });
Promise.resolve(promise).then(
(value) => finish(resolve, value),
(error) => finish(reject, error),
);
if (signal.aborted) onAbort();
});
}
function canClaimInteractionReply(message, pending) {
return pending.questions[pending.index]
&& nonEmptyString(message?.senderId) === pending.actor
@ -133,6 +221,8 @@ export function createQqBridgeStatus() {
messagesReceived: 0,
messagesReplied: 0,
messagesRejected: 0,
artifactsSent: 0,
artifactSendErrors: 0,
lastMessageAt: null,
lastReplyAt: null,
lastRejectedAt: null,
@ -150,6 +240,7 @@ export class QqHarnessBridge {
#replyTimeoutMs;
#signal;
#fetchImpl;
#fileUploadTimeoutMs;
#queues = new Map();
#pendingInteractions = new Map();
#interactionKeys = new Map();
@ -168,11 +259,15 @@ export class QqHarnessBridge {
replyTimeoutMs = 600_000,
signal,
fetchImpl = fetch,
fileUploadTimeoutMs = DEFAULT_FILE_UPLOAD_TIMEOUT_MS,
}) {
if (!bot || typeof bot.sendText !== 'function') throw new TypeError('QQ bot client is required');
if (!ownerUserOpenid) throw new TypeError('QQ scanner identity is required');
if (!harness || !state) throw new TypeError('Harness client and state store are required');
if (typeof fetchImpl !== 'function') throw new TypeError('fetchImpl must be a function');
if (!Number.isInteger(fileUploadTimeoutMs) || fileUploadTimeoutMs < 1) {
throw new TypeError('fileUploadTimeoutMs must be a positive integer');
}
this.#bot = bot;
this.#ownerUserOpenid = ownerUserOpenid;
this.#harness = harness;
@ -182,6 +277,7 @@ export class QqHarnessBridge {
this.#replyTimeoutMs = replyTimeoutMs;
this.#signal = signal;
this.#fetchImpl = fetchImpl;
this.#fileUploadTimeoutMs = Math.min(fileUploadTimeoutMs, DEFAULT_FILE_UPLOAD_TIMEOUT_MS);
this.#approvals = new HarnessApprovalQueue({ label: 'qq', logger });
}
@ -343,6 +439,87 @@ export class QqHarnessBridge {
this.#status.lastError = null;
}
async #deliverArtifacts(target, replyTo, artifacts = [], baseReceipt = null) {
if (artifacts.length === 0) {
return { receipt: baseReceipt, failureNoticeVisible: false };
}
const receipts = baseReceipt ? [baseReceipt] : [];
let failureNoticeVisible = false;
for (const artifact of artifacts) {
this.#signal?.throwIfAborted();
try {
if (typeof this.#bot.sendFile !== 'function') {
const unavailable = new Error('QQ file delivery is unavailable');
unavailable.code = 'artifact-provider-unavailable';
throw unavailable;
}
const file = await materializeOutboundArtifact(artifact, {
signal: this.#signal,
});
this.#signal?.throwIfAborted();
let result;
try {
const timeout = AbortSignal.timeout(this.#fileUploadTimeoutMs);
const waitSignal = this.#signal ? AbortSignal.any([this.#signal, timeout]) : timeout;
const pending = this.#bot.sendFile(
target,
{ buffer: file.bytes },
{
fileName: file.fileName,
onProgress: () => this.#signal?.throwIfAborted(),
},
);
trackOutboundArtifactProviderPromise(file, pending);
result = await waitWithSignal(pending, waitSignal);
} catch (error) {
if (this.#signal?.aborted) throw abortReason(this.#signal);
throw qqArtifactError(error, { dispatched: true });
}
this.#signal?.throwIfAborted();
const messageId = nonEmptyString(result?.message?.id);
receipts.push(createDeliveryReceipt({
deliveryId: file.deliveryKey,
presentation: 'qq-file',
providerMessageIds: messageId ? [messageId] : [],
artifacts: [{ artifactId: file.artifactId, outcome: 'sent' }],
}));
this.#status.artifactsSent = (this.#status.artifactsSent ?? 0) + 1;
} catch (rawError) {
if (this.#signal?.aborted) throw rawError;
const error = qqArtifactError(rawError);
this.#status.artifactSendErrors = (this.#status.artifactSendErrors ?? 0) + 1;
this.#logger.warn?.(
`[dsh-im:qq] result file delivery failed (${error?.code ?? error?.name ?? 'unknown'})`,
);
let providerMessageIds = [];
try {
const notice = await this.#bot.sendText(target, artifactFailureText(artifact?.fileName, error));
failureNoticeVisible = true;
providerMessageIds = providerMessageIdsFor(notice);
} catch (noticeError) {
if (this.#signal?.aborted) throw noticeError;
this.#logger.warn?.('[dsh-im:qq] unable to send the safe result-file failure notice');
}
receipts.push(createArtifactFailureReceipt({
artifactId: artifact?.artifactId ?? 'unknown',
deliveryId: artifact?.deliveryKey ?? artifact?.artifactId ?? 'unknown',
error,
providerMessageIds,
}));
} finally {
releaseOutboundArtifact(artifact);
}
}
return {
receipt: mergeDeliveryReceipts({
deliveryId: replyTo,
presentation: baseReceipt ? 'qq-text-and-files' : 'qq-files',
receipts,
}),
failureNoticeVisible,
};
}
async #process(message, key, { alreadyRecorded = false } = {}) {
if (this.#signal?.aborted) return;
const messageId = nonEmptyString(message?.messageId);
@ -427,8 +604,9 @@ export class QqHarnessBridge {
}
}
let answer;
let artifacts = [];
try {
({ answer } = await askInWorkspaceSession({
({ answer, artifacts = [] } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
@ -462,21 +640,50 @@ export class QqHarnessBridge {
this.#approvals.closeRoute(key),
]);
}
if (stream) {
try {
await stream.update(answer);
await stream.complete();
streamFinished = true;
} catch (error) {
stream.cancel?.();
this.#logger.warn?.('[dsh-im:qq] QQ stream finalization failed; using a text reply:', error);
this.#signal?.throwIfAborted();
const displayAnswer = answerTextForDelivery(answer, artifacts);
let textReceipt = null;
let textSendError = null;
try {
if (stream) {
try {
await stream.update(displayAnswer);
await stream.complete();
streamFinished = true;
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'qq-text',
providerMessageIds: providerMessageIdsFor(stream),
});
} catch (error) {
stream.cancel?.();
this.#logger.warn?.('[dsh-im:qq] QQ stream finalization failed; using a text reply:', error);
}
}
if (!streamFinished) {
const sent = await this.#bot.sendText(target, displayAnswer);
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'qq-text',
providerMessageIds: providerMessageIdsFor(sent),
});
}
} catch (error) {
textSendError = error;
this.#logger.warn?.('[dsh-im:qq] final text delivery failed; continuing with result files:', error);
}
const delivery = await this.#deliverArtifacts(target, messageId, artifacts, textReceipt);
const artifactDispatched = delivery.receipt?.artifacts?.some(
({ outcome }) => outcome === 'sent' || outcome === 'unknown',
);
if (textSendError && !artifactDispatched && !delivery.failureNoticeVisible) {
throw textSendError;
}
if (!streamFinished) await this.#bot.sendText(target, answer);
await this.#state.markSeen(messageId);
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
return delivery.receipt;
} catch (error) {
if (error?.code === 'turn-stopped') {
if (stream) {

View file

@ -21,9 +21,11 @@ export function createEditableMessageStream({
create,
edit,
sendRemainder,
messageIdForResult = () => null,
logger = console,
}) {
let messageId;
const providerMessageIds = [];
let pending = null;
let timer = null;
let inFlight = null;
@ -49,8 +51,17 @@ export function createEditableMessageStream({
};
return {
get messageId() {
return messageId;
},
get providerMessageIds() {
return [...providerMessageIds];
},
async start() {
messageId = await create(initialText);
if ((typeof messageId === 'string' && messageId.trim()) || Number.isSafeInteger(messageId)) {
providerMessageIds.push(String(messageId));
}
return this;
},
update(text) {
@ -69,7 +80,13 @@ export function createEditableMessageStream({
const first = chunks[0] ?? '处理完成。';
if (first !== lastSent) await edit(messageId, first);
lastSent = first;
for (const chunk of chunks.slice(1)) await sendRemainder(chunk);
for (const chunk of chunks.slice(1)) {
const result = await sendRemainder(chunk);
const id = messageIdForResult(result);
if ((typeof id === 'string' && id.trim()) || Number.isSafeInteger(id)) {
providerMessageIds.push(String(id));
}
}
},
cancel() {
closed = true;

View file

@ -3,6 +3,7 @@ import { randomUUID } from 'node:crypto';
import { isAbsolute } from 'node:path';
import { adoptRegisteredWorkspaceSession } from './harness-session-binding.mjs';
import { outboundArtifactRegistry } from './semantic/artifact.mjs';
// Every channel plugin runs in the same Host process. Sharing ownership by
// Harness origin prevents two channel-specific clients bound to one Session
@ -1025,6 +1026,7 @@ export class HarnessClient {
const timeoutMs = options.timeoutMs ?? 600_000;
const signal = options.signal;
const onUpdate = typeof options.onUpdate === 'function' ? options.onUpdate : null;
const onArtifact = typeof options.onArtifact === 'function' ? options.onArtifact : null;
const onInteraction = typeof options.onInteraction === 'function'
? options.onInteraction
: undefined;
@ -1071,11 +1073,34 @@ export class HarnessClient {
}
: null;
let interactionTask = null;
let artifactsDelivered = false;
let deliveredArtifactCount = 0;
const deliverArtifacts = async () => {
if (!onArtifact || artifactsDelivered || tracker.turn === null) {
return deliveredArtifactCount;
}
artifactsDelivered = true;
const artifacts = outboundArtifactRegistry.take(sessionId, tracker.turn, { signal });
for (const artifact of artifacts) {
try {
await onArtifact(artifact);
deliveredArtifactCount += 1;
} catch (error) {
outboundArtifactRegistry.release(artifact);
console.warn(`[${this.#logPrefix}] ignored an artifact handoff failure:`, error.message);
}
}
return deliveredArtifactCount;
};
if (ownership) {
this.#registerInteractionOwnership(sessionId, ownership);
this.#registerControlOwnership(ownership);
}
// This is resource ownership, not a feature Gate: it lets the Host retain
// this Turn's snapshots until the channel has polled and claimed them.
const closeArtifactConsumer = outboundArtifactRegistry.openConsumer(sessionId, promptRpcId);
try {
if (interactionSignal) {
@ -1134,7 +1159,15 @@ export class HarnessClient {
}
}
if (!tracker.finished) continue;
if (tracker.answer) return tracker.answer;
// An accepted /stop revokes attachment delivery even when Harness
// preserved a useful partial text answer for the existing UX.
const artifactCount = ownership?.stopRequested
? 0
: await deliverArtifacts();
if (tracker.answer) {
return tracker.answer;
}
if (artifactCount > 0) return '';
if (ownership?.stopRequested) throw turnStoppedError();
throw new Error(
`Harness turn ended without a text reply${tracker.reason ? ` (${JSON.stringify(tracker.reason)})` : ''}`,
@ -1145,15 +1178,19 @@ export class HarnessClient {
// Once cancellation was accepted, transport/poll failures and timeouts
// describe the convergence of that stop, not an unrelated ask failure.
if (!ownership?.stopRequested) throw error;
if (tracker.answer) return tracker.answer;
if (tracker.answer) {
return tracker.answer;
}
if (error?.code === 'turn-stopped') throw error;
throw turnStoppedError();
}
} finally {
closeArtifactConsumer();
if (ownership) {
this.#unregisterControlOwnership(ownership);
this.#unregisterInteractionOwnership(sessionId, ownership);
}
if (tracker.turn !== null) outboundArtifactRegistry.discard(sessionId, tracker.turn);
interactionController?.abort(new DOMException('Harness turn finished', 'AbortError'));
if (interactionTask) await interactionTask.catch(() => undefined);
}

View file

@ -0,0 +1,748 @@
import { createHash, randomUUID } from 'node:crypto';
import { constants as fsConstants } from 'node:fs';
import { copyFile, lstat, mkdtemp, open, realpath, unlink } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import { basename, extname, isAbsolute, join, resolve } from 'node:path';
export const OUTBOUND_ARTIFACT_TOOL = 'dsh_im_return_file';
const ARTIFACT_KIND = 'dsh-im-outbound-artifact';
const ARTIFACT_READ_CHUNK_BYTES = 64 * 1024;
const MIME_BY_EXTENSION = new Map([
['.csv', 'text/csv'],
['.doc', 'application/msword'],
['.docx', 'application/vnd.openxmlformats-officedocument.wordprocessingml.document'],
['.gif', 'image/gif'],
['.html', 'text/html'],
['.jpeg', 'image/jpeg'],
['.jpg', 'image/jpeg'],
['.json', 'application/json'],
['.md', 'text/markdown'],
['.pdf', 'application/pdf'],
['.png', 'image/png'],
['.rar', 'application/vnd.rar'],
['.txt', 'text/plain'],
['.webp', 'image/webp'],
['.xls', 'application/vnd.ms-excel'],
['.xlsx', 'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet'],
['.xml', 'application/xml'],
['.zip', 'application/zip'],
]);
const artifactStorage = new WeakMap();
const materializedArtifactSources = new WeakMap();
const artifactProviderSettlements = new WeakMap();
let managedSnapshotDirectoryPromise;
function managedSnapshotDirectory() {
managedSnapshotDirectoryPromise ??= mkdtemp(join(tmpdir(), 'dsh-im-outbound-'));
return managedSnapshotDirectoryPromise;
}
function artifactError(code, message) {
const error = new Error(message);
error.code = code;
return error;
}
function currentTurn(agent) {
const events = agent?.session?.events;
if (!Array.isArray(events)) return null;
let turn = null;
for (const event of events) {
if (event?.type === 'turn/start') {
turn = Number.isInteger(event.data?.turn) ? event.data.turn : null;
} else if (event?.type === 'turn/end' && event.data?.turn === turn) {
turn = null;
}
}
return Number.isInteger(turn) && turn >= 0 ? turn : null;
}
function sessionIdOf(session) {
const sessionId = session?.id ?? session?.header?.id;
return typeof sessionId === 'string' && sessionId ? sessionId : null;
}
function turnKey(sessionId, turn) {
return `${sessionId}\u0000${turn}`;
}
function promptKey(sessionId, promptRpcId) {
return `${sessionId}\u0000${promptRpcId}`;
}
function safeFileName(value) {
const cleaned = String(value ?? '')
.replace(/[\u0000-\u001f\u007f]/g, ' ')
.replace(/\p{Cf}/gu, '')
.replace(/\s+/g, ' ')
.trim();
if (!cleaned) return 'result.bin';
return [...cleaned].slice(0, 255).join('');
}
function mediaTypeFor(name) {
return MIME_BY_EXTENSION.get(extname(name).toLowerCase()) ?? 'application/octet-stream';
}
function sha256(bytes) {
return createHash('sha256').update(bytes).digest('hex');
}
function sameIdentity(left, right) {
return typeof left?.dev === 'bigint'
&& typeof left?.ino === 'bigint'
&& typeof left?.size === 'bigint'
&& typeof left?.mtimeNs === 'bigint'
&& typeof left?.ctimeNs === 'bigint'
&& left.dev === right?.dev
&& left.ino === right?.ino
&& left.size === right?.size
&& left.mtimeNs === right?.mtimeNs
&& left.ctimeNs === right?.ctimeNs;
}
async function hashFile(path, signal) {
let handle;
try {
handle = await open(path, fsConstants.O_RDONLY);
const hash = createHash('sha256');
const stream = handle.createReadStream({
autoClose: false,
highWaterMark: ARTIFACT_READ_CHUNK_BYTES,
...(signal ? { signal } : {}),
});
for await (const chunk of stream) {
signal?.throwIfAborted();
hash.update(chunk);
}
return hash.digest('hex');
} finally {
await handle?.close().catch(() => undefined);
}
}
/** Read exactly the observed file size and prove EOF. No project-level size or time limit applies. */
export async function readExactArtifactFile(handle, expectedSize, {
signal,
errorCode = 'artifact-changed',
errorMessage = 'The result file changed while it was being read.',
} = {}) {
const changed = () => artifactError(errorCode, errorMessage);
if (!Number.isSafeInteger(expectedSize) || expectedSize < 0) throw changed();
signal?.throwIfAborted();
let bytes;
try {
bytes = Buffer.allocUnsafe(expectedSize);
} catch {
throw changed();
}
if (typeof handle?.createReadStream === 'function') {
const stream = handle.createReadStream({
autoClose: false,
start: 0,
// Node's end offset is inclusive, so this reads one byte beyond the
// observed size when the file grows and catches the change.
end: expectedSize,
highWaterMark: ARTIFACT_READ_CHUNK_BYTES,
...(signal ? { signal } : {}),
});
let offset = 0;
for await (const chunk of stream) {
signal?.throwIfAborted();
const part = Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk);
if (offset + part.length > expectedSize) throw changed();
part.copy(bytes, offset);
offset += part.length;
}
if (offset !== expectedSize) throw changed();
return bytes;
}
let offset = 0;
while (offset < expectedSize) {
signal?.throwIfAborted();
const length = Math.min(ARTIFACT_READ_CHUNK_BYTES, expectedSize - offset);
const result = await handle.read(bytes, offset, length, offset);
const bytesRead = Number(result?.bytesRead ?? 0);
if (!Number.isInteger(bytesRead) || bytesRead <= 0 || bytesRead > length) throw changed();
offset += bytesRead;
}
signal?.throwIfAborted();
const probe = Buffer.allocUnsafe(1);
if ((await handle.read(probe, 0, 1, expectedSize))?.bytesRead !== 0) throw changed();
return bytes;
}
async function snapshotFile(workspace, requestedPath, signal) {
signal?.throwIfAborted();
let storagePath;
try {
const candidate = isAbsolute(requestedPath)
? requestedPath
: resolve(workspace, requestedPath);
const canonicalPath = await realpath(candidate);
const source = await lstat(canonicalPath, { bigint: true });
if (!source.isFile()) {
throw artifactError('artifact-not-file', 'The requested path is not a file.');
}
const directory = await managedSnapshotDirectory();
storagePath = join(directory, `${randomUUID()}.artifact`);
await copyFile(canonicalPath, storagePath, fsConstants.COPYFILE_EXCL);
signal?.throwIfAborted();
const snapshot = await lstat(storagePath, { bigint: true });
const size = Number(snapshot.size);
if (!snapshot.isFile() || !Number.isSafeInteger(size) || size < 0) {
throw artifactError('artifact-unavailable', 'The file could not be prepared for delivery.');
}
const fileName = safeFileName(basename(candidate));
return Object.freeze({
fileName,
mediaType: mediaTypeFor(fileName),
size,
digest: await hashFile(storagePath, signal),
storagePath,
});
} catch (error) {
if (storagePath) await unlink(storagePath).catch(() => undefined);
if (error?.code?.startsWith?.('artifact-')) throw error;
if (signal?.aborted) throw signal.reason ?? error;
throw artifactError('artifact-unavailable', 'The requested file is unavailable.');
}
}
function publicArtifact(artifact) {
return Object.freeze({
artifactId: artifact.artifactId,
fileName: artifact.fileName,
size: artifact.size,
});
}
function artifactKey(artifact) {
return artifact.artifactId;
}
function cleanupArtifactStorage(artifact) {
const storage = artifactStorage.get(artifact);
if (!storage) return;
storage.releaseRequested = true;
if (storage.materializing > 0) {
storage.cleanupRequested = true;
return;
}
const pending = [...artifactProviderSettlements.get(artifact) ?? []];
if (pending.length > 0) {
if (!storage.cleanupDeferred) {
storage.cleanupDeferred = true;
void Promise.allSettled(pending).finally(() => {
storage.cleanupDeferred = false;
cleanupArtifactStorage(artifact);
});
}
return;
}
artifactStorage.delete(artifact);
artifactProviderSettlements.delete(artifact);
storage.onCleanup?.();
void unlink(storage.path).catch(() => undefined);
}
/**
* Holds successful tool results until the channel that owns the Session Turn
* collects them. It does not decide which files a user may send.
*/
export class OutboundArtifactRegistry {
#turns = new Map();
#stagedTurns = new Map();
#claimedTurns = new Map();
#claimSignals = new Map();
#signalClaims = new Map();
#consumersByPrompt = new Map();
#consumersByTurn = new Map();
#openTurns = new Map();
#uuid;
constructor({ uuid = randomUUID } = {}) {
if (typeof uuid !== 'function') throw new TypeError('uuid must be a function');
this.#uuid = uuid;
}
/**
* Bind one channel request to the Turn it starts. This owns cleanup only:
* it never changes tool visibility or decides whether a file may be sent.
*/
openConsumer(sessionId, promptRpcId) {
if (typeof sessionId !== 'string' || !sessionId
|| typeof promptRpcId !== 'string' || !promptRpcId) {
throw new TypeError('sessionId and promptRpcId are required');
}
const key = promptKey(sessionId, promptRpcId);
const consumer = {
sessionId,
promptRpcId,
turn: null,
released: false,
};
this.#consumersByPrompt.set(key, consumer);
return () => {
if (consumer.released) return;
consumer.released = true;
if (this.#consumersByPrompt.get(key) === consumer) {
this.#consumersByPrompt.delete(key);
}
if (consumer.turn !== null) {
const keyForTurn = turnKey(sessionId, consumer.turn);
if (this.#consumersByTurn.get(keyForTurn) === consumer) {
this.#consumersByTurn.delete(keyForTurn);
}
// Claimed artifacts have already crossed into the provider pipeline and
// are released there. discard() only removes unclaimed work.
this.discard(sessionId, consumer.turn);
}
};
}
/** Observe durable Session events solely to terminate unclaimed snapshots. */
observeSessionEvent(session, event) {
const sessionId = sessionIdOf(session);
if (!sessionId || !event || typeof event !== 'object') return;
if (event.type === 'turn/start') {
const turn = event.data?.turn;
if (Number.isInteger(turn) && turn >= 0) this.#openTurns.set(sessionId, turn);
return;
}
if (event.type === 'user/message') {
const rpcId = event.data?.source?.rpcId;
const turn = this.#openTurns.get(sessionId) ?? currentTurn({ session });
if (typeof rpcId !== 'string' || !rpcId || !Number.isInteger(turn)) return;
this.#openTurns.set(sessionId, turn);
const consumer = this.#consumersByPrompt.get(promptKey(sessionId, rpcId));
if (!consumer || consumer.released) return;
consumer.turn = turn;
this.#consumersByTurn.set(turnKey(sessionId, turn), consumer);
return;
}
if (event.type !== 'turn/end') return;
const turn = event.data?.turn;
if (!Number.isInteger(turn) || turn < 0) return;
if (this.#openTurns.get(sessionId) === turn) this.#openTurns.delete(sessionId);
const consumer = this.#consumersByTurn.get(turnKey(sessionId, turn));
if (!consumer || consumer.released) this.discard(sessionId, turn);
}
/** A disposed Session cannot have another channel consumer claim its files. */
disposeSession(session) {
const sessionId = sessionIdOf(session);
if (!sessionId) return;
const prefix = `${sessionId}\u0000`;
const artifacts = new Set();
for (const entries of [this.#turns, this.#stagedTurns]) {
for (const [key, turnArtifacts] of entries) {
if (!key.startsWith(prefix)) continue;
for (const artifact of turnArtifacts.values()) artifacts.add(artifact);
entries.delete(key);
}
}
for (const [key, consumer] of this.#consumersByPrompt) {
if (!key.startsWith(prefix)) continue;
consumer.released = true;
this.#consumersByPrompt.delete(key);
}
for (const key of this.#consumersByTurn.keys()) {
if (key.startsWith(prefix)) this.#consumersByTurn.delete(key);
}
this.#openTurns.delete(sessionId);
for (const artifact of artifacts) cleanupArtifactStorage(artifact);
}
async stage(args, exec) {
const requestedPath = args?.path;
if (typeof requestedPath !== 'string' || !requestedPath.trim()) {
throw new TypeError('A file path is required.');
}
const agent = exec?.agent;
const sessionId = agent?.session?.header?.id;
const workspace = agent?.session?.header?.cwd;
const turn = currentTurn(agent);
if (typeof sessionId !== 'string' || !sessionId
|| typeof workspace !== 'string' || !workspace || turn === null) {
throw artifactError(
'artifact-context-required',
'A live Harness Session is required to return a file.',
);
}
const snapshot = await snapshotFile(workspace, requestedPath, exec?.signal);
const { storagePath, ...snapshotMetadata } = snapshot;
const artifact = Object.freeze({
kind: ARTIFACT_KIND,
schemaVersion: 1,
artifactId: this.#uuid(),
deliveryKey: this.#uuid(),
...snapshotMetadata,
source: 'managed-temp',
registeredBy: Object.freeze({
kind: 'tool-result',
eventId: typeof exec?.callId === 'string' ? exec.callId : 'unknown',
toolName: OUTBOUND_ARTIFACT_TOOL,
}),
origin: Object.freeze({
sessionId,
turn,
callId: typeof exec?.callId === 'string' ? exec.callId : null,
}),
createdAt: Date.now(),
});
artifactStorage.set(artifact, {
path: storagePath,
materializing: 0,
materialized: false,
releaseRequested: false,
cleanupRequested: false,
cleanupDeferred: false,
onCleanup: () => this.#forgetArtifact(artifact),
});
const keyForTurn = turnKey(sessionId, turn);
let staged = this.#stagedTurns.get(keyForTurn);
if (!staged) {
staged = new Map();
this.#stagedTurns.set(keyForTurn, staged);
}
staged.set(artifactKey(artifact), artifact);
return artifact;
}
commit(artifact) {
if (artifact?.kind !== ARTIFACT_KIND || !artifactStorage.has(artifact)) return null;
const keyForTurn = turnKey(artifact.origin.sessionId, artifact.origin.turn);
const key = artifactKey(artifact);
const staged = this.#stagedTurns.get(keyForTurn);
if (staged?.get(key) !== artifact) return null;
staged.delete(key);
if (staged.size === 0) this.#stagedTurns.delete(keyForTurn);
let committed = this.#turns.get(keyForTurn);
if (!committed) {
committed = new Map();
this.#turns.set(keyForTurn, committed);
}
committed.set(key, artifact);
return publicArtifact(artifact);
}
release(artifact) {
if (artifact?.kind !== ARTIFACT_KIND) return;
const keyForTurn = turnKey(artifact.origin.sessionId, artifact.origin.turn);
const key = artifactKey(artifact);
const staged = this.#stagedTurns.get(keyForTurn);
if (staged?.get(key) === artifact) staged.delete(key);
if (staged?.size === 0) this.#stagedTurns.delete(keyForTurn);
cleanupArtifactStorage(artifact);
}
take(sessionId, turn, { signal } = {}) {
const keyForTurn = turnKey(sessionId, turn);
const committed = this.#turns.get(keyForTurn);
this.#turns.delete(keyForTurn);
if (!committed) return [];
if (signal?.aborted) {
for (const artifact of committed.values()) cleanupArtifactStorage(artifact);
return [];
}
const claimed = this.#claimedTurns.get(keyForTurn) ?? new Map();
this.#claimedTurns.set(keyForTurn, claimed);
const artifacts = [];
for (const [key, artifact] of committed) {
if (!artifactStorage.has(artifact)) continue;
claimed.set(key, artifact);
artifacts.push(artifact);
}
if (claimed.size === 0) this.#claimedTurns.delete(keyForTurn);
else if (!this.#bindClaimSignal(keyForTurn, signal)) return [];
return artifacts;
}
discard(sessionId, turn) {
if (typeof sessionId !== 'string' || !Number.isInteger(turn)) return;
const keyForTurn = turnKey(sessionId, turn);
const artifacts = new Set([
...this.#turns.get(keyForTurn)?.values() ?? [],
...this.#stagedTurns.get(keyForTurn)?.values() ?? [],
]);
this.#turns.delete(keyForTurn);
this.#stagedTurns.delete(keyForTurn);
for (const artifact of artifacts) cleanupArtifactStorage(artifact);
}
clear() {
const artifacts = new Set();
for (const entries of this.#turns.values()) {
for (const artifact of entries.values()) artifacts.add(artifact);
}
for (const entries of this.#stagedTurns.values()) {
for (const artifact of entries.values()) artifacts.add(artifact);
}
for (const entries of this.#claimedTurns.values()) {
for (const artifact of entries.values()) artifacts.add(artifact);
}
for (const artifact of artifacts) cleanupArtifactStorage(artifact);
this.#turns.clear();
this.#stagedTurns.clear();
this.#claimedTurns.clear();
for (const turnKey of this.#claimSignals.keys()) this.#releaseClaimSignal(turnKey);
for (const consumer of this.#consumersByPrompt.values()) consumer.released = true;
this.#consumersByPrompt.clear();
this.#consumersByTurn.clear();
this.#openTurns.clear();
}
#bindClaimSignal(turnKey, signal) {
if (!signal) return true;
this.#releaseClaimSignal(turnKey);
let claim = this.#signalClaims.get(signal);
if (!claim) {
claim = { turnKeys: new Set(), onAbort: null };
claim.onAbort = () => {
for (const claimedTurnKey of [...claim.turnKeys]) {
for (const artifact of this.#claimedTurns.get(claimedTurnKey)?.values() ?? []) {
cleanupArtifactStorage(artifact);
}
}
};
this.#signalClaims.set(signal, claim);
signal.addEventListener('abort', claim.onAbort, { once: true });
}
claim.turnKeys.add(turnKey);
this.#claimSignals.set(turnKey, signal);
if (signal.aborted) {
claim.onAbort();
return false;
}
return true;
}
#releaseClaimSignal(turnKey) {
const signal = this.#claimSignals.get(turnKey);
if (!signal) return;
this.#claimSignals.delete(turnKey);
const claim = this.#signalClaims.get(signal);
claim?.turnKeys.delete(turnKey);
if (claim?.turnKeys.size === 0) {
signal.removeEventListener('abort', claim.onAbort);
this.#signalClaims.delete(signal);
}
}
#forgetArtifact(artifact) {
const keyForTurn = turnKey(artifact.origin.sessionId, artifact.origin.turn);
const key = artifactKey(artifact);
const committed = this.#turns.get(keyForTurn);
if (committed?.get(key) === artifact) committed.delete(key);
if (committed?.size === 0) this.#turns.delete(keyForTurn);
const staged = this.#stagedTurns.get(keyForTurn);
if (staged?.get(key) === artifact) staged.delete(key);
if (staged?.size === 0) this.#stagedTurns.delete(keyForTurn);
const claimed = this.#claimedTurns.get(keyForTurn);
if (claimed?.get(key) === artifact) claimed.delete(key);
if (claimed?.size === 0) {
this.#claimedTurns.delete(keyForTurn);
this.#releaseClaimSignal(keyForTurn);
}
}
}
export const outboundArtifactRegistry = new OutboundArtifactRegistry();
/**
* Build a two-phase tool: execute stages a file; the authoritative tools/result
* observer commits only a successful native call or successful Code Mode parent.
*/
export function createOutboundArtifactTool({ registry = outboundArtifactRegistry } = {}) {
const staged = new WeakMap();
const pendingByParent = new Map();
const appendPending = (parent, artifact) => {
let pending = pendingByParent.get(parent);
if (!pending) {
pending = { artifacts: [] };
pendingByParent.set(parent, pending);
}
pending.artifacts.push(artifact);
};
const definition = Object.freeze({
name: OUTBOUND_ARTIFACT_TOOL,
description: 'Send a readable file to the user through the current conversation. Existing and newly created files are both valid.',
parameters: {
type: 'object',
additionalProperties: false,
properties: {
path: {
type: 'string',
description: 'Absolute path, or a path relative to the current workspace.',
},
},
required: ['path'],
},
output: {
schema: {
type: 'object',
additionalProperties: false,
properties: {
artifactId: { type: 'string' },
fileName: { type: 'string' },
size: { type: 'number' },
},
required: ['artifactId', 'fileName', 'size'],
},
render: (_args, value) => [{
type: 'text',
text: `Registered ${value.fileName} (${value.size} bytes) for IM delivery.`,
}],
},
async execute(args, exec) {
const artifact = await registry.stage(args, exec);
staged.set(exec, artifact);
return publicArtifact(artifact);
},
});
const onResult = (exec, result) => {
if (exec?.name === OUTBOUND_ARTIFACT_TOOL) {
const artifact = staged.get(exec);
if (!artifact) return;
staged.delete(exec);
if (result?.isError) {
registry.release(artifact);
return;
}
if (exec.parent === undefined) registry.commit(artifact);
else appendPending(exec.parent, artifact);
return;
}
const pending = pendingByParent.get(exec?.token);
if (!pending) return;
pendingByParent.delete(exec.token);
if (result?.isError) {
for (const artifact of pending.artifacts) registry.release(artifact);
return;
}
if (exec.parent !== undefined) {
for (const artifact of pending.artifacts) appendPending(exec.parent, artifact);
return;
}
for (const artifact of pending.artifacts) registry.commit(artifact);
};
return Object.freeze({ definition, onResult });
}
/** Register the file-return tool without a per-request Gate. */
export function installOutboundArtifactTool(ctx, { registry = outboundArtifactRegistry } = {}) {
if (typeof ctx?.tools?.register !== 'function'
|| typeof ctx?.systemPrompt?.section !== 'function'
|| typeof ctx?.on !== 'function') return false;
const tool = createOutboundArtifactTool({ registry });
ctx.tools.register(tool.definition);
ctx.on('tools/result', tool.onResult);
ctx.on('session/event', (session, event) => registry.observeSessionEvent(session, event));
ctx.on('session/disposed', (session) => registry.disposeSession(session));
ctx.systemPrompt.section({
name: 'dsh-im:return-file',
order: 115,
text: `When the user asks to receive a file, call ${OUTBOUND_ARTIFACT_TOOL} with its path. Existing files can be sent directly; do not recreate or rename a file solely for delivery.`,
});
return true;
}
/** Materialize the registered snapshot for the channel provider. */
export async function materializeOutboundArtifact(artifact, {
signal,
} = {}) {
if (signal?.aborted) {
cleanupArtifactStorage(artifact);
signal.throwIfAborted();
}
if (artifact?.kind !== ARTIFACT_KIND || artifact.schemaVersion !== 1
|| !artifactStorage.has(artifact)
|| typeof artifact.digest !== 'string'
|| !Number.isSafeInteger(artifact.size) || artifact.size < 0) {
throw artifactError('artifact-invalid', 'The file registration is invalid.');
}
const storage = artifactStorage.get(artifact);
storage.materializing += 1;
let handle;
let materialized = false;
try {
const noFollow = Number.isInteger(fsConstants.O_NOFOLLOW) ? fsConstants.O_NOFOLLOW : 0;
handle = await open(storage.path, fsConstants.O_RDONLY | noFollow);
const before = await handle.stat({ bigint: true });
if (!before.isFile() || before.size !== BigInt(artifact.size)) {
throw artifactError('artifact-invalid', 'The file registration is invalid.');
}
const bytes = await readExactArtifactFile(handle, artifact.size, {
signal,
errorCode: 'artifact-invalid',
errorMessage: 'The file registration is invalid.',
});
const after = await handle.stat({ bigint: true });
if (!sameIdentity(before, after)
|| bytes.byteLength !== artifact.size
|| sha256(bytes) !== artifact.digest) {
throw artifactError('artifact-invalid', 'The file registration is invalid.');
}
signal?.throwIfAborted();
materialized = true;
const file = Object.freeze({
artifactId: artifact.artifactId,
deliveryKey: artifact.deliveryKey,
fileName: artifact.fileName,
mediaType: artifact.mediaType,
size: artifact.size,
bytes,
});
materializedArtifactSources.set(file, artifact);
storage.materialized = true;
return file;
} catch (error) {
if (error?.code?.startsWith?.('artifact-')) throw error;
if (signal?.aborted) throw signal.reason ?? error;
throw artifactError('artifact-invalid', 'The file registration is invalid.');
} finally {
await handle?.close().catch(() => undefined);
storage.materializing -= 1;
if (!materialized || storage.cleanupRequested) {
cleanupArtifactStorage(artifact);
}
}
}
/** Keep a materialized snapshot alive until an unabortable provider call settles. */
export function trackOutboundArtifactProviderPromise(file, promise) {
const artifact = materializedArtifactSources.get(file);
if (!artifact || !artifactStorage.has(artifact)
|| !promise || typeof promise.then !== 'function') return promise;
const settlements = artifactProviderSettlements.get(artifact) ?? new Set();
const settlement = Promise.resolve(promise).then(
() => undefined,
() => undefined,
);
settlements.add(settlement);
artifactProviderSettlements.set(artifact, settlements);
void settlement.finally(() => {
settlements.delete(settlement);
if (settlements.size === 0) artifactProviderSettlements.delete(artifact);
});
return promise;
}
/** Release a claimed snapshot after its provider send reaches a terminal result. */
export function releaseOutboundArtifact(artifact) {
const storage = artifactStorage.get(artifact);
if (storage) storage.releaseRequested = true;
cleanupArtifactStorage(artifact);
}

View file

@ -0,0 +1,153 @@
export const DELIVERY_RECEIPT_SCHEMA_VERSION = 1;
const ARTIFACT_OUTCOMES = new Set(['sent', 'rejected', 'failed', 'unknown']);
const REJECTED_ARTIFACT_ERRORS = new Set([
'artifact-changed',
'artifact-context-required',
'artifact-empty',
'artifact-invalid',
'artifact-not-file',
'artifact-permission-required',
'artifact-provider-rejected',
'artifact-too-large',
'artifact-unavailable',
]);
export function providerMessageIdsFor(value) {
if (!value || typeof value !== 'object') return [];
const ids = Array.isArray(value.providerMessageIds)
? value.providerMessageIds
.filter((candidate) => (
(typeof candidate === 'string' && candidate.trim())
|| Number.isSafeInteger(candidate)
))
.map(String)
: [];
const candidates = [
value.message_id,
value.messageId,
value.id,
value.ts,
value.message?.message_id,
value.message?.messageId,
value.message?.id,
value.message?.ts,
value.key?.id,
value.data?.message_id,
];
const id = candidates.find((candidate) => (
(typeof candidate === 'string' && candidate.trim())
|| Number.isSafeInteger(candidate)
));
if (id !== undefined) ids.push(String(id));
return [...new Set(ids)];
}
function requiredString(value, name) {
if (typeof value !== 'string' || !value.trim()) {
throw new TypeError(`${name} must be a non-empty string`);
}
return value;
}
function providerIds(values) {
if (!Array.isArray(values)) throw new TypeError('providerMessageIds must be an array');
const ids = [];
const seen = new Set();
for (const value of values) {
const id = requiredString(value, 'providerMessageId');
if (seen.has(id)) continue;
seen.add(id);
ids.push(id);
}
return Object.freeze(ids);
}
function artifactResults(values) {
if (!Array.isArray(values)) throw new TypeError('artifacts must be an array');
return Object.freeze(values.map((value) => {
if (!value || typeof value !== 'object') {
throw new TypeError('artifact result must be an object');
}
const artifactId = requiredString(value.artifactId, 'artifactId');
if (!ARTIFACT_OUTCOMES.has(value.outcome)) {
throw new TypeError('artifact outcome must be sent, rejected, failed, or unknown');
}
const reason = value.reason === undefined
? undefined
: requiredString(value.reason, 'artifact reason');
return Object.freeze({
artifactId,
outcome: value.outcome,
...(reason === undefined ? {} : { reason }),
});
}));
}
export function artifactOutcomeForError(error) {
const code = typeof error === 'string' ? error : error?.code;
if (code === 'artifact-delivery-uncertain') return 'unknown';
if (REJECTED_ARTIFACT_ERRORS.has(code)) return 'rejected';
return 'failed';
}
export function createDeliveryReceipt({
deliveryId,
presentation,
providerMessageIds = [],
artifacts = [],
}) {
return Object.freeze({
schemaVersion: DELIVERY_RECEIPT_SCHEMA_VERSION,
deliveryId: requiredString(deliveryId, 'deliveryId'),
presentation: requiredString(presentation, 'presentation'),
providerMessageIds: providerIds(providerMessageIds),
artifacts: artifactResults(artifacts),
});
}
export function createArtifactFailureReceipt({
artifactId,
deliveryId,
error,
presentation = 'text-fallback',
providerMessageIds = [],
}) {
const code = typeof error === 'string' ? error : error?.code;
const reason = typeof code === 'string' && code
? code
: 'artifact-provider-failed';
return createDeliveryReceipt({
deliveryId,
presentation,
providerMessageIds,
artifacts: [{
artifactId,
outcome: artifactOutcomeForError(error),
reason,
}],
});
}
export function mergeDeliveryReceipts({ deliveryId, presentation, receipts }) {
if (!Array.isArray(receipts) || receipts.length === 0) {
throw new TypeError('receipts must contain at least one delivery receipt');
}
const messageIds = [];
const artifacts = new Map();
for (const receipt of receipts) {
if (!receipt || receipt.schemaVersion !== DELIVERY_RECEIPT_SCHEMA_VERSION) {
throw new TypeError('receipt must use DeliveryReceipt schema version 1');
}
messageIds.push(...(receipt.providerMessageIds ?? []));
for (const artifact of receipt.artifacts ?? []) {
artifacts.set(artifact.artifactId, artifact);
}
}
return createDeliveryReceipt({
deliveryId,
presentation,
providerMessageIds: messageIds,
artifacts: [...artifacts.values()],
});
}

View file

@ -28,8 +28,19 @@ import {
harnessQuestionText,
validHarnessQuestion,
} from './harness-question.mjs';
import {
materializeOutboundArtifact,
releaseOutboundArtifact,
} from './semantic/artifact.mjs';
import {
createArtifactFailureReceipt,
createDeliveryReceipt,
mergeDeliveryReceipts,
providerMessageIdsFor,
} from './semantic/delivery.mjs';
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
const FILE_ONLY_COMPLETION_TEXT = '任务已完成。';
function cleanText(value) {
return typeof value === 'string' ? value.trim() : '';
@ -42,6 +53,42 @@ function canClaimInteractionReply(message, pending, senderId) {
&& Boolean(cleanText(message.content));
}
function artifactFailureText(fileName, error, descriptor) {
const name = String(fileName ?? '结果文件')
.replace(/[\r\n]+/g, ' ')
.trim()
.slice(0, 255) || '结果文件';
switch (error?.code) {
case 'artifact-delivery-uncertain':
return `结果文件「${name}」的发送结果未能确认,请先检查聊天内是否已收到,不要立即重试。`;
case 'artifact-permission-required':
if (descriptor?.key === 'slack') {
return `结果文件「${name}」已生成,但 Slack 应用缺少 files:write 权限。请更新 Manifest、重新安装应用并重新连接机器人后重试。`;
}
if (descriptor?.key === 'discord') {
return `结果文件「${name}」已生成,但机器人缺少 Discord 的 Send Messages、Attach Files 或 Read Message History 权限。`;
}
if (descriptor?.key === 'telegram') {
return `结果文件「${name}」已生成,但 Telegram 不允许机器人在当前聊天发送文档,请检查聊天权限。`;
}
return `结果文件「${name}」已生成,但当前机器人没有文件发送权限,请检查渠道权限。`;
case 'artifact-too-large':
return `结果文件「${name}」超过当前渠道大小上限,未发送。`;
case 'artifact-empty':
return `结果文件「${name}」为空,未发送。`;
case 'artifact-invalid':
case 'artifact-changed':
case 'artifact-unavailable':
return `结果文件「${name}」暂时无法读取或准备发送,请确认文件仍可访问后重试。`;
case 'artifact-rate-limited':
return `结果文件「${name}」暂时被当前渠道限流,未能发送,请稍后重试。`;
case 'artifact-provider-rejected':
return `结果文件「${name}」已生成,但当前渠道拒绝了该文件或文件消息。`;
default:
return `结果文件「${name}」已生成,但当前渠道暂时未能发送,请稍后重试。`;
}
}
export function createTextBridgeStatus() {
return {
messagesReceived: 0,
@ -285,6 +332,76 @@ export class TextHarnessBridge {
});
}
async #deliverArtifacts(target, replyTo, artifacts = [], baseReceipt) {
const receipts = baseReceipt ? [baseReceipt] : [];
let userVisible = Boolean(baseReceipt);
for (const artifact of artifacts) {
this.#signal?.throwIfAborted();
try {
if (typeof this.#bot.sendFile !== 'function') {
const unavailable = new Error('Native file delivery is unavailable');
unavailable.code = 'artifact-provider-unavailable';
throw unavailable;
}
const file = await materializeOutboundArtifact(artifact, {
signal: this.#signal,
});
this.#signal?.throwIfAborted();
const result = await this.#bot.sendFile(target, file);
receipts.push(createDeliveryReceipt({
deliveryId: file.deliveryKey,
presentation: `${this.#descriptor.key}-file`,
providerMessageIds: providerMessageIdsFor(result),
artifacts: [{ artifactId: file.artifactId, outcome: 'sent' }],
}));
userVisible = true;
this.#status.artifactsSent = (this.#status.artifactsSent ?? 0) + 1;
} catch (error) {
if (this.#signal?.aborted) throw error;
this.#status.artifactSendErrors = (this.#status.artifactSendErrors ?? 0) + 1;
this.#logger.warn?.(
`[dsh-im:${this.#descriptor.key}] result file delivery failed (${error?.code ?? 'unknown'})`,
);
let providerMessageIds = [];
let noticeSent = false;
try {
const notice = await this.#bot.sendText(
target,
artifactFailureText(artifact?.fileName, error, this.#descriptor),
);
providerMessageIds = providerMessageIdsFor(notice);
noticeSent = true;
} catch {
this.#logger.warn?.(
`[dsh-im:${this.#descriptor.key}] unable to send the safe result-file failure notice`,
);
}
const failureReceipt = createArtifactFailureReceipt({
artifactId: artifact?.artifactId ?? 'unknown',
deliveryId: artifact?.deliveryKey ?? artifact?.artifactId ?? 'unknown',
error,
providerMessageIds,
});
receipts.push(failureReceipt);
if (noticeSent || failureReceipt.artifacts[0]?.outcome === 'unknown') userVisible = true;
} finally {
releaseOutboundArtifact(artifact);
}
}
const receipt = receipts.length === 0
? null
: receipts.length === 1
? receipts[0]
: mergeDeliveryReceipts({
deliveryId: replyTo,
presentation: baseReceipt
? `${this.#descriptor.key}-text-and-files`
: `${this.#descriptor.key}-files`,
receipts,
});
return { receipt, userVisible };
}
async #process(message, messageId, senderId, conversationKey, {
alreadyRecorded = false,
} = {}) {
@ -386,7 +503,7 @@ export class TextHarnessBridge {
const content = hasImages
? await promptContentForMessage(message, { signal: this.#signal })
: undefined;
const { answer } = await askInWorkspaceSession({
const { answer, artifacts = [] } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key: conversationKey,
@ -412,10 +529,23 @@ export class TextHarnessBridge {
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
},
});
const visibleAnswer = !cleanText(answer) && artifacts.length > 0
? FILE_ONLY_COMPLETION_TEXT
: answer;
let textDeliveryError = null;
let textReceipt = null;
if (stream) {
try {
await stream.finish(answer);
const result = await stream.finish(visibleAnswer);
streamFinished = true;
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: `${this.#descriptor.key}-stream`,
providerMessageIds: [
...providerMessageIdsFor(stream),
...providerMessageIdsFor(result),
],
});
} catch (error) {
stream.cancel?.();
this.#logger.warn?.(
@ -424,10 +554,26 @@ export class TextHarnessBridge {
);
}
}
if (!streamFinished) await this.#bot.sendText(target, answer);
if (!streamFinished) {
try {
const result = await this.#bot.sendText(target, visibleAnswer);
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: `${this.#descriptor.key}-text`,
providerMessageIds: providerMessageIdsFor(result),
});
} catch (error) {
textDeliveryError = error;
}
}
// A failed final text must not discard an already registered result file.
// Settle the independent attachment path before surfacing the text error.
const delivery = await this.#deliverArtifacts(target, messageId, artifacts, textReceipt);
if (textDeliveryError && !delivery.userVisible) throw textDeliveryError;
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
return delivery.receipt;
} catch (error) {
if (error?.code === 'turn-stopped') {
if (stream) {

View file

@ -59,9 +59,23 @@ export async function askInWorkspaceSession({
return { sessionId, session };
});
if (!binding) continue;
const artifacts = [];
const originalOnArtifact = typeof askOptions === 'object'
&& typeof askOptions?.onArtifact === 'function'
? askOptions.onArtifact
: null;
const artifactOptions = typeof askOptions === 'number'
? { timeoutMs: askOptions }
: { ...askOptions };
artifactOptions.onArtifact = async (artifact) => {
artifacts.push(artifact);
await originalOnArtifact?.(artifact);
};
const answer = await binding.session.ask(content ?? text, artifactOptions);
return {
sessionId: binding.sessionId,
answer: await binding.session.ask(content ?? text, askOptions),
answer,
...(artifacts.length > 0 ? { artifacts } : {}),
};
} catch (error) {
if (error?.code !== WORKSPACE_SESSION_STALE) throw error;

View file

@ -18,6 +18,7 @@ oauth_config:
- app_mentions:read
- chat:write
- files:read
- files:write
- im:history
settings:
event_subscriptions:

View file

@ -5,6 +5,8 @@ const SLACK_FILE_HOST = 'files.slack.com';
const LEGACY_SLACK_FILE_HOST = 'slack.com';
const SLACK_FILE_PATH_PREFIX = '/files-pri/';
const SLACK_FILE_HOSTS = Object.freeze([SLACK_FILE_HOST]);
const SLACK_UPLOAD_PATH_PREFIX = '/upload/';
const DEFAULT_FILE_UPLOAD_TIMEOUT_MS = 120_000;
function cleanString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
@ -15,6 +17,47 @@ function requestSignal(signal, timeoutMs) {
return signal ? AbortSignal.any([signal, timeout]) : timeout;
}
function abortReason(signal) {
return signal?.reason instanceof Error
? signal.reason
: new DOMException('The operation was aborted', 'AbortError');
}
function positiveTimeout(value, name) {
if (!Number.isInteger(value) || value < 1) throw new TypeError(`${name} must be a positive integer`);
return value;
}
function preserveProviderMetadata(target, source) {
if (source?.providerCode !== undefined) target.providerCode = source.providerCode;
if (Number.isInteger(source?.status)) target.status = source.status;
return target;
}
function slackArtifactPreparationError(cause) {
let code = 'artifact-provider-failed';
let message = 'Slack file upload preparation failed';
if (cause?.code === 'slack-file_upload_size_restricted') {
code = 'artifact-too-large';
message = 'Slack rejected the result file because it is too large';
} else if (cause?.code === 'slack-missing-scope') {
code = 'artifact-permission-required';
message = 'Slack file delivery requires the files:write scope';
} else if (cause?.code?.startsWith?.('slack-')) {
code = 'artifact-provider-rejected';
message = 'Slack rejected file upload preparation';
}
const error = new Error(message, { cause });
error.code = code;
return preserveProviderMetadata(error, cause);
}
function uncertainSlackDelivery(cause) {
const error = new Error('Slack file completion result is uncertain', { cause });
error.code = 'artifact-delivery-uncertain';
return preserveProviderMetadata(error, cause);
}
function isRedirectStatus(status) {
return Number.isInteger(status) && status >= 300 && status < 400;
}
@ -43,6 +86,17 @@ function secureSlackFileUrl(value) {
return url;
}
function secureSlackUploadUrl(value) {
const url = new URL(value);
if (url.protocol !== 'https:' || url.username || url.password
|| (url.port && url.port !== '443') || url.hostname !== SLACK_FILE_HOST
|| !url.pathname.startsWith(SLACK_UPLOAD_PATH_PREFIX)) {
throw new Error('Slack returned an unsafe file upload URL');
}
url.hash = '';
return url;
}
function redirectUrl(response, source) {
const location = response?.headers?.get?.('location');
if (!location) return null;
@ -111,6 +165,7 @@ function safeOutgoingText(value, { trim = true } = {}) {
function apiFailure(method, payload, tokenKind) {
const reason = cleanString(payload?.error) ?? 'unknown_error';
const error = new Error(`Slack ${method} failed: ${reason.replaceAll('_', ' ')}`);
error.providerCode = reason;
if (['invalid_auth', 'not_authed', 'token_revoked', 'account_inactive'].includes(reason)) {
error.code = tokenKind === 'app' ? 'slack-invalid-app-token' : 'slack-invalid-bot-token';
} else if (reason === 'missing_scope') {
@ -141,8 +196,15 @@ export class SlackApi {
#fetch;
#baseUrl;
#botScopes = null;
#fileUploadTimeoutMs;
constructor({ botToken, appToken, fetchImpl = fetch, baseUrl = DEFAULT_BASE_URL }) {
constructor({
botToken,
appToken,
fetchImpl = fetch,
baseUrl = DEFAULT_BASE_URL,
fileUploadTimeoutMs = DEFAULT_FILE_UPLOAD_TIMEOUT_MS,
}) {
if (botToken !== undefined && !validSlackBotToken(botToken)) {
throw new TypeError('Slack Bot Token is invalid');
}
@ -155,6 +217,7 @@ export class SlackApi {
this.#appToken = appToken?.trim();
this.#fetch = fetchImpl;
this.#baseUrl = new URL(baseUrl);
this.#fileUploadTimeoutMs = positiveTimeout(fileUploadTimeoutMs, 'fileUploadTimeoutMs');
}
authTest(options = {}) {
@ -239,6 +302,99 @@ export class SlackApi {
});
}
async uploadFile({ channelId, threadTs, file, signal }) {
if (!file || typeof file !== 'object'
|| typeof file.fileName !== 'string' || !file.fileName
|| !Buffer.isBuffer(file.bytes)) {
throw new TypeError('A Slack file is required');
}
if (signal?.aborted) throw abortReason(signal);
const uploadSignal = requestSignal(signal, this.#fileUploadTimeoutMs);
let ticket;
try {
ticket = await this.#request('files.getUploadURLExternal', {
tokenKind: 'bot',
signal: uploadSignal,
timeoutMs: this.#fileUploadTimeoutMs,
body: { filename: file.fileName, length: file.bytes.byteLength },
});
} catch (error) {
if (signal?.aborted) throw abortReason(signal);
throw slackArtifactPreparationError(error);
}
let fileId;
let uploadUrl;
try {
fileId = requiredString(ticket?.file_id, 'file id');
uploadUrl = secureSlackUploadUrl(ticket?.upload_url);
} catch (error) {
throw slackArtifactPreparationError(error);
}
let uploaded;
try {
uploaded = await this.#fetch(uploadUrl, {
method: 'POST',
headers: { 'content-type': file.mediaType ?? 'application/octet-stream' },
body: file.bytes,
signal: uploadSignal,
redirect: 'error',
});
} catch (error) {
if (signal?.aborted) throw abortReason(signal);
throw slackArtifactPreparationError(error);
}
if (!Number.isInteger(uploaded?.status)) {
throw slackArtifactPreparationError(new Error('Slack file upload returned an invalid response'));
}
if (uploaded.status !== 200) {
await cancelResponseBody(uploaded);
const error = new Error(`Slack file upload failed with HTTP ${uploaded.status}`);
error.status = uploaded.status;
if (uploaded.status === 413) {
error.code = 'artifact-too-large';
} else if (uploaded.status === 429) {
error.code = 'artifact-rate-limited';
} else {
error.code = 'artifact-provider-rejected';
}
throw error;
}
await cancelResponseBody(uploaded);
let completionBody;
try {
completionBody = {
files: [{ id: fileId, title: file.fileName }],
channel_id: slackId(channelId, 'channel id'),
...(threadTs ? { thread_ts: requiredString(threadTs, 'thread timestamp') } : {}),
};
} catch (error) {
throw slackArtifactPreparationError(error);
}
if (uploadSignal.aborted) {
if (signal?.aborted) throw abortReason(signal);
throw slackArtifactPreparationError(uploadSignal.reason);
}
try {
const completed = await this.#request('files.completeUploadExternal', {
tokenKind: 'bot',
signal: uploadSignal,
timeoutMs: this.#fileUploadTimeoutMs,
retry: false,
body: completionBody,
});
if (!Array.isArray(completed?.files)
|| !completed.files.some((entry) => cleanString(entry?.id) === fileId)) {
throw new Error('Slack file completion did not confirm the uploaded file');
}
return completed;
} catch (error) {
if (signal?.aborted) throw abortReason(signal);
if (error?.code === 'slack-missing-scope') throw slackArtifactPreparationError(error);
throw uncertainSlackDelivery(error);
}
}
async downloadFile({ url, signal, maxBytes }) {
if (!this.#botToken) throw new TypeError('Slack bot token is required for file download');
const target = secureSlackFileUrl(url);
@ -276,18 +432,21 @@ export class SlackApi {
}) {
const token = tokenKind === 'app' ? this.#appToken : this.#botToken;
if (!token) throw new TypeError(`Slack ${tokenKind} token is required for ${method}`);
const formEncoded = method === 'files.getUploadURLExternal';
let response;
try {
response = await this.#fetch(new URL(method, this.#baseUrl), {
method: 'POST',
headers: {
authorization: `Bearer ${token}`,
'content-type': body === undefined
'content-type': body === undefined || formEncoded
? 'application/x-www-form-urlencoded;charset=utf-8'
: 'application/json;charset=utf-8',
'user-agent': 'DeepSeek-Harness-dsh-im (https://github.com/xmanrui/dsh-im, 0.2.2)',
},
...(body === undefined ? {} : { body: JSON.stringify(body) }),
...(body === undefined ? {} : {
body: formEncoded ? new URLSearchParams(body).toString() : JSON.stringify(body),
}),
signal: requestSignal(signal, timeoutMs),
redirect: 'error',
});
@ -312,7 +471,11 @@ export class SlackApi {
await delay(Math.min(10_000, Math.max(100, seconds * 1_000)), signal);
return this.#request(method, { tokenKind, body, signal, timeoutMs, retry: false });
}
if (!response.ok || payload?.ok !== true) throw apiFailure(method, payload, tokenKind);
if (!response.ok || payload?.ok !== true) {
const error = apiFailure(method, payload, tokenKind);
error.status = response.status;
throw error;
}
return payload;
}
}

View file

@ -120,6 +120,7 @@ async function createSlackMessageStream({ api, target, signal, logger }) {
let inFlight = null;
let broken = false;
let closed = false;
const providerMessageIds = [ts];
const appendLatest = async (text) => {
const next = splitMessageText(text, SLACK_MESSAGE_LIMIT)[0] ?? '';
@ -150,6 +151,10 @@ async function createSlackMessageStream({ api, target, signal, logger }) {
};
return {
messageId: ts,
get providerMessageIds() {
return [...providerMessageIds];
},
update(text) {
if (closed || broken || typeof text !== 'string' || !text.trim() || isToolProgress(text)) return;
pending = text;
@ -173,12 +178,13 @@ async function createSlackMessageStream({ api, target, signal, logger }) {
await api.updateMessage({ channelId: target.channelId, ts, text: first, signal });
}
for (const chunk of chunks.slice(1)) {
await api.postMessage({
const result = await api.postMessage({
channelId: target.channelId,
threadTs: target.threadTs,
text: chunk,
signal,
});
if (typeof result?.ts === 'string' && result.ts) providerMessageIds.push(result.ts);
}
},
cancel() {
@ -191,7 +197,7 @@ async function createSlackMessageStream({ api, target, signal, logger }) {
};
}
class SlackBotClient {
export class SlackBotClient {
#api;
#signal;
#logger;
@ -204,16 +210,17 @@ class SlackBotClient {
async sendText(target, text) {
const chunks = splitMessageText(text, SLACK_MESSAGE_LIMIT);
let result = null;
const providerMessageIds = [];
for (const chunk of chunks) {
result = await this.#api.postMessage({
const result = await this.#api.postMessage({
channelId: target.channelId,
threadTs: target.threadTs,
text: chunk,
signal: this.#signal,
});
if (typeof result?.ts === 'string' && result.ts) providerMessageIds.push(result.ts);
}
return result;
return { providerMessageIds };
}
openStream(target) {
@ -224,6 +231,15 @@ class SlackBotClient {
logger: this.#logger,
});
}
sendFile(target, file) {
return this.#api.uploadFile({
channelId: target.channelId,
threadTs: target.threadTs,
file,
signal: this.#signal,
});
}
}
export function createSlackRuntimeStatus() {

View file

@ -2,6 +2,7 @@ import { fetchImageBuffer } from '../shared/image-prompt.mjs';
const DEFAULT_BASE_URL = 'https://api.telegram.org/';
const TELEGRAM_FILE_HOSTS = Object.freeze(['api.telegram.org']);
const DEFAULT_FILE_UPLOAD_TIMEOUT_MS = 120_000;
function cleanString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
@ -12,6 +13,58 @@ function requestSignal(signal, timeoutMs) {
return signal ? AbortSignal.any([signal, timeout]) : timeout;
}
function abortReason(signal) {
return signal?.reason instanceof Error
? signal.reason
: new DOMException('The operation was aborted', 'AbortError');
}
function positiveTimeout(value, name) {
if (!Number.isInteger(value) || value < 1) throw new TypeError(`${name} must be a positive integer`);
return value;
}
function preserveProviderMetadata(target, source) {
if (source?.providerCode !== undefined) target.providerCode = source.providerCode;
if (source?.retry_after !== undefined) {
target.retry_after = source.retry_after;
target.retryAfter = source.retry_after;
}
if (Number.isInteger(source?.status)) target.status = source.status;
return target;
}
function telegramArtifactProviderError(cause) {
const providerCode = Number(cause?.providerCode);
const status = Number(cause?.status);
const message = cleanString(cause?.message) ?? '';
let code = 'artifact-provider-rejected';
let summary = 'Telegram rejected the document.';
if (providerCode === 401 || providerCode === 403 || status === 401 || status === 403) {
code = 'artifact-permission-required';
summary = 'Telegram denied permission to send the document.';
} else if (providerCode === 413 || status === 413
|| /(?:file|request|entity).{0,20}(?:too (?:big|large)|size limit)|too (?:big|large)/i.test(message)) {
code = 'artifact-too-large';
summary = 'The document exceeds Telegram\'s size limit.';
} else if (providerCode === 429 || status === 429) {
code = 'artifact-rate-limited';
summary = 'Telegram rate-limited document delivery.';
} else if (providerCode >= 500 || status >= 500) {
code = 'artifact-delivery-uncertain';
summary = 'Telegram document delivery result is uncertain.';
}
const error = new Error(summary, { cause });
error.code = code;
return preserveProviderMetadata(error, cause);
}
function uncertainTelegramDelivery(cause) {
const error = new Error('Telegram document delivery result is uncertain', { cause });
error.code = 'artifact-delivery-uncertain';
return preserveProviderMetadata(error, cause);
}
export function validTelegramToken(value) {
return typeof value === 'string' && /^\d{5,20}:[A-Za-z0-9_-]{20,}$/.test(value.trim());
}
@ -32,13 +85,20 @@ export class TelegramApi {
#token;
#fetch;
#baseUrl;
#fileUploadTimeoutMs;
constructor({ token, fetchImpl = fetch, baseUrl = DEFAULT_BASE_URL }) {
constructor({
token,
fetchImpl = fetch,
baseUrl = DEFAULT_BASE_URL,
fileUploadTimeoutMs = DEFAULT_FILE_UPLOAD_TIMEOUT_MS,
}) {
if (!validTelegramToken(token)) throw new TypeError('Telegram Bot Token is invalid');
if (typeof fetchImpl !== 'function') throw new TypeError('TelegramApi requires fetch');
this.#token = token.trim();
this.#fetch = fetchImpl;
this.#baseUrl = new URL(baseUrl);
this.#fileUploadTimeoutMs = positiveTimeout(fileUploadTimeoutMs, 'fileUploadTimeoutMs');
}
async getMe(options = {}) {
@ -108,6 +168,43 @@ export class TelegramApi {
}, { signal });
}
async sendDocument({ chatId, file, replyToMessageId, messageThreadId, signal }) {
if (!file || typeof file !== 'object'
|| typeof file.fileName !== 'string' || !file.fileName
|| !Buffer.isBuffer(file.bytes)) {
throw new TypeError('A Telegram document is required');
}
const payload = new FormData();
payload.append('chat_id', String(chatId));
payload.append(
'document',
new Blob([file.bytes], { type: file.mediaType ?? 'application/octet-stream' }),
file.fileName,
);
if (replyToMessageId) {
payload.append('reply_parameters', JSON.stringify({
message_id: replyToMessageId,
allow_sending_without_reply: true,
}));
}
if (messageThreadId) payload.append('message_thread_id', String(messageThreadId));
if (signal?.aborted) throw abortReason(signal);
const uploadSignal = requestSignal(signal, this.#fileUploadTimeoutMs);
try {
return await this.#call('sendDocument', payload, {
signal: uploadSignal,
timeoutMs: this.#fileUploadTimeoutMs,
multipart: true,
});
} catch (error) {
if (signal?.aborted) throw abortReason(signal);
if (error?.code?.startsWith?.('telegram-')) {
throw telegramArtifactProviderError(error);
}
throw uncertainTelegramDelivery(error);
}
}
async editMessageText({ chatId, messageId, text, signal }) {
return this.#call('editMessageText', {
chat_id: chatId,
@ -144,15 +241,15 @@ export class TelegramApi {
return this.#call('setChatMenuButton', { menu_button: menuButton }, { signal });
}
async #call(method, payload, { signal, timeoutMs = 15_000 } = {}) {
async #call(method, payload, { signal, timeoutMs = 15_000, multipart = false } = {}) {
const url = new URL(this.#baseUrl);
url.pathname = `${url.pathname.replace(/\/$/, '')}/bot${this.#token}/${method}`;
let response;
try {
response = await this.#fetch(url, {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify(payload),
...(multipart ? {} : { headers: { 'content-type': 'application/json' } }),
body: multipart ? payload : JSON.stringify(payload),
signal: requestSignal(signal, timeoutMs),
redirect: 'error',
});
@ -164,12 +261,21 @@ export class TelegramApi {
try {
body = await response.json();
} catch {
throw new Error(`Telegram ${method} returned invalid JSON`);
const error = new Error(`Telegram ${method} returned invalid JSON`);
error.status = response?.status;
throw error;
}
if (!response.ok || body?.ok !== true) {
const description = cleanString(body?.description);
const error = new Error(description ?? `Telegram ${method} failed`);
error.code = Number.isInteger(body?.error_code) ? `telegram-${body.error_code}` : 'telegram-api-error';
error.status = response.status;
if (Number.isInteger(body?.error_code)) error.providerCode = body.error_code;
const retryAfter = Number(body?.parameters?.retry_after);
if (Number.isFinite(retryAfter) && retryAfter >= 0) {
error.retry_after = retryAfter;
error.retryAfter = retryAfter;
}
throw error;
}
return body.result;

View file

@ -146,7 +146,7 @@ export function telegramInboundAllowed(message, {
&& allowedPrivateUserIds.has(String(message.senderId));
}
class TelegramBotClient {
export class TelegramBotClient {
#api;
#signal;
@ -157,17 +157,20 @@ class TelegramBotClient {
async sendText(target, text) {
const chunks = splitMessageText(text, 4_000);
let result = null;
const providerMessageIds = [];
for (const [index, chunk] of chunks.entries()) {
result = await this.#api.sendMessage({
const result = await this.#api.sendMessage({
chatId: target.chatId,
text: chunk,
replyToMessageId: index === 0 ? target.replyToMessageId : undefined,
messageThreadId: target.messageThreadId,
signal: this.#signal,
});
if (Number.isSafeInteger(result?.message_id)) {
providerMessageIds.push(String(result.message_id));
}
}
return result;
return { providerMessageIds };
}
sendTyping(target) {
@ -178,6 +181,16 @@ class TelegramBotClient {
});
}
sendFile(target, file) {
return this.#api.sendDocument({
chatId: target.chatId,
file,
replyToMessageId: target.replyToMessageId,
messageThreadId: target.messageThreadId,
signal: this.#signal,
});
}
async openStream(target) {
const stream = createEditableMessageStream({
limit: 4_000,
@ -203,6 +216,7 @@ class TelegramBotClient {
messageThreadId: target.messageThreadId,
signal: this.#signal,
}),
messageIdForResult: (message) => message?.message_id,
});
return stream.start();
}

View file

@ -27,6 +27,18 @@ import {
promptContentForMessage,
} from '../shared/image-prompt.mjs';
import { rememberConnectionTestTarget } from '../shared/connection-test.mjs';
import {
materializeOutboundArtifact,
releaseOutboundArtifact,
trackOutboundArtifactProviderPromise,
} from '../shared/semantic/artifact.mjs';
import {
createArtifactFailureReceipt,
createDeliveryReceipt,
mergeDeliveryReceipts,
} from '../shared/semantic/delivery.mjs';
const DEFAULT_FILE_UPLOAD_TIMEOUT_MS = 120_000;
const HELP_TEXT = [
'企业微信机器人已连接 DeepSeek Harness。',
@ -212,6 +224,88 @@ function progressText(update) {
return update?.text;
}
function artifactFailureText(fileName, error) {
const name = String(fileName ?? '结果文件').replace(/[\r\n]+/g, ' ').trim() || '结果文件';
switch (error?.code) {
case 'artifact-delivery-uncertain':
return `结果文件「${name}」的发送结果未能确认,请先检查聊天内是否已收到,不要立即重试。`;
case 'artifact-permission-required':
return `结果文件「${name}」已生成,但企业微信智能机器人缺少素材上传或文件消息能力,请检查机器人权限。`;
case 'artifact-too-large':
return `结果文件「${name}」超过当前企业微信机器人可发送的文件大小,未发送。`;
case 'artifact-empty':
return `结果文件「${name}」为空,企业微信不允许发送空文件。`;
case 'artifact-changed':
case 'artifact-invalid':
case 'artifact-unavailable':
return `结果文件「${name}」暂时无法读取或准备发送,请确认文件仍可访问后重试。`;
case 'artifact-rate-limited':
return `结果文件「${name}」暂时被企业微信限流,未能发送,请稍后重试。`;
case 'artifact-provider-rejected':
return `结果文件「${name}」已生成,但企业微信拒绝了该文件或文件消息。`;
default:
return `结果文件「${name}」已生成,但暂时未能通过企业微信发送,请稍后重试。`;
}
}
function abortReason(signal) {
return signal?.reason instanceof Error
? signal.reason
: new DOMException('The operation was aborted', 'AbortError');
}
function waitWithSignal(promise, signal) {
if (!signal) return promise;
signal.throwIfAborted();
return new Promise((resolve, reject) => {
let settled = false;
const finish = (callback, value) => {
if (settled) return;
settled = true;
signal.removeEventListener('abort', onAbort);
callback(value);
};
const onAbort = () => finish(reject, abortReason(signal));
signal.addEventListener('abort', onAbort, { once: true });
Promise.resolve(promise).then(
(value) => finish(resolve, value),
(error) => finish(reject, error),
);
if (signal.aborted) onAbort();
});
}
function wecomArtifactError(error, { dispatched = false } = {}) {
if (error?.code?.startsWith?.('artifact-')) return error;
const status = Number(error?.httpStatus ?? error?.status ?? error?.response?.status);
const providerCode = Number(error?.providerCode ?? error?.errcode ?? error?.body?.errcode);
const wrapped = new Error('Enterprise WeChat file delivery failed', { cause: error });
if (status === 401 || status === 403 || providerCode === 48002) {
wrapped.code = 'artifact-permission-required';
} else if (status === 413) {
wrapped.code = 'artifact-too-large';
} else if (status === 429 || providerCode === 45009) {
wrapped.code = 'artifact-rate-limited';
} else if (Number.isFinite(providerCode) && providerCode !== 0) {
wrapped.code = 'artifact-provider-rejected';
} else {
wrapped.code = dispatched ? 'artifact-delivery-uncertain' : 'artifact-provider-failed';
}
if (Number.isFinite(status)) wrapped.status = status;
if (Number.isFinite(providerCode)) wrapped.providerCode = providerCode;
return wrapped;
}
function answerTextForDelivery(answer, artifacts) {
if (typeof answer === 'string' && answer.trim()) return answer;
return artifacts.length > 0 ? '结果文件已生成。' : '任务已完成,但没有生成可显示的文本。';
}
function providerMessageId(result) {
return nonEmptyString(result?.body?.msgid)
?? nonEmptyString(result?.body?.message_id);
}
function canClaimInteractionReply(frame, pending) {
return pending.questions[pending.index]
&& nonEmptyString(bodyOf(frame).from?.userid) === pending.actor
@ -223,6 +317,8 @@ export function createWecomBridgeStatus() {
messagesReceived: 0,
messagesReplied: 0,
messagesRejected: 0,
artifactsSent: 0,
artifactSendErrors: 0,
lastMessageAt: null,
lastReplyAt: null,
lastRejectedAt: null,
@ -239,6 +335,7 @@ export class WecomHarnessBridge {
#replyTimeoutMs;
#generateReqId;
#signal;
#fileUploadTimeoutMs;
#queues = new Map();
#pendingInteractions = new Map();
#interactionKeys = new Map();
@ -256,12 +353,16 @@ export class WecomHarnessBridge {
logger = console,
replyTimeoutMs = 600_000,
generateStreamId = generateReqId,
fileUploadTimeoutMs = DEFAULT_FILE_UPLOAD_TIMEOUT_MS,
signal,
}) {
if (!client || typeof client.replyStream !== 'function' || typeof client.sendMessage !== 'function') {
throw new TypeError('Enterprise WeChat client is required');
}
if (!harness || !state) throw new TypeError('Harness client and state store are required');
if (!Number.isInteger(fileUploadTimeoutMs) || fileUploadTimeoutMs < 1) {
throw new TypeError('fileUploadTimeoutMs must be a positive integer');
}
this.#client = client;
this.#harness = harness;
this.#state = state;
@ -269,6 +370,7 @@ export class WecomHarnessBridge {
this.#logger = logger;
this.#replyTimeoutMs = replyTimeoutMs;
this.#generateReqId = generateStreamId;
this.#fileUploadTimeoutMs = Math.min(fileUploadTimeoutMs, DEFAULT_FILE_UPLOAD_TIMEOUT_MS);
this.#signal = signal;
this.#approvals = new HarnessApprovalQueue({ label: 'wecom', logger });
}
@ -458,10 +560,17 @@ export class WecomHarnessBridge {
}
async #sendActive(chatId, text) {
const providerMessageIds = [];
for (const chunk of splitUtf8(text)) {
this.#signal?.throwIfAborted();
await this.#client.sendMessage(chatId, { msgtype: 'markdown', markdown: { content: chunk } });
const result = await this.#client.sendMessage(
chatId,
{ msgtype: 'markdown', markdown: { content: chunk } },
);
const messageId = providerMessageId(result);
if (messageId) providerMessageIds.push(messageId);
}
return providerMessageIds;
}
async #sendImmediate(frame, chatId, text) {
@ -478,6 +587,105 @@ export class WecomHarnessBridge {
}
}
async #deliverArtifacts(chatId, replyTo, artifacts = [], baseReceipt = null) {
if (artifacts.length === 0) {
return { receipt: baseReceipt, failureNoticeVisible: false };
}
const receipts = baseReceipt ? [baseReceipt] : [];
let failureNoticeVisible = false;
for (const artifact of artifacts) {
this.#signal?.throwIfAborted();
try {
if (typeof this.#client.uploadMedia !== 'function'
|| typeof this.#client.sendMediaMessage !== 'function') {
const unavailable = new Error('Enterprise WeChat file delivery is unavailable');
unavailable.code = 'artifact-provider-unavailable';
throw unavailable;
}
const file = await materializeOutboundArtifact(artifact, {
signal: this.#signal,
});
this.#signal?.throwIfAborted();
const timeout = AbortSignal.timeout(this.#fileUploadTimeoutMs);
const waitSignal = this.#signal ? AbortSignal.any([this.#signal, timeout]) : timeout;
let uploaded;
try {
const pending = this.#client.uploadMedia(file.bytes, {
type: 'file',
filename: file.fileName,
});
trackOutboundArtifactProviderPromise(file, pending);
uploaded = await waitWithSignal(pending, waitSignal);
} catch (error) {
if (this.#signal?.aborted) throw abortReason(this.#signal);
throw wecomArtifactError(error);
}
this.#signal?.throwIfAborted();
const mediaId = nonEmptyString(uploaded?.media_id);
if (!mediaId) {
const rejected = new Error('Enterprise WeChat upload returned no media id');
rejected.code = 'artifact-provider-rejected';
throw rejected;
}
let sent;
try {
const pending = this.#client.sendMediaMessage(chatId, 'file', mediaId);
trackOutboundArtifactProviderPromise(file, pending);
sent = await waitWithSignal(pending, waitSignal);
} catch (error) {
if (this.#signal?.aborted) throw abortReason(this.#signal);
throw wecomArtifactError(error, { dispatched: true });
}
this.#signal?.throwIfAborted();
const providerCode = Number(sent?.body?.errcode ?? sent?.errcode);
if (Number.isFinite(providerCode) && providerCode !== 0) {
throw wecomArtifactError({ providerCode });
}
const messageId = providerMessageId(sent);
receipts.push(createDeliveryReceipt({
deliveryId: file.deliveryKey,
presentation: 'wecom-file',
providerMessageIds: messageId ? [messageId] : [],
artifacts: [{ artifactId: file.artifactId, outcome: 'sent' }],
}));
this.#status.artifactsSent = (this.#status.artifactsSent ?? 0) + 1;
} catch (error) {
if (this.#signal?.aborted) throw error;
this.#status.artifactSendErrors = (this.#status.artifactSendErrors ?? 0) + 1;
this.#logger.warn?.(
`[dsh-im:wecom] result file delivery failed (${error?.code ?? 'unknown'})`,
);
let providerMessageIds = [];
try {
providerMessageIds = await this.#sendActive(
chatId,
artifactFailureText(artifact?.fileName, error),
);
failureNoticeVisible = true;
} catch (noticeError) {
if (this.#signal?.aborted) throw noticeError;
this.#logger.warn?.('[dsh-im:wecom] unable to send the safe result-file failure notice');
}
receipts.push(createArtifactFailureReceipt({
artifactId: artifact?.artifactId ?? 'unknown',
deliveryId: artifact?.deliveryKey ?? artifact?.artifactId ?? 'unknown',
error,
providerMessageIds,
}));
} finally {
releaseOutboundArtifact(artifact);
}
}
return {
receipt: mergeDeliveryReceipts({
deliveryId: replyTo,
presentation: baseReceipt ? 'wecom-text-and-files' : 'wecom-files',
receipts,
}),
failureNoticeVisible,
};
}
async #process(frame, { alreadyRecorded = false, preparedMessage } = {}) {
if (this.#signal?.aborted) return;
const body = bodyOf(frame);
@ -556,7 +764,7 @@ export class WecomHarnessBridge {
const content = hasImages
? await promptContentForMessage(message, { signal: this.#signal })
: undefined;
const { answer } = await askInWorkspaceSession({
const { answer, artifacts = [] } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
@ -584,24 +792,64 @@ export class WecomHarnessBridge {
},
});
const chunks = splitUtf8(answer || '任务已完成,但没有生成可显示的文本。');
this.#signal?.throwIfAborted();
const displayAnswer = answerTextForDelivery(answer, artifacts);
const chunks = splitUtf8(displayAnswer);
let finalSent = false;
if (streamStarted && chunks.length > 0) {
try {
await this.#client.replyStream(frame, streamId, chunks[0], true);
for (const chunk of chunks.slice(1)) {
await this.#client.sendMessage(chatId, { msgtype: 'markdown', markdown: { content: chunk } });
let textReceipt = null;
let textSendError = null;
try {
if (streamStarted && chunks.length > 0) {
try {
const providerMessageIds = [];
const streamed = await this.#client.replyStream(frame, streamId, chunks[0], true);
const streamedMessageId = providerMessageId(streamed);
if (streamedMessageId) providerMessageIds.push(streamedMessageId);
for (const chunk of chunks.slice(1)) {
const sent = await this.#client.sendMessage(
chatId,
{ msgtype: 'markdown', markdown: { content: chunk } },
);
const messageId = providerMessageId(sent);
if (messageId) providerMessageIds.push(messageId);
}
finalSent = true;
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'wecom-text',
providerMessageIds,
});
} catch (error) {
this.#logger.warn?.('[dsh-im:wecom] stream finalization failed; using an active reply:', error);
}
finalSent = true;
} catch (error) {
this.#logger.warn?.('[dsh-im:wecom] stream finalization failed; using an active reply:', error);
}
if (!finalSent) {
const providerMessageIds = await this.#sendActive(chatId, displayAnswer);
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'wecom-text',
providerMessageIds,
});
}
} catch (error) {
textSendError = error;
this.#logger.warn?.(
'[dsh-im:wecom] final text delivery failed; continuing with result files:',
error,
);
}
const delivery = await this.#deliverArtifacts(chatId, messageId, artifacts, textReceipt);
const artifactDispatched = delivery.receipt?.artifacts?.some(
({ outcome }) => outcome === 'sent' || outcome === 'unknown',
);
if (textSendError && !artifactDispatched && !delivery.failureNoticeVisible) {
throw textSendError;
}
if (!finalSent) await this.#sendActive(chatId, answer);
await this.#state.markSeen(messageId);
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
return delivery.receipt;
} catch (error) {
if (error?.code === 'turn-stopped') {
if (streamStarted && streamId) {

View file

@ -1,4 +1,10 @@
import { createDecipheriv, randomBytes, randomUUID } from 'node:crypto';
import {
createCipheriv,
createDecipheriv,
createHash,
randomBytes,
randomUUID,
} from 'node:crypto';
import { fetchImageBuffer } from '../shared/image-prompt.mjs';
@ -13,6 +19,7 @@ const ILINK_APP_ID = 'bot';
const ILINK_CLIENT_VERSION = (2 << 16) | (4 << 8) | 6;
const DEFAULT_TIMEOUT_MS = 15_000;
const DEFAULT_LONG_POLL_TIMEOUT_MS = 35_000;
const WEIXIN_CDN_UPLOAD_RETRIES = 3;
const LOGIN_STATUSES = new Set([
'wait',
'scaned',
@ -30,6 +37,7 @@ export class WeixinApiError extends Error {
this.name = 'WeixinApiError';
this.code = code;
this.status = options.status;
this.providerCode = options.providerCode;
}
}
@ -37,6 +45,70 @@ function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
function safeProviderCode(value) {
const code = value === undefined || value === null ? null : String(value).trim();
return code && /^-?[A-Za-z0-9_.:-]{1,160}$/.test(code) ? code : undefined;
}
function preserveArtifactMetadata(target, source) {
if (Number.isInteger(source?.status)) target.status = source.status;
if (source?.providerCode !== undefined) target.providerCode = source.providerCode;
return target;
}
function weixinArtifactError(cause, { fallback = 'artifact-provider-rejected' } = {}) {
if (cause?.code?.startsWith?.('artifact-')) return cause;
const status = Number(cause?.status);
const providerCode = safeProviderCode(cause?.providerCode);
const providerText = providerCode ?? '';
let code = fallback;
let message = 'Weixin could not prepare the file for delivery.';
if (status === 401 || status === 403 || providerCode === '401' || providerCode === '403'
|| /(?:permission|forbidden|unauthor|access.?denied)/i.test(providerText)) {
code = 'artifact-permission-required';
message = 'Weixin denied permission to send the file.';
} else if (status === 413 || providerCode === '413'
|| /(?:too.?large|size.?limit)/i.test(providerText)) {
code = 'artifact-too-large';
message = 'The file exceeds Weixin\'s size limit.';
} else if (status === 429 || providerCode === '429'
|| /(?:rate.?limit|too.?many)/i.test(providerText)) {
code = 'artifact-rate-limited';
message = 'Weixin rate-limited file delivery.';
} else if (fallback === 'artifact-provider-rejected') {
message = 'Weixin rejected the file message.';
}
const error = new Error(message, { cause });
error.code = code;
return preserveArtifactMetadata(error, cause);
}
function uncertainWeixinDelivery(cause) {
const error = new Error('Weixin file delivery result is uncertain', { cause });
error.code = 'artifact-delivery-uncertain';
return preserveArtifactMetadata(error, cause);
}
function rejectedProviderResponse(value) {
if (!value || typeof value !== 'object') return null;
for (const field of ['ret', 'errcode']) {
if (value[field] !== undefined && value[field] !== 0 && value[field] !== '0') {
return safeProviderCode(value[field]) ?? 'rejected';
}
}
return null;
}
function classifyWeixinFinalDeliveryError(error, signal) {
if (signal?.aborted) throw abortError(signal);
const status = Number(error?.status);
if (error?.code === 'network-error' || error?.code === 'timeout'
|| error?.code === 'invalid-response' || (status >= 500 && status < 600)) {
return uncertainWeixinDelivery(error);
}
return weixinArtifactError(error);
}
function strictBase64(value) {
const text = nonEmptyString(value);
if (!text || text.length % 4 !== 0 || !/^[A-Za-z0-9+/]+={0,2}$/.test(text)) return null;
@ -188,6 +260,89 @@ function baseInfo() {
};
}
function aesEcbPaddedSize(size) {
return Math.ceil((size + 1) / 16) * 16;
}
function trustedWeixinCdnUploadUrl(value) {
let url;
try {
url = new URL(value);
} catch {
throw new WeixinApiError('invalid-upload-url', '微信服务返回了无效的文件上传地址。');
}
if (url.protocol !== 'https:' || url.hostname !== WEIXIN_CDN_HOST
|| (url.port && url.port !== '443') || url.pathname !== '/c2c/upload'
|| url.username || url.password) {
throw new WeixinApiError('untrusted-upload-url', '微信服务返回了不受信任的文件上传地址。');
}
url.hash = '';
return url;
}
function weixinCdnUploadUrl(response, fileKey) {
const fullUrl = nonEmptyString(response?.upload_full_url);
if (fullUrl) return trustedWeixinCdnUploadUrl(fullUrl);
const uploadParam = nonEmptyString(response?.upload_param);
if (!uploadParam) {
throw new WeixinApiError('missing-upload-url', '微信服务没有返回文件上传地址。');
}
const url = new URL(`${WEIXIN_CDN_BASE_URL}/upload`);
url.searchParams.set('encrypted_query_param', uploadParam);
url.searchParams.set('filekey', fileKey);
return trustedWeixinCdnUploadUrl(url);
}
function encryptWeixinUpload(bytes, key) {
const cipher = createCipheriv('aes-128-ecb', key, null);
return Buffer.concat([cipher.update(bytes), cipher.final()]);
}
async function uploadWeixinCdn(fetchImpl, url, ciphertext, { signal } = {}) {
let lastError;
for (let attempt = 1; attempt <= WEIXIN_CDN_UPLOAD_RETRIES; attempt += 1) {
signal?.throwIfAborted();
try {
const response = await fetchImpl(url, {
method: 'POST',
headers: { 'content-type': 'application/octet-stream' },
body: ciphertext,
signal: signal
? AbortSignal.any([signal, AbortSignal.timeout(60_000)])
: AbortSignal.timeout(60_000),
redirect: 'error',
});
if (response.status >= 400 && response.status < 500) {
throw new WeixinApiError(
'upload-rejected',
`微信文件上传被拒绝(HTTP ${response.status})。`,
{ status: response.status },
);
}
if (response.status !== 200) {
throw new WeixinApiError(
'upload-failed',
`微信文件上传失败(HTTP ${response.status})。`,
{ status: response.status },
);
}
const downloadParam = nonEmptyString(response.headers.get('x-encrypted-param'));
await response.body?.cancel?.().catch(() => undefined);
if (!downloadParam) {
throw new WeixinApiError('invalid-upload-response', '微信文件上传响应缺少下载参数。');
}
return downloadParam;
} catch (error) {
if (signal?.aborted) throw abortError(signal);
if (error instanceof WeixinApiError
&& (error.code === 'upload-rejected' || error.status < 500)) throw error;
lastError = error;
}
}
if (lastError instanceof WeixinApiError) throw lastError;
throw new WeixinApiError('upload-failed', '微信文件上传失败。', { cause: lastError });
}
function abortError(signal) {
if (signal?.reason instanceof Error) return signal.reason;
return new DOMException('The operation was aborted', 'AbortError');
@ -349,6 +504,117 @@ export function createWeixinApi({ fetchImpl = fetch } = {}) {
return true;
},
async sendFile({ baseUrl, token, toUserId, file, contextToken, runId, signal }) {
const recipient = nonEmptyString(toUserId);
if (!recipient || !file || typeof file !== 'object'
|| typeof file.fileName !== 'string' || !file.fileName
|| !Buffer.isBuffer(file.bytes)) {
throw new TypeError('toUserId and a file are required');
}
signal?.throwIfAborted();
const fileKey = randomBytes(16).toString('hex');
const aesKey = randomBytes(16);
const rawMd5 = createHash('md5').update(file.bytes).digest('hex');
let upload;
try {
upload = await requestJson(fetchImpl, {
method: 'POST',
baseUrl,
endpoint: 'ilink/bot/getuploadurl',
token,
signal,
body: {
filekey: fileKey,
media_type: 3,
to_user_id: recipient,
rawsize: file.bytes.byteLength,
rawfilemd5: rawMd5,
filesize: aesEcbPaddedSize(file.bytes.byteLength),
no_need_thumb: true,
aeskey: aesKey.toString('hex'),
base_info: baseInfo(),
},
});
} catch (error) {
if (signal?.aborted) throw abortError(signal);
const status = Number(error?.status);
const fallback = error?.code === 'http-error' && status >= 400 && status < 500
? 'artifact-provider-rejected'
: 'artifact-provider-failed';
throw weixinArtifactError(error, { fallback });
}
const uploadRejection = rejectedProviderResponse(upload);
if (uploadRejection) {
throw weixinArtifactError(new WeixinApiError(
'upload-url-rejected',
'微信服务拒绝了文件上传请求。',
{ providerCode: uploadRejection },
));
}
const uploadUrl = weixinCdnUploadUrl(upload, fileKey);
const ciphertext = encryptWeixinUpload(file.bytes, aesKey);
let downloadParam;
try {
downloadParam = await uploadWeixinCdn(fetchImpl, uploadUrl, ciphertext, { signal });
} catch (error) {
if (signal?.aborted) throw abortError(signal);
const status = Number(error?.status);
const fallback = error?.code === 'upload-rejected' || (status >= 400 && status < 500)
? 'artifact-provider-rejected'
: 'artifact-provider-failed';
throw weixinArtifactError(error, { fallback });
}
signal?.throwIfAborted();
const deliverySeed = nonEmptyString(file.deliveryKey) ?? nonEmptyString(file.artifactId)
?? randomUUID();
const clientId = `dsh-weixin-${createHash('sha256').update(deliverySeed).digest('hex').slice(0, 32)}`;
let response;
try {
response = await requestJson(fetchImpl, {
method: 'POST',
baseUrl,
endpoint: 'ilink/bot/sendmessage',
token,
signal,
body: {
msg: {
from_user_id: '',
to_user_id: recipient,
client_id: clientId,
message_type: 2,
message_state: 2,
item_list: [{
type: 4,
file_item: {
media: {
encrypt_query_param: downloadParam,
aes_key: Buffer.from(aesKey.toString('hex')).toString('base64'),
encrypt_type: 1,
},
file_name: file.fileName,
len: String(file.bytes.byteLength),
},
}],
...(nonEmptyString(contextToken) ? { context_token: contextToken.trim() } : {}),
...(nonEmptyString(runId) ? { run_id: runId.trim() } : {}),
},
base_info: baseInfo(),
},
});
} catch (error) {
throw classifyWeixinFinalDeliveryError(error, signal);
}
const sendRejection = rejectedProviderResponse(response);
if (sendRejection) {
throw weixinArtifactError(new WeixinApiError(
'send-rejected',
'微信服务拒绝了文件消息。',
{ providerCode: sendRejection },
));
}
return { messageId: clientId };
},
async notifyStart({ baseUrl, token, signal }) {
const response = await requestJson(fetchImpl, {
method: 'POST',

View file

@ -32,6 +32,16 @@ import {
promptContentForMessage,
} from '../shared/image-prompt.mjs';
import { rememberConnectionTestTarget } from '../shared/connection-test.mjs';
import {
materializeOutboundArtifact,
releaseOutboundArtifact,
} from '../shared/semantic/artifact.mjs';
import {
createArtifactFailureReceipt,
createDeliveryReceipt,
mergeDeliveryReceipts,
providerMessageIdsFor,
} from '../shared/semantic/delivery.mjs';
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
const GENERIC_PROCESSING_ERROR = '消息处理失败,请稍后重试。';
@ -89,6 +99,28 @@ function safeMessageError(error, userMessage = GENERIC_PROCESSING_ERROR) {
};
}
function artifactFailureText(fileName, error) {
const name = String(fileName ?? '结果文件').replace(/[\r\n]+/g, ' ').trim() || '结果文件';
switch (error?.code) {
case 'artifact-delivery-uncertain':
return `结果文件「${name}」发送结果未能确认,请先检查聊天内是否已收到,不要立即重试。`;
case 'artifact-permission-required':
return `结果文件「${name}」已生成,但微信机器人当前没有文件消息发送权限,请检查机器人文件消息能力。`;
case 'artifact-too-large':
return `结果文件「${name}」超过当前微信会话可发送的文件大小,未发送。`;
case 'artifact-rate-limited':
return `结果文件「${name}」暂时被微信限流,未能发送,请稍后重试。`;
case 'artifact-provider-rejected':
return `结果文件「${name}」已生成,但微信拒绝了该文件消息。`;
case 'artifact-invalid':
case 'artifact-changed':
case 'artifact-unavailable':
return `结果文件「${name}」暂时无法读取或准备发送,请确认文件仍可访问后重试。`;
default:
return `结果文件「${name}」已生成,但暂时未能通过微信发送,请稍后重试。`;
}
}
export function createWeixinBridgeStatus() {
return {
messagesReceived: 0,
@ -394,8 +426,9 @@ export class WeixinHarnessBridge {
? await promptContentForMessage(promptMessage, { signal: this.#signal })
: undefined;
let answer;
let artifacts = [];
try {
({ answer } = await askInWorkspaceSession({
({ answer, artifacts = [] } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
@ -421,12 +454,35 @@ export class WeixinHarnessBridge {
this.#approvals.closeRoute(key),
]);
}
await this.#send(sender, answer, contextToken, runId);
const answerText = typeof answer === 'string' && answer.trim()
? answer
: artifacts.length > 0 ? '结果文件已生成。' : answer;
let textDeliveryError = null;
let textReceipt = null;
try {
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'weixin-text',
providerMessageIds: await this.#send(sender, answerText, contextToken, runId),
});
} catch (error) {
textDeliveryError = error;
}
const delivery = await this.#deliverArtifacts(
sender,
messageId,
artifacts,
contextToken,
runId,
textReceipt,
);
if (textDeliveryError && !delivery.userVisible) throw textDeliveryError;
await this.#state.markSeen(messageId);
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
this.#status.lastMessageError = null;
return delivery.receipt;
} catch (error) {
if (error?.code === 'turn-stopped') {
await this.#state.markSeen(messageId);
@ -762,15 +818,90 @@ export class WeixinHarnessBridge {
}
async #send(toUserId, text, contextToken, runId) {
const providerMessageIds = [];
for (const chunk of splitWeixinText(text, this.#maxMessageChars)) {
await this.#api.sendText({
const result = await this.#api.sendText({
baseUrl: this.#baseUrl,
token: this.#token,
toUserId,
text: chunk,
contextToken,
runId,
signal: this.#signal,
});
providerMessageIds.push(...providerMessageIdsFor(result));
}
return providerMessageIds;
}
async #deliverArtifacts(toUserId, replyTo, artifacts, contextToken, runId, baseReceipt) {
const receipts = baseReceipt ? [baseReceipt] : [];
let userVisible = Boolean(baseReceipt);
for (const artifact of artifacts) {
this.#signal?.throwIfAborted();
try {
if (typeof this.#api.sendFile !== 'function') {
const unavailable = new Error('Weixin file delivery is unavailable');
unavailable.code = 'artifact-provider-unavailable';
throw unavailable;
}
const file = await materializeOutboundArtifact(artifact, {
signal: this.#signal,
});
const result = await this.#api.sendFile({
baseUrl: this.#baseUrl,
token: this.#token,
toUserId,
file,
contextToken,
runId,
signal: this.#signal,
});
receipts.push(createDeliveryReceipt({
deliveryId: file.deliveryKey,
presentation: 'weixin-file',
providerMessageIds: providerMessageIdsFor(result),
artifacts: [{ artifactId: file.artifactId, outcome: 'sent' }],
}));
userVisible = true;
this.#status.artifactsSent = (this.#status.artifactsSent ?? 0) + 1;
} catch (error) {
if (this.#signal?.aborted) throw error;
this.#status.artifactSendErrors = (this.#status.artifactSendErrors ?? 0) + 1;
this.#logger.warn?.(
`[dsh-weixin] result file delivery failed (${error?.code ?? 'unknown'})`,
);
let noticeSent = false;
const providerMessageIds = await this.#send(
toUserId,
artifactFailureText(artifact?.fileName, error),
contextToken,
runId,
).then((ids) => {
noticeSent = true;
return ids;
}).catch(() => []);
const failureReceipt = createArtifactFailureReceipt({
artifactId: artifact?.artifactId ?? 'unknown',
deliveryId: artifact?.deliveryKey ?? artifact?.artifactId ?? 'unknown',
error,
providerMessageIds,
});
receipts.push(failureReceipt);
if (noticeSent || failureReceipt.artifacts[0]?.outcome === 'unknown') userVisible = true;
} finally {
releaseOutboundArtifact(artifact);
}
}
const receipt = receipts.length === 0
? null
: receipts.length === 1
? receipts[0]
: mergeDeliveryReceipts({
deliveryId: replyTo,
presentation: baseReceipt ? 'weixin-text-and-files' : 'weixin-files',
receipts,
});
return { receipt, userVisible };
}
}

View file

@ -1,3 +1,5 @@
import { createHash } from 'node:crypto';
import {
areJidsSameUser,
downloadMediaMessage,
@ -6,12 +8,14 @@ import {
import { splitMessageText } from '../shared/editable-message-stream.mjs';
import { ImagePromptError } from '../shared/image-prompt.mjs';
import { trackOutboundArtifactProviderPromise } from '../shared/semantic/artifact.mjs';
import { createWhatsappBridgeStatus, WhatsappHarnessBridge } from './whatsapp-bridge.mjs';
import { createWhatsappWebSession } from './whatsapp-web-session.mjs';
const IMAGE_MEDIA_TYPES = new Set(['image/jpeg', 'image/png', 'image/webp', 'image/gif']);
const DEFAULT_MAX_IMAGE_BYTES = 5 * 1024 * 1024;
const IMAGE_DOWNLOAD_TIMEOUT_MS = 15_000;
const WHATSAPP_MEDIA_UPLOAD_TIMEOUT_MS = 120_000;
const MESSAGE_WRAPPER_KEYS = [
'ephemeralMessage',
'viewOnceMessage',
@ -225,27 +229,114 @@ class RecentWhatsappOutboundIds {
}
}
class WhatsappBotClient {
function abortReason(signal) {
return signal?.reason ?? new DOMException('The operation was aborted.', 'AbortError');
}
function waitWithSignal(promise, signal) {
if (!signal) return promise;
signal.throwIfAborted();
return new Promise((resolve, reject) => {
let settled = false;
const finish = (callback, value) => {
if (settled) return;
settled = true;
signal.removeEventListener('abort', onAbort);
callback(value);
};
const onAbort = () => finish(reject, abortReason(signal));
signal.addEventListener('abort', onAbort, { once: true });
Promise.resolve(promise).then(
(value) => finish(resolve, value),
(error) => finish(reject, error),
);
if (signal.aborted) onAbort();
});
}
function uncertainArtifactDelivery(error) {
if (error?.code === 'artifact-delivery-uncertain') return error;
const uncertain = new Error('WhatsApp could not confirm result file delivery.');
uncertain.code = 'artifact-delivery-uncertain';
uncertain.cause = error;
return uncertain;
}
export class WhatsappBotClient {
#socket;
#outboundIds;
#signal;
#mediaUploadTimeoutMs;
#typingTimers = new Map();
constructor(socket, outboundIds) {
constructor(socket, outboundIds, {
signal,
mediaUploadTimeoutMs = WHATSAPP_MEDIA_UPLOAD_TIMEOUT_MS,
} = {}) {
this.#socket = socket;
this.#outboundIds = outboundIds;
this.#signal = signal;
this.#mediaUploadTimeoutMs = mediaUploadTimeoutMs;
}
async sendText(target, text) {
await this.#stopTyping(target.jid);
let result = null;
const providerMessageIds = [];
for (const [index, chunk] of splitMessageText(text, 4_000).entries()) {
result = await this.#socket.sendMessage(
const result = await this.#socket.sendMessage(
target.jid,
{ text: chunk },
index === 0 && target.quoted ? { quoted: target.quoted } : undefined,
);
this.#outboundIds.remember(result?.key?.id);
if (typeof result?.key?.id === 'string' && result.key.id) {
providerMessageIds.push(result.key.id);
}
}
return { providerMessageIds };
}
async sendFile(target, file) {
this.#signal?.throwIfAborted();
await this.#stopTyping(target.jid);
this.#signal?.throwIfAborted();
const deliverySeed = typeof file.deliveryKey === 'string' && file.deliveryKey
? file.deliveryKey
: file.artifactId;
const messageId = typeof deliverySeed === 'string' && deliverySeed
? createHash('sha256').update(deliverySeed).digest('hex').slice(0, 20).toUpperCase()
: undefined;
const options = {
...(target.quoted ? { quoted: target.quoted } : {}),
...(messageId ? { messageId } : {}),
mediaUploadTimeoutMs: this.#mediaUploadTimeoutMs,
};
// Baileys can emit a self-chat echo before sendMessage settles. Reserve the
// deterministic id before dispatch so that echo cannot re-enter the bridge.
this.#outboundIds.remember(messageId);
let result;
try {
const pending = this.#socket.sendMessage(
target.jid,
{
document: file.bytes,
mimetype: file.mediaType ?? 'application/octet-stream',
fileName: file.fileName,
},
options,
);
trackOutboundArtifactProviderPromise(file, pending);
const timeout = AbortSignal.timeout(this.#mediaUploadTimeoutMs);
const waitSignal = this.#signal
? AbortSignal.any([this.#signal, timeout])
: timeout;
result = await waitWithSignal(pending, waitSignal);
} catch (error) {
if (this.#signal?.aborted) throw abortReason(this.#signal);
throw uncertainArtifactDelivery(error);
}
this.#signal?.throwIfAborted();
this.#outboundIds.remember(result?.key?.id);
return result;
}
@ -298,6 +389,7 @@ export class WhatsappRuntime {
#logger;
#replyTimeoutMs;
#connectTimeoutMs;
#mediaUploadTimeoutMs;
#createSession;
#status = createWhatsappRuntimeStatus();
#abortController = null;
@ -314,6 +406,7 @@ export class WhatsappRuntime {
logger = console,
replyTimeoutMs = 600_000,
connectTimeoutMs = 30_000,
mediaUploadTimeoutMs = WHATSAPP_MEDIA_UPLOAD_TIMEOUT_MS,
createSession = createWhatsappWebSession,
}) {
if (!config || !authDir || !harness || !state || typeof createSession !== 'function') {
@ -326,6 +419,13 @@ export class WhatsappRuntime {
this.#logger = logger;
this.#replyTimeoutMs = replyTimeoutMs;
this.#connectTimeoutMs = connectTimeoutMs;
if (!Number.isSafeInteger(mediaUploadTimeoutMs) || mediaUploadTimeoutMs <= 0) {
throw new TypeError('mediaUploadTimeoutMs must be a positive safe integer');
}
this.#mediaUploadTimeoutMs = Math.min(
mediaUploadTimeoutMs,
WHATSAPP_MEDIA_UPLOAD_TIMEOUT_MS,
);
this.#createSession = createSession;
}
@ -390,7 +490,10 @@ export class WhatsappRuntime {
if (!areJidsSameUser(identity.accountJid, this.#config.accountJid)) {
throw new Error('WhatsApp linked account does not match the saved bot');
}
const client = new WhatsappBotClient(session.socket, outboundIds);
const client = new WhatsappBotClient(session.socket, outboundIds, {
signal: controller.signal,
mediaUploadTimeoutMs: this.#mediaUploadTimeoutMs,
});
this.#client = client;
this.#bridge = new WhatsappHarnessBridge({
bot: client,