mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-09 20:13:21 +08:00
Merge pull request #24 from NIU-001-LIU/feat/feishu-watch-push
feat(feishu): persistent session watches and completion pushes
This commit is contained in:
commit
6cafba863e
8 changed files with 1243 additions and 159 deletions
263
lib/index.js
263
lib/index.js
File diff suppressed because one or more lines are too long
|
|
@ -32,11 +32,14 @@ import { runWorkspaceCommand, resolveSessionListWorkspace, workspacePathSnapshot
|
|||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
import {
|
||||
MENU_PAGE_SIZE,
|
||||
completionCard,
|
||||
menuCard,
|
||||
menuHelpText,
|
||||
sessionListCard,
|
||||
watchListCard,
|
||||
workspaceListCard,
|
||||
} from './feishu-cards.mjs';
|
||||
import { MAX_WATCHES_PER_KEY } from './state-store.mjs';
|
||||
|
||||
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
|
||||
const RESOLVED_REPLY_TTL_MS = 30 * 60_000;
|
||||
|
|
@ -44,6 +47,9 @@ const RESOLVED_REPLY_TTL_MS = 30 * 60_000;
|
|||
const MENU_COMMAND = /^\/m(?:enu)?$/i;
|
||||
const REPAIR_COMMAND_PREFIX = /^\/repair(?:\s|$)/i;
|
||||
const REPAIR_COMMAND = /^\/repair(?:\s+(qr|status|cancel|verify))?\s*$/i;
|
||||
const WATCH_COMMAND = /^\/watch(?:\s+([^\s]+))?$/i;
|
||||
const UNWATCH_COMMAND = /^\/unwatch(?:\s+([^\s]+))?$/i;
|
||||
const WATCHLIST_COMMAND = /^\/watchlist$/i;
|
||||
const SESSION_LIST_PREFIX = /^\/sessionlist(?:\s|$)/i;
|
||||
const WORKSPACE_LIST_COMMAND = /^\/workspacelist$/i;
|
||||
const NUMBER_REPLY = /^\d{1,2}$/;
|
||||
|
|
@ -83,9 +89,15 @@ const HELP_TEXT = [
|
|||
'/status 检查连接状态',
|
||||
'/repair 修复卡片按钮回调',
|
||||
'/m(或 /menu) 打开交互卡片菜单',
|
||||
'/watch [Session ID 或序号] 关注会话,任务完成自动推送',
|
||||
'/unwatch [Session ID 或序号] 取消关注',
|
||||
'/watchlist 查看关注列表',
|
||||
'/archived on|off 会话列表是否包含归档会话',
|
||||
'/help 显示本帮助',
|
||||
].join('\n');
|
||||
|
||||
const ARCHIVED_COMMAND = /^\/archived(?:\s+(on|off))?$/i;
|
||||
|
||||
/** Safe user-facing text for bind/workspace failures (no raw messages). */
|
||||
function safeErrorText(error) {
|
||||
switch (error?.code) {
|
||||
|
|
@ -106,6 +118,13 @@ function nonEmptyString(value) {
|
|||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
function orderedHistoryEvents(history) {
|
||||
return (Array.isArray(history?.events) ? history.events : [])
|
||||
.map((entry) => entry?.event ?? entry)
|
||||
.filter((entry) => entry && typeof entry === 'object' && Number.isFinite(entry.seq))
|
||||
.sort((left, right) => left.seq - right.seq);
|
||||
}
|
||||
|
||||
function senderOpenId(event) {
|
||||
return nonEmptyString(event?.sender?.sender_id?.open_id)
|
||||
?? nonEmptyString(event?.sender?.sender_id?.user_id);
|
||||
|
|
@ -240,6 +259,12 @@ export class FeishuHarnessBridge {
|
|||
#menus = new Map();
|
||||
/** Interactive-card message id → route context for button callbacks. */
|
||||
#cardKeys = new Map();
|
||||
/** The global event-mux watcher (one per bridge). */
|
||||
#eventWatcher = null;
|
||||
/** Serializes live completions and reconnect compensation. */
|
||||
#eventTail = Promise.resolve();
|
||||
/** Earliest completion that still needs delivery for each watch. */
|
||||
#failedWatchSeqs = new Map();
|
||||
|
||||
constructor({
|
||||
client,
|
||||
|
|
@ -295,6 +320,11 @@ export class FeishuHarnessBridge {
|
|||
this.#approvals = new HarnessApprovalQueue({ label: 'Feishu', logger });
|
||||
this.#signal = signal;
|
||||
ensureStatus(this.#status);
|
||||
// Persisted watches must resume at runtime start, not on the first
|
||||
// message. Older hosts without the mux watcher simply skip this.
|
||||
if (typeof this.#harness?.watchHarnessEvents === 'function') {
|
||||
queueMicrotask(() => this.#ensureEventWatcher());
|
||||
}
|
||||
}
|
||||
|
||||
accept(event) {
|
||||
|
|
@ -519,6 +549,7 @@ export class FeishuHarnessBridge {
|
|||
)),
|
||||
...this.#interactionTasks,
|
||||
...this.#commandTasks,
|
||||
this.#eventTail,
|
||||
]);
|
||||
}
|
||||
|
||||
|
|
@ -604,6 +635,36 @@ export class FeishuHarnessBridge {
|
|||
await this.#showWorkspaces({ chatId: event.message.chat_id, key });
|
||||
return;
|
||||
}
|
||||
if (WATCH_COMMAND.test(commandText)) {
|
||||
const target = (WATCH_COMMAND.exec(commandText)?.[1] ?? '').trim() || null;
|
||||
await this.#runWatch(key, event.message.chat_id, target);
|
||||
return;
|
||||
}
|
||||
if (UNWATCH_COMMAND.test(commandText)) {
|
||||
const target = (UNWATCH_COMMAND.exec(commandText)?.[1] ?? '').trim() || null;
|
||||
await this.#runUnwatch(key, event.message.chat_id, target);
|
||||
return;
|
||||
}
|
||||
if (WATCHLIST_COMMAND.test(commandText)) {
|
||||
await this.#showWatchList(key, event.message.chat_id);
|
||||
return;
|
||||
}
|
||||
if (ARCHIVED_COMMAND.test(commandText)) {
|
||||
const match = ARCHIVED_COMMAND.exec(commandText);
|
||||
const value = match[1]?.toLowerCase();
|
||||
if (value !== 'on' && value !== 'off') {
|
||||
await this.#send(event.message.chat_id, '用法:/archived on(包含归档会话)或 /archived off(隐藏归档会话)');
|
||||
return;
|
||||
}
|
||||
if (typeof this.#state?.setIncludeArchivedSessions === 'function') {
|
||||
await this.#state.setIncludeArchivedSessions(value === 'on');
|
||||
}
|
||||
await this.#send(
|
||||
event.message.chat_id,
|
||||
value === 'on' ? '已开启:会话列表包含归档会话。' : '已关闭:会话列表隐藏归档会话。',
|
||||
);
|
||||
return;
|
||||
}
|
||||
if (NUMBER_REPLY.test(commandText)) {
|
||||
const menu = this.#takeMenu(key);
|
||||
if (menu) {
|
||||
|
|
@ -991,7 +1052,15 @@ export class FeishuHarnessBridge {
|
|||
if (!action) return Promise.resolve();
|
||||
const messageId = nonEmptyString(event?.context?.open_message_id);
|
||||
const entry = messageId ? this.#cardKeys.get(messageId) : null;
|
||||
if (!entry) return Promise.resolve();
|
||||
if (!entry) {
|
||||
// The card predates this process (the in-memory mapping resets on
|
||||
// restart) or never came from us: nudge instead of staying silent.
|
||||
const chatId = nonEmptyString(event?.context?.open_chat_id);
|
||||
if (chatId) {
|
||||
this.#send(chatId, '这个菜单已过期,请回复 /m 重新打开。').catch(() => undefined);
|
||||
}
|
||||
return Promise.resolve();
|
||||
}
|
||||
// The promise is returned so tests (and future callers) can await the
|
||||
// action; the runtime dispatcher ignores it.
|
||||
return this.#handleCardAction(action, entry).catch((error) => {
|
||||
|
|
@ -1009,6 +1078,10 @@ export class FeishuHarnessBridge {
|
|||
await this.#showWorkspaces({ chatId, key });
|
||||
return;
|
||||
}
|
||||
if (action === 'watchlist') {
|
||||
await this.#showWatchList(key, chatId);
|
||||
return;
|
||||
}
|
||||
if (action === 'new') {
|
||||
await this.#state.clearSession(key);
|
||||
await this.#send(chatId, '已开启全新 Harness 会话。');
|
||||
|
|
@ -1029,6 +1102,14 @@ export class FeishuHarnessBridge {
|
|||
}
|
||||
if (action.startsWith('workspace:')) {
|
||||
await this.#switchWorkspace(key, chatId, action.slice('workspace:'.length));
|
||||
return;
|
||||
}
|
||||
if (action.startsWith('unwatch:')) {
|
||||
await this.#runUnwatch(key, chatId, action.slice('unwatch:'.length));
|
||||
return;
|
||||
}
|
||||
if (action.startsWith('watch:')) {
|
||||
await this.#runWatch(key, chatId, action.slice('watch:'.length));
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1053,7 +1134,7 @@ export class FeishuHarnessBridge {
|
|||
|
||||
async #handleMenuPick(menu, number, { chatId, key, event }) {
|
||||
if (menu.kind === 'menu') {
|
||||
const action = ['sessions', 'workspaces', 'new', 'status', 'help', 'repair'][number - 1];
|
||||
const action = ['sessions', 'workspaces', 'new', 'status', 'help', 'repair', 'watchlist'][number - 1];
|
||||
if (!action) {
|
||||
await this.#send(chatId, '菜单没有这个编号,回复 /m 重新打开。');
|
||||
return;
|
||||
|
|
@ -1071,6 +1152,7 @@ export class FeishuHarnessBridge {
|
|||
await this.#send(chatId, `本页只有 ${menu.sessions.length} 个会话,回复 /sessionlist 重新查看。`);
|
||||
return;
|
||||
}
|
||||
// The number label sits on the session (bind) button of the row.
|
||||
await this.#handleCardAction(`use:${session.sessionId}`, { chatId, key });
|
||||
return;
|
||||
}
|
||||
|
|
@ -1081,7 +1163,24 @@ export class FeishuHarnessBridge {
|
|||
return;
|
||||
}
|
||||
await this.#handleCardAction(`workspace:${workspace}`, { chatId, key });
|
||||
return;
|
||||
}
|
||||
if (menu.kind === 'watches') {
|
||||
const entry = menu.entries[number - 1];
|
||||
if (!entry?.sessionId) {
|
||||
await this.#send(chatId, `关注列表只有 ${menu.entries.length} 个会话。`);
|
||||
return;
|
||||
}
|
||||
await this.#handleCardAction(`unwatch:${entry.sessionId}`, { chatId, key });
|
||||
}
|
||||
}
|
||||
|
||||
/** The sessions visible under the bot's archived policy. */
|
||||
#visibleSessions(sessions) {
|
||||
if (this.#state?.includesArchivedSessions?.() === false) {
|
||||
return sessions.filter((session) => session.archived !== true);
|
||||
}
|
||||
return sessions;
|
||||
}
|
||||
|
||||
async #showSessions({ chatId, key }, selector, page = 0) {
|
||||
|
|
@ -1092,7 +1191,7 @@ export class FeishuHarnessBridge {
|
|||
return;
|
||||
}
|
||||
const listed = await this.#harness.listWorkspaceSessions(resolved.workspace);
|
||||
const sessions = Array.isArray(listed?.sessions) ? listed.sessions : [];
|
||||
const sessions = this.#visibleSessions(Array.isArray(listed?.sessions) ? listed.sessions : []);
|
||||
const workspace = listed?.workspace ?? resolved.workspace;
|
||||
if (sessions.length === 0) {
|
||||
await this.#send(chatId, `工作区:${workspace}\n该工作区暂无会话。`);
|
||||
|
|
@ -1100,16 +1199,24 @@ export class FeishuHarnessBridge {
|
|||
}
|
||||
const pageCount = Math.ceil(sessions.length / MENU_PAGE_SIZE);
|
||||
const safePage = Number.isSafeInteger(page) && page > 0 ? Math.min(page, pageCount - 1) : 0;
|
||||
const watchedSet = new Set(
|
||||
(this.#state.watchEntries?.(key) ?? []).map((entry) => entry.sessionId),
|
||||
);
|
||||
const pageSlice = sessions.slice(safePage * MENU_PAGE_SIZE, (safePage + 1) * MENU_PAGE_SIZE);
|
||||
this.#rememberMenu(key, {
|
||||
kind: 'sessions',
|
||||
sessions: sessions.slice(safePage * MENU_PAGE_SIZE, (safePage + 1) * MENU_PAGE_SIZE),
|
||||
});
|
||||
await this.#sendCard(chatId, sessionListCard(workspace, sessions, safePage, sessions.length), {
|
||||
key,
|
||||
// Keep the canonical selector result for later page callbacks. The
|
||||
// list response's workspace is display data and is not authoritative.
|
||||
sessionWorkspace: resolved.workspace,
|
||||
sessions: pageSlice.map((session) => ({ ...session, watched: watchedSet.has(session.sessionId) })),
|
||||
});
|
||||
await this.#sendCard(
|
||||
chatId,
|
||||
sessionListCard(workspace, sessions, safePage, sessions.length, watchedSet),
|
||||
{
|
||||
key,
|
||||
// Keep the canonical selector result for later page callbacks. The
|
||||
// list response's workspace is display data and is not authoritative.
|
||||
sessionWorkspace: resolved.workspace,
|
||||
},
|
||||
);
|
||||
} catch (error) {
|
||||
this.#logger.warn?.('[dsh-feishu] session list failed:', error.message);
|
||||
await this.#send(chatId, '暂时无法获取会话列表,请稍后重试。');
|
||||
|
|
@ -1171,6 +1278,256 @@ export class FeishuHarnessBridge {
|
|||
return messageId;
|
||||
}
|
||||
|
||||
// ── Watches: read-only session tracking + completion pushes ─────────────
|
||||
|
||||
#ensureEventWatcher() {
|
||||
if (this.#eventWatcher) return;
|
||||
if (typeof this.#harness?.watchHarnessEvents !== 'function') return;
|
||||
if (this.#signal?.aborted) return;
|
||||
const signal = this.#signal ?? new AbortController().signal;
|
||||
try {
|
||||
this.#eventWatcher = this.#harness.watchHarnessEvents({
|
||||
signal,
|
||||
onSessionEvent: (payload) => this.#onHarnessEvent(payload),
|
||||
onReconnect: () => {
|
||||
void this.#queueEventTask(() => this.#compensateMissedEvents());
|
||||
},
|
||||
});
|
||||
Promise.resolve(this.#eventWatcher).catch((error) => {
|
||||
if (!signal.aborted) {
|
||||
this.#logger.warn?.('[dsh-feishu] event watcher stopped:', error.message);
|
||||
}
|
||||
});
|
||||
} catch (error) {
|
||||
this.#eventWatcher = null;
|
||||
this.#logger.warn?.('[dsh-feishu] event watcher failed to start:', error.message);
|
||||
}
|
||||
}
|
||||
|
||||
#queueEventTask(task) {
|
||||
const next = this.#eventTail.then(task, task).catch((error) => {
|
||||
if (!this.#signal?.aborted) {
|
||||
this.#logger.warn?.('[dsh-feishu] completion event failed:', error.message);
|
||||
}
|
||||
});
|
||||
this.#eventTail = next;
|
||||
return next;
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a /watch target READ-ONLY: a session id is validated against
|
||||
* the registered workspaces' listings, an index against the current
|
||||
* workspace. Nothing is bound and no workspace is switched.
|
||||
*/
|
||||
async #resolveWatchTarget(target) {
|
||||
if (typeof target !== 'string' || target === '') {
|
||||
return { error: '用法:/watch <Session ID 或当前工作区序号>' };
|
||||
}
|
||||
const numeric = /^\d{1,4}$/.test(target) ? Number(target) : null;
|
||||
const currentPath = typeof this.#harness?.currentWorkspace === 'function'
|
||||
? this.#harness.currentWorkspace()
|
||||
: null;
|
||||
const listSessions = async (workspace) => {
|
||||
const listed = await this.#harness.listWorkspaceSessions(workspace);
|
||||
return Array.isArray(listed?.sessions) ? listed.sessions : [];
|
||||
};
|
||||
if (numeric !== null) {
|
||||
if (!currentPath) return { error: '当前机器人没有可用的工作区,无法按序号解析会话。' };
|
||||
const sessions = this.#visibleSessions(await listSessions(currentPath));
|
||||
const session = sessions[numeric - 1];
|
||||
if (!session?.sessionId) {
|
||||
return { error: `当前工作区只有 ${sessions.length} 个会话。` };
|
||||
}
|
||||
return { sessionId: session.sessionId, title: session.title ?? '暂无标题' };
|
||||
}
|
||||
const extraPaths = typeof this.#harness?.listWorkspaces === 'function'
|
||||
? (await this.#harness.listWorkspaces()).filter((path) => path !== currentPath)
|
||||
: [];
|
||||
const paths = [currentPath, ...extraPaths].filter(Boolean);
|
||||
for (const workspace of paths) {
|
||||
const sessions = await listSessions(workspace);
|
||||
const session = sessions.find((candidate) => candidate.sessionId === target);
|
||||
if (session) return { sessionId: target, title: session.title ?? '暂无标题' };
|
||||
}
|
||||
return { error: '没有找到这个会话,请用 /sessionlist 查看可用会话。' };
|
||||
}
|
||||
|
||||
async #latestSessionSeq(sessionId) {
|
||||
if (typeof this.#harness?.rpc !== 'function') return null;
|
||||
const history = await this.#harness.rpc(
|
||||
'session.history',
|
||||
{ sessionId, maxMessages: 20 },
|
||||
30_000,
|
||||
{ signal: this.#signal },
|
||||
);
|
||||
return orderedHistoryEvents(history).at(-1)?.seq ?? -1;
|
||||
}
|
||||
|
||||
async #runWatch(key, chatId, target) {
|
||||
this.#ensureEventWatcher();
|
||||
if (typeof this.#state?.setWatch !== 'function') {
|
||||
await this.#send(chatId, '当前状态存储不支持关注。');
|
||||
return;
|
||||
}
|
||||
let resolved;
|
||||
try {
|
||||
resolved = await this.#resolveWatchTarget(target);
|
||||
} catch (error) {
|
||||
await this.#send(chatId, `无法解析会话:${safeErrorText(error)}`);
|
||||
return;
|
||||
}
|
||||
if (resolved.error) {
|
||||
await this.#send(chatId, resolved.error);
|
||||
return;
|
||||
}
|
||||
const existing = this.#state.watchEntries?.(key) ?? [];
|
||||
const existingEntry = existing.find((entry) => entry.sessionId === resolved.sessionId);
|
||||
if (!existingEntry && existing.length >= MAX_WATCHES_PER_KEY) {
|
||||
await this.#send(chatId, `每个聊天最多关注 ${MAX_WATCHES_PER_KEY} 个会话。`);
|
||||
return;
|
||||
}
|
||||
try {
|
||||
const lastSeq = typeof existingEntry?.lastSeq === 'number'
|
||||
? existingEntry.lastSeq
|
||||
: await this.#latestSessionSeq(resolved.sessionId);
|
||||
await this.#state.setWatch(key, {
|
||||
sessionId: resolved.sessionId,
|
||||
title: resolved.title,
|
||||
chatId,
|
||||
lastSeq,
|
||||
});
|
||||
await this.#send(chatId, `已关注会话「${String(resolved.title).replace(/\s+/gu, ' ')}」,任务完成会推送结果。`);
|
||||
await this.#queueEventTask(() => this.#compensateSession(resolved.sessionId));
|
||||
} catch (error) {
|
||||
await this.#send(chatId, `关注失败:${safeErrorText(error)}`);
|
||||
}
|
||||
}
|
||||
|
||||
async #runUnwatch(key, chatId, target) {
|
||||
if (typeof this.#state?.removeWatch !== 'function') return;
|
||||
const entries = this.#state.watchEntries?.(key) ?? [];
|
||||
const entry = typeof target === 'string' && /^\d{1,4}$/.test(target)
|
||||
? entries[Number(target) - 1]
|
||||
: entries.find((candidate) => candidate.sessionId === target);
|
||||
if (!entry) {
|
||||
await this.#send(chatId, '关注列表里没有这个会话,回复 /watchlist 查看。');
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await this.#state.removeWatch(key, entry.sessionId);
|
||||
this.#failedWatchSeqs.delete(`${key}\0${entry.sessionId}`);
|
||||
await this.#send(chatId, `已取消关注「${String(entry.title ?? '').replace(/\s+/gu, ' ')}」。`);
|
||||
} catch (error) {
|
||||
await this.#send(chatId, `取消失败:${safeErrorText(error)}`);
|
||||
}
|
||||
}
|
||||
|
||||
async #showWatchList(key, chatId) {
|
||||
const entries = this.#state.watchEntries?.(key) ?? [];
|
||||
this.#rememberMenu(key, { kind: 'watches', entries });
|
||||
await this.#sendCard(chatId, watchListCard(entries), { key });
|
||||
}
|
||||
|
||||
/** Queue live turn completions behind any reconnect compensation. */
|
||||
#onHarnessEvent({ sessionId, event }) {
|
||||
if (this.#signal?.aborted
|
||||
|| !sessionId
|
||||
|| !event
|
||||
|| typeof event !== 'object'
|
||||
|| event.type !== 'turn/end'
|
||||
|| !Number.isFinite(event.seq)) return;
|
||||
void this.#queueEventTask(async () => {
|
||||
const hasFailedDelivery = (this.#state.keysWatching?.(sessionId) ?? [])
|
||||
.some((key) => this.#failedWatchSeqs.has(`${key}\0${sessionId}`));
|
||||
if (hasFailedDelivery) await this.#compensateSession(sessionId);
|
||||
await this.#deliverCompletion(sessionId, event);
|
||||
});
|
||||
}
|
||||
|
||||
async #deliverCompletion(sessionId, event) {
|
||||
if (this.#signal?.aborted || typeof this.#state?.keysWatching !== 'function') return;
|
||||
const reason = event?.data?.reason?.kind ?? event?.data?.reason ?? null;
|
||||
for (const key of this.#state.keysWatching(sessionId)) {
|
||||
if (this.#signal?.aborted) return;
|
||||
const entry = this.#state.watchEntry?.(key, sessionId);
|
||||
const deliveryKey = `${key}\0${sessionId}`;
|
||||
let failedSeq = this.#failedWatchSeqs.get(deliveryKey);
|
||||
if (typeof failedSeq === 'number'
|
||||
&& typeof entry?.lastSeq === 'number'
|
||||
&& entry.lastSeq >= failedSeq) {
|
||||
this.#failedWatchSeqs.delete(deliveryKey);
|
||||
failedSeq = undefined;
|
||||
}
|
||||
if (!entry?.chatId
|
||||
|| (typeof entry.lastSeq === 'number' && entry.lastSeq >= event.seq)
|
||||
|| (typeof failedSeq === 'number' && event.seq > failedSeq)) continue;
|
||||
try {
|
||||
await this.#sendCard(
|
||||
entry.chatId,
|
||||
completionCard(sessionId, entry.title, reason),
|
||||
{ key },
|
||||
);
|
||||
const current = this.#state.watchEntry?.(key, sessionId);
|
||||
if (!current
|
||||
|| current.chatId !== entry.chatId
|
||||
|| (typeof current.lastSeq === 'number' && current.lastSeq >= event.seq)) continue;
|
||||
await this.#state.setWatch(key, { ...current, lastSeq: event.seq });
|
||||
if (failedSeq === event.seq) this.#failedWatchSeqs.delete(deliveryKey);
|
||||
} catch (error) {
|
||||
this.#failedWatchSeqs.set(
|
||||
deliveryKey,
|
||||
typeof failedSeq === 'number' ? Math.min(failedSeq, event.seq) : event.seq,
|
||||
);
|
||||
this.#logger.warn?.('[dsh-feishu] completion push failed:', error.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async #compensateSession(sessionId) {
|
||||
if (this.#signal?.aborted || typeof this.#harness?.rpc !== 'function') return;
|
||||
try {
|
||||
const history = await this.#harness.rpc(
|
||||
'session.history',
|
||||
{ sessionId, maxMessages: 20 },
|
||||
30_000,
|
||||
{ signal: this.#signal },
|
||||
);
|
||||
const events = orderedHistoryEvents(history);
|
||||
const latestSeq = events.at(-1)?.seq ?? -1;
|
||||
const keys = typeof this.#state?.keysWatching === 'function'
|
||||
? this.#state.keysWatching(sessionId)
|
||||
: [];
|
||||
|
||||
// Watches created by older versions have no baseline. Establish one
|
||||
// without replaying completions that predate the watch.
|
||||
for (const key of keys) {
|
||||
const entry = this.#state.watchEntry?.(key, sessionId);
|
||||
if (entry && typeof entry.lastSeq !== 'number') {
|
||||
await this.#state.setWatch(key, { ...entry, lastSeq: latestSeq });
|
||||
}
|
||||
}
|
||||
|
||||
for (const event of events) {
|
||||
if (event.type === 'turn/end') await this.#deliverCompletion(sessionId, event);
|
||||
}
|
||||
} catch (error) {
|
||||
if (!this.#signal?.aborted) {
|
||||
this.#logger.warn?.(`[dsh-feishu] watch compensation failed for ${sessionId}:`, error.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Replay recent turn completions missed while the mux was disconnected. */
|
||||
async #compensateMissedEvents() {
|
||||
const sessionIds = typeof this.#state?.watchedSessionIds === 'function'
|
||||
? this.#state.watchedSessionIds()
|
||||
: [];
|
||||
for (const sessionId of sessionIds) {
|
||||
if (this.#signal?.aborted) return;
|
||||
await this.#compensateSession(sessionId);
|
||||
}
|
||||
}
|
||||
|
||||
#interactionAskOptions(event, key) {
|
||||
return {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
|
|
|
|||
|
|
@ -4,10 +4,14 @@
|
|||
* `im.message.create` API expects as `content` for `msg_type: interactive`
|
||||
* (card schema 2.0; callback buttons live inside a column_set/column layout).
|
||||
*
|
||||
* Buttons carry a small `{ action }` callback behavior that
|
||||
* `card.action.trigger` events echo back (when the app subscribes that
|
||||
* callback); every button also carries a numeric label so the number-reply
|
||||
* fallback stays usable without button callbacks.
|
||||
* The session list lays each row out as a `column_set` (fixed-width ⭐
|
||||
* watch-toggle column + weighted session-button column), which is how V2
|
||||
* expresses a row of buttons.
|
||||
*
|
||||
* Buttons carry a small `{ action }` value object that `card.action.trigger`
|
||||
* events echo back (when the app subscribes that event); every numbered
|
||||
* button also has a numeric label so the number-reply fallback stays usable
|
||||
* without button callbacks.
|
||||
*/
|
||||
|
||||
export const MENU_PAGE_SIZE = 10;
|
||||
|
|
@ -39,6 +43,17 @@ function button(content, actionValue) {
|
|||
};
|
||||
}
|
||||
|
||||
/** The raw button element (without the full-width column_set wrapper). */
|
||||
function buttonElement(content, actionValue) {
|
||||
return {
|
||||
tag: 'button',
|
||||
text: plainText(content),
|
||||
type: 'default',
|
||||
width: 'fill',
|
||||
behaviors: [{ type: 'callback', value: { action: actionValue } }],
|
||||
};
|
||||
}
|
||||
|
||||
function safeTitle(value) {
|
||||
const title = String(value ?? '').replace(/[\p{Cc}\p{Cf}\p{Zl}\p{Zp}]+/gu, ' ').replace(/\s+/gu, ' ').trim();
|
||||
return title || '暂无标题';
|
||||
|
|
@ -65,6 +80,7 @@ export function menuCard() {
|
|||
// have card.action.trigger yet, so rendering it as a callback button would
|
||||
// send the user straight back to Feishu's broken callback setup popup.
|
||||
{ tag: 'div', text: markdown('**6 · 修复卡片按钮**(请直接回复数字 **6**)') },
|
||||
button('7 · 关注列表', 'watchlist'),
|
||||
]);
|
||||
}
|
||||
|
||||
|
|
@ -101,23 +117,43 @@ export function cardActionProbeCard(nonce) {
|
|||
}
|
||||
|
||||
/**
|
||||
* One page of the workspace's sessions. Each row is a bind button; the
|
||||
* number label equals the reply-number for the same action (fallback).
|
||||
* One page of the workspace's sessions. Each row is a `column_set` pair:
|
||||
* the fixed-width ⭐ watch toggle (`⭐关注` / `⭐取关` for already-watched
|
||||
* sessions) followed by the session button that carries the page-local
|
||||
* number label (reply-number fallback = bind). Archived sessions are marked
|
||||
* in the label. `watchedSessionIds` is a Set-like of ids this conversation
|
||||
* already watches.
|
||||
*/
|
||||
export function sessionListCard(workspace, sessions, page, total) {
|
||||
export function sessionListCard(workspace, sessions, page, total, watchedSessionIds = new Set()) {
|
||||
const start = page * MENU_PAGE_SIZE;
|
||||
const slice = sessions.slice(start, start + MENU_PAGE_SIZE);
|
||||
const pageCount = Math.max(1, Math.ceil(total / MENU_PAGE_SIZE));
|
||||
const watched = (id) => typeof watchedSessionIds?.has === 'function' && watchedSessionIds.has(id);
|
||||
/** One row: fixed 90px watch toggle + the session button filling the rest. */
|
||||
const row = (watchButton, sessionButton) => ({
|
||||
tag: 'column_set',
|
||||
flex_mode: 'none',
|
||||
horizontal_spacing: 'default',
|
||||
columns: [
|
||||
{ tag: 'column', width: '90px', vertical_align: 'center', elements: [watchButton] },
|
||||
{ tag: 'column', width: 'weighted', weight: 1, vertical_align: 'center', elements: [sessionButton] },
|
||||
],
|
||||
});
|
||||
const elements = [
|
||||
{ tag: 'div', text: markdown(`**工作区**:\`${workspace}\`\n共 **${total}** 个会话${total > MENU_PAGE_SIZE ? `(第 ${page + 1}/${pageCount} 页)` : ''}`) },
|
||||
...slice.map((session, offset) => button(
|
||||
`${offset + 1}. ${safeTitle(session.title)}`,
|
||||
`use:${session.sessionId}`,
|
||||
)),
|
||||
...slice.map((session, offset) => {
|
||||
// Page-local numbering: number replies resolve against this page.
|
||||
const label = `${offset + 1}. ${safeTitle(session.title)}${session.archived === true ? '(已归档)' : ''}`;
|
||||
const watching = watched(session.sessionId);
|
||||
return row(
|
||||
buttonElement(watching ? '⭐取关' : '⭐关注', watching ? `unwatch:${session.sessionId}` : `watch:${session.sessionId}`),
|
||||
buttonElement(label, `use:${session.sessionId}`),
|
||||
);
|
||||
}),
|
||||
];
|
||||
if (page > 0) elements.push(button('◀ 上一页', `sessions:${page - 1}`));
|
||||
if (page + 1 < pageCount) elements.push(button('下一页 ▶', `sessions:${page + 1}`));
|
||||
elements.push({ tag: 'div', text: markdown('回复数字(1~N)同样可以绑定本页会话。') });
|
||||
elements.push({ tag: 'div', text: markdown('回复数字(1~N)绑定本页会话。') });
|
||||
return cardWith('📂 会话列表', elements);
|
||||
}
|
||||
|
||||
|
|
@ -146,10 +182,49 @@ export function menuHelpText() {
|
|||
'4 · /status 连接状态',
|
||||
'5 · /help 本帮助',
|
||||
'6 · /repair 修复卡片按钮(请回复数字 6)',
|
||||
'7 · /watchlist 关注列表',
|
||||
'',
|
||||
'直接发送文字/图片即继续当前会话。',
|
||||
'/session ID 或序号 绑定已有会话',
|
||||
'/watch ID 或序号 关注会话(完成后推送)',
|
||||
'/compact 压缩上下文',
|
||||
'/workspace 绝对路径 切换工作区',
|
||||
].join('\n');
|
||||
}
|
||||
|
||||
/** The watch list for one conversation (unwatch buttons + reply fallback). */
|
||||
export function watchListCard(entries) {
|
||||
const elements = entries.length === 0
|
||||
? [{ tag: 'div', text: markdown('当前没有关注的会话。\n`/watch <ID|序号>` 关注后,任务完成会自动推送。') }]
|
||||
: [
|
||||
{ tag: 'div', text: markdown('任务完成会自动推送,回复数字或点按钮取消关注:') },
|
||||
...entries.map((entry, index) => button(
|
||||
`${index + 1}. ${safeTitle(entry.title)}`,
|
||||
`unwatch:${entry.sessionId}`,
|
||||
)),
|
||||
];
|
||||
return cardWith('👁 关注列表', elements);
|
||||
}
|
||||
|
||||
/**
|
||||
* The completion push card. `title` is the session title, `reason` the
|
||||
* turn-end kind (completed / stopped / aborted).
|
||||
*/
|
||||
export function completionCard(sessionId, title, reason) {
|
||||
const reasonText = reason === 'completed'
|
||||
? '已完成'
|
||||
: reason === 'stopped'
|
||||
? '已停止'
|
||||
: reason === 'aborted'
|
||||
? '已中止'
|
||||
: reason === 'cancelled'
|
||||
? '已取消'
|
||||
: '已结束';
|
||||
return cardWith('✅ 任务完成', [
|
||||
{ tag: 'div', text: markdown(`**${safeTitle(title)}**\n\`${sessionId}\``) },
|
||||
{ tag: 'div', text: markdown(`**状态**:${reasonText}`) },
|
||||
button('打开会话列表', 'sessions'),
|
||||
button('工作区', 'workspaces'),
|
||||
{ tag: 'div', text: markdown('绑定该会话后可继续追问,输入文字即可。') },
|
||||
]);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,7 +1,24 @@
|
|||
import { mkdir, readFile, rename, writeFile } from 'node:fs/promises';
|
||||
import { dirname } from 'node:path';
|
||||
|
||||
const EMPTY_STATE = Object.freeze({ version: 1, sessions: {}, seenMessageIds: [] });
|
||||
const EMPTY_STATE = Object.freeze({
|
||||
version: 1,
|
||||
sessions: {},
|
||||
seenMessageIds: [],
|
||||
watches: {},
|
||||
includeArchivedSessions: false,
|
||||
});
|
||||
|
||||
/** One conversation key may watch at most this many sessions. */
|
||||
export const MAX_WATCHES_PER_KEY = 20;
|
||||
|
||||
/** A persisted watch entry: the watched session plus its delivery target. */
|
||||
function validWatchEntry(value) {
|
||||
return value
|
||||
&& typeof value === 'object'
|
||||
&& typeof value.sessionId === 'string' && value.sessionId.length > 0
|
||||
&& typeof value.chatId === 'string' && value.chatId.length > 0;
|
||||
}
|
||||
|
||||
export class StateStore {
|
||||
#path;
|
||||
|
|
@ -19,6 +36,10 @@ export class StateStore {
|
|||
version: 1,
|
||||
sessions: parsed.sessions && typeof parsed.sessions === 'object' ? parsed.sessions : {},
|
||||
seenMessageIds: Array.isArray(parsed.seenMessageIds) ? parsed.seenMessageIds.slice(-1000) : [],
|
||||
watches: parsed.watches && typeof parsed.watches === 'object' ? parsed.watches : {},
|
||||
includeArchivedSessions: typeof parsed.includeArchivedSessions === 'boolean'
|
||||
? parsed.includeArchivedSessions
|
||||
: false,
|
||||
};
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
|
|
@ -63,6 +84,69 @@ export class StateStore {
|
|||
return structuredClone(this.#state);
|
||||
}
|
||||
|
||||
// ── Watches (persisted: surviving restarts) ─────────────────────────────
|
||||
|
||||
watchEntries(key) {
|
||||
const list = this.#state.watches[key];
|
||||
return Array.isArray(list) ? list.filter(validWatchEntry) : [];
|
||||
}
|
||||
|
||||
watchEntry(key, sessionId) {
|
||||
return this.watchEntries(key).find((entry) => entry.sessionId === sessionId) ?? null;
|
||||
}
|
||||
|
||||
async setWatch(key, entry) {
|
||||
const list = this.#state.watches[key] ?? [];
|
||||
const index = list.findIndex((existing) => existing.sessionId === entry.sessionId);
|
||||
if (index === -1) {
|
||||
if (list.length >= MAX_WATCHES_PER_KEY) list.shift();
|
||||
list.push(entry);
|
||||
} else {
|
||||
list[index] = entry;
|
||||
}
|
||||
this.#state.watches[key] = list;
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
async removeWatch(key, sessionId) {
|
||||
const list = this.#state.watches[key] ?? [];
|
||||
this.#state.watches[key] = list.filter((entry) => entry.sessionId !== sessionId);
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
async clearWatches(key) {
|
||||
delete this.#state.watches[key];
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
/** Every conversation key currently watching the given session. */
|
||||
keysWatching(sessionId) {
|
||||
return Object.entries(this.#state.watches)
|
||||
.filter(([, list]) => Array.isArray(list) && list.some((entry) => validWatchEntry(entry) && entry.sessionId === sessionId))
|
||||
.map(([key]) => key);
|
||||
}
|
||||
|
||||
/** Unique watched session ids across all keys (restart compensation). */
|
||||
watchedSessionIds() {
|
||||
const ids = new Set();
|
||||
for (const list of Object.values(this.#state.watches)) {
|
||||
if (!Array.isArray(list)) continue;
|
||||
for (const entry of list) if (validWatchEntry(entry)) ids.add(entry.sessionId);
|
||||
}
|
||||
return [...ids];
|
||||
}
|
||||
|
||||
// ── Session-list archived policy (per bot) ───────────────────
|
||||
|
||||
includesArchivedSessions() {
|
||||
return this.#state.includeArchivedSessions === true;
|
||||
}
|
||||
|
||||
async setIncludeArchivedSessions(include) {
|
||||
this.#state.includeArchivedSessions = include === true;
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
async #persist() {
|
||||
const snapshot = JSON.stringify(this.#state, null, 2) + '\n';
|
||||
this.#writeQueue = this.#writeQueue.then(async () => {
|
||||
|
|
|
|||
|
|
@ -1223,6 +1223,125 @@ export class HarnessClient {
|
|||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Watch the global Harness event mux (all sessions) until `signal`
|
||||
* aborts, reconnecting on drop. The Desktop host serves the mux as a
|
||||
* WebSocket downlink; frames are `server-request` envelopes whose payload
|
||||
* is a `session/event` — only those are forwarded. `onReconnect` (when
|
||||
* provided) fires after every (re)connection so callers can compensate
|
||||
* for events missed while offline.
|
||||
*/
|
||||
async watchHarnessEvents({ signal, onSessionEvent, onReconnect } = {}) {
|
||||
if (typeof onSessionEvent !== 'function') {
|
||||
throw new TypeError('watchHarnessEvents requires onSessionEvent');
|
||||
}
|
||||
if (!signal || typeof signal.addEventListener !== 'function') {
|
||||
throw new TypeError('watchHarnessEvents requires an AbortSignal');
|
||||
}
|
||||
if (onReconnect !== undefined && typeof onReconnect !== 'function') {
|
||||
throw new TypeError('onReconnect must be a function');
|
||||
}
|
||||
const url = new URL('/api/events.mux', this.#baseUrl);
|
||||
url.protocol = url.protocol === 'https:' ? 'wss:' : 'ws:';
|
||||
while (!signal.aborted) {
|
||||
try {
|
||||
await this.#watchHarnessEventSocket(url.toString(), {
|
||||
signal,
|
||||
onSessionEvent,
|
||||
onReconnect,
|
||||
});
|
||||
} catch (error) {
|
||||
if (signal.aborted) return;
|
||||
console.warn(`[${this.#logPrefix}] Harness event mux disconnected:`, error.message);
|
||||
}
|
||||
if (signal.aborted) return;
|
||||
try {
|
||||
await sleep(this.#interactionReconnectDelayMs, signal);
|
||||
} catch {
|
||||
if (signal.aborted) return;
|
||||
throw new Error('Harness event mux reconnect wait failed');
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#watchHarnessEventSocket(url, { signal, onSessionEvent, onReconnect }) {
|
||||
return new Promise((resolve, reject) => {
|
||||
let socket;
|
||||
try {
|
||||
socket = this.#createWebSocket(url);
|
||||
} catch (error) {
|
||||
reject(error);
|
||||
return;
|
||||
}
|
||||
let opened = false;
|
||||
let finished = false;
|
||||
const close = () => {
|
||||
try {
|
||||
socket.close();
|
||||
} catch {
|
||||
// Already closed.
|
||||
}
|
||||
};
|
||||
const finish = (error) => {
|
||||
if (finished) return;
|
||||
finished = true;
|
||||
socket.removeEventListener('open', handleOpen);
|
||||
socket.removeEventListener('message', handleMessage);
|
||||
socket.removeEventListener('close', handleClose);
|
||||
socket.removeEventListener('error', handleError);
|
||||
signal.removeEventListener('abort', handleAbort);
|
||||
if (error) reject(error);
|
||||
else resolve();
|
||||
};
|
||||
const handleOpen = () => {
|
||||
opened = true;
|
||||
try {
|
||||
onReconnect?.();
|
||||
} catch (error) {
|
||||
console.warn(`[${this.#logPrefix}] mux reconnect hook failed:`, error.message);
|
||||
}
|
||||
};
|
||||
const handleMessage = (event) => {
|
||||
try {
|
||||
if (typeof event.data !== 'string') return;
|
||||
const envelope = JSON.parse(event.data);
|
||||
const payload = envelope?.payload;
|
||||
if (envelope?.type !== 'server-request'
|
||||
|| !payload
|
||||
|| typeof payload !== 'object'
|
||||
|| envelope.method !== payload.type
|
||||
|| payload.type !== 'session/event'
|
||||
|| typeof payload.sessionId !== 'string'
|
||||
|| !payload.event
|
||||
|| typeof payload.event !== 'object') return;
|
||||
onSessionEvent({ sessionId: payload.sessionId, event: payload.event });
|
||||
} catch (error) {
|
||||
console.warn(`[${this.#logPrefix}] ignored a malformed global mux frame:`, error.message);
|
||||
}
|
||||
};
|
||||
const handleClose = () => finish(opened ? null : new Error(
|
||||
'Harness event mux WebSocket closed before opening',
|
||||
));
|
||||
const handleError = () => {
|
||||
finish(new Error(opened
|
||||
? 'Harness event mux WebSocket failed'
|
||||
: 'Harness event mux 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');
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2045,7 +2045,9 @@ test('session list paginates by page number across 25 sessions', async () => {
|
|||
const firstButton = firstLayout?.columns?.[0]?.elements?.[0];
|
||||
assert.equal(firstButton?.tag, 'button');
|
||||
assert.equal(Object.hasOwn(firstButton, 'value'), false, 'V2 buttons must not use the legacy value field');
|
||||
assert.equal(callbackAction(firstButton), 'use:session-01');
|
||||
assert.equal(callbackAction(firstButton), 'watch:session-01', 'the watch toggle leads each row');
|
||||
const sessionButton = firstLayout?.columns?.[1]?.elements?.[0];
|
||||
assert.equal(callbackAction(sessionButton), 'use:session-01');
|
||||
assert.equal(useActionsFromCard(page0).length, 10);
|
||||
assert.equal(useActionsFromCard(page0)[0], 'session-01');
|
||||
|
||||
|
|
@ -2081,7 +2083,8 @@ test('number replies on a later session page use page-local labels', async () =>
|
|||
await bridge.waitForIdle();
|
||||
|
||||
const page2Buttons = buttonsFromCard(cards(sent).at(-1).content);
|
||||
assert.match(page2Buttons[0].text.content, /^1\. Session 21$/);
|
||||
const sessionButtons = page2Buttons.filter((candidate) => (callbackAction(candidate) ?? '').startsWith('use:'));
|
||||
assert.match(sessionButtons[0].text.content, /^1\. Session 21$/);
|
||||
|
||||
await bridge.accept(event('sessions-number-pick', '1', { senderOpenId: 'ou_owner' }));
|
||||
await bridge.waitForIdle();
|
||||
|
|
@ -2127,7 +2130,8 @@ test('session pagination preserves an explicitly selected workspace', async () =
|
|||
await bridge.onCardAction(cardActionEvent('om_card_1', 'sessions:1', 'ou_owner'));
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(useActionsFromCard(cards(sent).at(-1).content)[0], 'selected-11');
|
||||
assert.match(JSON.stringify(cards(sent).at(-1).content), new RegExp(workspaceB.replace(/[.*+?^${}()|[\]\\]/g, '\\$&')));
|
||||
const header = cards(sent).at(-1).content.body.elements[0].text.content;
|
||||
assert.match(header, new RegExp(workspaceB.replace(/[.*+?^${}()|[\]\\]/g, '\\$&')));
|
||||
});
|
||||
|
||||
const REPAIR_APP_ID = 'cli_repair_test';
|
||||
|
|
@ -2404,3 +2408,369 @@ test('repair monitor reports expiry without claiming that the callback was fixed
|
|||
.find((text) => text.includes('授权链接已过期'));
|
||||
assert.doesNotMatch(terminal, /修复完成/);
|
||||
});
|
||||
|
||||
// ── Watches: read-only tracking, persistence, compensation, dedup ─────────
|
||||
|
||||
import { StateStore } from '../../../src/channels/feishu/state-store.mjs';
|
||||
|
||||
function watchHarness({ sessionsByWorkspace = { 'C:/work': [] }, current = 'C:/work', history = [] } = {}) {
|
||||
const listeners = [];
|
||||
let currentHistory = history;
|
||||
return {
|
||||
ensureRunning: async () => true,
|
||||
currentWorkspace: () => current,
|
||||
listWorkspaces: async () => Object.keys(sessionsByWorkspace),
|
||||
listWorkspaceSessions: async (workspace) => ({ workspace, sessions: sessionsByWorkspace[workspace] ?? [] }),
|
||||
bindWorkspaceSession: async (_key, sessionId) => ({ sessionId, title: `Title ${sessionId}` }),
|
||||
switchWorkspace: async (path) => path,
|
||||
rpc: async (method, params) => (method === 'session.history' ? { events: currentHistory } : null),
|
||||
watchHarnessEvents: ({ signal, onSessionEvent, onReconnect }) => {
|
||||
listeners.push({ signal, onSessionEvent, onReconnect });
|
||||
return new Promise((resolve) => {
|
||||
if (signal.aborted) resolve();
|
||||
else signal.addEventListener('abort', resolve, { once: true });
|
||||
});
|
||||
},
|
||||
_listeners: listeners,
|
||||
_setHistory: (next) => { currentHistory = next; },
|
||||
};
|
||||
}
|
||||
|
||||
async function watchStoreFixture(seedSessions = []) {
|
||||
const path = join(tmpdir(), `dsh-im-watch-test-${Math.random().toString(36).slice(2)}.json`);
|
||||
const store = new StateStore(path);
|
||||
await store.load();
|
||||
for (const [key, sessionId] of seedSessions) await store.setSession(key, sessionId);
|
||||
return { path, store, state: store };
|
||||
}
|
||||
|
||||
test('/watch resolves read-only: no binding, no workspace switch', async () => {
|
||||
const { state } = await watchStoreFixture([['p2p:ou_owner', 'bound-session']]);
|
||||
let bindCalls = 0;
|
||||
let switchCalls = 0;
|
||||
const harness = watchHarness({
|
||||
sessionsByWorkspace: { 'C:/work': [{ sessionId: 'target-session', title: 'Target' }] },
|
||||
});
|
||||
harness.bindWorkspaceSession = async () => { bindCalls += 1; throw new Error('must not bind'); };
|
||||
harness.switchWorkspace = async () => { switchCalls += 1; throw new Error('must not switch'); };
|
||||
const sent = [];
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: textClient(async ({ text }) => sent.push(text)),
|
||||
channel: {},
|
||||
harness,
|
||||
state,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
});
|
||||
|
||||
await bridge.accept(event('watch-1', '/watch 1', { senderOpenId: 'ou_owner' }));
|
||||
await bridge.waitForIdle();
|
||||
assert.match(sent.at(-1), /已关注会话「Target」/);
|
||||
assert.equal(bindCalls, 0, 'watch must not bind the conversation');
|
||||
assert.equal(switchCalls, 0, 'watch must not switch workspaces');
|
||||
assert.equal(state.sessionFor('p2p:ou_owner'), 'bound-session', 'existing binding unchanged');
|
||||
const entry = state.watchEntry('p2p:ou_owner', 'target-session');
|
||||
assert.ok(entry, 'watch entry persisted');
|
||||
assert.equal(entry.chatId, 'oc_chat');
|
||||
});
|
||||
|
||||
test('/watch finds a session in another workspace without switching', async () => {
|
||||
const { state } = await watchStoreFixture();
|
||||
const harness = watchHarness({
|
||||
current: 'C:/work',
|
||||
sessionsByWorkspace: {
|
||||
'C:/work': [],
|
||||
'D:/other': [{ sessionId: 'other-session', title: 'Other Session' }],
|
||||
},
|
||||
});
|
||||
const sent = [];
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: textClient(async ({ text }) => sent.push(text)),
|
||||
channel: {},
|
||||
harness,
|
||||
state,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
});
|
||||
|
||||
await bridge.accept(event('watch-x', '/watch other-session', { senderOpenId: 'ou_owner' }));
|
||||
await bridge.waitForIdle();
|
||||
assert.match(sent.at(-1), /已关注会话「Other Session」/);
|
||||
assert.equal(state.sessionFor('p2p:ou_owner'), null, 'cross-workspace watch must not bind');
|
||||
assert.ok(state.watchEntry('p2p:ou_owner', 'other-session'));
|
||||
});
|
||||
|
||||
test('persisted watches resume the event watcher at runtime start', async () => {
|
||||
const { path, state } = await watchStoreFixture();
|
||||
await state.setWatch('p2p:ou_owner', { sessionId: 'kept-session', title: 'Kept', chatId: 'oc_chat', lastSeq: 3 });
|
||||
const reloadedState = await new StateStore(path).load();
|
||||
const harness = watchHarness();
|
||||
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: textClient(async () => {}),
|
||||
channel: {},
|
||||
harness,
|
||||
state: reloadedState,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
});
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
assert.equal(harness._listeners.length, 1, 'watcher must restart from persisted state');
|
||||
assert.ok(bridge);
|
||||
});
|
||||
|
||||
test('reconnect compensation replays missed turn/end and dedups duplicates', async () => {
|
||||
const { state } = await watchStoreFixture();
|
||||
await state.setWatch('p2p:ou_owner', { sessionId: 'watched-session', title: 'Watched', chatId: 'oc_chat', lastSeq: 9 });
|
||||
const harness = watchHarness({
|
||||
history: [
|
||||
{ event: { type: 'turn/end', seq: 11, data: { turn: 't2', reason: { kind: 'stopped' } } } },
|
||||
{ event: { type: 'turn/end', seq: 10, data: { turn: 't1', reason: { kind: 'completed' } } } },
|
||||
],
|
||||
});
|
||||
const cards = [];
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: cardClient(async ({ msgType, content }) => {
|
||||
if (msgType === 'interactive') cards.push(content);
|
||||
}),
|
||||
channel: {},
|
||||
harness,
|
||||
state,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
});
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
assert.equal(harness._listeners.length, 1);
|
||||
|
||||
// Real history wraps events and may return them out of order.
|
||||
harness._listeners[0].onReconnect();
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(cards.length, 2);
|
||||
assert.match(JSON.stringify(cards[0]), /已完成/);
|
||||
assert.match(JSON.stringify(cards[1]), /已停止/);
|
||||
assert.equal(state.watchEntry('p2p:ou_owner', 'watched-session').lastSeq, 11);
|
||||
|
||||
// Reconnect and an overlapping live frame are both deduplicated by lastSeq.
|
||||
harness._listeners[0].onReconnect();
|
||||
harness._listeners[0].onSessionEvent({
|
||||
sessionId: 'watched-session',
|
||||
event: { type: 'turn/end', seq: 11, data: { turn: 't2', reason: { kind: 'stopped' } } },
|
||||
});
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(cards.length, 2);
|
||||
});
|
||||
|
||||
test('/watch baselines existing history and completion-card buttons keep their route', async () => {
|
||||
const { state } = await watchStoreFixture();
|
||||
const work = realpathSync(tmpdir());
|
||||
const oldCompletion = {
|
||||
event: { type: 'turn/end', seq: 10, data: { turn: 'old', reason: { kind: 'completed' } } },
|
||||
};
|
||||
const harness = watchHarness({
|
||||
current: work,
|
||||
sessionsByWorkspace: { [work]: [{ sessionId: 'watched-session', title: 'Watched' }] },
|
||||
history: [oldCompletion],
|
||||
});
|
||||
const sent = [];
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: cardClient(async (outgoing) => sent.push(outgoing)),
|
||||
channel: {},
|
||||
harness,
|
||||
state,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
});
|
||||
|
||||
await bridge.accept(event('watch-baseline', '/watch 1', { senderOpenId: 'ou_owner' }));
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(state.watchEntry('p2p:ou_owner', 'watched-session').lastSeq, 10);
|
||||
assert.equal(cards(sent).length, 0, 'a new watch must not replay an older completion');
|
||||
|
||||
harness._listeners[0].onReconnect();
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(cards(sent).length, 0);
|
||||
|
||||
harness._setHistory([
|
||||
oldCompletion,
|
||||
{ event: { type: 'turn/end', seq: 11, data: { turn: 'new', reason: { kind: 'completed' } } } },
|
||||
]);
|
||||
harness._listeners[0].onReconnect();
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(cards(sent).length, 1);
|
||||
|
||||
// The text confirmation is om_card_1, so the completion is om_card_2.
|
||||
await bridge.onCardAction(cardActionEvent('om_card_2', 'sessions', 'ou_owner'));
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(cards(sent).length, 2);
|
||||
assert.equal(cards(sent).at(-1).content.header.title.content, '📂 会话列表');
|
||||
});
|
||||
|
||||
test('a failed completion push keeps its watermark and later activity retries it', async () => {
|
||||
const { state } = await watchStoreFixture();
|
||||
await state.setWatch('p2p:ou_owner', {
|
||||
sessionId: 'cross-workspace-session',
|
||||
title: 'Cross Workspace Title',
|
||||
chatId: 'oc_chat',
|
||||
lastSeq: 10,
|
||||
});
|
||||
const completion = {
|
||||
event: { type: 'turn/end', seq: 11, data: { turn: 'retry', reason: { kind: 'completed' } } },
|
||||
};
|
||||
const laterCompletion = {
|
||||
event: { type: 'turn/end', seq: 12, data: { turn: 'later', reason: { kind: 'completed' } } },
|
||||
};
|
||||
const harness = watchHarness({ history: [laterCompletion, completion] });
|
||||
const cardsSent = [];
|
||||
let failNext = true;
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: cardClient(async ({ msgType, content }) => {
|
||||
if (msgType !== 'interactive') return;
|
||||
if (failNext) {
|
||||
failNext = false;
|
||||
throw new Error('temporary Feishu failure');
|
||||
}
|
||||
cardsSent.push(content);
|
||||
}),
|
||||
channel: {},
|
||||
harness,
|
||||
state,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
logger: { warn: () => undefined },
|
||||
});
|
||||
await eventually(() => harness._listeners.length === 1);
|
||||
|
||||
harness._listeners[0].onSessionEvent({
|
||||
sessionId: 'cross-workspace-session',
|
||||
event: completion.event,
|
||||
});
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(state.watchEntry('p2p:ou_owner', 'cross-workspace-session').lastSeq, 10);
|
||||
assert.equal(cardsSent.length, 0);
|
||||
|
||||
// A later live completion recovers the earlier failure through history;
|
||||
// no socket reconnect is required to unstick this watch.
|
||||
harness._listeners[0].onSessionEvent({
|
||||
sessionId: 'cross-workspace-session',
|
||||
event: laterCompletion.event,
|
||||
});
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(state.watchEntry('p2p:ou_owner', 'cross-workspace-session').lastSeq, 12);
|
||||
assert.equal(cardsSent.length, 2);
|
||||
assert.match(JSON.stringify(cardsSent[0]), /Cross Workspace Title/);
|
||||
});
|
||||
|
||||
test('legacy watches establish a baseline without replaying old completions', async () => {
|
||||
const { state } = await watchStoreFixture();
|
||||
await state.setWatch('p2p:ou_owner', {
|
||||
sessionId: 'legacy-session',
|
||||
title: 'Legacy',
|
||||
chatId: 'oc_chat',
|
||||
lastSeq: null,
|
||||
});
|
||||
const harness = watchHarness({
|
||||
history: [{ event: { type: 'turn/end', seq: 20, data: { turn: 'old' } } }],
|
||||
});
|
||||
const cardsSent = [];
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: cardClient(async ({ msgType, content }) => {
|
||||
if (msgType === 'interactive') cardsSent.push(content);
|
||||
}),
|
||||
channel: {},
|
||||
harness,
|
||||
state,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
});
|
||||
await eventually(() => harness._listeners.length === 1);
|
||||
|
||||
harness._listeners[0].onReconnect();
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(cardsSent.length, 0);
|
||||
assert.equal(state.watchEntry('p2p:ou_owner', 'legacy-session').lastSeq, 20);
|
||||
});
|
||||
|
||||
test('runtime abort stops the old event watcher before a new bridge starts', async () => {
|
||||
const firstHarness = watchHarness();
|
||||
const firstController = new AbortController();
|
||||
const { state: firstState } = await watchStoreFixture();
|
||||
new FeishuHarnessBridge({
|
||||
client: textClient(async () => {}),
|
||||
channel: {},
|
||||
harness: firstHarness,
|
||||
state: firstState,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
signal: firstController.signal,
|
||||
});
|
||||
await eventually(() => firstHarness._listeners.length === 1);
|
||||
assert.equal(firstHarness._listeners[0].signal.aborted, false);
|
||||
firstController.abort();
|
||||
assert.equal(firstHarness._listeners[0].signal.aborted, true);
|
||||
|
||||
const secondHarness = watchHarness();
|
||||
const secondController = new AbortController();
|
||||
const { state: secondState } = await watchStoreFixture();
|
||||
new FeishuHarnessBridge({
|
||||
client: textClient(async () => {}),
|
||||
channel: {},
|
||||
harness: secondHarness,
|
||||
state: secondState,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
signal: secondController.signal,
|
||||
});
|
||||
await eventually(() => secondHarness._listeners.length === 1);
|
||||
assert.equal(secondHarness._listeners[0].signal.aborted, false);
|
||||
secondController.abort();
|
||||
});
|
||||
|
||||
test('archived sessions are hidden by default; /archived on reveals them', async () => {
|
||||
const { state } = await watchStoreFixture();
|
||||
const workRaw = join(tmpdir(), 'dsh-im-archived-test-work');
|
||||
mkdirSync(workRaw, { recursive: true });
|
||||
const work = realpathSync(workRaw);
|
||||
const harness = watchHarness({
|
||||
current: work,
|
||||
sessionsByWorkspace: {
|
||||
[work]: [
|
||||
{ sessionId: 'live-session', title: 'Live', archived: false },
|
||||
{ sessionId: 'old-session', title: 'Old', archived: true },
|
||||
],
|
||||
},
|
||||
});
|
||||
const sent = [];
|
||||
const cards = [];
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: cardClient(async ({ msgType, content }) => {
|
||||
if (msgType === 'interactive') cards.push(content);
|
||||
}),
|
||||
channel: {},
|
||||
harness,
|
||||
state,
|
||||
status: bridgeStatus(),
|
||||
allowedSenderOpenIds: new Set(['ou_owner']),
|
||||
});
|
||||
bridge._sent = sent;
|
||||
|
||||
// Default: hidden. The explicit toggle re-enables inclusion for the card check.
|
||||
assert.equal(state.includesArchivedSessions(), false);
|
||||
|
||||
await bridge.accept(event('sessions-arch', '/sessionlist', { senderOpenId: 'ou_owner' }));
|
||||
await bridge.waitForIdle();
|
||||
const use = useActionsFromCard(cards.at(-1));
|
||||
assert.deepEqual(use, ['live-session'], 'archived session hidden');
|
||||
|
||||
// The numeric watch index must resolve against the filtered list too.
|
||||
await bridge.accept(event('watch-arch', '/watch 1', { senderOpenId: 'ou_owner' }));
|
||||
await bridge.waitForIdle();
|
||||
assert.ok(state.watchEntry('p2p:ou_owner', 'live-session'));
|
||||
assert.equal(state.watchEntry('p2p:ou_owner', 'old-session'), null);
|
||||
|
||||
await bridge.accept(event('archived-on', '/archived on', { senderOpenId: 'ou_owner' }));
|
||||
await bridge.waitForIdle();
|
||||
await bridge.accept(event('sessions-arch-2', '/sessionlist', { senderOpenId: 'ou_owner' }));
|
||||
await bridge.waitForIdle();
|
||||
const useOn = useActionsFromCard(cards.at(-1));
|
||||
assert.deepEqual(useOn, ['live-session', 'old-session'], '/archived on restores archived sessions');
|
||||
});
|
||||
|
|
|
|||
|
|
@ -16,13 +16,13 @@ function buttons(value, result = []) {
|
|||
return result;
|
||||
}
|
||||
|
||||
test('menu exposes repair as number-only text instead of a callback button', () => {
|
||||
test('menu appends watchlist and keeps repair number-only', () => {
|
||||
const card = JSON.parse(menuCard());
|
||||
assert.match(JSON.stringify(card), /6 · 修复卡片按钮/);
|
||||
const actions = buttons(card).flatMap((button) => (
|
||||
button.behaviors?.map((behavior) => behavior?.value?.action) ?? []
|
||||
));
|
||||
assert.deepEqual(actions, ['sessions', 'workspaces', 'new', 'status', 'help']);
|
||||
assert.deepEqual(actions, ['sessions', 'workspaces', 'new', 'status', 'help', 'watchlist']);
|
||||
assert.equal(actions.includes('repair'), false);
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -159,6 +159,84 @@ test('interaction watcher uses the real Harness wire protocol and leaves approva
|
|||
assert.equal(socket.readyState, 3);
|
||||
});
|
||||
|
||||
test('global event watcher reconnects without resolving until abort', async () => {
|
||||
const sockets = [];
|
||||
const socketUrls = [];
|
||||
const client = new HarnessClient({
|
||||
baseUrl: 'http://127.0.0.1:3080/base',
|
||||
workspace: '/tmp/dsh-feishu-workspace',
|
||||
interactionReconnectDelayMs: 0,
|
||||
createWebSocket: (url) => {
|
||||
const socket = new FakeSocket();
|
||||
sockets.push(socket);
|
||||
socketUrls.push(url);
|
||||
queueMicrotask(() => socket.open());
|
||||
return socket;
|
||||
},
|
||||
});
|
||||
const controller = new AbortController();
|
||||
const events = [];
|
||||
let reconnects = 0;
|
||||
let settled = false;
|
||||
const watching = client.watchHarnessEvents({
|
||||
signal: controller.signal,
|
||||
onReconnect: () => { reconnects += 1; },
|
||||
onSessionEvent: (payload) => events.push(payload),
|
||||
});
|
||||
watching.finally(() => { settled = true; });
|
||||
|
||||
await eventually(() => reconnects === 1);
|
||||
sockets[0].frame({
|
||||
type: 'server-request',
|
||||
rpcId: 'event-one',
|
||||
method: 'session/event',
|
||||
payload: {
|
||||
type: 'session/event',
|
||||
sessionId: 'session-one',
|
||||
event: { type: 'turn/end', seq: 1 },
|
||||
},
|
||||
});
|
||||
sockets[0].frame({
|
||||
type: 'server-request',
|
||||
rpcId: 'invalid-method',
|
||||
method: 'different/method',
|
||||
payload: {
|
||||
type: 'session/event',
|
||||
sessionId: 'ignored',
|
||||
event: { type: 'turn/end', seq: 2 },
|
||||
},
|
||||
});
|
||||
await eventually(() => events.length === 1);
|
||||
|
||||
sockets[0].close();
|
||||
await eventually(() => reconnects === 2);
|
||||
assert.equal(settled, false, 'a dropped socket must not complete the watcher');
|
||||
sockets[1].frame({
|
||||
type: 'server-request',
|
||||
rpcId: 'event-two',
|
||||
method: 'session/event',
|
||||
payload: {
|
||||
type: 'session/event',
|
||||
sessionId: 'session-two',
|
||||
event: { type: 'turn/end', seq: 3 },
|
||||
},
|
||||
});
|
||||
await eventually(() => events.length === 2);
|
||||
|
||||
assert.deepEqual(socketUrls, [
|
||||
'ws://127.0.0.1:3080/api/events.mux',
|
||||
'ws://127.0.0.1:3080/api/events.mux',
|
||||
]);
|
||||
assert.deepEqual(events.map(({ sessionId, event }) => [sessionId, event.seq]), [
|
||||
['session-one', 1],
|
||||
['session-two', 3],
|
||||
]);
|
||||
controller.abort();
|
||||
await watching;
|
||||
assert.equal(sockets[1].readyState, 3);
|
||||
assert.equal(settled, true);
|
||||
});
|
||||
|
||||
test('HarnessClient lists only absolute workspace paths', async () => {
|
||||
const client = new HarnessClient({
|
||||
baseUrl: 'http://127.0.0.1:3080',
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue