mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-11 00:32:43 +08:00
feat(weixin): show typing indicator during turns
This commit is contained in:
parent
ad3c08ace3
commit
cb34e3f8de
10 changed files with 948 additions and 200 deletions
|
|
@ -637,6 +637,53 @@ export function createWeixinApi({ fetchImpl = fetch } = {}) {
|
|||
}
|
||||
},
|
||||
|
||||
async getConfig({ baseUrl, token, toUserId, contextToken, signal }) {
|
||||
const recipient = nonEmptyString(toUserId);
|
||||
if (!recipient) throw new TypeError('toUserId is required');
|
||||
const response = await requestJson(fetchImpl, {
|
||||
method: 'POST',
|
||||
baseUrl,
|
||||
endpoint: 'ilink/bot/getconfig',
|
||||
token,
|
||||
signal,
|
||||
timeoutMs: 10_000,
|
||||
body: {
|
||||
ilink_user_id: recipient,
|
||||
...(nonEmptyString(contextToken) ? { context_token: contextToken.trim() } : {}),
|
||||
base_info: baseInfo(),
|
||||
},
|
||||
});
|
||||
if (response?.ret !== undefined && response.ret !== 0) {
|
||||
throw new WeixinApiError('config-rejected', '微信服务拒绝了机器人配置请求。');
|
||||
}
|
||||
return { typingTicket: nonEmptyString(response?.typing_ticket) };
|
||||
},
|
||||
|
||||
async sendTyping({ baseUrl, token, toUserId, typingTicket, status, signal }) {
|
||||
const recipient = nonEmptyString(toUserId);
|
||||
const ticket = nonEmptyString(typingTicket);
|
||||
if (!recipient || !ticket) throw new TypeError('toUserId and typingTicket are required');
|
||||
if (status !== 1 && status !== 2) throw new TypeError('typing status must be 1 or 2');
|
||||
const response = await requestJson(fetchImpl, {
|
||||
method: 'POST',
|
||||
baseUrl,
|
||||
endpoint: 'ilink/bot/sendtyping',
|
||||
token,
|
||||
signal,
|
||||
timeoutMs: 10_000,
|
||||
body: {
|
||||
ilink_user_id: recipient,
|
||||
typing_ticket: ticket,
|
||||
status,
|
||||
base_info: baseInfo(),
|
||||
},
|
||||
});
|
||||
if (response?.ret !== undefined && response.ret !== 0) {
|
||||
throw new WeixinApiError('typing-rejected', '微信服务拒绝了输入状态请求。');
|
||||
}
|
||||
return true;
|
||||
},
|
||||
|
||||
async sendText({ baseUrl, token, toUserId, text, contextToken, runId, signal }) {
|
||||
const recipient = nonEmptyString(toUserId);
|
||||
const content = nonEmptyString(text);
|
||||
|
|
|
|||
|
|
@ -58,6 +58,8 @@ import {
|
|||
import { t } from '../shared/i18n.mjs';
|
||||
|
||||
const INTERACTION_RESOLVED_TEXT = () => t('这个问题已在其他客户端处理,无需再次回答。');
|
||||
const DEFAULT_TYPING_KEEPALIVE_MS = 5_000;
|
||||
const TYPING_RETRY_DELAY_MS = 60_000;
|
||||
|
||||
const HELP_TEXT = () => [
|
||||
t('微信已连接 DeepSeek Harness。'),
|
||||
|
|
@ -181,6 +183,7 @@ export class WeixinHarnessBridge {
|
|||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#maxMessageChars;
|
||||
#typingKeepaliveMs;
|
||||
#signal;
|
||||
#queues = new Map();
|
||||
#pendingInteractions = new Map();
|
||||
|
|
@ -190,6 +193,14 @@ export class WeixinHarnessBridge {
|
|||
#commandTasks = new Set();
|
||||
#approvals;
|
||||
#batchInputs = new BatchInputManager();
|
||||
#typingTicket = null;
|
||||
#typingTarget = null;
|
||||
#typingTimer = null;
|
||||
#typingGeneration = 0;
|
||||
#typingTail = Promise.resolve();
|
||||
#typingClosed = false;
|
||||
#typingRetryAt = 0;
|
||||
#typingTicketStale = false;
|
||||
|
||||
constructor({
|
||||
api,
|
||||
|
|
@ -202,11 +213,15 @@ export class WeixinHarnessBridge {
|
|||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
maxMessageChars = DEFAULT_WEIXIN_MAX_MESSAGE_CHARS,
|
||||
typingKeepaliveMs = DEFAULT_TYPING_KEEPALIVE_MS,
|
||||
signal,
|
||||
}) {
|
||||
if (!api || typeof api.sendText !== 'function') throw new TypeError('Weixin API is required');
|
||||
if (!baseUrl || !token || !ownerUserId) throw new TypeError('Weixin account credentials are required');
|
||||
if (!harness || !state) throw new TypeError('Harness client and state store are required');
|
||||
if (!Number.isFinite(typingKeepaliveMs) || typingKeepaliveMs <= 0) {
|
||||
throw new TypeError('typingKeepaliveMs must be a positive number');
|
||||
}
|
||||
this.#api = api;
|
||||
this.#baseUrl = baseUrl;
|
||||
this.#token = token;
|
||||
|
|
@ -217,6 +232,7 @@ export class WeixinHarnessBridge {
|
|||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#maxMessageChars = maxMessageChars;
|
||||
this.#typingKeepaliveMs = typingKeepaliveMs;
|
||||
this.#signal = signal;
|
||||
this.#approvals = new HarnessApprovalQueue({ label: 'weixin', logger });
|
||||
}
|
||||
|
|
@ -263,6 +279,7 @@ export class WeixinHarnessBridge {
|
|||
return this.#finishBatchResult(
|
||||
message,
|
||||
messageId,
|
||||
key,
|
||||
sender,
|
||||
contextToken,
|
||||
runId,
|
||||
|
|
@ -294,7 +311,7 @@ export class WeixinHarnessBridge {
|
|||
`[dsh-weixin] failed to process a command [${failure.referenceId}]:`,
|
||||
error,
|
||||
);
|
||||
return this.#send(sender, messageFailureText(failure), contextToken, runId)
|
||||
return this.#sendOutOfBand(key, sender, messageFailureText(failure), contextToken, runId)
|
||||
.catch(() => undefined);
|
||||
}).finally(() => {
|
||||
this.#acceptedMessageIds.delete(messageId);
|
||||
|
|
@ -384,7 +401,7 @@ export class WeixinHarnessBridge {
|
|||
return current;
|
||||
}
|
||||
|
||||
#finishBatchResult(message, messageId, sender, contextToken, runId, result) {
|
||||
#finishBatchResult(message, messageId, key, sender, contextToken, runId, result) {
|
||||
let task;
|
||||
task = Promise.resolve().then(async () => {
|
||||
if (this.#state.hasSeen(messageId)) return;
|
||||
|
|
@ -392,7 +409,7 @@ export class WeixinHarnessBridge {
|
|||
this.#status.messagesReceived += 1;
|
||||
this.#status.lastMessageAt = new Date().toISOString();
|
||||
if (result.message) {
|
||||
await this.#send(sender, result.message, contextToken, runId);
|
||||
await this.#sendOutOfBand(key, sender, result.message, contextToken, runId);
|
||||
}
|
||||
this.#status.lastError = null;
|
||||
}).catch(async (error) => {
|
||||
|
|
@ -403,7 +420,7 @@ export class WeixinHarnessBridge {
|
|||
`[dsh-weixin] failed to process a batch input message [${failure.referenceId}]:`,
|
||||
error,
|
||||
);
|
||||
await this.#send(sender, messageFailureText(failure), contextToken, runId)
|
||||
await this.#sendOutOfBand(key, sender, messageFailureText(failure), contextToken, runId)
|
||||
.catch(() => undefined);
|
||||
}).finally(() => {
|
||||
this.#acceptedMessageIds.delete(messageId);
|
||||
|
|
@ -421,9 +438,15 @@ export class WeixinHarnessBridge {
|
|||
)),
|
||||
...this.#approvalTasks,
|
||||
...this.#commandTasks,
|
||||
this.#typingTail,
|
||||
]);
|
||||
}
|
||||
|
||||
async close() {
|
||||
this.#typingClosed = true;
|
||||
await this.#stopTyping({ signal: AbortSignal.timeout(5_000) });
|
||||
}
|
||||
|
||||
async #processFastCommand(
|
||||
message,
|
||||
messageId,
|
||||
|
|
@ -454,7 +477,12 @@ export class WeixinHarnessBridge {
|
|||
]);
|
||||
}
|
||||
for (const reply of result?.messages ?? [result?.message]) {
|
||||
if (reply) await this.#send(sender, reply, contextToken, runId);
|
||||
if (!reply) continue;
|
||||
if (result?.stopped) {
|
||||
await this.#send(sender, reply, contextToken, runId);
|
||||
} else {
|
||||
await this.#sendOutOfBand(key, sender, reply, contextToken, runId);
|
||||
}
|
||||
}
|
||||
this.#status.lastError = null;
|
||||
}
|
||||
|
|
@ -537,12 +565,13 @@ export class WeixinHarnessBridge {
|
|||
return;
|
||||
}
|
||||
|
||||
const content = hasImages
|
||||
? await promptContentForMessage(promptMessage, { signal: this.#signal })
|
||||
: undefined;
|
||||
let answer;
|
||||
let artifacts = [];
|
||||
await this.#startTyping(sender, contextToken);
|
||||
try {
|
||||
const content = hasImages
|
||||
? await promptContentForMessage(promptMessage, { signal: this.#signal })
|
||||
: undefined;
|
||||
await this.#state.markSeen(messageId);
|
||||
promptRecorded = true;
|
||||
({ answer, artifacts = [] } = await askInWorkspaceSession({
|
||||
|
|
@ -556,13 +585,17 @@ export class WeixinHarnessBridge {
|
|||
timeoutMs: this.#replyTimeoutMs,
|
||||
signal: this.#signal,
|
||||
control: { owner: this, key },
|
||||
onUpdate: () => this.#resumeTyping(key, sender, contextToken),
|
||||
onInteraction: (interaction) => this.#handleInteraction(interaction, {
|
||||
key,
|
||||
actor: sender,
|
||||
contextToken,
|
||||
runId,
|
||||
}),
|
||||
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
|
||||
onInteractionResolved: async (resolution) => {
|
||||
await this.#handleInteractionResolved(resolution);
|
||||
await this.#resumeTyping(key, sender, contextToken);
|
||||
},
|
||||
files: promptMessage.files,
|
||||
},
|
||||
}));
|
||||
|
|
@ -572,6 +605,7 @@ export class WeixinHarnessBridge {
|
|||
}
|
||||
} finally {
|
||||
await Promise.allSettled([
|
||||
this.#stopTyping(),
|
||||
this.#cancelPendingInteraction(key),
|
||||
this.#approvals.closeRoute(key),
|
||||
]);
|
||||
|
|
@ -966,6 +1000,7 @@ export class WeixinHarnessBridge {
|
|||
}
|
||||
|
||||
async #send(toUserId, text, contextToken, runId) {
|
||||
await this.#stopTyping();
|
||||
const providerMessageIds = [];
|
||||
for (const chunk of splitWeixinText(text, this.#maxMessageChars)) {
|
||||
const result = await this.#api.sendText({
|
||||
|
|
@ -982,6 +1017,164 @@ export class WeixinHarnessBridge {
|
|||
return providerMessageIds;
|
||||
}
|
||||
|
||||
async #sendOutOfBand(key, toUserId, text, contextToken, runId) {
|
||||
try {
|
||||
return await this.#send(toUserId, text, contextToken, runId);
|
||||
} finally {
|
||||
if (this.#queues.has(key)) {
|
||||
await this.#resumeTyping(key, toUserId, contextToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#queueTyping(operation) {
|
||||
const task = this.#typingTail
|
||||
.catch(() => undefined)
|
||||
.then(operation);
|
||||
this.#typingTail = task.catch(() => undefined);
|
||||
return task;
|
||||
}
|
||||
|
||||
async #startTyping(toUserId, contextToken) {
|
||||
const target = nonEmptyString(toUserId);
|
||||
if (!target
|
||||
|| this.#typingClosed
|
||||
|| this.#signal?.aborted
|
||||
|| Date.now() < this.#typingRetryAt
|
||||
|| typeof this.#api.getConfig !== 'function'
|
||||
|| typeof this.#api.sendTyping !== 'function') return false;
|
||||
if (this.#typingTarget === target) return true;
|
||||
|
||||
const generation = ++this.#typingGeneration;
|
||||
try {
|
||||
return await this.#queueTyping(async () => {
|
||||
if (generation !== this.#typingGeneration || this.#typingClosed) return false;
|
||||
let ticket = this.#typingTicket;
|
||||
if (!ticket) {
|
||||
const config = await this.#api.getConfig({
|
||||
baseUrl: this.#baseUrl,
|
||||
token: this.#token,
|
||||
toUserId: target,
|
||||
contextToken,
|
||||
signal: this.#signal,
|
||||
});
|
||||
ticket = nonEmptyString(config?.typingTicket);
|
||||
if (!ticket) {
|
||||
this.#typingRetryAt = Date.now() + TYPING_RETRY_DELAY_MS;
|
||||
return false;
|
||||
}
|
||||
if (generation !== this.#typingGeneration || this.#typingClosed) return false;
|
||||
this.#typingTicket = ticket;
|
||||
this.#typingRetryAt = 0;
|
||||
}
|
||||
|
||||
this.#typingTarget = target;
|
||||
await this.#api.sendTyping({
|
||||
baseUrl: this.#baseUrl,
|
||||
token: this.#token,
|
||||
toUserId: target,
|
||||
typingTicket: ticket,
|
||||
status: 1,
|
||||
signal: this.#signal,
|
||||
});
|
||||
if (generation !== this.#typingGeneration || this.#typingClosed) return false;
|
||||
this.#typingTicketStale = false;
|
||||
this.#typingRetryAt = 0;
|
||||
this.#scheduleTyping(generation);
|
||||
return true;
|
||||
});
|
||||
} catch (error) {
|
||||
if (generation === this.#typingGeneration) {
|
||||
this.#clearTypingTimer();
|
||||
this.#typingTicketStale = Boolean(this.#typingTarget && this.#typingTicket);
|
||||
this.#typingRetryAt = Date.now() + TYPING_RETRY_DELAY_MS;
|
||||
}
|
||||
if (!this.#signal?.aborted && !this.#typingClosed) {
|
||||
this.#logger.warn?.('[dsh-weixin] typing indicator failed:', error);
|
||||
}
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
async #stopTyping({ signal = AbortSignal.timeout(5_000) } = {}) {
|
||||
++this.#typingGeneration;
|
||||
this.#clearTypingTimer();
|
||||
const target = this.#typingTarget;
|
||||
const ticket = this.#typingTicket;
|
||||
const stale = this.#typingTicketStale;
|
||||
this.#typingTarget = null;
|
||||
if (typeof this.#api.sendTyping !== 'function') return false;
|
||||
try {
|
||||
return await this.#queueTyping(async () => {
|
||||
if (!target || !ticket) return false;
|
||||
try {
|
||||
await this.#api.sendTyping({
|
||||
baseUrl: this.#baseUrl,
|
||||
token: this.#token,
|
||||
toUserId: target,
|
||||
typingTicket: ticket,
|
||||
status: 2,
|
||||
signal,
|
||||
});
|
||||
} finally {
|
||||
if (stale && this.#typingTicket === ticket) this.#typingTicket = null;
|
||||
this.#typingTicketStale = false;
|
||||
}
|
||||
return true;
|
||||
});
|
||||
} catch (error) {
|
||||
if (!signal?.aborted) {
|
||||
this.#logger.warn?.('[dsh-weixin] typing cancellation failed:', error);
|
||||
}
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
#scheduleTyping(generation) {
|
||||
this.#clearTypingTimer();
|
||||
if (generation !== this.#typingGeneration || !this.#typingTarget || this.#typingClosed) return;
|
||||
const timer = setTimeout(() => {
|
||||
if (this.#typingTimer === timer) this.#typingTimer = null;
|
||||
void this.#queueTyping(async () => {
|
||||
if (generation !== this.#typingGeneration || !this.#typingTarget || this.#typingClosed) {
|
||||
return false;
|
||||
}
|
||||
await this.#api.sendTyping({
|
||||
baseUrl: this.#baseUrl,
|
||||
token: this.#token,
|
||||
toUserId: this.#typingTarget,
|
||||
typingTicket: this.#typingTicket,
|
||||
status: 1,
|
||||
signal: this.#signal,
|
||||
});
|
||||
return true;
|
||||
}).then((sent) => {
|
||||
if (sent) this.#scheduleTyping(generation);
|
||||
}).catch((error) => {
|
||||
if (generation !== this.#typingGeneration) return;
|
||||
this.#typingTicketStale = Boolean(this.#typingTarget && this.#typingTicket);
|
||||
this.#typingRetryAt = Date.now() + TYPING_RETRY_DELAY_MS;
|
||||
if (!this.#signal?.aborted && !this.#typingClosed) {
|
||||
this.#logger.warn?.('[dsh-weixin] typing keepalive failed:', error);
|
||||
}
|
||||
});
|
||||
}, this.#typingKeepaliveMs);
|
||||
timer.unref?.();
|
||||
this.#typingTimer = timer;
|
||||
}
|
||||
|
||||
#clearTypingTimer() {
|
||||
if (this.#typingTimer) clearTimeout(this.#typingTimer);
|
||||
this.#typingTimer = null;
|
||||
}
|
||||
|
||||
#resumeTyping(key, toUserId, contextToken) {
|
||||
if (this.#pendingInteractions.has(key) || this.#approvals.hasPending(key)) {
|
||||
return Promise.resolve(false);
|
||||
}
|
||||
return this.#startTyping(toUserId, contextToken);
|
||||
}
|
||||
|
||||
async #deliverArtifacts(toUserId, replyTo, artifacts, contextToken, runId, baseReceipt) {
|
||||
const sendArtifact = (method, file) => this.#api[method]({
|
||||
baseUrl: this.#baseUrl,
|
||||
|
|
|
|||
|
|
@ -285,6 +285,7 @@ export class WeixinRuntime {
|
|||
this.#abortController?.abort();
|
||||
this.#abortController = null;
|
||||
this.#monitor = null;
|
||||
await bridge?.close?.();
|
||||
await monitor?.catch(() => undefined);
|
||||
await bridge?.waitForIdle();
|
||||
this.#bridge = null;
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue