fix: support Harness interactions in DingTalk and Feishu

This commit is contained in:
xmanrui 2026-08-18 13:03:15 +08:00
parent c9ec84cc67
commit 3e7e1dcdd7
13 changed files with 5061 additions and 881 deletions

View file

@ -3,11 +3,17 @@ import {
splitDingtalkText,
} from './dingtalk-api.mjs';
import { createDingTalkCardStream } from './dingtalk-card-stream.mjs';
import {
harnessAnswerForQuestion,
harnessQuestionText,
validHarnessQuestion,
} from '../shared/harness-question.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 = '消息处理失败,请稍后重试。';
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
const HELP_TEXT = [
'钉钉机器人已连接 DeepSeek Harness。',
@ -55,6 +61,19 @@ function progressText(update) {
return `_${nonEmptyString(update?.text) ?? '正在处理…'}_`;
}
function canClaimInteractionReply(message, pending, sender) {
if (pending.needsPresentation || !pending.questions[pending.index]) return false;
if (pending.actor !== sender) return false;
if (String(message?.conversationType) === '2' && message?.isInAtList !== true) return false;
if (message?.msgtype !== 'text' || !nonEmptyString(message?.text?.content)) return false;
try {
normalizeDingtalkSessionWebhook(message.sessionWebhook);
return true;
} catch {
return false;
}
}
function ensureStats(status) {
status.stats ??= {};
for (const key of ['messagesReceived', 'messagesReplied', 'messagesRejected', 'messagesIgnored']) {
@ -102,6 +121,8 @@ export class DingtalkHarnessBridge {
#maxMessageChars;
#signal;
#queues = new Map();
#pendingInteractions = new Map();
#interactionKeys = new Map();
#acceptedMessageIds = new Set();
constructor({
@ -157,12 +178,49 @@ export class DingtalkHarnessBridge {
this.#status.lastRejectedAt = new Date().toISOString();
return Promise.resolve();
}
const pending = this.#pendingInteractions.get(key);
if (pending && pending.actor !== sender) {
return this.#enqueueMessage(message, messageId, sender, key);
}
// Once one valid answer has been claimed, later messages are subsequent
// prompts even if the network submission eventually needs a retry. Invalid
// replies do not claim the question, so the next valid answer can still
// pass through this interaction queue.
if (pending?.submitting || pending?.claimedReplyMessageId) {
return this.#enqueueMessage(message, messageId, sender, key);
}
if (pending) {
if (canClaimInteractionReply(message, pending, sender)) {
pending.claimedReplyMessageId = messageId;
}
const previous = pending.queue ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#processInteractionReply(message, messageId, sender, key, pending))
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
if (pending.queue === current) pending.queue = null;
});
pending.queue = current;
return current;
}
return this.#enqueueMessage(message, messageId, sender, key);
}
#enqueueMessage(message, messageId, sender, key, {
releaseMessageId = true,
alreadyRecorded = false,
} = {}) {
const previous = this.#queues.get(key) ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#process(message, messageId, sender, key))
.then(() => this.#process(message, messageId, sender, key, { alreadyRecorded }))
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
if (this.#queues.get(key) === current) this.#queues.delete(key);
});
this.#queues.set(key, current);
@ -170,15 +228,22 @@ export class DingtalkHarnessBridge {
}
async waitForIdle() {
await Promise.allSettled([...this.#queues.values()]);
await Promise.allSettled([
...this.#queues.values(),
...[...this.#pendingInteractions.values()].flatMap((pending) => (
pending.queue ? [pending.queue] : []
)),
]);
}
async #process(message, messageId, sender, key) {
async #process(message, messageId, sender, key, { alreadyRecorded = false } = {}) {
this.#signal?.throwIfAborted();
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
increment(this.#status, 'messagesReceived');
this.#status.lastMessageAt = new Date().toISOString();
if (!alreadyRecorded) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
increment(this.#status, 'messagesReceived');
this.#status.lastMessageAt = new Date().toISOString();
}
if (String(message.conversationType) === '2' && message.isInAtList !== true) {
increment(this.#status, 'messagesIgnored');
@ -253,6 +318,13 @@ export class DingtalkHarnessBridge {
onUpdate: cardStarted
? (update) => cardStream.push(progressText(update))
: undefined,
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key,
actor: sender,
sessionWebhook,
requiresMention: String(message.conversationType) === '2',
}),
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
},
});
const streamed = cardStarted && await cardStream.finish(answer);
@ -270,6 +342,286 @@ export class DingtalkHarnessBridge {
} catch {
this.#logger.error?.('[dsh-dingtalk] failed to send the safe error reply');
}
} finally {
await this.#cancelPendingInteraction(key);
}
}
async #processInteractionReply(message, messageId, sender, key, expected) {
this.#signal?.throwIfAborted();
const current = this.#pendingInteractions.get(key);
const claimed = expected.claimedReplyMessageId === messageId;
if (!current || current !== expected || current.submitting) {
if (claimed && (!current || current !== expected)) {
return this.#discardResolvedInteractionReply(message, messageId);
}
return this.#enqueueMessage(message, messageId, sender, key, { releaseMessageId: false });
}
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
increment(this.#status, 'messagesReceived');
this.#status.lastMessageAt = new Date().toISOString();
if (String(message.conversationType) === '2' && message.isInAtList !== true) {
increment(this.#status, 'messagesIgnored');
return;
}
let sessionWebhook;
try {
sessionWebhook = normalizeDingtalkSessionWebhook(message.sessionWebhook);
} catch {
increment(this.#status, 'messagesRejected');
this.#status.lastRejectedAt = new Date().toISOString();
this.#status.lastError = '钉钉消息没有安全的回复地址。';
return;
}
const text = message?.msgtype === 'text' ? nonEmptyString(message?.text?.content) : null;
if (!text) {
try {
await this.#send(sessionWebhook, '请用文字回答当前问题。');
} catch {
this.#logger.error?.('[dsh-dingtalk] failed to reject a non-text interaction reply');
}
return;
}
const pending = this.#pendingInteractions.get(key);
if (!pending || pending !== expected || pending.submitting) {
if (claimed && (!pending || pending !== expected)) {
try {
await this.#send(sessionWebhook, INTERACTION_RESOLVED_TEXT);
} catch {
this.#logger.error?.('[dsh-dingtalk] failed to send an expired interaction notice');
}
return;
}
return this.#enqueueMessage(message, messageId, sender, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
pending.sessionWebhook = sessionWebhook;
if (pending.needsPresentation) {
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = '钉钉交互问题发送失败。';
this.#logger.error?.('[dsh-dingtalk] failed to retry an interaction question');
pending.interaction.reconnect?.();
}
return;
}
const question = pending.questions[pending.index];
if (!question) return;
pending.answers.push(harnessAnswerForQuestion(question, text));
pending.index += 1;
if (pending.index < pending.questions.length) {
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
pending.needsPresentation = true;
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = '钉钉交互问题发送失败。';
this.#logger.error?.('[dsh-dingtalk] failed to send the next interaction question');
pending.interaction.reconnect?.();
}
return;
}
pending.submitting = true;
try {
await pending.interaction.respond({
ok: true,
value: {
sessionId: pending.sessionId,
answer: { answers: pending.answers },
},
});
this.#clearPendingInteraction(key, pending.interactionId);
this.#status.lastError = null;
} catch (error) {
if (this.#signal?.aborted) return;
if (this.#pendingInteractions.get(key) !== pending) return;
if (error?.code === 'interaction-not-pending') {
this.#clearPendingInteraction(key, pending.interactionId);
try {
await this.#send(sessionWebhook, INTERACTION_RESOLVED_TEXT);
} catch {
this.#logger.error?.('[dsh-dingtalk] failed to send an expired interaction notice');
}
return;
}
pending.submitting = false;
pending.answers.pop();
pending.index -= 1;
this.#status.lastError = '回答提交失败。';
this.#logger.error?.('[dsh-dingtalk] failed to answer a Harness interaction');
try {
await this.#send(sessionWebhook, '回答提交失败,请重新发送当前问题的答案。');
} catch {
this.#logger.error?.('[dsh-dingtalk] failed to send an interaction retry notice');
}
}
}
async #handleInteraction(interaction, {
key,
actor,
sessionWebhook,
requiresMention,
}) {
// The transport deliberately exposes every interaction kind. Approval is
// left unanswered until #5 adds its own policy and renderer.
if (interaction?.kind !== 'question') return;
const questions = interaction?.payload?.questions;
const interactionId = typeof interaction?.interactionId === 'string'
? interaction.interactionId
: interaction?.rpcId;
if (typeof interaction.rpcId !== 'string'
|| typeof interactionId !== 'string'
|| typeof interaction.sessionId !== 'string'
|| !Array.isArray(questions)
|| questions.length === 0
|| questions.some((question) => !validHarnessQuestion(question))) {
this.#logger.warn?.('[dsh-dingtalk] ignored an invalid Harness question interaction');
return;
}
if (interaction.recovered === true) {
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'DingTalk safely cancelled an interaction left by an earlier client.',
details: {},
},
});
try {
await this.#send(sessionWebhook, '检测到这个 Session 中遗留的待回答问题,已安全取消并继续处理你刚才的消息。');
} catch {
this.#logger.error?.('[dsh-dingtalk] failed to send an interaction recovery notice');
}
return;
}
const existing = this.#pendingInteractions.get(key);
if (existing?.interactionId === interactionId) {
existing.interaction = interaction;
if (existing.needsPresentation) await this.#presentInteraction(existing);
return;
}
if (this.#interactionKeys.has(interactionId)) return;
if (existing) {
this.#logger.warn?.('[dsh-dingtalk] cancelled a second pending Harness question');
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'DingTalk is already handling another user interaction.',
details: {},
},
});
return;
}
const pending = {
kind: 'question',
interactionId,
sessionId: interaction.sessionId,
interaction,
actor,
requiresMention,
questions,
answers: [],
index: 0,
sessionWebhook,
queue: null,
claimedReplyMessageId: null,
submitting: false,
needsPresentation: true,
};
this.#pendingInteractions.set(key, pending);
this.#interactionKeys.set(pending.interactionId, key);
await this.#presentInteraction(pending);
}
#handleInteractionResolved(resolution) {
const interactionId = resolution?.interactionId;
if (resolution?.kind !== 'question' || typeof interactionId !== 'string') return;
const key = this.#interactionKeys.get(interactionId);
if (!key) return;
this.#clearPendingInteraction(key, interactionId);
}
async #presentInteraction(pending) {
const question = pending.questions[pending.index];
if (!question) return;
await this.#send(
pending.sessionWebhook,
harnessQuestionText(
question,
pending.index,
pending.questions.length,
{ requiresMention: pending.requiresMention },
),
);
pending.needsPresentation = false;
}
async #discardResolvedInteractionReply(message, messageId) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
increment(this.#status, 'messagesReceived');
this.#status.lastMessageAt = new Date().toISOString();
let sessionWebhook;
try {
sessionWebhook = normalizeDingtalkSessionWebhook(message.sessionWebhook);
} catch {
increment(this.#status, 'messagesRejected');
this.#status.lastRejectedAt = new Date().toISOString();
return;
}
try {
await this.#send(sessionWebhook, INTERACTION_RESOLVED_TEXT);
} catch {
this.#logger.error?.('[dsh-dingtalk] failed to send an expired interaction notice');
}
}
#takePendingInteraction(key, interactionId) {
const pending = this.#pendingInteractions.get(key);
if (!pending
|| (interactionId !== undefined && pending.interactionId !== interactionId)) return null;
this.#pendingInteractions.delete(key);
this.#interactionKeys.delete(pending.interactionId);
return pending;
}
#clearPendingInteraction(key, interactionId) {
return this.#takePendingInteraction(key, interactionId) !== null;
}
async #cancelPendingInteraction(key) {
const pending = this.#takePendingInteraction(key);
if (!pending || pending.kind !== 'question') return;
try {
await pending.interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'The DingTalk interaction ended before the user answered.',
details: {},
},
}, { signal: AbortSignal.timeout(5_000) });
} catch (error) {
if (error?.code !== 'interaction-not-pending') {
this.#logger.warn?.('[dsh-dingtalk] failed to cancel a pending Harness interaction');
}
}
}

