mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-10 03:50:45 +08:00
feat: support per-bot workspace switching
This commit is contained in:
parent
adcace3512
commit
44dd50216a
64 changed files with 3878 additions and 712 deletions
|
|
@ -3,6 +3,8 @@ import {
|
|||
splitDingtalkText,
|
||||
} from './dingtalk-api.mjs';
|
||||
import { createDingTalkCardStream } from './dingtalk-card-stream.mjs';
|
||||
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
|
||||
const CARD_INITIAL_TEXT = '已连接 DeepSeek Harness,正在思考…';
|
||||
const CARD_ERROR_TEXT = '消息处理失败,请稍后重试。';
|
||||
|
|
@ -12,6 +14,7 @@ const HELP_TEXT = [
|
|||
'',
|
||||
'直接发送文字即可继续当前会话。',
|
||||
'/new 开启一个全新会话',
|
||||
'/workspace 工作区绝对路径 切换工作区',
|
||||
'/status 检查连接状态',
|
||||
'/help 显示本帮助',
|
||||
].join('\n');
|
||||
|
|
@ -213,12 +216,12 @@ export class DingtalkHarnessBridge {
|
|||
await this.#send(sessionWebhook, '已开启新会话。请发送你的问题。');
|
||||
return;
|
||||
}
|
||||
|
||||
let sessionId = this.#state.sessionFor(key);
|
||||
if (!sessionId || !(await this.#harness.sessionExists(sessionId, { signal: this.#signal }))) {
|
||||
sessionId = await this.#harness.createSession({ signal: this.#signal });
|
||||
await this.#state.setSession(key, sessionId);
|
||||
const workspaceCommand = await runWorkspaceCommand(text, this.#harness);
|
||||
if (workspaceCommand) {
|
||||
await this.#send(sessionWebhook, workspaceCommand.message);
|
||||
return;
|
||||
}
|
||||
|
||||
if (typeof this.#api.createAiCard === 'function'
|
||||
&& typeof this.#api.updateAiCard === 'function'
|
||||
&& typeof this.#api.finishAiCard === 'function') {
|
||||
|
|
@ -232,12 +235,20 @@ export class DingtalkHarnessBridge {
|
|||
});
|
||||
cardStarted = await cardStream.start(CARD_INITIAL_TEXT);
|
||||
}
|
||||
const answer = await this.#harness.ask(sessionId, text, {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
signal: this.#signal,
|
||||
onUpdate: cardStarted
|
||||
? (update) => cardStream.push(progressText(update))
|
||||
: undefined,
|
||||
const { answer } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
text,
|
||||
createOptions: { signal: this.#signal },
|
||||
existsOptions: { signal: this.#signal },
|
||||
askOptions: {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
signal: this.#signal,
|
||||
onUpdate: cardStarted
|
||||
? (update) => cardStream.push(progressText(update))
|
||||
: undefined,
|
||||
},
|
||||
});
|
||||
const streamed = cardStarted && await cardStream.finish(answer);
|
||||
if (!streamed) await this.#send(sessionWebhook, answer);
|
||||
|
|
|
|||
|
|
@ -610,7 +610,12 @@ export class DingtalkController {
|
|||
await this.#stopRuntime(identity.botId);
|
||||
if (previousConfig) await this.#configStore.save(previousConfig).catch(() => undefined);
|
||||
else if (this.#configStore.get(identity.botId)) {
|
||||
await this.#configStore.remove(identity.botId).catch(() => undefined);
|
||||
const removed = await this.#configStore.remove(identity.botId).catch(() => null);
|
||||
if (removed) {
|
||||
await this.#deleteState({ botId: identity.botId, config }).catch((cleanupError) => {
|
||||
this.#logger.warn?.('[dsh-dingtalk] failed to clean up cancelled bot state:', cleanupError);
|
||||
});
|
||||
}
|
||||
}
|
||||
await this.#restoreCredential(identity.secretRef, previousSecret);
|
||||
if (previousConfig && cleanString(previousSecret?.value)) {
|
||||
|
|
|
|||
|
|
@ -217,10 +217,11 @@ export class HarnessClient {
|
|||
}
|
||||
|
||||
async workspaceId(options = {}) {
|
||||
const { items } = await this.rpc('workspace.list', {}, 30_000, options);
|
||||
const existing = items.find((item) => item.path === this.#workspace);
|
||||
const { workspace = this.#workspace, ...rpcOptions } = options;
|
||||
const { items } = await this.rpc('workspace.list', {}, 30_000, rpcOptions);
|
||||
const existing = items.find((item) => item.path === workspace);
|
||||
if (existing) return existing.workspaceId;
|
||||
const created = await this.rpc('workspace.create', { path: this.#workspace }, 30_000, options);
|
||||
const created = await this.rpc('workspace.create', { path: workspace }, 30_000, rpcOptions);
|
||||
return created.workspace.workspaceId;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -117,6 +117,11 @@ export class DingtalkStateStore {
|
|||
await this.#persist();
|
||||
}
|
||||
|
||||
async clearSessions() {
|
||||
this.#state.sessions = {};
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
hasSeen(messageId) {
|
||||
const id = nonEmptyString(messageId);
|
||||
return Boolean(id && this.#state.seenMessageIds.includes(id));
|
||||
|
|
|
|||
|
|
@ -5,12 +5,15 @@ import {
|
|||
isBotSender,
|
||||
splitText,
|
||||
} from './message-utils.mjs';
|
||||
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
|
||||
const HELP_TEXT = [
|
||||
'北汇星河 AIOS 已连接 DeepSeek Harness。',
|
||||
'',
|
||||
'直接发送问题即可继续当前会话。',
|
||||
'/new 开启一个全新会话',
|
||||
'/workspace 工作区绝对路径 切换工作区',
|
||||
'/status 检查连接状态',
|
||||
'/help 显示本帮助',
|
||||
].join('\n');
|
||||
|
|
@ -106,25 +109,30 @@ export class FeishuHarnessBridge {
|
|||
await this.#send(event.message.chat_id, '飞书机器人与 DeepSeek Harness 连接正常。');
|
||||
return;
|
||||
}
|
||||
|
||||
let sessionId = this.#state.sessionFor(key);
|
||||
if (!sessionId || !(await this.#harness.sessionExists(sessionId))) {
|
||||
sessionId = await this.#harness.createSession();
|
||||
await this.#state.setSession(key, sessionId);
|
||||
const workspaceCommand = await runWorkspaceCommand(text, this.#harness);
|
||||
if (workspaceCommand) {
|
||||
await this.#send(event.message.chat_id, workspaceCommand.message);
|
||||
return;
|
||||
}
|
||||
|
||||
console.info(`[bridge] processing ${event.message.chat_type} message ${messageId} in ${sessionId}`);
|
||||
await this.#answerWithStream(event, sessionId, text);
|
||||
console.info(`[bridge] processing ${event.message.chat_type} message ${messageId}`);
|
||||
await this.#answerWithStream(event, key, text);
|
||||
this.#status.messagesReplied += 1;
|
||||
this.#status.lastReplyAt = new Date().toISOString();
|
||||
this.#status.lastError = null;
|
||||
}
|
||||
|
||||
async #answerWithStream(event, sessionId, text) {
|
||||
async #answerWithStream(event, key, text) {
|
||||
const chatId = event.message.chat_id;
|
||||
const messageId = event.message.message_id;
|
||||
if (!this.#channel?.stream) {
|
||||
const answer = await this.#harness.ask(sessionId, text, { timeoutMs: this.#replyTimeoutMs });
|
||||
const { answer } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
text,
|
||||
askOptions: { timeoutMs: this.#replyTimeoutMs },
|
||||
});
|
||||
for (const chunk of splitText(answer)) await this.#send(chatId, chunk);
|
||||
this.#status.streamFallbacks = (this.#status.streamFallbacks ?? 0) + 1;
|
||||
return;
|
||||
|
|
@ -136,13 +144,19 @@ export class FeishuHarnessBridge {
|
|||
await this.#channel.stream(chatId, {
|
||||
markdown: async (controller) => {
|
||||
promptStarted = true;
|
||||
completedAnswer = await this.#harness.ask(sessionId, text, {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
onUpdate: async (update) => {
|
||||
await controller.setContent(this.#progressText(update));
|
||||
this.#status.streamUpdates = (this.#status.streamUpdates ?? 0) + 1;
|
||||
({ answer: completedAnswer } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
text,
|
||||
askOptions: {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
onUpdate: async (update) => {
|
||||
await controller.setContent(this.#progressText(update));
|
||||
this.#status.streamUpdates = (this.#status.streamUpdates ?? 0) + 1;
|
||||
},
|
||||
},
|
||||
});
|
||||
}));
|
||||
await controller.setContent(completedAnswer);
|
||||
},
|
||||
}, { replyTo: messageId });
|
||||
|
|
@ -158,7 +172,13 @@ export class FeishuHarnessBridge {
|
|||
if (promptStarted) throw error;
|
||||
|
||||
console.warn('[bridge] native Feishu stream unavailable; using text fallback:', error.message);
|
||||
const answer = await this.#harness.ask(sessionId, text, { timeoutMs: this.#replyTimeoutMs });
|
||||
const { answer } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
text,
|
||||
askOptions: { timeoutMs: this.#replyTimeoutMs },
|
||||
});
|
||||
for (const chunk of splitText(answer)) await this.#send(chatId, chunk);
|
||||
this.#status.streamFallbacks = (this.#status.streamFallbacks ?? 0) + 1;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -197,17 +197,18 @@ export class HarnessClient {
|
|||
throw new Error(`Harness did not become ready: ${lastError?.message ?? 'timeout'}`);
|
||||
}
|
||||
|
||||
async workspaceId() {
|
||||
async workspaceId(options = {}) {
|
||||
const workspace = options.workspace ?? this.#workspace;
|
||||
const { items } = await this.rpc('workspace.list', {});
|
||||
const existing = items.find((item) => item.path === this.#workspace);
|
||||
const existing = items.find((item) => item.path === workspace);
|
||||
if (existing) return existing.workspaceId;
|
||||
const created = await this.rpc('workspace.create', { path: this.#workspace });
|
||||
const created = await this.rpc('workspace.create', { path: workspace });
|
||||
return created.workspace.workspaceId;
|
||||
}
|
||||
|
||||
async createSession() {
|
||||
async createSession(options = {}) {
|
||||
await this.ensureRunning();
|
||||
const workspaceId = await this.workspaceId();
|
||||
const workspaceId = await this.workspaceId(options);
|
||||
const created = await this.rpc('session.create', {
|
||||
workspaceId,
|
||||
agentPreset: this.#agentPreset,
|
||||
|
|
|
|||
|
|
@ -41,6 +41,11 @@ export class StateStore {
|
|||
await this.#persist();
|
||||
}
|
||||
|
||||
async clearSessions() {
|
||||
this.#state.sessions = {};
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
hasSeen(messageId) {
|
||||
return this.#state.seenMessageIds.includes(messageId);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,8 +1,12 @@
|
|||
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
|
||||
const HELP_TEXT = [
|
||||
'QQ 机器人已连接 DeepSeek Harness。',
|
||||
'',
|
||||
'直接发送文字即可继续当前会话。',
|
||||
'/new 开启一个全新会话',
|
||||
'/workspace 工作区绝对路径 切换工作区',
|
||||
'/status 检查连接状态',
|
||||
'/help 显示本帮助',
|
||||
].join('\n');
|
||||
|
|
@ -121,11 +125,11 @@ export class QqHarnessBridge {
|
|||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
|
||||
let sessionId = this.#state.sessionFor(key);
|
||||
if (!sessionId || !(await this.#harness.sessionExists(sessionId))) {
|
||||
sessionId = await this.#harness.createSession();
|
||||
await this.#state.setSession(key, sessionId);
|
||||
const workspaceCommand = await runWorkspaceCommand(text, this.#harness);
|
||||
if (workspaceCommand) {
|
||||
await this.#bot.sendText(target, workspaceCommand.message);
|
||||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
|
||||
let stream = null;
|
||||
|
|
@ -137,16 +141,22 @@ export class QqHarnessBridge {
|
|||
this.#logger.warn?.('[dsh-im:qq] unable to start a QQ stream; using a text reply:', error);
|
||||
}
|
||||
}
|
||||
const answer = await this.#harness.ask(sessionId, text, {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
onUpdate: stream ? async (update) => {
|
||||
const progress = update.type === 'text'
|
||||
? update.text
|
||||
: update.type === 'tool'
|
||||
? `正在使用${update.name}…`
|
||||
: update.text;
|
||||
if (progress) await stream.update(progress);
|
||||
} : undefined,
|
||||
const { answer } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
text,
|
||||
askOptions: {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
onUpdate: stream ? async (update) => {
|
||||
const progress = update.type === 'text'
|
||||
? update.text
|
||||
: update.type === 'tool'
|
||||
? `正在使用${update.name}…`
|
||||
: update.text;
|
||||
if (progress) await stream.update(progress);
|
||||
} : undefined,
|
||||
},
|
||||
});
|
||||
if (stream) {
|
||||
try {
|
||||
|
|
|
|||
|
|
@ -388,7 +388,14 @@ export class QqController {
|
|||
if (record.controller.signal.aborted) {
|
||||
await this.#stopRuntime(identity.botId);
|
||||
if (previousConfig) await this.#configStore.save(previousConfig).catch(() => undefined);
|
||||
else await this.#configStore.remove(identity.botId).catch(() => undefined);
|
||||
else {
|
||||
const removed = await this.#configStore.remove(identity.botId).catch(() => null);
|
||||
if (removed) {
|
||||
await this.#deleteState({ botId: identity.botId, config }).catch((cleanupError) => {
|
||||
this.#logger.warn?.('[dsh-im:qq] cancelled bot state cleanup failed:', cleanupError);
|
||||
});
|
||||
}
|
||||
}
|
||||
await this.#restoreCredential(identity.secretRef, previousSecret);
|
||||
throw error;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -56,6 +56,11 @@ export class QqStateStore {
|
|||
await this.#persist();
|
||||
}
|
||||
|
||||
async clearSessions() {
|
||||
this.#state.sessions = {};
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
hasSeen(messageId) {
|
||||
return this.#state.seenMessageIds.includes(messageId);
|
||||
}
|
||||
|
|
|
|||
584
src/channels/shared/bot-workspace-store.mjs
Normal file
584
src/channels/shared/bot-workspace-store.mjs
Normal file
|
|
@ -0,0 +1,584 @@
|
|||
import { mkdir, readFile, rename, stat, unlink, writeFile } from 'node:fs/promises';
|
||||
import { dirname, isAbsolute, resolve } from 'node:path';
|
||||
|
||||
import { WORKSPACE_SESSION_STALE } from './workspace-session.mjs';
|
||||
|
||||
const EMPTY_DOCUMENT = Object.freeze({ version: 1, workspaces: Object.freeze({}) });
|
||||
|
||||
function botIdOf(value) {
|
||||
if (typeof value !== 'string' || !/^[A-Za-z0-9_-]{1,128}$/.test(value)) {
|
||||
throw new TypeError('Invalid bot id');
|
||||
}
|
||||
return value;
|
||||
}
|
||||
|
||||
function normalizeDocument(value) {
|
||||
if (!value || value.version !== 1 || !value.workspaces
|
||||
|| typeof value.workspaces !== 'object' || Array.isArray(value.workspaces)) return null;
|
||||
const workspaces = {};
|
||||
for (const [botId, workspace] of Object.entries(value.workspaces)) {
|
||||
if (!/^[A-Za-z0-9_-]{1,128}$/.test(botId)
|
||||
|| typeof workspace !== 'string' || !isAbsolute(workspace)) return null;
|
||||
workspaces[botId] = resolve(workspace);
|
||||
}
|
||||
return { version: 1, workspaces };
|
||||
}
|
||||
|
||||
export async function validateWorkspacePath(value) {
|
||||
if (typeof value !== 'string' || !value.trim() || !isAbsolute(value.trim())) {
|
||||
const error = new Error('工作区必须是绝对路径。');
|
||||
error.code = 'workspace-not-absolute';
|
||||
throw error;
|
||||
}
|
||||
const workspace = resolve(value.trim());
|
||||
let info;
|
||||
try {
|
||||
info = await stat(workspace);
|
||||
} catch (cause) {
|
||||
const error = new Error('工作区路径不存在。', { cause });
|
||||
error.code = 'workspace-not-found';
|
||||
throw error;
|
||||
}
|
||||
if (!info.isDirectory()) {
|
||||
const error = new Error('工作区路径必须指向一个目录。');
|
||||
error.code = 'workspace-not-directory';
|
||||
throw error;
|
||||
}
|
||||
return workspace;
|
||||
}
|
||||
|
||||
export class BotWorkspaceStore {
|
||||
#path;
|
||||
#defaultWorkspace;
|
||||
#workspaces = {};
|
||||
#generations = new Map();
|
||||
#nextGeneration = 1;
|
||||
#incarnations = new Map();
|
||||
#nextIncarnation = 1;
|
||||
#removals = new Map();
|
||||
#removalDetails = new WeakMap();
|
||||
#dirtyRemovals = new Set();
|
||||
#writeQueue = Promise.resolve();
|
||||
#botQueues = new Map();
|
||||
|
||||
constructor(path, { defaultWorkspace = process.cwd() } = {}) {
|
||||
if (typeof path !== 'string' || !path) throw new TypeError('workspace store path is required');
|
||||
this.#path = path;
|
||||
this.#defaultWorkspace = resolve(defaultWorkspace);
|
||||
}
|
||||
|
||||
async load() {
|
||||
try {
|
||||
const normalized = normalizeDocument(JSON.parse(await readFile(this.#path, 'utf8')));
|
||||
if (!normalized) throw new Error('dsh-im workspace config is invalid');
|
||||
this.#workspaces = normalized.workspaces;
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
this.#workspaces = {};
|
||||
}
|
||||
this.#generations.clear();
|
||||
this.#nextGeneration = 1;
|
||||
this.#incarnations.clear();
|
||||
this.#nextIncarnation = 1;
|
||||
this.#removals.clear();
|
||||
this.#dirtyRemovals.clear();
|
||||
for (const botId of Object.keys(this.#workspaces)) {
|
||||
this.#generations.set(botId, this.#freshGeneration());
|
||||
this.#incarnations.set(botId, this.#freshIncarnation());
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
has(botId) {
|
||||
const id = botIdOf(botId);
|
||||
return Object.hasOwn(this.#workspaces, id) && !this.#removals.has(id);
|
||||
}
|
||||
|
||||
incarnationFor(botId) {
|
||||
return this.#incarnations.get(botIdOf(botId)) ?? null;
|
||||
}
|
||||
|
||||
workspaceFor(botId) {
|
||||
return this.#workspaces[botIdOf(botId)] ?? this.#defaultWorkspace;
|
||||
}
|
||||
|
||||
generationFor(botId) {
|
||||
return this.#generations.get(botIdOf(botId)) ?? null;
|
||||
}
|
||||
|
||||
async whenIdle() {
|
||||
await this.#writeQueue;
|
||||
}
|
||||
|
||||
async whenBotIdle(botId) {
|
||||
const id = botIdOf(botId);
|
||||
while (true) {
|
||||
const pending = this.#botQueues.get(id);
|
||||
if (!pending) return;
|
||||
await pending;
|
||||
if (this.#botQueues.get(id) === pending) return;
|
||||
}
|
||||
}
|
||||
|
||||
async ensure(botId, { workspace = this.#defaultWorkspace } = {}) {
|
||||
const id = botIdOf(botId);
|
||||
const initialWorkspace = resolve(workspace);
|
||||
return this.#enqueue(id, async () => {
|
||||
if (!this.#workspaces[id]) {
|
||||
this.#workspaces[id] = initialWorkspace;
|
||||
this.#generations.set(id, this.#freshGeneration());
|
||||
this.#incarnations.set(id, this.#freshIncarnation());
|
||||
try {
|
||||
await this.#persist();
|
||||
} catch (error) {
|
||||
delete this.#workspaces[id];
|
||||
this.#generations.delete(id);
|
||||
this.#incarnations.delete(id);
|
||||
throw error;
|
||||
}
|
||||
} else if (!this.#generations.has(id)) {
|
||||
this.#generations.set(id, this.#freshGeneration());
|
||||
}
|
||||
return this.#workspaces[id];
|
||||
});
|
||||
}
|
||||
|
||||
async setWorkspace(botId, value, { clearSessions, incarnation } = {}) {
|
||||
const id = botIdOf(botId);
|
||||
if (!this.has(id)
|
||||
|| (incarnation !== undefined && incarnation !== this.incarnationFor(id))) {
|
||||
const error = new Error('找不到要修改的机器人。');
|
||||
error.code = 'workspace-bot-not-found';
|
||||
throw error;
|
||||
}
|
||||
const workspace = await validateWorkspacePath(value);
|
||||
return this.#enqueue(id, async () => {
|
||||
if (!this.has(id)
|
||||
|| (incarnation !== undefined && incarnation !== this.incarnationFor(id))) {
|
||||
const error = new Error('找不到要修改的机器人。');
|
||||
error.code = 'workspace-bot-not-found';
|
||||
throw error;
|
||||
}
|
||||
if (workspace === this.workspaceFor(id)) return workspace;
|
||||
const previous = this.#workspaces[id];
|
||||
// Advance first so a session creation that started before this queued
|
||||
// transition can never be written back after the clear.
|
||||
this.#generations.set(id, this.#freshGeneration());
|
||||
// Clear the old session mapping before publishing the new workspace.
|
||||
// A crash can then lose conversation continuity, but can never pair the
|
||||
// new workspace with sessions created in the old one.
|
||||
await clearSessions?.();
|
||||
this.#workspaces[id] = workspace;
|
||||
try {
|
||||
await this.#persist();
|
||||
} catch (error) {
|
||||
this.#workspaces[id] = previous;
|
||||
throw error;
|
||||
}
|
||||
return workspace;
|
||||
});
|
||||
}
|
||||
|
||||
async invalidateSessions(botId, { clearSessions } = {}) {
|
||||
const id = botIdOf(botId);
|
||||
return this.#enqueue(id, async () => {
|
||||
this.#generations.set(id, this.#freshGeneration());
|
||||
await clearSessions?.();
|
||||
});
|
||||
}
|
||||
|
||||
/** Fence one lifecycle and return the opaque token required to abort/finish it. */
|
||||
async beginRemoval(botId, { clearSessions } = {}) {
|
||||
const id = botIdOf(botId);
|
||||
return this.#enqueue(id, async () => {
|
||||
const existing = this.#removals.get(id);
|
||||
if (existing) return existing;
|
||||
const transaction = Object.freeze({});
|
||||
this.#removals.set(id, transaction);
|
||||
this.#removalDetails.set(transaction, {
|
||||
botId: id,
|
||||
incarnation: this.incarnationFor(id),
|
||||
});
|
||||
this.#generations.set(id, this.#freshGeneration());
|
||||
try {
|
||||
await clearSessions?.();
|
||||
} catch (error) {
|
||||
if (this.#removals.get(id) === transaction) this.#removals.delete(id);
|
||||
throw error;
|
||||
}
|
||||
return transaction;
|
||||
});
|
||||
}
|
||||
|
||||
/** Re-open only the lifecycle represented by transaction; stale tokens are no-ops. */
|
||||
async abortRemoval(transaction) {
|
||||
const { botId: id } = this.#removalDetailsFor(transaction);
|
||||
return this.#enqueue(id, async () => {
|
||||
if (this.#removals.get(id) !== transaction) return false;
|
||||
this.#removals.delete(id);
|
||||
if (Object.hasOwn(this.#workspaces, id)) {
|
||||
this.#generations.set(id, this.#freshGeneration());
|
||||
if (!this.#incarnations.has(id)) {
|
||||
this.#incarnations.set(id, this.#freshIncarnation());
|
||||
}
|
||||
}
|
||||
return true;
|
||||
});
|
||||
}
|
||||
|
||||
/** Retire only the lifecycle represented by transaction; stale tokens are no-ops. */
|
||||
async finishRemoval(transaction) {
|
||||
const { botId: id, incarnation } = this.#removalDetailsFor(transaction);
|
||||
return this.#enqueue(id, async () => {
|
||||
if (this.#removals.get(id) !== transaction) {
|
||||
return { removed: false, persisted: true, error: null, stale: true };
|
||||
}
|
||||
if (this.incarnationFor(id) !== incarnation) {
|
||||
this.#removals.delete(id);
|
||||
return { removed: false, persisted: true, error: null, stale: true };
|
||||
}
|
||||
this.#removals.delete(id);
|
||||
return this.#retireCurrentIncarnation(id);
|
||||
});
|
||||
}
|
||||
|
||||
/** Commit the workspace lifecycle after the config store durably removed a bot. */
|
||||
async retireAfterConfigCommit(botId) {
|
||||
const id = botIdOf(botId);
|
||||
return this.#enqueue(id, async () => {
|
||||
this.#removals.delete(id);
|
||||
return this.#retireCurrentIncarnation(id);
|
||||
});
|
||||
}
|
||||
|
||||
async remove(botId) {
|
||||
const result = await this.retireAfterConfigCommit(botId);
|
||||
if (result.error) throw result.error;
|
||||
return result.removed;
|
||||
}
|
||||
|
||||
async reconcile(activeBotIds) {
|
||||
const active = new Set([...activeBotIds].map(botIdOf));
|
||||
const candidates = new Set([...Object.keys(this.#workspaces), ...this.#dirtyRemovals]);
|
||||
for (const botId of candidates) {
|
||||
if (!active.has(botId)) await this.remove(botId);
|
||||
}
|
||||
}
|
||||
|
||||
decorateStatus(status) {
|
||||
if (!status || typeof status !== 'object' || !Array.isArray(status.bots)) return status;
|
||||
return {
|
||||
...status,
|
||||
bots: status.bots.map((bot) => bot?.botId
|
||||
? { ...bot, workspace: this.workspaceFor(bot.botId) }
|
||||
: bot),
|
||||
};
|
||||
}
|
||||
|
||||
#freshGeneration() {
|
||||
const generation = this.#nextGeneration;
|
||||
this.#nextGeneration += 1;
|
||||
return generation;
|
||||
}
|
||||
|
||||
#freshIncarnation() {
|
||||
const incarnation = this.#nextIncarnation;
|
||||
this.#nextIncarnation += 1;
|
||||
return incarnation;
|
||||
}
|
||||
|
||||
#removalDetailsFor(transaction) {
|
||||
if (!transaction || typeof transaction !== 'object') {
|
||||
throw new TypeError('Invalid workspace removal transaction');
|
||||
}
|
||||
const details = this.#removalDetails.get(transaction);
|
||||
if (!details) throw new TypeError('Invalid workspace removal transaction');
|
||||
return details;
|
||||
}
|
||||
|
||||
async #retireCurrentIncarnation(id) {
|
||||
const hadWorkspace = Object.hasOwn(this.#workspaces, id);
|
||||
const needsCleanup = hadWorkspace || this.#dirtyRemovals.has(id);
|
||||
delete this.#workspaces[id];
|
||||
this.#generations.delete(id);
|
||||
this.#incarnations.delete(id);
|
||||
if (!needsCleanup) return {
|
||||
removed: false, persisted: true, error: null, stale: false,
|
||||
};
|
||||
try {
|
||||
await this.#persistCurrentDocument();
|
||||
return {
|
||||
removed: hadWorkspace, persisted: true, error: null, stale: false,
|
||||
};
|
||||
} catch (error) {
|
||||
this.#dirtyRemovals.add(id);
|
||||
return {
|
||||
removed: hadWorkspace, persisted: false, error, stale: false,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
async #enqueue(botId, operation) {
|
||||
const queued = this.#writeQueue.then(operation, operation);
|
||||
const settled = queued.then(() => undefined, () => undefined);
|
||||
this.#writeQueue = settled;
|
||||
this.#botQueues.set(botId, settled);
|
||||
void settled.finally(() => {
|
||||
if (this.#botQueues.get(botId) === settled) this.#botQueues.delete(botId);
|
||||
});
|
||||
return queued;
|
||||
}
|
||||
|
||||
async #persist() {
|
||||
const document = { version: 1, workspaces: this.#workspaces };
|
||||
await mkdir(dirname(this.#path), { recursive: true, mode: 0o700 });
|
||||
const temporary = `${this.#path}.tmp`;
|
||||
await writeFile(temporary, `${JSON.stringify(document, null, 2)}\n`, {
|
||||
encoding: 'utf8',
|
||||
mode: 0o600,
|
||||
});
|
||||
await rename(temporary, this.#path);
|
||||
this.#dirtyRemovals.clear();
|
||||
}
|
||||
|
||||
async #persistCurrentDocument() {
|
||||
if (Object.keys(this.#workspaces).length > 0) {
|
||||
await this.#persist();
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await unlink(this.#path);
|
||||
this.#dirtyRemovals.clear();
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
this.#dirtyRemovals.clear();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function decorateResult(workspaces, result) {
|
||||
return result && typeof result.then === 'function'
|
||||
? result.then((value) => workspaces.decorateStatus(value))
|
||||
: workspaces.decorateStatus(result);
|
||||
}
|
||||
|
||||
function targetStatus(controller) {
|
||||
return Promise.resolve(controller.status());
|
||||
}
|
||||
|
||||
/** Observe the config store's durable removal commit without changing its API. */
|
||||
export function observeBotWorkspaceRemovals(
|
||||
configStore,
|
||||
{ workspaces, method = 'remove', botIdFromRemoved = (removed) => removed?.botId },
|
||||
) {
|
||||
if (!configStore || !workspaces || typeof configStore[method] !== 'function') {
|
||||
throw new TypeError('configStore removal observer dependencies are required');
|
||||
}
|
||||
return new Proxy(configStore, {
|
||||
get(target, property) {
|
||||
const value = Reflect.get(target, property, target);
|
||||
if (property === method) {
|
||||
return async (...args) => {
|
||||
const removed = await value.apply(target, args);
|
||||
const botId = removed ? botIdFromRemoved(removed, args) : null;
|
||||
if (botId) await workspaces.retireAfterConfigCommit(botId);
|
||||
return removed;
|
||||
};
|
||||
}
|
||||
return typeof value === 'function' ? value.bind(target) : value;
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
export function createBotWorkspaceScope(harness, { botId, workspaces, state }) {
|
||||
if (!harness || !workspaces || !state) throw new TypeError('harness, workspaces, and state are required');
|
||||
const incarnation = workspaces.incarnationFor(botId);
|
||||
const isCurrentScope = () => workspaces.has(botId)
|
||||
&& workspaces.incarnationFor(botId) === incarnation;
|
||||
const sessionGenerations = new Map();
|
||||
const scopedHarness = new Proxy(harness, {
|
||||
get(target, property) {
|
||||
if (property === 'currentWorkspace') return () => workspaces.workspaceFor(botId);
|
||||
if (property === 'switchWorkspace') {
|
||||
return (workspace) => {
|
||||
if (!isCurrentScope()) {
|
||||
const error = new Error('找不到要修改的机器人。');
|
||||
error.code = 'workspace-bot-not-found';
|
||||
return Promise.reject(error);
|
||||
}
|
||||
return workspaces.setWorkspace(botId, workspace, {
|
||||
clearSessions: () => state.clearSessions(),
|
||||
incarnation,
|
||||
});
|
||||
};
|
||||
}
|
||||
if (property === 'createSession') {
|
||||
return async (options = {}) => {
|
||||
await workspaces.whenBotIdle(botId);
|
||||
if (!isCurrentScope()) {
|
||||
const error = new Error('找不到要修改的机器人。');
|
||||
error.code = 'workspace-bot-not-found';
|
||||
throw error;
|
||||
}
|
||||
const generation = workspaces.generationFor(botId);
|
||||
const sessionId = await target.createSession({
|
||||
...options,
|
||||
workspace: workspaces.workspaceFor(botId),
|
||||
});
|
||||
sessionGenerations.set(sessionId, generation);
|
||||
return sessionId;
|
||||
};
|
||||
}
|
||||
if (property === 'sessionExists') {
|
||||
return (sessionId, ...args) => {
|
||||
if (!isCurrentScope()) return false;
|
||||
const generation = sessionGenerations.get(sessionId);
|
||||
if (generation !== undefined && generation !== workspaces.generationFor(botId)) {
|
||||
sessionGenerations.delete(sessionId);
|
||||
return false;
|
||||
}
|
||||
return target.sessionExists(sessionId, ...args);
|
||||
};
|
||||
}
|
||||
if (property === 'ask') {
|
||||
return (sessionId, ...args) => {
|
||||
const generation = sessionGenerations.get(sessionId);
|
||||
sessionGenerations.delete(sessionId);
|
||||
if (!isCurrentScope()
|
||||
|| (generation !== undefined && generation !== workspaces.generationFor(botId))) {
|
||||
const error = new Error('The bot workspace changed before this prompt started.');
|
||||
error.code = WORKSPACE_SESSION_STALE;
|
||||
throw error;
|
||||
}
|
||||
return target.ask(sessionId, ...args);
|
||||
};
|
||||
}
|
||||
const value = Reflect.get(target, property, target);
|
||||
return typeof value === 'function' ? value.bind(target) : value;
|
||||
},
|
||||
});
|
||||
const scopedState = new Proxy(state, {
|
||||
get(target, property) {
|
||||
if (property === 'sessionFor') {
|
||||
return (key, ...args) => {
|
||||
if (!isCurrentScope()) return null;
|
||||
const sessionId = target.sessionFor(key, ...args);
|
||||
if (sessionId && !sessionGenerations.has(sessionId)) {
|
||||
sessionGenerations.set(sessionId, workspaces.generationFor(botId));
|
||||
}
|
||||
return sessionId;
|
||||
};
|
||||
}
|
||||
if (property === 'setSession') {
|
||||
return (key, sessionId, ...args) => {
|
||||
const generation = sessionGenerations.get(sessionId);
|
||||
if (!isCurrentScope()
|
||||
|| (generation !== undefined && generation !== workspaces.generationFor(botId))) {
|
||||
sessionGenerations.delete(sessionId);
|
||||
return false;
|
||||
}
|
||||
return target.setSession(key, sessionId, ...args);
|
||||
};
|
||||
}
|
||||
const value = Reflect.get(target, property, target);
|
||||
return typeof value === 'function' ? value.bind(target) : value;
|
||||
},
|
||||
});
|
||||
return Object.freeze({ harness: scopedHarness, state: scopedState });
|
||||
}
|
||||
|
||||
export function createBotScopedHarness(harness, options) {
|
||||
return createBotWorkspaceScope(harness, options).harness;
|
||||
}
|
||||
|
||||
export function createWorkspaceAwareController(controller, { workspaces, stateFor }) {
|
||||
if (!controller || !workspaces || typeof stateFor !== 'function') {
|
||||
throw new TypeError('controller, workspaces, and stateFor are required');
|
||||
}
|
||||
const transitions = new Map();
|
||||
const withBotTransition = (botId, operation) => {
|
||||
const previous = transitions.get(botId) ?? Promise.resolve();
|
||||
const current = previous.catch(() => undefined).then(operation);
|
||||
transitions.set(botId, current);
|
||||
return current.finally(() => {
|
||||
if (transitions.get(botId) === current) transitions.delete(botId);
|
||||
});
|
||||
};
|
||||
const updateWorkspace = (botId, workspace) => {
|
||||
// Capture at API invocation, before even waiting for an older outer
|
||||
// transition. A queued request still belongs to the incarnation that the
|
||||
// caller observed, not a deterministic same-id rebind that appears later.
|
||||
const incarnation = workspaces.incarnationFor(botId);
|
||||
return withBotTransition(botId, async () => {
|
||||
const snapshot = await controller.status();
|
||||
if (!snapshot?.bots?.some((bot) => bot?.botId === botId)) {
|
||||
const error = new Error('找不到要修改的机器人。');
|
||||
error.code = 'workspace-bot-not-found';
|
||||
throw error;
|
||||
}
|
||||
const state = await stateFor(botId);
|
||||
await workspaces.setWorkspace(botId, workspace, {
|
||||
clearSessions: () => state.clearSessions(),
|
||||
incarnation,
|
||||
});
|
||||
return workspaces.decorateStatus(await controller.status());
|
||||
});
|
||||
};
|
||||
const deleteWithWorkspace = (botId, invokeDelete) => withBotTransition(botId, async () => {
|
||||
// Fence the old runtime without changing the durable mapping. A crash
|
||||
// before the controller removes its config therefore keeps the bot's
|
||||
// workspace, while a crash after that commit is healed by startup
|
||||
// reconciliation.
|
||||
const removal = await workspaces.beginRemoval(botId, {
|
||||
clearSessions: async () => {
|
||||
try {
|
||||
const state = await stateFor(botId);
|
||||
if (!state || typeof state.clearSessions !== 'function') {
|
||||
throw new TypeError('bot state does not support session cleanup');
|
||||
}
|
||||
await state.clearSessions();
|
||||
} catch (error) {
|
||||
console.warn(
|
||||
`[dsh-im] ignored session cleanup failure while deleting bot ${botId}:`,
|
||||
error?.message ?? error,
|
||||
);
|
||||
}
|
||||
},
|
||||
});
|
||||
try {
|
||||
const result = await invokeDelete();
|
||||
await workspaces.finishRemoval(removal);
|
||||
return workspaces.decorateStatus(result);
|
||||
} catch (error) {
|
||||
const after = await targetStatus(controller).catch(() => null);
|
||||
const knownAbsent = Array.isArray(after?.bots)
|
||||
&& !after.bots.some((bot) => bot?.botId === botId);
|
||||
if (knownAbsent) await workspaces.finishRemoval(removal);
|
||||
else await workspaces.abortRemoval(removal);
|
||||
throw error;
|
||||
}
|
||||
});
|
||||
|
||||
return new Proxy(controller, {
|
||||
get(target, property) {
|
||||
if (property === 'updateWorkspace') return updateWorkspace;
|
||||
const value = Reflect.get(target, property, target);
|
||||
if (typeof value !== 'function') return value;
|
||||
if (property === 'deleteBot') {
|
||||
return (botId, ...args) => deleteWithWorkspace(
|
||||
botId,
|
||||
() => value.call(target, botId, ...args),
|
||||
);
|
||||
}
|
||||
if (property === 'disconnect') {
|
||||
return async (...args) => {
|
||||
const before = await target.status();
|
||||
const botId = before?.bots?.[0]?.botId;
|
||||
if (!botId) return decorateResult(workspaces, value.apply(target, args));
|
||||
return deleteWithWorkspace(botId, () => value.apply(target, args));
|
||||
};
|
||||
}
|
||||
return (...args) => decorateResult(workspaces, value.apply(target, args));
|
||||
},
|
||||
});
|
||||
}
|
||||
|
|
@ -57,6 +57,11 @@ export class ConversationStateStore {
|
|||
await this.#persist();
|
||||
}
|
||||
|
||||
async clearSessions() {
|
||||
this.#state.sessions = {};
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
hasSeen(messageId) {
|
||||
return this.#state.seenMessageIds.includes(messageId);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,3 +1,6 @@
|
|||
import { runWorkspaceCommand } from './workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from './workspace-session.mjs';
|
||||
|
||||
function cleanText(value) {
|
||||
return typeof value === 'string' ? value.trim() : '';
|
||||
}
|
||||
|
|
@ -97,6 +100,7 @@ export class TextHarnessBridge {
|
|||
'',
|
||||
'直接发送文字即可继续当前会话。',
|
||||
'/new 开启一个全新会话',
|
||||
'/workspace 工作区绝对路径 切换工作区',
|
||||
'/status 检查连接状态',
|
||||
'/help 显示本帮助',
|
||||
].join('\n'));
|
||||
|
|
@ -109,6 +113,12 @@ export class TextHarnessBridge {
|
|||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
const workspaceCommand = await runWorkspaceCommand(text, this.#harness);
|
||||
if (workspaceCommand) {
|
||||
await this.#bot.sendText(target, workspaceCommand.message);
|
||||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
const conversationKey = `${message.kind}:${message.conversationId}`;
|
||||
if (command === '/new') {
|
||||
await this.#state.clearSession(conversationKey);
|
||||
|
|
@ -117,12 +127,6 @@ export class TextHarnessBridge {
|
|||
return;
|
||||
}
|
||||
|
||||
let sessionId = this.#state.sessionFor(conversationKey);
|
||||
if (!sessionId || !(await this.#harness.sessionExists(sessionId))) {
|
||||
sessionId = await this.#harness.createSession();
|
||||
await this.#state.setSession(conversationKey, sessionId);
|
||||
}
|
||||
|
||||
await this.#bot.sendTyping?.(target).catch((error) => {
|
||||
this.#logger.warn?.(`[dsh-im:${this.#descriptor.key}] typing indicator failed:`, error);
|
||||
});
|
||||
|
|
@ -138,13 +142,19 @@ export class TextHarnessBridge {
|
|||
);
|
||||
}
|
||||
}
|
||||
const answer = await this.#harness.ask(sessionId, text, {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
onUpdate: stream ? async (update) => {
|
||||
const progress = update.type === 'text' ? update.text
|
||||
: update.type === 'tool' ? `正在使用${update.name}…` : update.text;
|
||||
if (progress) await stream.update(progress);
|
||||
} : undefined,
|
||||
const { answer } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key: conversationKey,
|
||||
text,
|
||||
askOptions: {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
onUpdate: stream ? async (update) => {
|
||||
const progress = update.type === 'text' ? update.text
|
||||
: update.type === 'tool' ? `正在使用${update.name}…` : update.text;
|
||||
if (progress) await stream.update(progress);
|
||||
} : undefined,
|
||||
},
|
||||
});
|
||||
if (stream) {
|
||||
try {
|
||||
|
|
|
|||
26
src/channels/shared/workspace-command.mjs
Normal file
26
src/channels/shared/workspace-command.mjs
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
const WORKSPACE_COMMAND = /^\/workspace(?:\s+([\s\S]+))?$/i;
|
||||
|
||||
export async function runWorkspaceCommand(text, harness) {
|
||||
if (typeof text !== 'string') return null;
|
||||
const match = WORKSPACE_COMMAND.exec(text.trim());
|
||||
if (!match) return null;
|
||||
const workspace = match[1]?.trim();
|
||||
if (!workspace) {
|
||||
return { handled: true, message: '用法:/workspace 工作区绝对路径' };
|
||||
}
|
||||
if (typeof harness?.switchWorkspace !== 'function') {
|
||||
return { handled: true, message: '当前机器人暂不支持切换工作区。' };
|
||||
}
|
||||
try {
|
||||
const current = await harness.switchWorkspace(workspace);
|
||||
return { handled: true, message: `工作区已切换为:${current}` };
|
||||
} catch (error) {
|
||||
if (['workspace-not-absolute', 'workspace-not-found', 'workspace-not-directory'].includes(error?.code)) {
|
||||
return { handled: true, message: `${error.message}\n用法:/workspace 工作区绝对路径` };
|
||||
}
|
||||
if (error?.code === 'workspace-bot-not-found') {
|
||||
return { handled: true, message: '机器人正在移除或已重新接入,无法切换原会话的工作区。' };
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
44
src/channels/shared/workspace-session.mjs
Normal file
44
src/channels/shared/workspace-session.mjs
Normal file
|
|
@ -0,0 +1,44 @@
|
|||
export const WORKSPACE_SESSION_STALE = 'workspace-session-stale';
|
||||
|
||||
async function sessionExists(harness, sessionId, options) {
|
||||
return options === undefined
|
||||
? harness.sessionExists(sessionId)
|
||||
: harness.sessionExists(sessionId, options);
|
||||
}
|
||||
|
||||
async function createSession(harness, options) {
|
||||
return options === undefined
|
||||
? harness.createSession()
|
||||
: harness.createSession(options);
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve, persist, and ask through a session that belongs to the bot's
|
||||
* current workspace. A concurrent workspace switch invalidates the scoped
|
||||
* session and retries before any prompt is sent to the stale session.
|
||||
*/
|
||||
export async function askInWorkspaceSession({
|
||||
harness,
|
||||
state,
|
||||
key,
|
||||
text,
|
||||
createOptions,
|
||||
existsOptions,
|
||||
askOptions,
|
||||
}) {
|
||||
while (true) {
|
||||
let sessionId = state.sessionFor(key);
|
||||
if (!sessionId || !(await sessionExists(harness, sessionId, existsOptions))) {
|
||||
sessionId = await createSession(harness, createOptions);
|
||||
if (await state.setSession(key, sessionId) === false) continue;
|
||||
}
|
||||
try {
|
||||
return {
|
||||
sessionId,
|
||||
answer: await harness.ask(sessionId, text, askOptions),
|
||||
};
|
||||
} catch (error) {
|
||||
if (error?.code !== WORKSPACE_SESSION_STALE) throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -54,6 +54,11 @@ export class WecomStateStore {
|
|||
await this.#persist();
|
||||
}
|
||||
|
||||
async clearSessions() {
|
||||
this.#state.sessions = {};
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
hasSeen(messageId) {
|
||||
return this.#state.seenMessageIds.includes(messageId);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,10 +1,13 @@
|
|||
import { generateReqId } from '@wecom/aibot-node-sdk';
|
||||
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
|
||||
const HELP_TEXT = [
|
||||
'企业微信机器人已连接 DeepSeek Harness。',
|
||||
'',
|
||||
'直接发送文字即可继续当前会话。',
|
||||
'/new 开启一个全新会话',
|
||||
'/workspace 工作区绝对路径 切换工作区',
|
||||
'/status 检查连接状态',
|
||||
'/help 显示本帮助',
|
||||
].join('\n');
|
||||
|
|
@ -182,11 +185,11 @@ export class WecomHarnessBridge {
|
|||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
|
||||
let sessionId = this.#state.sessionFor(key);
|
||||
if (!sessionId || !(await this.#harness.sessionExists(sessionId))) {
|
||||
sessionId = await this.#harness.createSession();
|
||||
await this.#state.setSession(key, sessionId);
|
||||
const workspaceCommand = await runWorkspaceCommand(text, this.#harness);
|
||||
if (workspaceCommand) {
|
||||
await this.#sendImmediate(frame, chatId, workspaceCommand.message);
|
||||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
|
||||
streamId = this.#generateReqId('stream');
|
||||
|
|
@ -197,14 +200,20 @@ export class WecomHarnessBridge {
|
|||
this.#logger.warn?.('[dsh-im:wecom] unable to start a stream; using an active reply:', error);
|
||||
}
|
||||
|
||||
const answer = await this.#harness.ask(sessionId, text, {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
onUpdate: streamStarted && typeof this.#client.replyStreamNonBlocking === 'function'
|
||||
? async (update) => {
|
||||
const progress = splitUtf8(progressText(update))[0];
|
||||
if (progress) await this.#client.replyStreamNonBlocking(frame, streamId, progress, false);
|
||||
}
|
||||
: undefined,
|
||||
const { answer } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
text,
|
||||
askOptions: {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
onUpdate: streamStarted && typeof this.#client.replyStreamNonBlocking === 'function'
|
||||
? async (update) => {
|
||||
const progress = splitUtf8(progressText(update))[0];
|
||||
if (progress) await this.#client.replyStreamNonBlocking(frame, streamId, progress, false);
|
||||
}
|
||||
: undefined,
|
||||
},
|
||||
});
|
||||
|
||||
const chunks = splitUtf8(answer || '任务已完成,但没有生成可显示的文本。');
|
||||
|
|
|
|||
|
|
@ -400,7 +400,14 @@ export class WecomController {
|
|||
if (record.controller.signal.aborted) {
|
||||
await this.#stopRuntime(identity.botId);
|
||||
if (previousConfig) await this.#configStore.save(previousConfig).catch(() => undefined);
|
||||
else await this.#configStore.remove(identity.botId).catch(() => undefined);
|
||||
else {
|
||||
const removed = await this.#configStore.remove(identity.botId).catch(() => null);
|
||||
if (removed) {
|
||||
await this.#deleteState({ botId: identity.botId, config }).catch((cleanupError) => {
|
||||
this.#logger.warn?.('[dsh-im:wecom] cancelled bot state cleanup failed:', cleanupError);
|
||||
});
|
||||
}
|
||||
}
|
||||
await this.#restoreCredential(identity.secretRef, previousSecret);
|
||||
throw error;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -187,17 +187,18 @@ export class HarnessClient {
|
|||
throw new Error(`Harness did not become ready: ${lastError?.message ?? 'timeout'}`);
|
||||
}
|
||||
|
||||
async workspaceId() {
|
||||
async workspaceId(options = {}) {
|
||||
const workspace = options.workspace ?? this.#workspace;
|
||||
const { items } = await this.rpc('workspace.list', {});
|
||||
const existing = items.find((item) => item.path === this.#workspace);
|
||||
const existing = items.find((item) => item.path === workspace);
|
||||
if (existing) return existing.workspaceId;
|
||||
const created = await this.rpc('workspace.create', { path: this.#workspace });
|
||||
const created = await this.rpc('workspace.create', { path: workspace });
|
||||
return created.workspace.workspaceId;
|
||||
}
|
||||
|
||||
async createSession() {
|
||||
async createSession(options = {}) {
|
||||
await this.ensureRunning();
|
||||
const workspaceId = await this.workspaceId();
|
||||
const workspaceId = await this.workspaceId(options);
|
||||
const created = await this.rpc('session.create', {
|
||||
workspaceId,
|
||||
agentPreset: this.#agentPreset,
|
||||
|
|
|
|||
|
|
@ -62,6 +62,11 @@ export class WeixinStateStore {
|
|||
await this.#persist();
|
||||
}
|
||||
|
||||
async clearSessions() {
|
||||
this.#state.sessions = {};
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
hasSeen(messageId) {
|
||||
return this.#state.seenMessageIds.includes(messageId);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -3,12 +3,15 @@ import {
|
|||
splitWeixinText,
|
||||
weixinMessageId,
|
||||
} from './weixin-api.mjs';
|
||||
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
|
||||
const HELP_TEXT = [
|
||||
'微信已连接 DeepSeek Harness。',
|
||||
'',
|
||||
'直接发送文字或带文字识别结果的语音即可继续当前会话。',
|
||||
'/new 开启一个全新会话',
|
||||
'/workspace 工作区绝对路径 切换工作区',
|
||||
'/status 检查连接状态',
|
||||
'/help 显示本帮助',
|
||||
].join('\n');
|
||||
|
|
@ -133,14 +136,21 @@ export class WeixinHarnessBridge {
|
|||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
const workspaceCommand = await runWorkspaceCommand(text, this.#harness);
|
||||
if (workspaceCommand) {
|
||||
await this.#send(sender, workspaceCommand.message, contextToken, runId);
|
||||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
|
||||
const key = conversationKey(sender);
|
||||
let sessionId = this.#state.sessionFor(key);
|
||||
if (!sessionId || !(await this.#harness.sessionExists(sessionId))) {
|
||||
sessionId = await this.#harness.createSession();
|
||||
await this.#state.setSession(key, sessionId);
|
||||
}
|
||||
const answer = await this.#harness.ask(sessionId, text, { timeoutMs: this.#replyTimeoutMs });
|
||||
const { answer } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
text,
|
||||
askOptions: { timeoutMs: this.#replyTimeoutMs },
|
||||
});
|
||||
await this.#send(sender, answer, contextToken, runId);
|
||||
await this.#state.markSeen(messageId);
|
||||
this.#status.messagesReplied += 1;
|
||||
|
|
|
|||
|
|
@ -468,7 +468,12 @@ export class WeixinController {
|
|||
await this.#stopRuntime(identity.botId);
|
||||
if (previousConfig) await this.#configStore.save(previousConfig).catch(() => undefined);
|
||||
else if (this.#configStore.get(identity.botId)) {
|
||||
await this.#configStore.remove(identity.botId).catch(() => undefined);
|
||||
const removed = await this.#configStore.remove(identity.botId).catch(() => null);
|
||||
if (removed) {
|
||||
await this.#deleteState({ botId: identity.botId, config }).catch((cleanupError) => {
|
||||
this.#logger.warn?.('[dsh-weixin] failed to clean up cancelled bot state:', cleanupError);
|
||||
});
|
||||
}
|
||||
}
|
||||
await this.#restoreCredential(identity.tokenRef, previousToken);
|
||||
if (previousConfig && previousToken?.value) {
|
||||
|
|
|
|||
|
|
@ -311,7 +311,14 @@ export class WhatsappController {
|
|||
if (record.controller.signal.aborted || this.#closed || error?.name === 'AbortError') {
|
||||
await this.#stopRuntime(botId);
|
||||
if (previous) await this.#configStore.save(previous).catch(() => undefined);
|
||||
else await this.#configStore.remove(botId).catch(() => undefined);
|
||||
else {
|
||||
const removed = await this.#configStore.remove(botId).catch(() => null);
|
||||
if (removed) {
|
||||
await this.#deleteState({ botId, config }).catch((cleanupError) => {
|
||||
this.#logger.warn?.('[dsh-im:whatsapp] cancelled bot state cleanup failed:', cleanupError);
|
||||
});
|
||||
}
|
||||
}
|
||||
await this.#deleteAuth(record.authDirectory).catch(() => undefined);
|
||||
if (previous) await this.#startRuntime(previous).catch(() => undefined);
|
||||
record.state = 'cancelled';
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue