feat: add model and turn control commands

This commit is contained in:
xmanrui 2026-08-20 01:28:51 +08:00
parent fc8f120e13
commit 99ce38e9a4
38 changed files with 3451 additions and 169 deletions

View file

@ -10,6 +10,14 @@ import {
} from '../shared/harness-question.mjs';
import { HarnessApprovalQueue } from '../shared/harness-approval.mjs';
import { runCompactCommand } from '../shared/compact-command.mjs';
import {
isControlCommand,
runControlCommand,
} from '../shared/control-command.mjs';
import {
isModelCommand,
runModelCommand,
} from '../shared/model-command.mjs';
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
import {
@ -33,6 +41,10 @@ const HELP_TEXT = [
'/workspacelist 列出工作区绝对路径',
'/sessionlist [工作区序号或绝对路径] 列出会话 ID 和标题',
'/session Session ID 或当前工作区序号 将当前聊天绑定到指定会话',
'/models 列出所有可用模型',
'/model [模型ID] 查看或切换当前会话模型',
'/stop 停止当前任务',
'/steer 补充指令 纠偏当前任务',
'/status 检查连接状态',
'/help 显示本帮助',
].join('\n');
@ -246,6 +258,7 @@ export class DingtalkHarnessBridge {
#pendingInteractions = new Map();
#interactionKeys = new Map();
#interactionTasks = new Set();
#commandTasks = new Set();
#acceptedMessageIds = new Set();
#approvals;
@ -314,6 +327,37 @@ export class DingtalkHarnessBridge {
rememberConnectionTestTarget(this.#state, { sessionWebhook });
}
const pending = this.#pendingInteractions.get(key);
const promptMessage = dingtalkInboundMessage(message, {
api: this.#api,
clientId: this.#clientId,
clientSecret: this.#clientSecret,
});
const commandText = nonEmptyString(promptMessage.content) ?? '';
const commandRunner = isControlCommand(commandText)
? runControlCommand
: (isModelCommand(commandText) ? runModelCommand : null);
const addressed = String(message.conversationType) !== '2' || message?.isInAtList === true;
if (commandRunner && sessionWebhook && addressed) {
let task;
task = this.#processFastCommand(
message,
messageId,
key,
sessionWebhook,
promptMessage,
commandRunner,
).catch((error) => {
if (error?.code === 'turn-stopped' || this.#signal?.aborted) return;
this.#status.lastError = '钉钉命令处理失败。';
this.#logger.error?.('[dsh-dingtalk] failed to process a command', safeErrorDiagnostic(error));
return this.#send(sessionWebhook, CARD_ERROR_TEXT).catch(() => undefined);
}).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
const approvalReply = this.#approvals.claimReply({
key,
actor: sender,
@ -412,9 +456,41 @@ export class DingtalkHarnessBridge {
pending.queue ? [pending.queue] : []
)),
...this.#interactionTasks,
...this.#commandTasks,
]);
}
async #processFastCommand(message, messageId, key, sessionWebhook, prompt, runner) {
this.#signal?.throwIfAborted();
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
increment(this.#status, 'messagesReceived');
this.#status.lastMessageAt = new Date().toISOString();
const result = await runner(
nonEmptyString(prompt.content) ?? '',
this.#harness,
this.#state,
key,
{
signal: this.#signal,
hasImages: hasInboundImages(prompt),
pendingInteraction: this.#pendingInteractions.has(key)
|| this.#approvals.hasPending(key),
control: { owner: this, key },
},
);
if (result?.stopped) {
await Promise.allSettled([
this.#cancelPendingInteraction(key),
this.#approvals.closeRoute(key),
]);
}
for (const reply of result?.messages ?? [result?.message]) {
if (reply) await this.#send(sessionWebhook, reply);
}
this.#status.lastError = null;
}
async #process(message, messageId, sender, key, { alreadyRecorded = false } = {}) {
this.#signal?.throwIfAborted();
if (!alreadyRecorded) {
@ -519,6 +595,7 @@ export class DingtalkHarnessBridge {
askOptions: {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
control: { owner: this, key },
onUpdate: cardStarted
? (update) => cardStream.push(progressText(update))
: undefined,
@ -537,6 +614,10 @@ export class DingtalkHarnessBridge {
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} catch (error) {
if (error?.code === 'turn-stopped') {
if (cardStarted) await cardStream.finish('已停止。').catch(() => undefined);
return;
}
if (this.#signal?.aborted) return;
this.#status.lastError = '钉钉消息处理失败。';
this.#logger.error?.(

View file

@ -18,7 +18,15 @@ import {
} from '../shared/harness-question.mjs';
import { HarnessApprovalQueue } from '../shared/harness-approval.mjs';
import { runCompactCommand } from '../shared/compact-command.mjs';
import {
isControlCommand,
runControlCommand,
} from '../shared/control-command.mjs';
import { rememberConnectionTestTarget } from '../shared/connection-test.mjs';
import {
isModelCommand,
runModelCommand,
} from '../shared/model-command.mjs';
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
@ -35,6 +43,10 @@ const HELP_TEXT = [
'/workspacelist 列出工作区绝对路径',
'/sessionlist [工作区序号或绝对路径] 列出会话 ID 和标题',
'/session Session ID 或当前工作区序号 将当前聊天绑定到指定会话',
'/models 列出所有可用模型',
'/model [模型ID] 查看或切换当前会话模型',
'/stop 停止当前任务',
'/steer 补充指令 纠偏当前任务',
'/status 检查连接状态',
'/help 显示本帮助',
].join('\n');
@ -77,6 +89,7 @@ export class FeishuHarnessBridge {
#resolvedQuestionReplies = new Map();
#acceptedMessageIds = new Set();
#interactionTasks = new Set();
#commandTasks = new Set();
#approvals;
#status;
#allowedSenderOpenIds;
@ -141,6 +154,37 @@ export class FeishuHarnessBridge {
this.#acceptedMessageIds.add(messageId);
const processingReaction = this.#addReaction(messageId, 'OnIt');
const commandMessage = extractInboundMessage(event, this.#client);
const commandText = nonEmptyString(commandMessage.content) ?? '';
const commandRunner = isControlCommand(commandText)
? runControlCommand
: (isModelCommand(commandText) ? runModelCommand : null);
const addressed = event?.message?.chat_type === 'p2p'
|| (Array.isArray(event?.message?.mentions) && event.message.mentions.length > 0);
if (commandRunner && addressed) {
const processing = this.#processFastCommand(
event,
messageId,
key,
commandMessage,
commandRunner,
);
let current;
current = processing
.then(() => this.#finishReaction(messageId, processingReaction, 'DONE'))
.catch((error) => this.#handleMessageFailure(
event,
messageId,
processingReaction,
error,
))
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(current);
});
this.#commandTasks.add(current);
return current;
}
if (this.#isResolvedQuestionReply(event, key)) {
const current = Promise.resolve()
.then(() => this.#discardResolvedInteractionReply(event, messageId))
@ -276,6 +320,10 @@ export class FeishuHarnessBridge {
}
async #handleMessageFailure(event, messageId, processingReaction, error) {
if (error?.code === 'turn-stopped') {
await this.#removeProcessingReaction(messageId, processingReaction);
return;
}
if (this.#signal?.aborted) {
await this.#removeProcessingReaction(messageId, processingReaction);
return;
@ -297,9 +345,41 @@ export class FeishuHarnessBridge {
pending.queue ? [pending.queue] : []
)),
...this.#interactionTasks,
...this.#commandTasks,
]);
}
async #processFastCommand(event, messageId, key, message, runner) {
this.#signal?.throwIfAborted();
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.lastMessageAt = new Date().toISOString();
this.#status.messagesReceived += 1;
const result = await runner(
nonEmptyString(message.content) ?? '',
this.#harness,
this.#state,
key,
{
signal: this.#signal,
hasImages: hasInboundImages(message),
pendingInteraction: this.#pendingInteractions.has(key)
|| this.#approvals.hasPending(key),
control: { owner: this, key },
},
);
if (result?.stopped) {
await Promise.allSettled([
this.#cancelPendingInteraction(key),
this.#approvals.closeRoute(key),
]);
}
for (const reply of result?.messages ?? [result?.message]) {
if (reply) await this.#send(event.message.chat_id, reply);
}
this.#status.lastError = null;
}
async #handle(event, key, { alreadyRecorded = false } = {}) {
this.#signal?.throwIfAborted();
const messageId = event.message.message_id;
@ -372,6 +452,7 @@ export class FeishuHarnessBridge {
return {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
control: { owner: this, key },
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key,
actor: senderOpenId(event),

View file

@ -1,5 +1,9 @@
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import { runCompactCommand } from '../shared/compact-command.mjs';
import {
isControlCommand,
runControlCommand,
} from '../shared/control-command.mjs';
import { rememberConnectionTestTarget } from '../shared/connection-test.mjs';
import {
harnessAnswerForQuestion,
@ -7,6 +11,10 @@ import {
validHarnessQuestion,
} from '../shared/harness-question.mjs';
import { HarnessApprovalQueue } from '../shared/harness-approval.mjs';
import {
isModelCommand,
runModelCommand,
} from '../shared/model-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
import {
fetchImageBuffer,
@ -38,6 +46,10 @@ const HELP_TEXT = [
'/workspacelist 列出工作区绝对路径',
'/sessionlist [工作区序号或绝对路径] 列出会话 ID 和标题',
'/session Session ID 或当前工作区序号 将当前聊天绑定到指定会话',
'/models 列出所有可用模型',
'/model [模型ID] 查看或切换当前会话模型',
'/stop 停止当前任务',
'/steer 补充指令 纠偏当前任务',
'/status 检查连接状态',
'/help 显示本帮助',
].join('\n');
@ -134,6 +146,7 @@ export class QqHarnessBridge {
#interactionKeys = new Map();
#acceptedMessageIds = new Set();
#approvalTasks = new Set();
#commandTasks = new Set();
#approvals;
constructor({
@ -184,6 +197,34 @@ export class QqHarnessBridge {
rememberConnectionTestTarget(this.#state, message.replyTarget);
}
const pending = this.#pendingInteractions.get(key);
const commandText = safeText(message);
const commandRunner = isControlCommand(commandText)
? runControlCommand
: (isModelCommand(commandText) ? runModelCommand : null);
const allowed = this.#ownerUserOpenid === '*' || sender === this.#ownerUserOpenid;
const addressed = message.kind !== 'group'
|| message.rawEventType === 'GROUP_AT_MESSAGE_CREATE';
if (commandRunner && allowed && addressed) {
let task;
task = this.#processFastCommand(
message,
messageId,
key,
commandText,
commandRunner,
).catch((error) => {
if (error?.code === 'turn-stopped' || this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-im:qq] failed to process a command:', error);
return this.#bot.sendText(message.replyTarget, '消息处理失败,请稍后重试。')
.catch(() => undefined);
}).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
const approval = this.#approvals.claimReply({
key,
actor: sender,
@ -262,9 +303,35 @@ export class QqHarnessBridge {
pending.queue ? [pending.queue] : []
)),
...this.#approvalTasks,
...this.#commandTasks,
]);
}
async #processFastCommand(message, messageId, key, text, runner) {
this.#signal?.throwIfAborted();
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
const result = await runner(text, this.#harness, this.#state, key, {
signal: this.#signal,
hasImages: hasQqImageAttachments(message),
pendingInteraction: this.#pendingInteractions.has(key)
|| this.#approvals.hasPending(key),
control: { owner: this, key },
});
if (result?.stopped) {
await Promise.allSettled([
this.#cancelPendingInteraction(key),
this.#approvals.closeRoute(key),
]);
}
for (const reply of result?.messages ?? [result?.message]) {
if (reply) await this.#bot.sendText(message.replyTarget, reply);
}
this.#status.lastError = null;
}
async #process(message, key, { alreadyRecorded = false } = {}) {
if (this.#signal?.aborted) return;
const messageId = nonEmptyString(message?.messageId);
@ -360,6 +427,7 @@ export class QqHarnessBridge {
askOptions: {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
control: { owner: this, key },
onUpdate: stream ? async (update) => {
const progress = update.type === 'text'
? update.text
@ -399,6 +467,22 @@ export class QqHarnessBridge {
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} catch (error) {
if (error?.code === 'turn-stopped') {
if (stream) {
try {
await stream.cancel?.();
} catch (streamError) {
this.#logger.warn?.('[dsh-im:qq] unable to cancel a stopped QQ stream:', streamError);
}
try {
await this.#bot.sendText(target, '已停止。');
} catch (sendError) {
this.#logger.warn?.('[dsh-im:qq] unable to announce a stopped QQ turn:', sendError);
}
}
await this.#state.markSeen(messageId);
return;
}
stream?.cancel?.();
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);

View file

@ -511,7 +511,9 @@ export function createBotWorkspaceScope(harness, { botId, workspaces, state }) {
}
};
}
if ((property === 'listWorkspaces' || property === 'listWorkspaceSessions')
if ((property === 'listWorkspaces'
|| property === 'listWorkspaceSessions'
|| property === 'listModels')
&& typeof target[property] === 'function') {
return async (...args) => {
if (!isCurrentScope()) {
@ -627,6 +629,30 @@ export function createBotWorkspaceScope(harness, { botId, workspaces, state }) {
sessionGenerations.delete(sessionId);
const isCurrentSession = () => isCurrentScope()
&& generation === workspaces.generationFor(botId);
const invokeCurrentSession = async (method, args, action) => {
if (!isCurrentSession()) {
throw workspaceSessionStale(
`The bot workspace changed before this ${action} started.`,
);
}
const result = await target[method](sessionId, ...args);
if (!isCurrentSession()) {
throw workspaceSessionStale(
`The bot workspace changed while this ${action} was running.`,
);
}
return result;
};
const invokeStartedSessionMutation = async (method, args, action) => {
if (!isCurrentSession()) {
throw workspaceSessionStale(
`The bot workspace changed before this ${action} started.`,
);
}
// Once an irreversible control mutation has started, preserve its
// actual outcome even if a workspace switch commits concurrently.
return target[method](sessionId, ...args);
};
return Object.freeze({
sessionId,
async sessionExists(...args) {
@ -634,6 +660,24 @@ export function createBotWorkspaceScope(harness, { botId, workspaces, state }) {
const exists = await target.sessionExists(sessionId, ...args);
return isCurrentSession() && exists;
},
models(...args) {
return invokeCurrentSession('getSessionModels', args, 'model listing');
},
selectModel(...args) {
return invokeCurrentSession('selectSessionModel', args, 'model selection');
},
isRunning(...args) {
return invokeCurrentSession('isSessionRunning', args, 'run-state check');
},
hasActiveTurn(...args) {
return invokeCurrentSession('hasActiveTurn', args, 'turn ownership check');
},
stopActiveTurn(...args) {
return invokeStartedSessionMutation('stopActiveTurn', args, 'turn stop');
},
steerActiveTurn(...args) {
return invokeStartedSessionMutation('steerActiveTurn', args, 'turn steering');
},
ask(...args) {
if (!isCurrentSession()) {
throw workspaceSessionStale(

View file

@ -0,0 +1,84 @@
const CONTROL_COMMAND = /^\/(?:stop|steer)(?=$|\s)/iu;
const STOP_COMMAND = /^\/stop(?=$|\s)/iu;
const STOP_USAGE = '用法:/stop(不带参数)';
const STEER_USAGE = '用法:/steer <补充指令>';
const TEXT_ONLY = '控制命令仅支持纯文字,请移除图片后重试。';
function commandResult(message, extra = {}) {
return { message, ...extra };
}
function requestOptions(signal) {
return signal ? { signal } : {};
}
function boundSession(harness, state, key) {
if (typeof state?.sessionFor !== 'function') return null;
const sessionId = state.sessionFor(key);
if (typeof sessionId !== 'string' || !sessionId) return null;
if (typeof harness?.workspaceSession !== 'function') {
throw new TypeError('Harness does not support workspace sessions');
}
const session = harness.workspaceSession(sessionId);
if (!session || typeof session !== 'object') {
throw new TypeError('Harness returned an invalid workspace session');
}
return session;
}
export function isControlCommand(text) {
return typeof text === 'string' && CONTROL_COMMAND.test(text.trim());
}
export async function runControlCommand(text, harness, state, key, {
signal,
hasImages = false,
pendingInteraction = false,
control,
} = {}) {
if (!isControlCommand(text)) return null;
const command = text.trim();
const stop = STOP_COMMAND.test(command);
if (hasImages) return commandResult(TEXT_ONLY);
if (stop) {
if (!/^\/stop$/iu.test(command)) return commandResult(STOP_USAGE);
const session = boundSession(harness, state, key);
if (!session) return commandResult('当前聊天没有正在运行的任务。');
if (typeof session.stopActiveTurn !== 'function') {
throw new TypeError('Harness session does not support stopping active turns');
}
const stopped = await session.stopActiveTurn(control, requestOptions(signal));
return stopped
? commandResult('已请求停止当前任务。', { stopped: true })
: commandResult('当前聊天没有正在运行的任务。');
}
const match = /^\/steer(?:\s+([\s\S]*))?$/iu.exec(command);
const instruction = match?.[1]?.trim() ?? '';
if (!instruction) return commandResult(STEER_USAGE);
if (pendingInteraction) {
return commandResult([
'当前任务正在等待你的回答或审批。',
'',
'请先处理当前请求,或者发送 /stop 停止任务。',
].join('\n'));
}
const session = boundSession(harness, state, key);
if (!session) {
return commandResult('当前聊天没有正在运行的任务,请直接发送普通消息。');
}
if (typeof session.steerActiveTurn !== 'function') {
throw new TypeError('Harness session does not support steering active turns');
}
const steered = await session.steerActiveTurn(
instruction,
control,
requestOptions(signal),
);
return steered
? commandResult('已提交补充指令,Agent 会在下一步读取。')
: commandResult('当前聊天没有正在运行的任务,请直接发送普通消息。');
}

View file

@ -110,6 +110,10 @@ export class HarnessApprovalQueue {
this.#logger = logger;
}
hasPending(key) {
return this.#routes.get(key)?.items.some((pending) => !pending.inactive) === true;
}
claimReply({
key,
actor,

View file

@ -12,12 +12,77 @@ const interactionRegistries = new Map();
function interactionRegistry(origin) {
let registry = interactionRegistries.get(origin);
if (!registry) {
registry = { ownerships: new Map(), claims: new Map(), nextOrder: 0 };
registry = {
ownerships: new Map(),
claims: new Map(),
controls: new WeakMap(),
nextOrder: 0,
};
interactionRegistries.set(origin, registry);
}
return registry;
}
function normalizeControl(control) {
const ownerType = typeof control?.owner;
if ((ownerType !== 'object' && ownerType !== 'function')
|| control.owner === null
|| typeof control.key !== 'string'
|| !control.key) return null;
return { owner: control.owner, key: control.key };
}
function validModelSelection(value) {
return value !== null
&& typeof value === 'object'
&& typeof value.provider === 'string'
&& Boolean(value.provider)
&& typeof value.model === 'string'
&& Boolean(value.model)
&& (value.reasoningEffort === undefined || typeof value.reasoningEffort === 'string');
}
function validateModelCatalog(value, method, { session = false } = {}) {
if (!value || typeof value !== 'object'
|| !Array.isArray(value.groups)
|| !Array.isArray(value.failures)) {
throw new Error(`Harness returned an invalid response for ${method}`);
}
for (const group of value.groups) {
if (!group || typeof group !== 'object'
|| typeof group.id !== 'string' || !group.id
|| typeof group.name !== 'string' || !group.name
|| !Array.isArray(group.models)) {
throw new Error(`Harness returned an invalid response for ${method}`);
}
for (const model of group.models) {
if (!model || typeof model !== 'object'
|| typeof model.id !== 'string' || !model.id
|| typeof model.name !== 'string' || !model.name) {
throw new Error(`Harness returned an invalid response for ${method}`);
}
}
}
for (const failure of value.failures) {
if (!failure || typeof failure !== 'object'
|| typeof failure.id !== 'string' || !failure.id
|| typeof failure.name !== 'string' || !failure.name
|| (failure.message !== undefined && typeof failure.message !== 'string')) {
throw new Error(`Harness returned an invalid response for ${method}`);
}
}
if (session && (!validModelSelection(value.current) || typeof value.routable !== 'boolean')) {
throw new Error(`Harness returned an invalid response for ${method}`);
}
return value;
}
function turnStoppedError() {
const error = new Error('Harness turn was stopped before producing a text reply');
error.code = 'turn-stopped';
return error;
}
function workspacePaths(value) {
if (!Array.isArray(value?.items)) return [];
return value.items.flatMap((item) => (
@ -319,10 +384,13 @@ export class HarnessClient {
#rpcIdPrefix;
#logPrefix;
#commandExecutor;
#controlExecutor;
#sessionMaintenanceExecutor;
#managedProcess = null;
#interactionRegistry;
#interactionOwnerships;
#interactionClaims;
#controlOwnerships;
constructor({
baseUrl,
@ -336,6 +404,8 @@ export class HarnessClient {
rpcIdPrefix = 'im',
logPrefix = 'dsh-im',
commandExecutor,
controlExecutor,
sessionMaintenanceExecutor,
}) {
if (typeof createWebSocket !== 'function') {
throw new TypeError('createWebSocket must be a function');
@ -352,6 +422,13 @@ export class HarnessClient {
if (commandExecutor !== undefined && typeof commandExecutor !== 'function') {
throw new TypeError('commandExecutor must be a function');
}
if (controlExecutor !== undefined && typeof controlExecutor !== 'function') {
throw new TypeError('controlExecutor must be a function');
}
if (sessionMaintenanceExecutor !== undefined
&& typeof sessionMaintenanceExecutor !== 'function') {
throw new TypeError('sessionMaintenanceExecutor must be a function');
}
this.#baseUrl = new URL(baseUrl);
this.#workspace = workspace;
// Keep an omitted preset absent so session.create resolves the Host's current default.
@ -364,9 +441,12 @@ export class HarnessClient {
this.#rpcIdPrefix = rpcIdPrefix.trim();
this.#logPrefix = logPrefix.trim();
this.#commandExecutor = commandExecutor;
this.#controlExecutor = controlExecutor;
this.#sessionMaintenanceExecutor = sessionMaintenanceExecutor;
this.#interactionRegistry = interactionRegistry(this.#baseUrl.origin);
this.#interactionOwnerships = this.#interactionRegistry.ownerships;
this.#interactionClaims = this.#interactionRegistry.claims;
this.#controlOwnerships = this.#interactionRegistry.controls;
}
async rpc(method, payload = {}, timeoutMs = 30_000, options = {}) {
@ -443,6 +523,61 @@ export class HarnessClient {
return workspaceSessions(workspace, workspaceList.archivedSessionIds, sessionList);
}
async listModels(options = {}) {
await this.ensureRunning(options);
const value = await this.rpc('llm.models', {}, 30_000, options);
return validateModelCatalog(value, 'llm.models');
}
async getSessionModels(sessionId, options = {}) {
if (typeof sessionId !== 'string' || !sessionId) throw new TypeError('sessionId is required');
await this.ensureRunning(options);
const value = await this.rpc('session.models', { sessionId }, 30_000, options);
return validateModelCatalog(value, 'session.models', { session: true });
}
async selectSessionModel(sessionId, selection, options = {}) {
if (typeof sessionId !== 'string' || !sessionId) throw new TypeError('sessionId is required');
if (!validModelSelection(selection)) {
throw new TypeError('A provider and model are required');
}
await this.ensureRunning(options);
const operation = (maintenanceSignal) => {
const signal = maintenanceSignal && options.signal
? AbortSignal.any([maintenanceSignal, options.signal])
: (maintenanceSignal ?? options.signal);
return this.rpc('session.selectModel', {
sessionId,
provider: selection.provider,
model: selection.model,
}, 30_000, signal ? { ...options, signal } : options);
};
const value = this.#sessionMaintenanceExecutor
? await this.#sessionMaintenanceExecutor({ sessionId, operation })
: await operation();
if (!value || typeof value !== 'object' || !validModelSelection(value.selected)) {
throw new Error('Harness returned an invalid response for session.selectModel');
}
return value;
}
async isSessionRunning(sessionId, options = {}) {
if (typeof sessionId !== 'string' || !sessionId) throw new TypeError('sessionId is required');
await this.ensureRunning(options);
const value = await this.rpc('session.list', {}, 30_000, options);
if (!value || typeof value !== 'object' || !Array.isArray(value.items)) {
throw new Error('Harness returned an invalid response for session.list');
}
for (const item of value.items) {
if (!item || typeof item !== 'object'
|| typeof item.sessionId !== 'string' || !item.sessionId
|| typeof item.running !== 'boolean') {
throw new Error('Harness returned an invalid response for session.list');
}
}
return value.items.find((item) => item.sessionId === sessionId)?.running ?? false;
}
async adoptWorkspaceSession(value, options = {}) {
return adoptRegisteredWorkspaceSession(this, value, options);
}
@ -585,6 +720,154 @@ export class HarnessClient {
}
}
#registerControlOwnership(ownership) {
if (!ownership.control) return;
let routes = this.#controlOwnerships.get(ownership.control.owner);
if (!routes) {
routes = new Map();
this.#controlOwnerships.set(ownership.control.owner, routes);
}
const owners = routes.get(ownership.control.key) ?? new Set();
owners.add(ownership);
routes.set(ownership.control.key, owners);
}
#unregisterControlOwnership(ownership) {
if (!ownership.control) return;
const routes = this.#controlOwnerships.get(ownership.control.owner);
const owners = routes?.get(ownership.control.key);
owners?.delete(ownership);
if (owners?.size === 0) routes.delete(ownership.control.key);
if (routes?.size === 0) this.#controlOwnerships.delete(ownership.control.owner);
}
#controlCandidates(sessionId, control) {
const normalized = normalizeControl(control);
if (!normalized) return [];
const owners = this.#controlOwnerships.get(normalized.owner)?.get(normalized.key);
return [...(owners ?? [])]
.filter((ownership) => ownership.sessionId === sessionId)
.sort((left, right) => left.order - right.order);
}
#activeControlOwnership(sessionId, control) {
return this.#controlCandidates(sessionId, control).find((ownership) => (
ownership.started
&& ownership.active
&& !ownership.completed
&& ownership.turn !== null
)) ?? null;
}
async #refreshControlOwnership(sessionId, control, options) {
const candidates = this.#controlCandidates(sessionId, control);
// An exact local owner is mandatory before even observing the Session.
// This prevents an unrelated chat bound to the same Session from using
// run state as authority to cancel or steer somebody else's turn.
if (candidates.length === 0) return null;
try {
const history = await this.rpc(
'session.history',
{ sessionId, maxMessages: 50 },
30_000,
options,
);
this.#consumeInteractionOwnerships(sessionId, history.events ?? []);
} catch (error) {
if (!(error instanceof HarnessRpcError) || error.code !== 'session-not-found') throw error;
for (const ownership of candidates) {
ownership.active = false;
ownership.completed = true;
}
return null;
}
return this.#activeControlOwnership(sessionId, control);
}
async hasActiveTurn(sessionId, control, options = {}) {
if (typeof sessionId !== 'string' || !sessionId) throw new TypeError('sessionId is required');
return Boolean(await this.#refreshControlOwnership(sessionId, control, options));
}
async stopActiveTurn(sessionId, control, options = {}) {
if (typeof sessionId !== 'string' || !sessionId) throw new TypeError('sessionId is required');
const ownership = await this.#refreshControlOwnership(sessionId, control, options);
if (!ownership) return false;
if (ownership.stopRequested) return true;
// Re-check after the refresh await. The owning ask may have completed and
// unregistered while session.history was in flight.
if (this.#activeControlOwnership(sessionId, control) !== ownership) return false;
ownership.stopRequested = true;
try {
if (this.#controlExecutor) {
const accepted = this.#controlExecutor({
sessionId,
expectedTurn: ownership.turn,
promptRpcId: ownership.promptRpcId,
action: 'stop',
});
if (accepted && typeof accepted.then === 'function') {
throw new TypeError('controlExecutor must return synchronously');
}
if (accepted !== undefined) {
if (typeof accepted !== 'boolean') {
throw new TypeError('controlExecutor must return a boolean or undefined');
}
if (!accepted) ownership.stopRequested = false;
return accepted;
}
}
await this.rpc(
'session.cancel',
{ sessionId, keepInbox: true },
30_000,
options,
);
return true;
} catch (error) {
if (this.#activeControlOwnership(sessionId, control) === ownership) {
ownership.stopRequested = false;
}
if (error instanceof HarnessRpcError && error.code === 'session-not-found') return false;
throw error;
}
}
async steerActiveTurn(sessionId, text, control, options = {}) {
if (typeof sessionId !== 'string' || !sessionId) throw new TypeError('sessionId is required');
if (typeof text !== 'string' || !text.trim()) {
throw new TypeError('Steering text is required');
}
const ownership = await this.#refreshControlOwnership(sessionId, control, options);
if (!ownership || ownership.stopRequested) return false;
if (this.#activeControlOwnership(sessionId, control) !== ownership) return false;
if (this.#controlExecutor) {
const accepted = this.#controlExecutor({
sessionId,
expectedTurn: ownership.turn,
promptRpcId: ownership.promptRpcId,
action: 'steer',
text,
});
if (accepted && typeof accepted.then === 'function') {
throw new TypeError('controlExecutor must return synchronously');
}
if (accepted !== undefined) {
if (typeof accepted !== 'boolean') {
throw new TypeError('controlExecutor must return a boolean or undefined');
}
return accepted;
}
}
await this.rpc('session.prompt', {
sessionId,
mode: 'steer',
content: [{ type: 'text', text }],
clientTimeZone: Intl.DateTimeFormat().resolvedOptions().timeZone,
}, 30_000, options);
return true;
}
#consumeInteractionOwnerships(sessionId, entries) {
for (const ownership of this.#interactionOwnerships.get(sessionId) ?? []) {
consumeInteractionOwnership(ownership, entries);
@ -634,6 +917,7 @@ export class HarnessClient {
const onInteractionResolved = typeof options.onInteractionResolved === 'function'
? options.onInteractionResolved
: undefined;
const control = normalizeControl(options.control);
await this.ensureRunning({ signal });
const before = await this.rpc(
'session.history',
@ -655,23 +939,29 @@ export class HarnessClient {
// The mux is host-global. A prompt RPC becomes the owner only when its
// durable user/message starts a turn, so two chats bound to one Session
// cannot answer each other's questions or approvals.
const ownership = interactionController
const ownership = interactionController || control
? {
sessionId,
promptRpcId,
active: false,
started: false,
completed: false,
stopRequested: false,
turn: null,
openTurn: null,
lastSeq: baselineSeq,
reconnect: null,
order: -1,
toolCalls: new Map(),
control,
}
: null;
let interactionTask = null;
if (ownership) this.#registerInteractionOwnership(sessionId, ownership);
if (ownership) {
this.#registerInteractionOwnership(sessionId, ownership);
this.#registerControlOwnership(ownership);
}
try {
if (interactionSignal) {
@ -706,39 +996,52 @@ export class HarnessClient {
clientTimeZone: Intl.DateTimeFormat().resolvedOptions().timeZone,
}, 30_000, { rpcId: promptRpcId, signal });
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
await sleep(300, signal);
const history = await this.rpc(
'session.history',
{ sessionId, maxMessages: 50 },
30_000,
{ signal },
);
const wasActive = ownership?.active === true;
if (ownership) {
this.#consumeInteractionOwnerships(sessionId, history.events ?? []);
if (!wasActive && ownership.active) ownership.reconnect?.();
}
const update = tracker.consume(history.events ?? []);
if (update && onUpdate) {
try {
await onUpdate(update);
} catch (error) {
console.warn(`[${this.#logPrefix}] ignored a progress update failure:`, error.message);
try {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
await sleep(300, signal);
const history = await this.rpc(
'session.history',
{ sessionId, maxMessages: 50 },
30_000,
{ signal },
);
const wasActive = ownership?.active === true;
if (ownership) {
this.#consumeInteractionOwnerships(sessionId, history.events ?? []);
if (!wasActive && ownership.active) ownership.reconnect?.();
}
const update = tracker.consume(history.events ?? []);
if (update && onUpdate) {
try {
await onUpdate(update);
} catch (error) {
console.warn(`[${this.#logPrefix}] ignored a progress update failure:`, error.message);
}
}
if (!tracker.finished) continue;
if (tracker.answer) return tracker.answer;
if (ownership?.stopRequested) throw turnStoppedError();
throw new Error(
`Harness turn ended without a text reply${tracker.reason ? ` (${JSON.stringify(tracker.reason)})` : ''}`,
);
}
if (!tracker.finished) continue;
throw new Error(`Harness reply timed out after ${Math.round(timeoutMs / 1_000)} seconds`);
} catch (error) {
// 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;
throw new Error(
`Harness turn ended without a text reply${tracker.reason ? ` (${JSON.stringify(tracker.reason)})` : ''}`,
);
if (error?.code === 'turn-stopped') throw error;
throw turnStoppedError();
}
throw new Error(`Harness reply timed out after ${Math.round(timeoutMs / 1_000)} seconds`);
} finally {
if (ownership) {
this.#unregisterControlOwnership(ownership);
this.#unregisterInteractionOwnership(sessionId, ownership);
}
interactionController?.abort(new DOMException('Harness turn finished', 'AbortError'));
if (interactionTask) await interactionTask.catch(() => undefined);
if (ownership) this.#unregisterInteractionOwnership(sessionId, ownership);
}
}
@ -775,6 +1078,10 @@ export class HarnessClient {
socket.removeEventListener('error', handleError);
signal.removeEventListener('abort', handleAbort);
if (ownership?.reconnect === close) ownership.reconnect = null;
if (signal.aborted) {
resolve();
return;
}
void callbackTail.then(() => {
const failure = error ?? callbackFailure;
if (failure) reject(failure);

View file

@ -0,0 +1,314 @@
import { splitWorkspaceCommandMessage } from './workspace-command.mjs';
import { WORKSPACE_SESSION_STALE } from './workspace-session.mjs';
import { withSessionBindingLock } from './session-binding-lock.mjs';
const MODEL_COMMAND = /^\/model(?=$|\s)/i;
const MODELS_COMMAND = /^\/models(?=$|\s)/i;
const MODEL_USAGE = '用法:/model <provider>/<model>';
const MODELS_USAGE = '用法:/models(不带参数)';
const SESSION_BINDING_CHANGED = 'session-binding-changed';
const UNSAFE_DISPLAY_TEXT_GLOBAL = /[\p{Cc}\p{Cf}\p{Zl}\p{Zp}]+/gu;
function commandResult(message) {
return {
handled: true,
message,
messages: splitWorkspaceCommandMessage(message),
};
}
function safeDisplayText(value) {
if (typeof value !== 'string') return '';
return value.replace(UNSAFE_DISPLAY_TEXT_GLOBAL, ' ').replace(/\s+/gu, ' ').trim();
}
function rpcOptions(signal) {
return signal ? { signal } : {};
}
function normalizeCatalog(value, { requireCurrent = false } = {}) {
if (!value || typeof value !== 'object'
|| !Array.isArray(value.groups) || !Array.isArray(value.failures)) {
throw new TypeError('Harness returned an invalid model catalog');
}
const groups = value.groups.map((group) => {
if (!group || typeof group !== 'object'
|| typeof group.id !== 'string' || !group.id
|| typeof group.name !== 'string' || !group.name
|| !Array.isArray(group.models)) {
throw new TypeError('Harness returned an invalid model provider group');
}
return {
id: group.id,
name: group.name,
models: group.models.map((model) => {
if (!model || typeof model !== 'object'
|| typeof model.id !== 'string' || !model.id
|| typeof model.name !== 'string' || !model.name) {
throw new TypeError('Harness returned an invalid model');
}
return { id: model.id, name: model.name };
}),
};
});
const failures = value.failures.map((failure) => {
if (!failure || typeof failure !== 'object'
|| typeof failure.id !== 'string' || !failure.id
|| typeof failure.name !== 'string' || !failure.name) {
throw new TypeError('Harness returned an invalid model provider failure');
}
return { id: failure.id, name: failure.name };
});
let current = null;
if (value.current !== undefined) {
if (!value.current || typeof value.current !== 'object'
|| typeof value.current.provider !== 'string' || !value.current.provider
|| typeof value.current.model !== 'string' || !value.current.model) {
throw new TypeError('Harness returned an invalid current model');
}
current = { provider: value.current.provider, model: value.current.model };
} else if (requireCurrent) {
throw new TypeError('Harness returned no current model');
}
return { groups, failures, current };
}
function modelId(provider, model) {
return `${provider}/${model}`;
}
function matchingModel(catalog, requested) {
for (const group of catalog.groups) {
for (const model of group.models) {
if (modelId(group.id, model.id) === requested) {
return { provider: group.id, model: model.id };
}
}
}
return null;
}
function formatCatalog(catalog) {
const currentId = catalog.current
? modelId(catalog.current.provider, catalog.current.model)
: null;
const lines = ['可用模型:'];
if (catalog.groups.length === 0) lines.push('', '当前没有可用模型。');
for (const group of catalog.groups) {
lines.push('', safeDisplayText(group.name) || safeDisplayText(group.id));
for (const model of group.models) {
const id = modelId(group.id, model.id);
lines.push(`- ${safeDisplayText(id)}${id === currentId ? '(当前)' : ''}`);
}
}
if (catalog.failures.length > 0) {
lines.push('', '以下模型提供方暂时不可用:');
for (const failure of catalog.failures) {
lines.push(`- ${safeDisplayText(failure.name) || safeDisplayText(failure.id)}`);
}
}
return lines.join('\n');
}
function currentModelMessage(current) {
return [
'当前模型:',
modelId(current.provider, current.model),
'',
'查看全部模型:/models',
'切换模型:/model <provider>/<model>',
].join('\n');
}
function noSessionMessage() {
return [
'当前聊天还没有会话。',
'',
'查看模型:/models',
'选择模型:/model <provider>/<model>',
].join('\n');
}
function errorCode(error) {
return error?.code ?? error?.failure?.code;
}
function modelErrorMessage(error, action) {
const code = errorCode(error);
if (code === 'agent-busy') {
return '当前任务正在运行,请等待完成或先发送 /stop。';
}
if (code === 'session-not-found') {
return '当前聊天绑定的会话已不存在,请重试。';
}
if (code === 'model-unavailable') {
return '无法切换到该模型。模型当前不可用,或不支持当前会话中的图片。';
}
if (code === WORKSPACE_SESSION_STALE || code === 'workspace-bot-not-found') {
return '工作区或机器人状态已发生变化,请重试。';
}
if (code === SESSION_BINDING_CHANGED) {
return '当前聊天绑定的会话已发生变化,请重试。';
}
if (code === 'cancelled' || error?.name === 'AbortError') {
return action === 'list' ? '获取模型列表已取消。' : '模型切换已取消。';
}
return action === 'list'
? '暂时无法获取模型列表,请稍后重试。'
: '模型切换失败,请稍后重试。';
}
async function boundSession(harness, state, key, options) {
if (typeof state?.sessionFor !== 'function') return null;
const sessionId = state.sessionFor(key);
if (typeof sessionId !== 'string' || !sessionId) return null;
if (typeof harness?.workspaceSession !== 'function') {
throw new TypeError('Harness does not support workspace sessions');
}
const session = harness.workspaceSession(sessionId);
if (!session || typeof session.sessionExists !== 'function') {
throw new TypeError('Harness returned an invalid workspace session');
}
if (await session.sessionExists(options)) return { sessionId, session };
if (typeof state.clearSession === 'function' && state.sessionFor(key) === sessionId) {
await state.clearSession(key);
}
return null;
}
async function sessionIsBusy(session, control, options) {
if (typeof session?.isRunning !== 'function'
|| typeof session?.hasActiveTurn !== 'function') {
throw new TypeError('Harness session does not expose run state');
}
if (await session.isRunning(options)) return true;
return Boolean(await session.hasActiveTurn(control, options));
}
async function listCatalog(harness, options) {
if (typeof harness?.listModels !== 'function') {
throw new TypeError('Harness does not support listing models');
}
return normalizeCatalog(await harness.listModels(options));
}
async function sessionCatalog(session, options) {
if (typeof session?.models !== 'function') {
throw new TypeError('Harness session does not support listing models');
}
return normalizeCatalog(await session.models(options), { requireCurrent: true });
}
function isModelsCommand(command) {
return MODELS_COMMAND.test(command);
}
export function isModelCommand(text) {
if (typeof text !== 'string') return false;
const command = text.trim();
return MODELS_COMMAND.test(command) || MODEL_COMMAND.test(command);
}
export async function runModelCommand(text, harness, state, key, options = {}) {
if (!isModelCommand(text)) return null;
const command = text.trim();
if (options.hasImages) {
return commandResult('模型命令仅支持纯文字,请移除图片后重试。');
}
const requestOptions = rpcOptions(options.signal);
if (isModelsCommand(command)) {
if (!/^\/models[ \t]*$/iu.test(command)) return commandResult(MODELS_USAGE);
try {
const bound = await boundSession(harness, state, key, requestOptions);
const catalog = bound
? await sessionCatalog(bound.session, requestOptions)
: await listCatalog(harness, requestOptions);
return commandResult(formatCatalog(catalog));
} catch (error) {
return commandResult(modelErrorMessage(error, 'list'));
}
}
const match = /^\/model(?:[ \t]+([^\s]+))?[ \t]*$/iu.exec(command);
if (!match) return commandResult(MODEL_USAGE);
const requested = match[1];
if (!requested) {
try {
const bound = await boundSession(harness, state, key, requestOptions);
if (!bound) return commandResult(noSessionMessage());
const catalog = await sessionCatalog(bound.session, requestOptions);
return commandResult(currentModelMessage(catalog.current));
} catch (error) {
return commandResult(modelErrorMessage(error, 'select'));
}
}
if (!requested.includes('/') || requested.startsWith('/') || requested.endsWith('/')) {
return commandResult(MODEL_USAGE);
}
if (options.pendingInteraction) {
return commandResult([
'当前任务正在等待你的回答或审批。',
'',
'请先处理当前请求,或者发送 /stop 停止任务。',
].join('\n'));
}
try {
return await withSessionBindingLock(state, key, async () => {
const bound = await boundSession(harness, state, key, requestOptions);
if (bound && await sessionIsBusy(bound.session, options.control, requestOptions)) {
return commandResult('当前任务正在运行,请等待完成或先发送 /stop。');
}
const catalog = bound
? await sessionCatalog(bound.session, requestOptions)
: await listCatalog(harness, requestOptions);
const selection = matchingModel(catalog, requested);
if (!selection) {
return commandResult([
`没有找到模型:${safeDisplayText(requested)}`,
'',
'请发送 /models 查看可用模型。',
].join('\n'));
}
if (bound) {
if (typeof bound.session.selectModel !== 'function') {
throw new TypeError('Harness session does not support model selection');
}
await bound.session.selectModel(selection, requestOptions);
} else {
if (typeof harness?.createSession !== 'function'
|| typeof harness?.workspaceSession !== 'function'
|| typeof state?.sessionFor !== 'function'
|| typeof state?.setSession !== 'function') {
throw new TypeError('Harness cannot create a conversation session');
}
const sessionId = await harness.createSession(requestOptions);
if (typeof sessionId !== 'string' || !sessionId) {
throw new TypeError('Harness returned an invalid session id');
}
const session = harness.workspaceSession(sessionId);
if (!session || typeof session.selectModel !== 'function') {
throw new TypeError('Harness session does not support model selection');
}
await session.selectModel(selection, requestOptions);
const currentSessionId = state.sessionFor(key);
if (typeof currentSessionId === 'string' && currentSessionId) {
const changed = new Error('Conversation binding changed during model selection');
changed.code = SESSION_BINDING_CHANGED;
throw changed;
}
if (await state.setSession(key, sessionId) === false) {
const stale = new Error('Workspace changed while binding the new session');
stale.code = WORKSPACE_SESSION_STALE;
throw stale;
}
}
return commandResult(`模型已切换为:\n${modelId(selection.provider, selection.model)}\n\n后续消息将使用该模型。`);
});
} catch (error) {
return commandResult(modelErrorMessage(error, 'select'));
}
}

View file

@ -0,0 +1,40 @@
const bindingLocks = new WeakMap();
function stateLocks(state) {
let locks = bindingLocks.get(state);
if (!locks) {
locks = new Map();
bindingLocks.set(state, locks);
}
return locks;
}
/**
* Serialize the short Session binding transaction for one conversation.
* The caller must keep long-running Session work, such as ask(), outside.
*/
export async function withSessionBindingLock(state, key, operation) {
const stateType = typeof state;
if ((stateType !== 'object' && stateType !== 'function') || state === null) {
throw new TypeError('state is required');
}
if (typeof key !== 'string' || !key) throw new TypeError('conversation key is required');
if (typeof operation !== 'function') throw new TypeError('operation is required');
const locks = stateLocks(state);
const previous = locks.get(key) ?? Promise.resolve();
let release;
const current = new Promise((resolve) => { release = resolve; });
locks.set(key, current);
await previous;
try {
return await operation();
} finally {
release();
if (locks.get(key) === current) {
locks.delete(key);
if (locks.size === 0) bindingLocks.delete(state);
}
}
}

View file

@ -1,9 +1,17 @@
import { runWorkspaceCommand } from './workspace-command.mjs';
import { runCompactCommand } from './compact-command.mjs';
import {
isControlCommand,
runControlCommand,
} from './control-command.mjs';
import {
rememberConnectionTestTarget,
sendRememberedConnectionTest,
} from './connection-test.mjs';
import {
isModelCommand,
runModelCommand,
} from './model-command.mjs';
import { askInWorkspaceSession } from './workspace-session.mjs';
import { HarnessApprovalQueue } from './harness-approval.mjs';
import {
@ -56,6 +64,7 @@ export class TextHarnessBridge {
#interactionKeys = new Map();
#acceptedMessageIds = new Set();
#approvalTasks = new Set();
#commandTasks = new Set();
#approvals;
constructor({
@ -110,6 +119,24 @@ export class TextHarnessBridge {
const key = `${kind}:${conversationId}`;
const pending = this.#pendingInteractions.get(key);
const text = cleanText(normalized.content);
const commandRunner = isControlCommand(text)
? runControlCommand
: (isModelCommand(text) ? runModelCommand : null);
if (commandRunner && (normalized.kind !== 'group' || normalized.addressed === true)) {
let task;
task = this.#processFastCommand(
normalized,
messageId,
key,
commandRunner,
).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
const approval = this.#approvals.claimReply({
key,
actor: senderId,
@ -201,9 +228,48 @@ export class TextHarnessBridge {
pending.queue ? [pending.queue] : []
)),
...this.#approvalTasks,
...this.#commandTasks,
]);
}
async #processFastCommand(message, messageId, key, runner) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
const target = message.replyTarget;
try {
const result = await runner(
cleanText(message.content),
this.#harness,
this.#state,
key,
{
signal: this.#signal,
hasImages: hasInboundImages(message),
pendingInteraction: this.#pendingInteractions.has(key)
|| this.#approvals.hasPending(key),
control: { owner: this, key },
},
);
if (result?.stopped) {
await Promise.allSettled([
this.#cancelPendingInteraction(key),
this.#approvals.closeRoute(key),
]);
}
for (const reply of result?.messages ?? [result?.message]) {
if (reply) await this.#bot.sendText(target, reply);
}
this.#status.lastError = null;
} catch (error) {
if (error?.code === 'turn-stopped' || this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.(`[dsh-im:${this.#descriptor.key}] failed to process a command:`, error);
await this.#bot.sendText(target, '消息处理失败,请稍后重试。').catch(() => undefined);
}
}
sendConnectionTest(text) {
return sendRememberedConnectionTest({
state: this.#state,
@ -250,6 +316,10 @@ export class TextHarnessBridge {
'/workspacelist 列出工作区绝对路径',
'/sessionlist [工作区序号或绝对路径] 列出会话 ID 和标题',
'/session Session ID 或当前工作区序号 将当前聊天绑定到指定会话',
'/models 列出所有可用模型',
'/model [模型ID] 查看或切换当前会话模型',
'/stop 停止当前任务',
'/steer 补充指令 纠偏当前任务',
'/status 检查连接状态',
'/help 显示本帮助',
].join('\n'));
@ -316,6 +386,7 @@ export class TextHarnessBridge {
askOptions: {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
control: { owner: this, key: conversationKey },
onUpdate: stream ? async (update) => {
const progress = update.type === 'text' ? update.text
: update.type === 'tool' ? `正在使用${update.name}…` : update.text;
@ -347,6 +418,16 @@ export class TextHarnessBridge {
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} catch (error) {
if (error?.code === 'turn-stopped') {
if (stream) {
try {
await stream.finish('已停止。');
} catch {
stream.cancel?.();
}
}
return;
}
stream?.cancel?.();
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);

View file

@ -1,3 +1,5 @@
import { withSessionBindingLock } from './session-binding-lock.mjs';
export const WORKSPACE_SESSION_STALE = 'workspace-session-stale';
function workspaceSession(harness, sessionId) {
@ -7,6 +9,12 @@ function workspaceSession(harness, sessionId) {
return Object.freeze({
sessionId,
sessionExists: (...args) => harness.sessionExists(sessionId, ...args),
models: (...args) => harness.getSessionModels(sessionId, ...args),
selectModel: (...args) => harness.selectSessionModel(sessionId, ...args),
isRunning: (...args) => harness.isSessionRunning(sessionId, ...args),
hasActiveTurn: (...args) => harness.hasActiveTurn(sessionId, ...args),
stopActiveTurn: (...args) => harness.stopActiveTurn(sessionId, ...args),
steerActiveTurn: (...args) => harness.steerActiveTurn(sessionId, ...args),
ask: (...args) => harness.ask(sessionId, ...args),
});
}
@ -40,16 +48,20 @@ export async function askInWorkspaceSession({
}) {
while (true) {
try {
let sessionId = state.sessionFor(key);
let session = sessionId ? workspaceSession(harness, sessionId) : null;
if (!session || !(await sessionExists(session, existsOptions))) {
sessionId = await createSession(harness, createOptions);
if (await state.setSession(key, sessionId) === false) continue;
session = workspaceSession(harness, sessionId);
}
const binding = await withSessionBindingLock(state, key, async () => {
let sessionId = state.sessionFor(key);
let session = sessionId ? workspaceSession(harness, sessionId) : null;
if (!session || !(await sessionExists(session, existsOptions))) {
sessionId = await createSession(harness, createOptions);
if (await state.setSession(key, sessionId) === false) return null;
session = workspaceSession(harness, sessionId);
}
return { sessionId, session };
});
if (!binding) continue;
return {
sessionId,
answer: await session.ask(content ?? text, askOptions),
sessionId: binding.sessionId,
answer: await binding.session.ask(content ?? text, askOptions),
};
} catch (error) {
if (error?.code !== WORKSPACE_SESSION_STALE) throw error;

View file

@ -6,6 +6,14 @@ import {
} from '../shared/harness-question.mjs';
import { HarnessApprovalQueue } from '../shared/harness-approval.mjs';
import { runCompactCommand } from '../shared/compact-command.mjs';
import {
isControlCommand,
runControlCommand,
} from '../shared/control-command.mjs';
import {
isModelCommand,
runModelCommand,
} from '../shared/model-command.mjs';
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
import {
@ -26,6 +34,10 @@ const HELP_TEXT = [
'/workspacelist 列出工作区绝对路径',
'/sessionlist [工作区序号或绝对路径] 列出会话 ID 和标题',
'/session Session ID 或当前工作区序号 将当前聊天绑定到指定会话',
'/models 列出所有可用模型',
'/model [模型ID] 查看或切换当前会话模型',
'/stop 停止当前任务',
'/steer 补充指令 纠偏当前任务',
'/status 检查连接状态',
'/help 显示本帮助',
].join('\n');
@ -223,6 +235,7 @@ export class WecomHarnessBridge {
#interactionKeys = new Map();
#acceptedMessageIds = new Set();
#approvalTasks = new Set();
#commandTasks = new Set();
#approvals;
#prefetchedImageCount = 0;
@ -274,6 +287,33 @@ export class WecomHarnessBridge {
rememberConnectionTestTarget(this.#state, { chatId });
}
const pending = this.#pendingInteractions.get(key);
const commandMessage = wecomInboundMessage(frame, this.#client);
const commandText = nonEmptyString(commandMessage.content) ?? '';
const commandRunner = isControlCommand(commandText)
? runControlCommand
: (isModelCommand(commandText) ? runModelCommand : null);
if (commandRunner) {
let task;
task = this.#processFastCommand(
frame,
messageId,
chatId,
key,
commandMessage,
commandRunner,
).catch((error) => {
if (error?.code === 'turn-stopped' || this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-im:wecom] failed to process a command');
return this.#sendImmediate(frame, chatId, '消息处理失败,请稍后重试。')
.catch(() => undefined);
}).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
const approval = this.#approvals.claimReply({
key,
actor: senderId,
@ -377,9 +417,35 @@ export class WecomHarnessBridge {
pending.queue ? [pending.queue] : []
)),
...this.#approvalTasks,
...this.#commandTasks,
]);
}
async #processFastCommand(frame, messageId, chatId, key, message, runner) {
this.#signal?.throwIfAborted();
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
const result = await runner(message.content, this.#harness, this.#state, key, {
signal: this.#signal,
hasImages: hasInboundImages(message),
pendingInteraction: this.#pendingInteractions.has(key)
|| this.#approvals.hasPending(key),
control: { owner: this, key },
});
if (result?.stopped) {
await Promise.allSettled([
this.#cancelPendingInteraction(key),
this.#approvals.closeRoute(key),
]);
}
for (const reply of result?.messages ?? [result?.message]) {
if (reply) await this.#sendImmediate(frame, chatId, reply);
}
this.#status.lastError = null;
}
async #sendActive(chatId, text) {
for (const chunk of splitUtf8(text)) {
this.#signal?.throwIfAborted();
@ -490,6 +556,7 @@ export class WecomHarnessBridge {
askOptions: {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
control: { owner: this, key },
onUpdate: streamStarted && typeof this.#client.replyStreamNonBlocking === 'function'
? async (update) => {
const progress = splitUtf8(progressText(update))[0];
@ -525,6 +592,14 @@ export class WecomHarnessBridge {
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} catch (error) {
if (error?.code === 'turn-stopped') {
if (streamStarted && streamId) {
await this.#client.replyStream(frame, streamId, '已停止。', true)
.catch(() => undefined);
}
await this.#state.markSeen(messageId);
return;
}
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-im:wecom] failed to process an inbound message');

View file

@ -11,6 +11,14 @@ import {
} from '../shared/harness-question.mjs';
import { HarnessApprovalQueue } from '../shared/harness-approval.mjs';
import { runCompactCommand } from '../shared/compact-command.mjs';
import {
isControlCommand,
runControlCommand,
} from '../shared/control-command.mjs';
import {
isModelCommand,
runModelCommand,
} from '../shared/model-command.mjs';
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
import {
@ -32,6 +40,10 @@ const HELP_TEXT = [
'/workspacelist 列出工作区绝对路径',
'/sessionlist [工作区序号或绝对路径] 列出会话 ID 和标题',
'/session Session ID 或当前工作区序号 将当前聊天绑定到指定会话',
'/models 列出所有可用模型',
'/model [模型ID] 查看或切换当前会话模型',
'/stop 停止当前任务',
'/steer 补充指令 纠偏当前任务',
'/status 检查连接状态',
'/help 显示本帮助',
].join('\n');
@ -85,6 +97,7 @@ export class WeixinHarnessBridge {
#interactionKeys = new Map();
#acceptedMessageIds = new Set();
#approvalTasks = new Set();
#commandTasks = new Set();
#approvals;
constructor({
@ -136,6 +149,34 @@ export class WeixinHarnessBridge {
const contextToken = nonEmptyString(message?.context_token) ?? undefined;
const runId = nonEmptyString(message?.run_id) ?? undefined;
const pending = this.#pendingInteractions.get(key);
const commandText = nonEmptyString(extractWeixinText(message)) ?? '';
const commandRunner = isControlCommand(commandText)
? runControlCommand
: (isModelCommand(commandText) ? runModelCommand : null);
if (commandRunner && sender === this.#ownerUserId) {
let task;
task = this.#processFastCommand(
message,
messageId,
key,
sender,
contextToken,
runId,
commandText,
commandRunner,
).catch((error) => {
if (error?.code === 'turn-stopped' || this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-weixin] failed to process a command:', error);
return this.#send(sender, '消息处理失败,请稍后重试。', contextToken, runId)
.catch(() => undefined);
}).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
const approval = this.#approvals.claimReply({
key,
actor: sender,
@ -211,9 +252,44 @@ export class WeixinHarnessBridge {
pending.queue ? [pending.queue] : []
)),
...this.#approvalTasks,
...this.#commandTasks,
]);
}
async #processFastCommand(
message,
messageId,
key,
sender,
contextToken,
runId,
text,
runner,
) {
this.#signal?.throwIfAborted();
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
const result = await runner(text, this.#harness, this.#state, key, {
signal: this.#signal,
hasImages: hasWeixinImageItems(message),
pendingInteraction: this.#pendingInteractions.has(key)
|| this.#approvals.hasPending(key),
control: { owner: this, key },
});
if (result?.stopped) {
await Promise.allSettled([
this.#cancelPendingInteraction(key),
this.#approvals.closeRoute(key),
]);
}
for (const reply of result?.messages ?? [result?.message]) {
if (reply) await this.#send(sender, reply, contextToken, runId);
}
this.#status.lastError = null;
}
async #process(message, key, { alreadyRecorded = false } = {}) {
this.#signal?.throwIfAborted();
const messageId = weixinMessageId(message);
@ -303,6 +379,7 @@ export class WeixinHarnessBridge {
askOptions: {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
control: { owner: this, key },
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key,
actor: sender,
@ -324,6 +401,10 @@ export class WeixinHarnessBridge {
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} catch (error) {
if (error?.code === 'turn-stopped') {
await this.#state.markSeen(messageId);
return;
}
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-weixin] failed to process an inbound message:', error);