View file

@ -1,369 +1,19 @@
import { spawn } from 'node:child_process';
import { randomUUID } from 'node:crypto';
import { isAbsolute } from 'node:path';
import {
HarnessClient as SharedHarnessClient,
} from '../shared/harness-client.mjs';
import { adoptRegisteredWorkspaceSession } from '../shared/harness-session-binding.mjs';
export {
HarnessInteractionError,
HarnessReplyTracker,
HarnessRpcError,
} from '../shared/harness-client.mjs';
function workspacePaths(value) {
if (!Array.isArray(value?.items)) return [];
return value.items.flatMap((item) => (
typeof item?.path === 'string' && isAbsolute(item.path) ? [item.path] : []
));
}
function workspaceFromList(workspacePath, workspaceList) {
if (!Array.isArray(workspaceList?.items)
|| !Array.isArray(workspaceList?.archivedSessionIds)) {
throw new Error('Harness returned an invalid response for workspace.list');
}
const workspace = workspaceList.items.find((item) => item?.path === workspacePath);
if (!workspace) return null;
if (!Array.isArray(workspace.sessionIds)
|| workspace.sessionIds.some((sessionId) => typeof sessionId !== 'string')) {
throw new Error('Harness returned invalid session IDs for workspace.list');
}
return workspace;
}
function workspaceSessions(workspace, archivedSessionIds, sessionList) {
if (!Array.isArray(sessionList?.items)) {
throw new Error('Harness returned an invalid response for session.list');
}
const archived = new Set(archivedSessionIds);
const summaries = new Map(sessionList.items.flatMap((item) => (
typeof item?.sessionId === 'string' ? [[item.sessionId, item]] : []
)));
return {
workspace: workspace.path,
sessions: workspace.sessionIds.map((sessionId) => {
const summary = summaries.get(sessionId);
const title = summary?.projections?.values?.title;
return {
sessionId,
title: typeof title === 'string' ? title : null,
archived: archived.has(sessionId),
blank: summary?.blank === true,
origin: summary?.origin === 'subagent' ? 'subagent' : null,
summaryAvailable: summary !== undefined,
};
}),
};
}
function sleep(ms, signal) {
return new Promise((resolve, reject) => {
if (signal?.aborted) {
reject(signal.reason ?? new DOMException('Aborted', 'AbortError'));
return;
}
const timer = setTimeout(() => {
signal?.removeEventListener('abort', onAbort);
resolve();
}, ms);
const onAbort = () => {
clearTimeout(timer);
reject(signal.reason ?? new DOMException('Aborted', 'AbortError'));
};
signal?.addEventListener('abort', onAbort, { once: true });
});
}
function assistantMessageText(event) {
return (event?.data?.message?.content ?? [])
.filter((part) => part.type === 'text' && typeof part.text === 'string')
.map((part) => part.text)
.join('\n')
.trim();
}
export class HarnessReplyTracker {
#promptRpcId;
#lastSeq;
#openTurn = null;
#targetTurn = null;
#stepText = new Map();
#latestText = '';
#finished = false;
#reason = null;
constructor({ promptRpcId, afterSeq = -1 }) {
this.#promptRpcId = promptRpcId;
this.#lastSeq = afterSeq;
}
get finished() {
return this.#finished;
}
get answer() {
return this.#latestText.trim();
}
get reason() {
return this.#reason;
}
consume(entries) {
let update = null;
const ordered = [...entries]
.map((entry) => entry?.event ?? entry)
.filter(Boolean)
.sort((left, right) => (left.seq ?? -1) - (right.seq ?? -1));
for (const event of ordered) {
const seq = event.seq ?? -1;
if (seq <= this.#lastSeq) continue;
this.#lastSeq = seq;
if (event.type === 'turn/start') this.#openTurn = event.data?.turn ?? null;
if (event.type === 'user/message' && event.data?.source?.rpcId === this.#promptRpcId) {
this.#targetTurn = this.#openTurn;
continue;
}
if (this.#targetTurn === null) continue;
if (event.type === 'turn/end') {
if (event.data?.turn !== this.#targetTurn) continue;
this.#finished = true;
this.#reason = event.data?.reason ?? null;
this.#openTurn = null;
continue;
}
if (event.data?.turn !== this.#targetTurn) continue;
if (event.type === 'assistant/chunk' && event.data?.chunk?.type === 'text-delta') {
const step = event.data?.step ?? 0;
const index = event.data.chunk.index ?? 0;
const key = `${step}:${index}`;
this.#stepText.set(key, (this.#stepText.get(key) ?? '') + event.data.chunk.text);
const prefix = `${step}:`;
const text = [...this.#stepText.entries()]
.filter(([partKey]) => partKey.startsWith(prefix))
.sort(([left], [right]) => Number(left.split(':')[1]) - Number(right.split(':')[1]))
.map(([, part]) => part)
.join('\n')
.trim();
if (text && text !== this.#latestText) {
this.#latestText = text;
update = { type: 'text', text };
}
continue;
}
if (event.type === 'assistant/message') {
const text = assistantMessageText(event);
if (text && text !== this.#latestText) {
this.#latestText = text;
update = { type: 'text', text };
}
continue;
}
if (event.type === 'tool/call') {
update = { type: 'tool', name: event.data?.name ?? '工具' };
} else if (event.type === 'tool/result') {
update = { type: 'status', text: '正在整理结果…' };
}
}
return update;
}
}
export class HarnessRpcError extends Error {
constructor(method, error) {
super(`${method}: ${error?.message ?? 'unknown Harness RPC error'}`);
this.name = 'HarnessRpcError';
this.method = method;
this.code = error?.code ?? 'internal';
this.details = error?.details ?? {};
}
}
export class HarnessClient {
#baseUrl;
#workspace;
#agentPreset;
#autostart;
#dshBin;
#fetch;
#managedProcess = null;
constructor({
baseUrl,
workspace,
agentPreset = 'standard',
autostart = false,
dshBin = 'dsh',
fetchImpl = fetch,
}) {
this.#baseUrl = new URL(baseUrl);
this.#workspace = workspace;
this.#agentPreset = agentPreset;
this.#autostart = autostart;
this.#dshBin = dshBin;
this.#fetch = fetchImpl;
}
async rpc(method, payload = {}, timeoutMs = 30_000, options = {}) {
const rpcId = options.rpcId ?? `dingtalk-${randomUUID()}`;
const timeoutSignal = AbortSignal.timeout(timeoutMs);
const signal = options.signal
? AbortSignal.any([options.signal, timeoutSignal])
: timeoutSignal;
const response = await this.#fetch(new URL(`/api/${method}`, this.#baseUrl), {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ type: 'client-request', rpcId, method, payload }),
signal,
export class HarnessClient extends SharedHarnessClient {
constructor(options) {
super({
...options,
rpcIdPrefix: 'dingtalk',
logPrefix: 'dsh-dingtalk',
});
if (!response.ok) throw new Error(`Harness transport ${method} failed: HTTP ${response.status}`);
const body = await response.json();
if (body?.type !== 'server-response' || body?.rpcId !== rpcId) {
throw new Error(`Harness returned an invalid response for ${method}`);
}
if (!body.result?.ok) throw new HarnessRpcError(method, body.result?.error);
return body.result.value;
}
async health(options = {}) {
await this.rpc('host.describe', {}, 5_000, options);
return true;
}
async ensureRunning(options = {}) {
try {
return await this.health(options);
} catch (firstError) {
if (!this.#autostart) throw firstError;
}
if (!this.#managedProcess || this.#managedProcess.exitCode !== null) {
const port = this.#baseUrl.port || (this.#baseUrl.protocol === 'https:' ? '443' : '80');
this.#managedProcess = spawn(this.#dshBin, [
'web', '--host', this.#baseUrl.hostname, '--port', port,
], {
cwd: this.#workspace,
env: process.env,
stdio: ['ignore', 'inherit', 'inherit'],
});
this.#managedProcess.on('error', (error) => {
console.error('[dsh-dingtalk] failed to start Harness:', error.message);
});
}
const deadline = Date.now() + 60_000;
let lastError;
while (Date.now() < deadline) {
await sleep(1_000, options.signal);
try {
return await this.health(options);
} catch (error) {
lastError = error;
}
}
throw new Error(`Harness did not become ready: ${lastError?.message ?? 'timeout'}`);
}
async listWorkspaces(options = {}) {
await this.ensureRunning(options);
return workspacePaths(await this.rpc('workspace.list', {}, 30_000, options));
}
async listWorkspaceSessions(workspacePath, options = {}) {
await this.ensureRunning(options);
const workspaceList = await this.rpc('workspace.list', {}, 30_000, options);
const workspace = workspaceFromList(workspacePath, workspaceList);
if (!workspace) return { workspace: workspacePath, sessions: [] };
const sessionList = await this.rpc('session.list', {}, 30_000, options);
return workspaceSessions(workspace, workspaceList.archivedSessionIds, sessionList);
}
async adoptWorkspaceSession(value, options = {}) {
return adoptRegisteredWorkspaceSession(this, value, options);
}
async workspaceId(options = {}) {
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: workspace }, 30_000, rpcOptions);
return created.workspace.workspaceId;
}
async createSession(options = {}) {
await this.ensureRunning(options);
const workspaceId = await this.workspaceId(options);
const created = await this.rpc('session.create', {
workspaceId,
agentPreset: this.#agentPreset,
}, 30_000, options);
return created.sessionId;
}
async sessionExists(sessionId, options = {}) {
try {
await this.rpc('session.history', { sessionId, maxMessages: 1 }, 30_000, options);
return true;
} catch (error) {
if (error instanceof HarnessRpcError && error.code === 'session-not-found') return false;
throw error;
}
}
async ask(sessionId, text, options = {}) {
if (typeof options === 'number') options = { timeoutMs: options };
const timeoutMs = options.timeoutMs ?? 600_000;
const signal = options.signal;
const onUpdate = typeof options.onUpdate === 'function' ? options.onUpdate : null;
await this.ensureRunning({ signal });
const before = await this.rpc(
'session.history',
{ sessionId, maxMessages: 1 },
30_000,
{ signal },
);
const baselineSeq = Math.max(-1, ...(before.events ?? []).map(({ event }) => event.seq ?? -1));
const promptRpcId = `dingtalk-${randomUUID()}`;
const tracker = new HarnessReplyTracker({ promptRpcId, afterSeq: baselineSeq });
await this.rpc('session.prompt', {
sessionId,
mode: 'queue',
content: [{ type: 'text', text }],
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 update = tracker.consume(history.events ?? []);
if (update && onUpdate) {
try {
await onUpdate(update);
} catch (error) {
console.warn('[dsh-dingtalk] ignored a progress update failure:', error.message);
}
}
if (!tracker.finished) continue;
if (tracker.answer) return tracker.answer;
throw new Error(
`Harness turn ended without a text reply${tracker.reason ? ` (${JSON.stringify(tracker.reason)})` : ''}`,
);
}
throw new Error(`Harness reply timed out after ${Math.round(timeoutMs / 1_000)} seconds`);
}
stopManagedProcess() {
if (this.#managedProcess?.exitCode === null) this.#managedProcess.kill('SIGTERM');
}
}

View file

@ -5,9 +5,17 @@ import {
isBotSender,
splitText,
} from './message-utils.mjs';
import {
harnessAnswerForQuestion,
harnessQuestionText,
validHarnessQuestion,
} from '../shared/harness-question.mjs';
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
const RESOLVED_REPLY_TTL_MS = 30 * 60_000;
const HELP_TEXT = [
'北汇星河 AIOS 已连接 DeepSeek Harness。',
'',
@ -21,16 +29,49 @@ const HELP_TEXT = [
'/help 显示本帮助',
].join('\n');
function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
function senderOpenId(event) {
return nonEmptyString(event?.sender?.sender_id?.open_id)
?? nonEmptyString(event?.sender?.sender_id?.user_id);
}
function canClaimInteractionReply(event, pending) {
return pending.needsPresentation !== true
&& pending.questions[pending.index]
&& senderOpenId(event) === pending.actor
&& event?.message?.message_type === 'text'
&& nonEmptyString(extractText(event));
}
function ensureStatus(status) {
for (const key of ['messagesReceived', 'messagesReplied', 'messagesRejected']) {
status[key] ??= 0;
}
status.lastMessageAt ??= null;
status.lastReplyAt ??= null;
status.lastRejectedAt ??= null;
status.lastError ??= null;
}
export class FeishuHarnessBridge {
#client;
#channel;
#harness;
#state;
#queues = new Map();
#pendingInteractions = new Map();
#interactionKeys = new Map();
#resolvedQuestionReplies = new Map();
#acceptedMessageIds = new Set();
#interactionTasks = new Set();
#status;
#allowedSenderOpenIds;
#replyTimeoutMs;
#logger;
#signal;
constructor({
client,
@ -39,8 +80,13 @@ export class FeishuHarnessBridge {
state,
status,
allowedSenderOpenIds = new Set(),
replyTimeoutMs = 600000,
replyTimeoutMs = 600_000,
logger = console,
signal,
}) {
if (!client || !harness || !state || !status) {
throw new TypeError('Feishu bridge dependencies are required');
}
this.#client = client;
this.#channel = channel;
this.#harness = harness;
@ -48,55 +94,171 @@ export class FeishuHarnessBridge {
this.#status = status;
this.#allowedSenderOpenIds = allowedSenderOpenIds;
this.#replyTimeoutMs = replyTimeoutMs;
this.#logger = logger;
this.#signal = signal;
ensureStatus(this.#status);
}
accept(event) {
const messageId = event?.message?.message_id;
if (!messageId || isBotSender(event) || event?.message?.message_type !== 'text') return;
if (this.#signal?.aborted) return Promise.resolve();
const messageId = nonEmptyString(event?.message?.message_id);
if (!messageId || isBotSender(event)) return Promise.resolve();
if (!isAllowedSender(event, this.#allowedSenderOpenIds)) {
this.#status.messagesRejected += 1;
this.#status.lastRejectedAt = new Date().toISOString();
console.warn('[bridge] ignored a message from a sender outside the allowlist');
return;
this.#logger.warn?.('[dsh-feishu] ignored a message from a sender outside the allowlist');
return Promise.resolve();
}
if (this.#state.hasSeen(messageId) || this.#acceptedMessageIds.has(messageId)) return;
if (this.#state.hasSeen(messageId) || this.#acceptedMessageIds.has(messageId)) {
return Promise.resolve();
}
let key;
try {
key = conversationKey(event);
} catch {
this.#status.messagesRejected += 1;
this.#status.lastRejectedAt = new Date().toISOString();
return Promise.resolve();
}
this.#acceptedMessageIds.add(messageId);
const processingReaction = this.#addReaction(messageId, 'OnIt');
if (this.#isResolvedQuestionReply(event, key)) {
const current = Promise.resolve()
.then(() => this.#discardResolvedInteractionReply(event, messageId))
.then(() => this.#finishReaction(messageId, processingReaction, 'DONE'))
.catch((error) => this.#handleMessageFailure(
event,
messageId,
processingReaction,
error,
))
.finally(() => this.#acceptedMessageIds.delete(messageId));
return current;
}
const pending = this.#pendingInteractions.get(key);
if (pending && senderOpenId(event) !== pending.actor) {
return this.#enqueueMessage(event, messageId, key, processingReaction);
}
if (pending?.submitting || pending?.claimedReplyMessageId) {
return this.#enqueueMessage(event, messageId, key, processingReaction);
}
if (pending) {
if (canClaimInteractionReply(event, pending)) pending.claimedReplyMessageId = messageId;
const previous = pending.queue ?? Promise.resolve();
const processing = previous
.catch(() => undefined)
.then(() => this.#processInteractionReply(
event,
messageId,
key,
pending,
processingReaction,
));
pending.queue = processing;
const key = conversationKey(event);
const releaseInteraction = () => {
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
if (pending.queue === processing) pending.queue = null;
};
let current;
current = processing
.then(
() => {
releaseInteraction();
return this.#finishReaction(messageId, processingReaction, 'DONE');
},
(error) => {
releaseInteraction();
return this.#handleMessageFailure(
event,
messageId,
processingReaction,
error,
);
},
)
.finally(() => {
releaseInteraction();
this.#acceptedMessageIds.delete(messageId);
this.#interactionTasks.delete(current);
});
this.#interactionTasks.add(current);
return current;
}
return this.#enqueueMessage(event, messageId, key, processingReaction);
}
#enqueueMessage(event, messageId, key, processingReaction, {
releaseMessageId = true,
alreadyRecorded = false,
finalize = true,
} = {}) {
const previous = this.#queues.get(key) ?? Promise.resolve();
const task = previous
const work = previous
.catch(() => undefined)
.then(() => this.#handle(event, key))
.then(() => this.#finishReaction(messageId, processingReaction, 'DONE'))
.catch(async (error) => {
console.error('[bridge] message handling failed:', error.message);
this.#status.lastError = error.message;
await this.#finishReaction(messageId, processingReaction, 'ERROR');
await this.#send(
event.message.chat_id,
'处理失败,请稍后重试。如果问题持续,请在 DeepSeek Harness 的飞书插件页面检查连接状态。',
).catch(() => undefined);
})
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
if (this.#queues.get(key) === task) this.#queues.delete(key);
});
this.#queues.set(key, task);
.then(() => this.#handle(event, key, { alreadyRecorded }));
const settled = finalize
? work
.then(() => this.#finishReaction(messageId, processingReaction, 'DONE'))
.catch((error) => this.#handleMessageFailure(
event,
messageId,
processingReaction,
error,
))
: work;
let current;
current = settled.finally(() => {
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
if (this.#queues.get(key) === current) this.#queues.delete(key);
});
this.#queues.set(key, current);
return current;
}
async #handleMessageFailure(event, messageId, processingReaction, error) {
if (this.#signal?.aborted) {
await this.#removeProcessingReaction(messageId, processingReaction);
return;
}
this.#logger.error?.('[dsh-feishu] message handling failed:', error?.message ?? String(error));
this.#status.lastError = error?.message ?? String(error);
await this.#finishReaction(messageId, processingReaction, 'ERROR');
await this.#send(
event.message.chat_id,
'处理失败,请稍后重试。如果问题持续,请在 DeepSeek Harness 的飞书插件页面检查连接状态。',
).catch(() => undefined);
}
async waitForIdle() {
await Promise.allSettled([...this.#queues.values()]);
await Promise.allSettled([
...this.#queues.values(),
...[...this.#pendingInteractions.values()].flatMap((pending) => (
pending.queue ? [pending.queue] : []
)),
...this.#interactionTasks,
]);
}
async #handle(event, key) {
async #handle(event, key, { alreadyRecorded = false } = {}) {
this.#signal?.throwIfAborted();
const messageId = event.message.message_id;
await this.#state.markSeen(messageId);
this.#status.lastMessageAt = new Date().toISOString();
this.#status.messagesReceived += 1;
if (!alreadyRecorded) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.lastMessageAt = new Date().toISOString();
this.#status.messagesReceived += 1;
}
const text = extractText(event);
if (!text) return;
if (!text) {
await this.#send(event.message.chat_id, '目前仅支持文字消息。');
return;
}
if (text === '/help') {
await this.#send(event.message.chat_id, HELP_TEXT);
@ -108,7 +270,7 @@ export class FeishuHarnessBridge {
return;
}
if (text === '/status') {
await this.#harness.ensureRunning();
await this.#harness.ensureRunning({ signal: this.#signal });
await this.#send(event.message.chat_id, '飞书机器人与 DeepSeek Harness 连接正常。');
return;
}
@ -120,11 +282,29 @@ export class FeishuHarnessBridge {
return;
}
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;
this.#logger.info?.(`[dsh-feishu] processing ${event.message.chat_type} message ${messageId}`);
try {
await this.#answerWithStream(event, key, text);
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} finally {
await this.#cancelPendingInteraction(key);
}
}
#interactionAskOptions(event, key) {
return {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key,
actor: senderOpenId(event),
chatId: event.message.chat_id,
requiresMention: event.message.chat_type !== 'p2p',
}),
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
};
}
async #answerWithStream(event, key, text) {
@ -136,7 +316,9 @@ export class FeishuHarnessBridge {
state: this.#state,
key,
text,
askOptions: { timeoutMs: this.#replyTimeoutMs },
createOptions: { signal: this.#signal },
existsOptions: { signal: this.#signal },
askOptions: this.#interactionAskOptions(event, key),
});
for (const chunk of splitText(answer)) await this.#send(chatId, chunk);
this.#status.streamFallbacks = (this.#status.streamFallbacks ?? 0) + 1;
@ -149,18 +331,21 @@ export class FeishuHarnessBridge {
await this.#channel.stream(chatId, {
markdown: async (controller) => {
promptStarted = true;
const askOptions = {
...this.#interactionAskOptions(event, key),
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;
},
},
createOptions: { signal: this.#signal },
existsOptions: { signal: this.#signal },
askOptions,
}));
await controller.setContent(completedAnswer);
},
@ -169,26 +354,308 @@ export class FeishuHarnessBridge {
} catch (error) {
this.#status.streamErrors = (this.#status.streamErrors ?? 0) + 1;
if (completedAnswer) {
console.warn('[bridge] native Feishu stream failed after generation; sending final text:', error.message);
this.#logger.warn?.(
'[dsh-feishu] native stream failed after generation; sending final text:',
error.message,
);
for (const chunk of splitText(completedAnswer)) await this.#send(chatId, chunk);
this.#status.streamFallbacks = (this.#status.streamFallbacks ?? 0) + 1;
return;
}
if (promptStarted) throw error;
console.warn('[bridge] native Feishu stream unavailable; using text fallback:', error.message);
this.#logger.warn?.('[dsh-feishu] native stream unavailable; using text fallback:', error.message);
const { answer } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
text,
askOptions: { timeoutMs: this.#replyTimeoutMs },
createOptions: { signal: this.#signal },
existsOptions: { signal: this.#signal },
askOptions: this.#interactionAskOptions(event, key),
});
for (const chunk of splitText(answer)) await this.#send(chatId, chunk);
this.#status.streamFallbacks = (this.#status.streamFallbacks ?? 0) + 1;
}
}
async #processInteractionReply(event, messageId, key, expected, processingReaction) {
this.#signal?.throwIfAborted();
const current = this.#pendingInteractions.get(key);
const claimed = expected.claimedReplyMessageId === messageId;
if (!current || current !== expected || current.submitting) {
if (this.#isResolvedQuestionReply(event, key)) {
return this.#discardResolvedInteractionReply(event, messageId);
}
if (claimed && (!current || current !== expected)) {
return this.#discardResolvedInteractionReply(event, messageId);
}
return this.#enqueueMessage(event, messageId, key, processingReaction, {
releaseMessageId: false,
finalize: false,
});
}
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.lastMessageAt = new Date().toISOString();
this.#status.messagesReceived += 1;
const text = extractText(event);
if (!text) {
await this.#send(event.message.chat_id, '请用文字回答当前问题。');
return;
}
const pending = this.#pendingInteractions.get(key);
if (!pending || pending !== expected || pending.submitting) {
if (this.#isResolvedQuestionReply(event, key)) {
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT).catch(() => undefined);
return;
}
if (claimed && (!pending || pending !== expected)) {
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT);
return;
}
return this.#enqueueMessage(event, messageId, key, processingReaction, {
releaseMessageId: false,
alreadyRecorded: true,
finalize: false,
});
}
pending.chatId = event.message.chat_id;
if (pending.needsPresentation) {
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = '飞书交互问题发送失败。';
this.#logger.error?.('[dsh-feishu] failed to retry an interaction question');
pending.interaction.reconnect?.();
}
return;
}
const question = pending.questions[pending.index];
if (!question) return;
pending.answers.push(harnessAnswerForQuestion(question, text));
pending.index += 1;
if (pending.index < pending.questions.length) {
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
pending.needsPresentation = true;
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = '飞书交互问题发送失败。';
this.#logger.error?.('[dsh-feishu] failed to send the next interaction question');
pending.interaction.reconnect?.();
}
return;
}
pending.submitting = true;
try {
await pending.interaction.respond({
ok: true,
value: {
sessionId: pending.sessionId,
answer: { answers: pending.answers },
},
});
this.#rememberResolvedInteraction(key, pending);
this.#clearPendingInteraction(key, pending.interactionId);
this.#status.lastError = null;
} catch (error) {
if (this.#signal?.aborted) return;
if (this.#pendingInteractions.get(key) !== pending) return;
if (error?.code === 'interaction-not-pending') {
this.#rememberResolvedInteraction(key, pending);
this.#clearPendingInteraction(key, pending.interactionId);
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT).catch(() => undefined);
return;
}
pending.submitting = false;
pending.answers.pop();
pending.index -= 1;
this.#status.lastError = '回答提交失败。';
this.#logger.error?.('[dsh-feishu] failed to answer a Harness interaction');
await this.#send(event.message.chat_id, '回答提交失败,请重新发送当前问题的答案。')
.catch(() => undefined);
}
}
async #handleInteraction(interaction, {
key,
actor,
chatId,
requiresMention,
}) {
// Approval is deliberately exposed by the transport but remains
// unanswered until #5 adds an authenticated policy and renderer.
if (interaction?.kind !== 'question') return;
const questions = interaction?.payload?.questions;
const interactionId = typeof interaction?.interactionId === 'string'
? interaction.interactionId
: interaction?.rpcId;
if (typeof interaction.rpcId !== 'string'
|| typeof interactionId !== 'string'
|| typeof interaction.sessionId !== 'string'
|| !Array.isArray(questions)
|| questions.length === 0
|| questions.some((question) => !validHarnessQuestion(question))) {
this.#logger.warn?.('[dsh-feishu] ignored an invalid Harness question interaction');
return;
}
if (interaction.recovered === true) {
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'Feishu safely cancelled an interaction left by an earlier client.',
details: {},
},
});
await this.#send(
chatId,
'检测到这个 Session 中遗留的待回答问题,已安全取消并继续处理你刚才的消息。',
).catch(() => undefined);
return;
}
const existing = this.#pendingInteractions.get(key);
if (existing?.interactionId === interactionId) {
existing.interaction = interaction;
if (existing.needsPresentation) await this.#presentInteraction(existing);
return;
}
if (this.#interactionKeys.has(interactionId)) return;
if (existing) {
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'Feishu is already handling another user interaction.',
details: {},
},
});
return;
}
const pending = {
kind: 'question',
interactionId,
sessionId: interaction.sessionId,
interaction,
key,
actor,
requiresMention,
questions,
answers: [],
index: 0,
chatId,
queue: null,
claimedReplyMessageId: null,
submitting: false,
needsPresentation: true,
questionMessageIds: new Set(),
inactive: false,
};
this.#pendingInteractions.set(key, pending);
this.#interactionKeys.set(pending.interactionId, key);
await this.#presentInteraction(pending);
}
#handleInteractionResolved(resolution) {
const interactionId = resolution?.interactionId;
if (resolution?.kind !== 'question' || typeof interactionId !== 'string') return;
const key = this.#interactionKeys.get(interactionId);
if (!key) return;
const pending = this.#pendingInteractions.get(key);
if (pending) this.#rememberResolvedInteraction(key, pending);
this.#clearPendingInteraction(key, interactionId);
}
async #presentInteraction(pending) {
const question = pending.questions[pending.index];
if (!question) return;
const messageId = await this.#send(
pending.chatId,
harnessQuestionText(
question,
pending.index,
pending.questions.length,
{ requiresMention: pending.requiresMention },
),
);
if (messageId) {
pending.questionMessageIds.add(messageId);
if (pending.inactive) this.#rememberResolvedInteraction(pending.key, pending);
}
pending.needsPresentation = false;
}
#rememberResolvedInteraction(key, pending) {
const expiresAt = Date.now() + RESOLVED_REPLY_TTL_MS;
for (const messageId of pending.questionMessageIds ?? []) {
this.#resolvedQuestionReplies.set(messageId, { key, expiresAt });
}
}
#isResolvedQuestionReply(event, key) {
const now = Date.now();
for (const [messageId, resolution] of this.#resolvedQuestionReplies) {
if (resolution.expiresAt <= now) this.#resolvedQuestionReplies.delete(messageId);
}
for (const reference of [event?.message?.parent_id, event?.message?.root_id]) {
const resolution = this.#resolvedQuestionReplies.get(reference);
if (resolution?.key === key && resolution.expiresAt > now) return true;
}
return false;
}
async #discardResolvedInteractionReply(event, messageId) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.lastMessageAt = new Date().toISOString();
this.#status.messagesReceived += 1;
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT).catch(() => undefined);
}
#takePendingInteraction(key, interactionId) {
const pending = this.#pendingInteractions.get(key);
if (!pending
|| (interactionId !== undefined && pending.interactionId !== interactionId)) return null;
this.#pendingInteractions.delete(key);
this.#interactionKeys.delete(pending.interactionId);
pending.inactive = true;
return pending;
}
#clearPendingInteraction(key, interactionId) {
return this.#takePendingInteraction(key, interactionId) !== null;
}
async #cancelPendingInteraction(key) {
const pending = this.#takePendingInteraction(key);
if (!pending || pending.kind !== 'question') return;
this.#rememberResolvedInteraction(key, pending);
try {
await pending.interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'The Feishu interaction ended before the user answered.',
details: {},
},
}, { signal: AbortSignal.timeout(5_000) });
} catch (error) {
if (error?.code !== 'interaction-not-pending') {
this.#logger.warn?.('[dsh-feishu] failed to cancel a pending Harness interaction');
}
}
}
#progressText(update) {
if (update.type === 'text' && update.text) return update.text;
if (update.type === 'tool') {
@ -206,12 +673,12 @@ export class FeishuHarnessBridge {
return reactionId;
} catch (error) {
this.#status.reactionErrors = (this.#status.reactionErrors ?? 0) + 1;
console.warn(`[bridge] unable to add ${emojiType} reaction:`, error.message);
this.#logger.warn?.(`[dsh-feishu] unable to add ${emojiType} reaction:`, error.message);
return null;
}
}
async #finishReaction(messageId, processingReaction, finalEmojiType) {
async #removeProcessingReaction(messageId, processingReaction) {
const reactionId = await processingReaction;
if (reactionId && this.#channel?.removeReaction) {
try {
@ -219,9 +686,13 @@ export class FeishuHarnessBridge {
this.#status.reactionsRemoved = (this.#status.reactionsRemoved ?? 0) + 1;
} catch (error) {
this.#status.reactionErrors = (this.#status.reactionErrors ?? 0) + 1;
console.warn('[bridge] unable to remove processing reaction:', error.message);
this.#logger.warn?.('[dsh-feishu] unable to remove processing reaction:', error.message);
}
}
}
async #finishReaction(messageId, processingReaction, finalEmojiType) {
await this.#removeProcessingReaction(messageId, processingReaction);
await this.#addReaction(messageId, finalEmojiType);
}
@ -237,5 +708,6 @@ export class FeishuHarnessBridge {
if (response?.code && response.code !== 0) {
throw new Error(`Feishu send failed: ${response.msg || response.code}`);
}
return nonEmptyString(response?.data?.message_id);
}
}

View file

@ -1,6 +1,26 @@
import { FeishuHarnessBridge } from './bridge.mjs';
import { VerifiedFeishuChannel } from './feishu-channel.mjs';
const DEFAULT_REQUEST_TIMEOUT_MS = 15_000;
function httpInstanceWithTimeout(httpInstance, timeoutMs) {
if (!httpInstance || typeof httpInstance.request !== 'function') return undefined;
const optionsWithTimeout = (options) => ({
...(options ?? {}),
timeout: options?.timeout ?? timeoutMs,
});
return {
request: (options) => httpInstance.request(optionsWithTimeout(options)),
get: (url, options) => httpInstance.get(url, optionsWithTimeout(options)),
delete: (url, options) => httpInstance.delete(url, optionsWithTimeout(options)),
head: (url, options) => httpInstance.head(url, optionsWithTimeout(options)),
options: (url, options) => httpInstance.options(url, optionsWithTimeout(options)),
post: (url, data, options) => httpInstance.post(url, data, optionsWithTimeout(options)),
put: (url, data, options) => httpInstance.put(url, data, optionsWithTimeout(options)),
patch: (url, data, options) => httpInstance.patch(url, data, optionsWithTimeout(options)),
};
}
export function createBridgeStatus({ allowedSenderCount = 1 } = {}) {
return {
startedAt: null,
@ -43,11 +63,13 @@ export class FeishuRuntime {
#state;
#replyTimeoutMs;
#connectTimeoutMs;
#requestTimeoutMs;
#logger;
#client = null;
#bridge = null;
#wsClient = null;
#starting = null;
#abortController = null;
#status;
constructor({
@ -61,6 +83,7 @@ export class FeishuRuntime {
state,
replyTimeoutMs = 600000,
connectTimeoutMs = 15000,
requestTimeoutMs = DEFAULT_REQUEST_TIMEOUT_MS,
logger = console,
}) {
if (!lark) throw new Error('FeishuRuntime requires the Feishu SDK');
@ -70,6 +93,9 @@ export class FeishuRuntime {
if (normalizedOwners.length === 0) throw new Error('FeishuRuntime requires at least one owner open_id');
if (!harness) throw new Error('FeishuRuntime requires a Harness client');
if (!state) throw new Error('FeishuRuntime requires a state store');
if (!Number.isFinite(requestTimeoutMs) || requestTimeoutMs <= 0) {
throw new TypeError('FeishuRuntime requestTimeoutMs must be a positive number');
}
this.#lark = lark;
this.#appId = appId;
@ -80,6 +106,7 @@ export class FeishuRuntime {
this.#state = state;
this.#replyTimeoutMs = replyTimeoutMs;
this.#connectTimeoutMs = connectTimeoutMs;
this.#requestTimeoutMs = requestTimeoutMs;
this.#logger = logger;
this.#status = createBridgeStatus({ allowedSenderCount: normalizedOwners.length });
}
@ -99,12 +126,15 @@ export class FeishuRuntime {
}
async #start() {
const abortController = new AbortController();
this.#abortController = abortController;
const { signal } = abortController;
this.#status.startedAt = new Date().toISOString();
this.#status.feishuLongConnectionState = 'connecting';
this.#status.lastError = null;
try {
await this.#harness.ensureRunning();
await this.#harness.ensureRunning({ signal });
this.#status.harnessReachable = true;
const sdkDomain = this.#domain === 'lark'
@ -115,6 +145,11 @@ export class FeishuRuntime {
appSecret: this.#appSecret,
domain: sdkDomain,
};
const httpInstance = httpInstanceWithTimeout(
this.#lark.defaultHttpInstance,
this.#requestTimeoutMs,
);
if (httpInstance) larkConfig.httpInstance = httpInstance;
this.#client = new this.#lark.Client(larkConfig);
const channel = new VerifiedFeishuChannel({
client: this.#client,
@ -128,6 +163,8 @@ export class FeishuRuntime {
status: this.#status,
allowedSenderOpenIds: new Set(this.#ownerOpenIds),
replyTimeoutMs: this.#replyTimeoutMs,
signal,
logger: this.#logger,
});
const dispatcher = new this.#lark.EventDispatcher({}).register({
@ -205,6 +242,9 @@ export class FeishuRuntime {
async stop({ preserveError = false } = {}) {
const error = preserveError ? this.#status.lastError : null;
const abortController = this.#abortController;
this.#abortController = null;
abortController?.abort(new DOMException('Feishu runtime stopped', 'AbortError'));
this.#status.ready = false;
if (this.#wsClient) {
this.#wsClient.close({ force: true });

View file

@ -1,338 +1,19 @@
import { spawn } from 'node:child_process';
import { randomUUID } from 'node:crypto';
import { isAbsolute } from 'node:path';
import {
HarnessClient as SharedHarnessClient,
} from '../shared/harness-client.mjs';
import { adoptRegisteredWorkspaceSession } from '../shared/harness-session-binding.mjs';
export {
HarnessInteractionError,
HarnessReplyTracker,
HarnessRpcError,
} from '../shared/harness-client.mjs';
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
function workspacePaths(value) {
if (!Array.isArray(value?.items)) return [];
return value.items.flatMap((item) => (
typeof item?.path === 'string' && isAbsolute(item.path) ? [item.path] : []
));
}
function workspaceFromList(workspacePath, workspaceList) {
if (!Array.isArray(workspaceList?.items)
|| !Array.isArray(workspaceList?.archivedSessionIds)) {
throw new Error('Harness returned an invalid response for workspace.list');
}
const workspace = workspaceList.items.find((item) => item?.path === workspacePath);
if (!workspace) return null;
if (!Array.isArray(workspace.sessionIds)
|| workspace.sessionIds.some((sessionId) => typeof sessionId !== 'string')) {
throw new Error('Harness returned invalid session IDs for workspace.list');
}
return workspace;
}
function workspaceSessions(workspace, archivedSessionIds, sessionList) {
if (!Array.isArray(sessionList?.items)) {
throw new Error('Harness returned an invalid response for session.list');
}
const archived = new Set(archivedSessionIds);
const summaries = new Map(sessionList.items.flatMap((item) => (
typeof item?.sessionId === 'string' ? [[item.sessionId, item]] : []
)));
return {
workspace: workspace.path,
sessions: workspace.sessionIds.map((sessionId) => {
const summary = summaries.get(sessionId);
const title = summary?.projections?.values?.title;
return {
sessionId,
title: typeof title === 'string' ? title : null,
archived: archived.has(sessionId),
blank: summary?.blank === true,
origin: summary?.origin === 'subagent' ? 'subagent' : null,
summaryAvailable: summary !== undefined,
};
}),
};
}
function messageText(event) {
return (event?.data?.message?.content ?? [])
.filter((part) => part.type === 'text' && typeof part.text === 'string')
.map((part) => part.text)
.join('\n')
.trim();
}
export class HarnessReplyTracker {
#promptRpcId;
#lastSeq;
#openTurn = null;
#targetTurn = null;
#stepText = new Map();
#latestText = '';
#finished = false;
#reason = null;
constructor({ promptRpcId, afterSeq = -1 }) {
this.#promptRpcId = promptRpcId;
this.#lastSeq = afterSeq;
}
get finished() {
return this.#finished;
}
get answer() {
return this.#latestText.trim();
}
get reason() {
return this.#reason;
}
consume(entries) {
let update = null;
const ordered = [...entries]
.map((entry) => entry?.event ?? entry)
.filter(Boolean)
.sort((left, right) => (left.seq ?? -1) - (right.seq ?? -1));
for (const event of ordered) {
const seq = event.seq ?? -1;
if (seq <= this.#lastSeq) continue;
this.#lastSeq = seq;
if (event.type === 'turn/start') {
this.#openTurn = event.data?.turn ?? null;
}
if (event.type === 'user/message'
&& event.data?.source?.rpcId === this.#promptRpcId) {
this.#targetTurn = this.#openTurn;
continue;
}
if (this.#targetTurn === null) continue;
if (event.type === 'turn/end') {
if (event.data?.turn !== this.#targetTurn) continue;
this.#finished = true;
this.#reason = event.data?.reason ?? null;
this.#openTurn = null;
continue;
}
if (event.data?.turn !== this.#targetTurn) continue;
if (event.type === 'assistant/chunk' && event.data?.chunk?.type === 'text-delta') {
const step = event.data?.step ?? 0;
const index = event.data.chunk.index ?? 0;
const key = `${step}:${index}`;
this.#stepText.set(key, (this.#stepText.get(key) ?? '') + event.data.chunk.text);
const stepPrefix = `${step}:`;
const text = [...this.#stepText.entries()]
.filter(([partKey]) => partKey.startsWith(stepPrefix))
.sort(([left], [right]) => Number(left.split(':')[1]) - Number(right.split(':')[1]))
.map(([, part]) => part)
.join('\n')
.trim();
if (text && text !== this.#latestText) {
this.#latestText = text;
update = { type: 'text', text };
}
continue;
}
if (event.type === 'assistant/message') {
const text = messageText(event);
if (text && text !== this.#latestText) {
this.#latestText = text;
update = { type: 'text', text };
}
continue;
}
if (event.type === 'tool/call') {
update = { type: 'tool', name: event.data?.name ?? '工具' };
} else if (event.type === 'tool/result') {
update = { type: 'status', text: '正在整理结果…' };
}
}
return update;
}
}
export class HarnessRpcError extends Error {
constructor(method, error) {
super(`${method}: ${error?.message ?? 'unknown Harness RPC error'}`);
this.name = 'HarnessRpcError';
this.method = method;
this.code = error?.code ?? 'internal';
this.details = error?.details ?? {};
}
}
export class HarnessClient {
#baseUrl;
#workspace;
#agentPreset;
#autostart;
#dshBin;
#managedProcess = null;
constructor({ baseUrl, workspace, agentPreset, autostart, dshBin }) {
this.#baseUrl = new URL(baseUrl);
this.#workspace = workspace;
this.#agentPreset = agentPreset;
this.#autostart = autostart;
this.#dshBin = dshBin;
}
async rpc(method, payload = {}, timeoutMs = 30000, options = {}) {
const rpcId = options.rpcId ?? `feishu-${randomUUID()}`;
const response = await fetch(new URL(`/api/${method}`, this.#baseUrl), {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ type: 'client-request', rpcId, method, payload }),
signal: AbortSignal.timeout(timeoutMs),
export class HarnessClient extends SharedHarnessClient {
constructor(options) {
super({
...options,
rpcIdPrefix: 'feishu',
logPrefix: 'dsh-feishu',
});
if (!response.ok) throw new Error(`Harness transport ${method} failed: HTTP ${response.status}`);
const body = await response.json();
if (body?.type !== 'server-response' || body?.rpcId !== rpcId) {
throw new Error(`Harness returned an invalid response for ${method}`);
}
if (!body.result?.ok) throw new HarnessRpcError(method, body.result?.error);
return body.result.value;
}
async health() {
await this.rpc('host.describe', {}, 5000);
return true;
}
async ensureRunning() {
try {
return await this.health();
} catch (firstError) {
if (!this.#autostart) throw firstError;
}
if (!this.#managedProcess || this.#managedProcess.exitCode !== null) {
const port = this.#baseUrl.port || (this.#baseUrl.protocol === 'https:' ? '443' : '80');
this.#managedProcess = spawn(this.#dshBin, [
'web',
'--host',
this.#baseUrl.hostname,
'--port',
port,
], {
cwd: this.#workspace,
env: process.env,
stdio: ['ignore', 'inherit', 'inherit'],
});
this.#managedProcess.on('error', (error) => {
console.error('[bridge] failed to start Harness:', error.message);
});
}
const deadline = Date.now() + 60000;
let lastError;
while (Date.now() < deadline) {
await sleep(1000);
try {
return await this.health();
} catch (error) {
lastError = error;
}
}
throw new Error(`Harness did not become ready: ${lastError?.message ?? 'timeout'}`);
}
async listWorkspaces(options = {}) {
await this.ensureRunning();
return workspacePaths(await this.rpc('workspace.list', {}, 30000, options));
}
async listWorkspaceSessions(workspacePath, options = {}) {
await this.ensureRunning();
const workspaceList = await this.rpc('workspace.list', {}, 30000, options);
const workspace = workspaceFromList(workspacePath, workspaceList);
if (!workspace) return { workspace: workspacePath, sessions: [] };
const sessionList = await this.rpc('session.list', {}, 30000, options);
return workspaceSessions(workspace, workspaceList.archivedSessionIds, sessionList);
}
async adoptWorkspaceSession(value, options = {}) {
return adoptRegisteredWorkspaceSession(this, value, options, 30000);
}
async workspaceId(options = {}) {
const workspace = options.workspace ?? this.#workspace;
const { items } = await this.rpc('workspace.list', {});
const existing = items.find((item) => item.path === workspace);
if (existing) return existing.workspaceId;
const created = await this.rpc('workspace.create', { path: workspace });
return created.workspace.workspaceId;
}
async createSession(options = {}) {
await this.ensureRunning();
const workspaceId = await this.workspaceId(options);
const created = await this.rpc('session.create', {
workspaceId,
agentPreset: this.#agentPreset,
});
return created.sessionId;
}
async sessionExists(sessionId) {
try {
await this.rpc('session.history', { sessionId, maxMessages: 1 });
return true;
} catch (error) {
if (error instanceof HarnessRpcError && error.code === 'session-not-found') return false;
throw error;
}
}
async ask(sessionId, text, options = {}) {
if (typeof options === 'number') options = { timeoutMs: options };
const timeoutMs = options.timeoutMs ?? 600000;
const onUpdate = typeof options.onUpdate === 'function' ? options.onUpdate : null;
await this.ensureRunning();
const before = await this.rpc('session.history', { sessionId, maxMessages: 1 });
const baselineSeq = Math.max(-1, ...(before.events ?? []).map(({ event }) => event.seq ?? -1));
const promptRpcId = `feishu-${randomUUID()}`;
const tracker = new HarnessReplyTracker({ promptRpcId, afterSeq: baselineSeq });
await this.rpc('session.prompt', {
sessionId,
mode: 'queue',
content: [{ type: 'text', text }],
clientTimeZone: Intl.DateTimeFormat().resolvedOptions().timeZone,
}, 30000, { rpcId: promptRpcId });
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
await sleep(300);
const history = await this.rpc('session.history', { sessionId, maxMessages: 50 });
const update = tracker.consume(history.events ?? []);
if (update && onUpdate) {
try {
await onUpdate(update);
} catch (error) {
console.warn('[bridge] ignored a progress update failure:', error.message);
}
}
if (!tracker.finished) continue;
if (tracker.answer) return tracker.answer;
throw new Error(`Harness turn ended without a text reply${tracker.reason ? ` (${JSON.stringify(tracker.reason)})` : ''}`);
}
throw new Error(`Harness reply timed out after ${Math.round(timeoutMs / 1000)} seconds`);
}
stopManagedProcess() {
if (this.#managedProcess?.exitCode === null) this.#managedProcess.kill('SIGTERM');
}
}

View file

@ -0,0 +1,825 @@
import { spawn } from 'node:child_process';
import { randomUUID } from 'node:crypto';
import { isAbsolute } from 'node:path';
import { adoptRegisteredWorkspaceSession } from './harness-session-binding.mjs';
// Every channel plugin runs in the same Host process. Sharing ownership by
// Harness origin prevents two channel-specific clients bound to one Session
// from claiming or cancelling each other's interactions.
const interactionRegistries = new Map();
function interactionRegistry(origin) {
let registry = interactionRegistries.get(origin);
if (!registry) {
registry = { ownerships: new Map(), claims: new Map(), nextOrder: 0 };
interactionRegistries.set(origin, registry);
}
return registry;
}
function workspacePaths(value) {
if (!Array.isArray(value?.items)) return [];
return value.items.flatMap((item) => (
typeof item?.path === 'string' && isAbsolute(item.path) ? [item.path] : []
));
}
function workspaceFromList(workspacePath, workspaceList) {
if (!Array.isArray(workspaceList?.items)
|| !Array.isArray(workspaceList?.archivedSessionIds)) {
throw new Error('Harness returned an invalid response for workspace.list');
}
const workspace = workspaceList.items.find((item) => item?.path === workspacePath);
if (!workspace) return null;
if (!Array.isArray(workspace.sessionIds)
|| workspace.sessionIds.some((sessionId) => typeof sessionId !== 'string')) {
throw new Error('Harness returned invalid session IDs for workspace.list');
}
return workspace;
}
function workspaceSessions(workspace, archivedSessionIds, sessionList) {
if (!Array.isArray(sessionList?.items)) {
throw new Error('Harness returned an invalid response for session.list');
}
const archived = new Set(archivedSessionIds);
const summaries = new Map(sessionList.items.flatMap((item) => (
typeof item?.sessionId === 'string' ? [[item.sessionId, item]] : []
)));
return {
workspace: workspace.path,
sessions: workspace.sessionIds.map((sessionId) => {
const summary = summaries.get(sessionId);
const title = summary?.projections?.values?.title;
return {
sessionId,
title: typeof title === 'string' ? title : null,
archived: archived.has(sessionId),
blank: summary?.blank === true,
origin: summary?.origin === 'subagent' ? 'subagent' : null,
summaryAvailable: summary !== undefined,
};
}),
};
}
function sleep(ms, signal) {
return new Promise((resolve, reject) => {
if (signal?.aborted) {
reject(signal.reason ?? new DOMException('Aborted', 'AbortError'));
return;
}
const timer = setTimeout(() => {
signal?.removeEventListener('abort', onAbort);
resolve();
}, ms);
const onAbort = () => {
clearTimeout(timer);
reject(signal.reason ?? new DOMException('Aborted', 'AbortError'));
};
signal?.addEventListener('abort', onAbort, { once: true });
});
}
function assistantMessageText(event) {
return (event?.data?.message?.content ?? [])
.filter((part) => part.type === 'text' && typeof part.text === 'string')
.map((part) => part.text)
.join('\n')
.trim();
}
function consumeInteractionOwnership(ownership, entries) {
const ordered = [...entries]
.map((entry) => entry?.event ?? entry)
.filter(Boolean)
.sort((left, right) => (left.seq ?? -1) - (right.seq ?? -1));
for (const event of ordered) {
const seq = event.seq ?? -1;
if (seq <= ownership.lastSeq) continue;
ownership.lastSeq = seq;
if (event.type === 'turn/start') {
const turn = event.data?.turn ?? null;
if (ownership.active && turn !== ownership.turn) ownership.active = false;
if (ownership.turn !== null && turn !== ownership.turn) ownership.completed = true;
ownership.openTurn = turn;
continue;
}
if (event.type === 'user/message' && event.data?.source?.rpcId === ownership.promptRpcId) {
ownership.active = true;
ownership.started = true;
ownership.completed = false;
ownership.turn = event.data?.turn ?? ownership.openTurn;
continue;
}
if (event.type === 'turn/end' && event.data?.turn === ownership.turn) {
ownership.active = false;
ownership.completed = true;
}
}
}
export class HarnessReplyTracker {
#promptRpcId;
#lastSeq;
#openTurn = null;
#targetTurn = null;
#stepText = new Map();
#latestText = '';
#finished = false;
#reason = null;
constructor({ promptRpcId, afterSeq = -1 }) {
this.#promptRpcId = promptRpcId;
this.#lastSeq = afterSeq;
}
get finished() {
return this.#finished;
}
get answer() {
return this.#latestText.trim();
}
get reason() {
return this.#reason;
}
get tracking() {
return this.#targetTurn !== null && !this.#finished;
}
get turn() {
return this.#targetTurn;
}
consume(entries) {
let update = null;
const ordered = [...entries]
.map((entry) => entry?.event ?? entry)
.filter(Boolean)
.sort((left, right) => (left.seq ?? -1) - (right.seq ?? -1));
for (const event of ordered) {
const seq = event.seq ?? -1;
if (seq <= this.#lastSeq) continue;
this.#lastSeq = seq;
if (event.type === 'turn/start') this.#openTurn = event.data?.turn ?? null;
if (event.type === 'user/message' && event.data?.source?.rpcId === this.#promptRpcId) {
this.#targetTurn = this.#openTurn;
continue;
}
if (this.#targetTurn === null) continue;
if (event.type === 'turn/end') {
if (event.data?.turn !== this.#targetTurn) continue;
this.#finished = true;
this.#reason = event.data?.reason ?? null;
this.#openTurn = null;
continue;
}
if (event.data?.turn !== this.#targetTurn) continue;
if (event.type === 'assistant/chunk' && event.data?.chunk?.type === 'text-delta') {
const step = event.data?.step ?? 0;
const index = event.data.chunk.index ?? 0;
const key = `${step}:${index}`;
this.#stepText.set(key, (this.#stepText.get(key) ?? '') + event.data.chunk.text);
const prefix = `${step}:`;
const text = [...this.#stepText.entries()]
.filter(([partKey]) => partKey.startsWith(prefix))
.sort(([left], [right]) => Number(left.split(':')[1]) - Number(right.split(':')[1]))
.map(([, part]) => part)
.join('\n')
.trim();
if (text && text !== this.#latestText) {
this.#latestText = text;
update = { type: 'text', text };
}
continue;
}
if (event.type === 'assistant/message') {
const text = assistantMessageText(event);
if (text && text !== this.#latestText) {
this.#latestText = text;
update = { type: 'text', text };
}
continue;
}
if (event.type === 'tool/call') {
update = { type: 'tool', name: event.data?.name ?? '工具' };
} else if (event.type === 'tool/result') {
update = { type: 'status', text: '正在整理结果…' };
}
}
return update;
}
}
export class HarnessRpcError extends Error {
constructor(method, error) {
super(`${method}: ${error?.message ?? 'unknown Harness RPC error'}`);
this.name = 'HarnessRpcError';
this.method = method;
this.code = error?.code ?? 'internal';
this.details = error?.details ?? {};
}
}
export class HarnessInteractionError extends Error {
constructor(code, message) {
super(message);
this.name = 'HarnessInteractionError';
this.code = code;
}
}
export class HarnessClient {
#baseUrl;
#workspace;
#agentPreset;
#autostart;
#dshBin;
#fetch;
#createWebSocket;
#interactionReconnectDelayMs;
#rpcIdPrefix;
#logPrefix;
#managedProcess = null;
#interactionRegistry;
#interactionOwnerships;
#interactionClaims;
constructor({
baseUrl,
workspace,
agentPreset = 'standard',
autostart = false,
dshBin = 'dsh',
fetchImpl = fetch,
createWebSocket = (url) => new WebSocket(url),
interactionReconnectDelayMs = 500,
rpcIdPrefix = 'im',
logPrefix = 'dsh-im',
}) {
if (typeof createWebSocket !== 'function') {
throw new TypeError('createWebSocket must be a function');
}
if (!Number.isFinite(interactionReconnectDelayMs) || interactionReconnectDelayMs < 0) {
throw new TypeError('interactionReconnectDelayMs must be a non-negative number');
}
if (typeof rpcIdPrefix !== 'string' || !rpcIdPrefix.trim()) {
throw new TypeError('rpcIdPrefix must be a non-empty string');
}
if (typeof logPrefix !== 'string' || !logPrefix.trim()) {
throw new TypeError('logPrefix must be a non-empty string');
}
this.#baseUrl = new URL(baseUrl);
this.#workspace = workspace;
this.#agentPreset = agentPreset;
this.#autostart = autostart;
this.#dshBin = dshBin;
this.#fetch = fetchImpl;
this.#createWebSocket = createWebSocket;
this.#interactionReconnectDelayMs = interactionReconnectDelayMs;
this.#rpcIdPrefix = rpcIdPrefix.trim();
this.#logPrefix = logPrefix.trim();
this.#interactionRegistry = interactionRegistry(this.#baseUrl.origin);
this.#interactionOwnerships = this.#interactionRegistry.ownerships;
this.#interactionClaims = this.#interactionRegistry.claims;
}
async rpc(method, payload = {}, timeoutMs = 30_000, options = {}) {
const rpcId = options.rpcId ?? `${this.#rpcIdPrefix}-${randomUUID()}`;
const timeoutSignal = AbortSignal.timeout(timeoutMs);
const signal = options.signal
? AbortSignal.any([options.signal, timeoutSignal])
: timeoutSignal;
const response = await this.#fetch(new URL(`/api/${method}`, this.#baseUrl), {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ type: 'client-request', rpcId, method, payload }),
signal,
});
if (!response.ok) throw new Error(`Harness transport ${method} failed: HTTP ${response.status}`);
const body = await response.json();
if (body?.type !== 'server-response' || body?.rpcId !== rpcId) {
throw new Error(`Harness returned an invalid response for ${method}`);
}
if (!body.result?.ok) throw new HarnessRpcError(method, body.result?.error);
return body.result.value;
}
async health(options = {}) {
await this.rpc('host.describe', {}, 5_000, options);
return true;
}
async ensureRunning(options = {}) {
try {
return await this.health(options);
} catch (firstError) {
if (!this.#autostart) throw firstError;
}
if (!this.#managedProcess || this.#managedProcess.exitCode !== null) {
const port = this.#baseUrl.port || (this.#baseUrl.protocol === 'https:' ? '443' : '80');
this.#managedProcess = spawn(this.#dshBin, [
'web', '--host', this.#baseUrl.hostname, '--port', port,
], {
cwd: this.#workspace,
env: process.env,
stdio: ['ignore', 'inherit', 'inherit'],
});
this.#managedProcess.on('error', (error) => {
console.error(`[${this.#logPrefix}] failed to start Harness:`, error.message);
});
}
const deadline = Date.now() + 60_000;
let lastError;
while (Date.now() < deadline) {
await sleep(1_000, options.signal);
try {
return await this.health(options);
} catch (error) {
lastError = error;
}
}
throw new Error(`Harness did not become ready: ${lastError?.message ?? 'timeout'}`);
}
async listWorkspaces(options = {}) {
await this.ensureRunning(options);
return workspacePaths(await this.rpc('workspace.list', {}, 30_000, options));
}
async listWorkspaceSessions(workspacePath, options = {}) {
await this.ensureRunning(options);
const workspaceList = await this.rpc('workspace.list', {}, 30_000, options);
const workspace = workspaceFromList(workspacePath, workspaceList);
if (!workspace) return { workspace: workspacePath, sessions: [] };
const sessionList = await this.rpc('session.list', {}, 30_000, options);
return workspaceSessions(workspace, workspaceList.archivedSessionIds, sessionList);
}
async adoptWorkspaceSession(value, options = {}) {
return adoptRegisteredWorkspaceSession(this, value, options);
}
async workspaceId(options = {}) {
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: workspace }, 30_000, rpcOptions);
return created.workspace.workspaceId;
}
async createSession(options = {}) {
await this.ensureRunning(options);
const workspaceId = await this.workspaceId(options);
const created = await this.rpc('session.create', {
workspaceId,
agentPreset: this.#agentPreset,
}, 30_000, options);
return created.sessionId;
}
async sessionExists(sessionId, options = {}) {
try {
await this.rpc('session.history', { sessionId, maxMessages: 1 }, 30_000, options);
return true;
} catch (error) {
if (error instanceof HarnessRpcError && error.code === 'session-not-found') return false;
throw error;
}
}
async respondInteraction(rpcId, result, options = {}) {
if (typeof rpcId !== 'string' || !rpcId) throw new TypeError('rpcId is required');
if (!result || typeof result !== 'object' || typeof result.ok !== 'boolean') {
throw new TypeError('A Harness RPC result is required');
}
const timeoutSignal = AbortSignal.timeout(options.timeoutMs ?? 30_000);
const signal = options.signal
? AbortSignal.any([options.signal, timeoutSignal])
: timeoutSignal;
const response = await this.#fetch(new URL('/api/respond', this.#baseUrl), {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ type: 'client-response', rpcId, result }),
signal,
});
if (!response.ok) {
throw new Error(`Harness transport respond failed: HTTP ${response.status}`);
}
const receipt = await response.json();
if (receipt?.accepted === true) return receipt;
if (receipt?.accepted !== false
|| (receipt.reason !== 'bad-response' && receipt.reason !== 'not-pending')) {
throw new Error('Harness returned an invalid interaction response receipt');
}
const reason = receipt.reason;
throw new HarnessInteractionError(
`interaction-${reason}`,
`Harness interaction response was rejected (${reason})`,
);
}
async watchInteractions(sessionId, {
signal,
onInteraction,
onResolved,
onOpen,
ownership,
} = {}) {
if (typeof sessionId !== 'string' || !sessionId) throw new TypeError('sessionId is required');
if (!signal || typeof signal.addEventListener !== 'function') {
throw new TypeError('watchInteractions requires an AbortSignal');
}
if (onInteraction !== undefined && typeof onInteraction !== 'function') {
throw new TypeError('onInteraction must be a function');
}
if (onResolved !== undefined && typeof onResolved !== 'function') {
throw new TypeError('onResolved must be a function');
}
if (onOpen !== undefined && typeof onOpen !== 'function') {
throw new TypeError('onOpen must be a function');
}
while (!signal.aborted) {
try {
await this.#watchInteractionSocket(sessionId, {
signal,
onInteraction,
onResolved,
onOpen,
ownership,
});
} catch (error) {
if (signal.aborted) return;
console.warn(`[${this.#logPrefix}] Harness interaction stream disconnected:`, error.message);
}
if (signal.aborted) return;
try {
await sleep(this.#interactionReconnectDelayMs, signal);
} catch {
if (signal.aborted) return;
throw new Error('Harness interaction reconnect wait failed');
}
}
}
#registerInteractionOwnership(sessionId, ownership) {
const owners = this.#interactionOwnerships.get(sessionId) ?? new Set();
ownership.order = this.#interactionRegistry.nextOrder;
this.#interactionRegistry.nextOrder += 1;
owners.add(ownership);
this.#interactionOwnerships.set(sessionId, owners);
}
#unregisterInteractionOwnership(sessionId, ownership) {
const owners = this.#interactionOwnerships.get(sessionId);
owners?.delete(ownership);
if (owners?.size === 0) this.#interactionOwnerships.delete(sessionId);
for (const [key, claim] of this.#interactionClaims) {
if (claim.ownership === ownership) this.#interactionClaims.delete(key);
}
}
#consumeInteractionOwnerships(sessionId, entries) {
for (const ownership of this.#interactionOwnerships.get(sessionId) ?? []) {
consumeInteractionOwnership(ownership, entries);
}
}
async #refreshInteractionOwnerships(sessionId, signal) {
const history = await this.rpc(
'session.history',
{ sessionId, maxMessages: 50 },
30_000,
{ signal },
);
this.#consumeInteractionOwnerships(sessionId, history.events ?? []);
}
#interactionOwner(sessionId, claimKey, kind) {
const claim = this.#interactionClaims.get(claimKey);
if (claim && this.#interactionOwnerships.get(sessionId)?.has(claim.ownership)) return claim;
const owners = [...(this.#interactionOwnerships.get(sessionId) ?? [])];
const active = owners
.filter((ownership) => ownership.active)
.sort((left, right) => left.order - right.order);
if (active.length > 0) return { ownership: active[0], recovered: false };
// Approval recovery must remain fail-closed until #5 can prove the actor
// and conversation that originally requested the decision.
if (kind !== 'question') return null;
// A newly attached IM conversation may encounter a question left by
// an earlier runtime before its queued prompt starts. Let the oldest such
// ask adopt that replay so the Session can recover instead of deadlocking.
const ownership = owners
.filter((ownership) => !ownership.started && !ownership.completed)
.sort((left, right) => left.order - right.order)[0] ?? null;
return ownership ? { ownership, recovered: true } : null;
}
async ask(sessionId, text, options = {}) {
if (typeof options === 'number') options = { timeoutMs: options };
const timeoutMs = options.timeoutMs ?? 600_000;
const signal = options.signal;
const onUpdate = typeof options.onUpdate === 'function' ? options.onUpdate : null;
const onInteraction = typeof options.onInteraction === 'function'
? options.onInteraction
: undefined;
const onInteractionResolved = typeof options.onInteractionResolved === 'function'
? options.onInteractionResolved
: undefined;
await this.ensureRunning({ signal });
const before = await this.rpc(
'session.history',
{ sessionId, maxMessages: 1 },
30_000,
{ signal },
);
const baselineSeq = Math.max(-1, ...(before.events ?? []).map(({ event }) => event.seq ?? -1));
const promptRpcId = `${this.#rpcIdPrefix}-${randomUUID()}`;
const tracker = new HarnessReplyTracker({ promptRpcId, afterSeq: baselineSeq });
const interactionController = onInteraction || onInteractionResolved
? new AbortController()
: null;
const interactionSignal = interactionController
? (signal
? AbortSignal.any([signal, interactionController.signal])
: interactionController.signal)
: null;
// 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
? {
promptRpcId,
active: false,
started: false,
completed: false,
turn: null,
openTurn: null,
lastSeq: baselineSeq,
reconnect: null,
order: -1,
}
: null;
let interactionTask = null;
if (ownership) this.#registerInteractionOwnership(sessionId, ownership);
try {
if (interactionSignal) {
let markOpen;
const opened = new Promise((resolve) => { markOpen = resolve; });
interactionTask = this.watchInteractions(sessionId, {
signal: interactionSignal,
onInteraction,
onResolved: onInteractionResolved,
onOpen: markOpen,
ownership,
});
void interactionTask.catch(() => undefined);
await Promise.race([
opened,
sleep(30_000, interactionSignal).then(() => {
throw new Error('Harness interaction stream did not open within 30 seconds');
}),
]);
}
await this.rpc('session.prompt', {
sessionId,
mode: 'queue',
content: [{ type: 'text', text }],
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);
}
}
if (!tracker.finished) continue;
if (tracker.answer) return tracker.answer;
throw new Error(
`Harness turn ended without a text reply${tracker.reason ? ` (${JSON.stringify(tracker.reason)})` : ''}`,
);
}
throw new Error(`Harness reply timed out after ${Math.round(timeoutMs / 1_000)} seconds`);
} finally {
interactionController?.abort(new DOMException('Harness turn finished', 'AbortError'));
if (interactionTask) await interactionTask.catch(() => undefined);
if (ownership) this.#unregisterInteractionOwnership(sessionId, ownership);
}
}
#watchInteractionSocket(sessionId, {
signal,
onInteraction,
onResolved,
onOpen,
ownership,
}) {
const url = new URL('/api/events.mux', this.#baseUrl);
url.protocol = url.protocol === 'https:' ? 'wss:' : 'ws:';
return new Promise((resolve, reject) => {
let socket;
try {
socket = this.#createWebSocket(url.toString());
} catch (error) {
reject(error);
return;
}
let opened = false;
let settled = false;
let callbackFailure = null;
let callbackTail = Promise.resolve();
let ownershipReady = ownership === undefined || ownership === null;
const bufferedEnvelopes = [];
const finish = (error) => {
if (settled) return;
settled = true;
socket.removeEventListener('open', handleOpen);
socket.removeEventListener('message', handleMessage);
socket.removeEventListener('close', handleClose);
socket.removeEventListener('error', handleError);
signal.removeEventListener('abort', handleAbort);
if (ownership?.reconnect === close) ownership.reconnect = null;
void callbackTail.then(() => {
const failure = error ?? callbackFailure;
if (failure) reject(failure);
else resolve();
}, reject);
};
const close = () => {
try {
if (socket.readyState === 0 || socket.readyState === 1) socket.close();
} catch {
// Cleanup must still settle the watcher if a WebSocket rejects close while connecting.
}
};
const handleOpen = () => {
opened = true;
if (ownership) ownership.reconnect = close;
try {
onOpen?.();
} catch (error) {
console.warn(`[${this.#logPrefix}] ignored an interaction open callback failure:`, error.message);
}
if (ownership) {
void this.#refreshInteractionOwnerships(sessionId, signal).then(() => {
if (settled) return;
ownershipReady = true;
for (const envelope of bufferedEnvelopes.splice(0)) processEnvelope(envelope);
}).catch((error) => {
callbackFailure ??= error;
close();
finish(error);
});
}
};
const dispatch = (callback, value) => {
if (!callback) return;
callbackTail = callbackTail
.then(() => callback(value))
.catch((error) => {
callbackFailure ??= error;
close();
finish(callbackFailure);
});
};
const processEnvelope = (envelope) => {
const payload = envelope.payload;
if (ownership && payload.type === 'session/event') {
this.#consumeInteractionOwnerships(sessionId, [payload.event]);
return;
}
if (payload.type === 'question/requested' || payload.type === 'approval/requested') {
const kind = payload.type === 'question/requested' ? 'question' : 'approval';
const interactionId = kind === 'question' ? envelope.rpcId : payload.approvalId;
const claimKey = `${kind}:${interactionId}`;
if (ownership) {
const claim = this.#interactionOwner(sessionId, claimKey, kind);
if (claim?.ownership !== ownership) return;
this.#interactionClaims.set(claimKey, claim);
}
dispatch(onInteraction, Object.freeze({
kind,
interactionId,
rpcId: envelope.rpcId,
sessionId,
payload,
recovered: ownership
? this.#interactionClaims.get(claimKey)?.recovered === true
: false,
reconnect: close,
respond: (result, options = {}) => this.respondInteraction(
envelope.rpcId,
result,
{ ...options, signal: options.signal ?? signal },
),
}));
return;
}
if (payload.type === 'question/resolved' || payload.type === 'approval/resolved') {
const kind = payload.type === 'question/resolved' ? 'question' : 'approval';
const interactionId = kind === 'question'
? payload.questionRpcId
: payload.approvalId;
const claimKey = `${kind}:${interactionId}`;
if (ownership) {
const claim = this.#interactionClaims.get(claimKey);
if (claim?.ownership !== ownership) return;
this.#interactionClaims.delete(claimKey);
}
dispatch(onResolved, Object.freeze({
kind,
interactionId,
sessionId,
outcome: payload.outcome,
payload,
}));
}
};
const handleMessage = (event) => {
try {
if (typeof event.data !== 'string') throw new Error('binary WebSocket frame');
const envelope = JSON.parse(event.data);
const payload = envelope?.payload;
if (envelope?.type !== 'server-request'
|| typeof envelope.rpcId !== 'string'
|| !payload || typeof payload !== 'object'
|| envelope.method !== payload.type) {
throw new Error('invalid server-request envelope');
}
if (payload.sessionId !== sessionId) return;
if (!ownershipReady) bufferedEnvelopes.push(envelope);
else processEnvelope(envelope);
} catch (error) {
console.warn(`[${this.#logPrefix}] ignored a malformed Harness interaction frame:`, error.message);
}
};
const handleClose = () => finish(opened ? null : new Error(
'Harness interaction WebSocket closed before opening',
));
const handleError = () => {
finish(new Error(opened
? 'Harness interaction WebSocket failed'
: 'Harness interaction WebSocket failed before opening'));
close();
};
const handleAbort = () => {
close();
finish();
};
socket.addEventListener('open', handleOpen);
socket.addEventListener('message', handleMessage);
socket.addEventListener('close', handleClose, { once: true });
socket.addEventListener('error', handleError, { once: true });
signal.addEventListener('abort', handleAbort, { once: true });
if (signal.aborted) handleAbort();
});
}
stopManagedProcess() {
if (this.#managedProcess?.exitCode === null) this.#managedProcess.kill('SIGTERM');
}
}

View file

@ -0,0 +1,85 @@
function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
export function validHarnessQuestion(question) {
return question && typeof question.id === 'string' && typeof question.question === 'string'
&& (question.header === undefined || typeof question.header === 'string')
&& (question.detail === undefined || typeof question.detail === 'string')
&& (question.multiSelect === undefined || typeof question.multiSelect === 'boolean')
&& (question.options === undefined || (Array.isArray(question.options)
&& question.options.every((option) => (
option && typeof option.label === 'string'
&& (option.description === undefined || typeof option.description === 'string')
))));
}
export function harnessQuestionText(question, index, total, { requiresMention = false } = {}) {
const lines = [];
const progress = total > 1 ? `(${index + 1}/${total})` : '';
lines.push(`DeepSeek Harness 需要你补充信息${progress}:`);
if (nonEmptyString(question.header)) lines.push('', question.header.trim());
lines.push('', nonEmptyString(question.question) ?? '请输入你的回答。');
if (nonEmptyString(question.detail)) lines.push('', question.detail.trim());
const options = Array.isArray(question.options) ? question.options : [];
if (options.length > 0) {
lines.push('');
options.forEach((option, optionIndex) => {
const label = typeof option?.label === 'string' ? option.label : '';
const description = nonEmptyString(option?.description);
lines.push(`${optionIndex + 1}. ${label}${description ? ` — ${description}` : ''}`);
});
lines.push('', question.multiSelect === true
? '请回复选项序号或文字;多选用逗号分隔,也可补充其他内容。'
: '请回复一个选项序号或文字,也可直接输入其他答案。');
} else {
lines.push('', '请直接回复你的答案。');
}
if (requiresMention) lines.push('', '群聊中请 @机器人 后发送答案。');
return lines.join('\n');
}
function optionLabel(token, options) {
const normalized = token.trim();
if (!normalized) return null;
if (/^\d+$/.test(normalized)) {
const option = options[Number(normalized) - 1];
return typeof option?.label === 'string' ? option.label : null;
}
const exact = options.find((option) => option?.label === normalized);
return typeof exact?.label === 'string' ? exact.label : null;
}
export function harnessAnswerForQuestion(question, text) {
const options = Array.isArray(question.options) ? question.options : [];
if (options.length === 0) {
return { id: question.id, selected: [], custom: text };
}
const wholeLabel = optionLabel(text, options);
if (question.multiSelect !== true) {
return wholeLabel
? { id: question.id, selected: [wholeLabel] }
: { id: question.id, selected: [], custom: text };
}
if (wholeLabel) return { id: question.id, selected: [wholeLabel] };
const selected = [];
const custom = [];
for (const token of text.split(/[,,、;;\n]+/)) {
const value = token.trim();
if (!value) continue;
const label = optionLabel(value, options);
if (label) {
if (!selected.includes(label)) selected.push(label);
} else {
custom.push(value);
}
}
return {
id: question.id,
selected,
...(custom.length > 0 ? { custom: custom.join('、') } : {}),
};
}