mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-09 05:20:46 +08:00
feat(feishu): persistent session watches and completion pushes
- /watch resolves the target READ-ONLY: sessions are validated against registered workspace listings (current or any other) without binding the conversation or switching workspaces. /unwatch and /watchlist round out the command set; the watchlist card supports button and number-reply unwatching. - Watches persist in the state store (sessionId + title + chatId + lastSeq) and resume at runtime start: the bridge starts the global event-mux watcher in its constructor when the harness supports it. - Completion pushes fire on turn/end, deduped by sessionId + turn id, with the seq watermark persisted per watch. After a mux reconnect the bridge replays each watched session's recent history (session.history) through the normal handler, so missed turn/end events are compensated without duplicates. - Bound-session push targets are persisted (chatTargets) instead of living only in memory, and refresh on every accepted message. - Watch failures map to safe user-facing messages.
This commit is contained in:
parent
0579c24b13
commit
beb269dc0e
6 changed files with 746 additions and 132 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,6 +89,9 @@ const HELP_TEXT = [
|
|||
'/status 检查连接状态',
|
||||
'/repair 修复卡片按钮回调',
|
||||
'/m(或 /menu) 打开交互卡片菜单',
|
||||
'/watch [Session ID 或序号] 关注会话,任务完成自动推送',
|
||||
'/unwatch [Session ID 或序号] 取消关注',
|
||||
'/watchlist 查看关注列表',
|
||||
'/help 显示本帮助',
|
||||
].join('\n');
|
||||
|
||||
|
|
@ -240,6 +249,13 @@ 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;
|
||||
#eventWatchSignal = null;
|
||||
/** Turn dedup: `sessionId:turnId` already processed. */
|
||||
#handledTurns = new Set();
|
||||
/** Session titles for completion cards (short in-memory cache). */
|
||||
#sessionTitleCache = new Map();
|
||||
|
||||
constructor({
|
||||
client,
|
||||
|
|
@ -295,6 +311,12 @@ export class FeishuHarnessBridge {
|
|||
this.#approvals = new HarnessApprovalQueue({ label: 'Feishu', logger });
|
||||
this.#signal = signal;
|
||||
ensureStatus(this.#status);
|
||||
// Persisted watches (and bound-session push targets) must resume at
|
||||
// runtime start, not on the first message. Guarded: harnesses without
|
||||
// the mux watcher (tests, older hosts) simply skip this.
|
||||
if (typeof this.#harness?.watchHarnessEvents === 'function') {
|
||||
queueMicrotask(() => this.#ensureEventWatcher());
|
||||
}
|
||||
}
|
||||
|
||||
accept(event) {
|
||||
|
|
@ -567,6 +589,8 @@ export class FeishuHarnessBridge {
|
|||
const text = message.content;
|
||||
const hasImages = hasInboundImages(message);
|
||||
const commandText = event.message.message_type === 'text' && !hasImages ? text : null;
|
||||
// Keep the persistent delivery target fresh for bound-session pushes.
|
||||
this.#rememberChatTarget(key, event.message.chat_id);
|
||||
if (!text && !hasImages) {
|
||||
await this.#send(event.message.chat_id, '目前支持文字和图片消息。');
|
||||
return;
|
||||
|
|
@ -604,6 +628,20 @@ 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 (NUMBER_REPLY.test(commandText)) {
|
||||
const menu = this.#takeMenu(key);
|
||||
if (menu) {
|
||||
|
|
@ -1029,6 +1067,10 @@ 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));
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1081,6 +1123,15 @@ 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 });
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1171,6 +1222,236 @@ export class FeishuHarnessBridge {
|
|||
return messageId;
|
||||
}
|
||||
|
||||
// ── Watches: read-only session tracking + completion pushes ─────────────
|
||||
|
||||
#ensureEventWatcher() {
|
||||
if (this.#eventWatcher) return;
|
||||
if (typeof this.#harness?.watchHarnessEvents !== 'function') return;
|
||||
this.#eventWatchSignal = new AbortController();
|
||||
try {
|
||||
this.#eventWatcher = this.#harness.watchHarnessEvents({
|
||||
signal: this.#eventWatchSignal.signal,
|
||||
onSessionEvent: (payload) => this.#onHarnessEvent(payload),
|
||||
onReconnect: () => void this.#compensateMissedEvents(),
|
||||
});
|
||||
this.#eventWatcher?.catch?.(() => undefined);
|
||||
} catch (error) {
|
||||
this.#eventWatcher = null;
|
||||
this.#logger.warn?.('[dsh-feishu] event watcher failed to start:', error.message);
|
||||
}
|
||||
}
|
||||
|
||||
/** Persist the chat delivery target for a conversation (bound pushes). */
|
||||
#rememberChatTarget(key, chatId) {
|
||||
if (typeof this.#state?.setChatTarget !== 'function' || !chatId) return;
|
||||
this.#state.setChatTarget(key, chatId).catch(() => undefined);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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 = 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 #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) ?? [];
|
||||
if (!existing.some((entry) => entry.sessionId === resolved.sessionId) && existing.length >= MAX_WATCHES_PER_KEY) {
|
||||
await this.#send(chatId, `每个聊天最多关注 ${MAX_WATCHES_PER_KEY} 个会话。`);
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await this.#state.setWatch(key, {
|
||||
sessionId: resolved.sessionId,
|
||||
title: resolved.title,
|
||||
chatId,
|
||||
lastSeq: null,
|
||||
});
|
||||
await this.#send(chatId, `已关注会话「${String(resolved.title).replace(/\s+/gu, ' ')}」,任务完成会推送结果。`);
|
||||
} 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);
|
||||
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 });
|
||||
}
|
||||
|
||||
/**
|
||||
* A global mux session event. Completion pushes fire on turn/end, with
|
||||
* dedup on `sessionId:turnId` and a seq watermark persisted per watch so
|
||||
* reconnect compensation can replay only what was missed.
|
||||
*/
|
||||
#onHarnessEvent({ sessionId, event }) {
|
||||
if (!sessionId || !event || typeof event !== 'object' || event.type !== 'turn/end') return;
|
||||
const turnId = nonEmptyString(event.data?.turn)
|
||||
?? nonEmptyString(event.data?.turnId)
|
||||
?? String(event.seq ?? '');
|
||||
const dedupKey = `${sessionId}:${turnId}`;
|
||||
if (this.#handledTurns.has(dedupKey)) return;
|
||||
this.#handledTurns.add(dedupKey);
|
||||
if (this.#handledTurns.size > 2000) {
|
||||
const oldest = this.#handledTurns.values().next().value;
|
||||
if (oldest !== undefined) this.#handledTurns.delete(oldest);
|
||||
}
|
||||
if (typeof event.seq === 'number' && typeof this.#state?.setWatch === 'function') {
|
||||
for (const key of (this.#state.keysWatching?.(sessionId) ?? [])) {
|
||||
const entry = this.#state.watchEntry?.(key, sessionId);
|
||||
if (entry && (typeof entry.lastSeq !== 'number' || entry.lastSeq < event.seq)) {
|
||||
this.#state.setWatch(key, { ...entry, lastSeq: event.seq }).catch(() => undefined);
|
||||
}
|
||||
}
|
||||
}
|
||||
void this.#sendCompletion(sessionId, event);
|
||||
}
|
||||
|
||||
async #sendCompletion(sessionId, event) {
|
||||
const reason = event?.data?.reason?.kind ?? event?.data?.reason ?? null;
|
||||
const targets = new Set();
|
||||
if (typeof this.#state?.keysWatching === 'function') {
|
||||
for (const key of this.#state.keysWatching(sessionId)) {
|
||||
const entry = this.#state.watchEntry?.(key, sessionId);
|
||||
if (entry?.chatId) targets.add(entry.chatId);
|
||||
}
|
||||
}
|
||||
if (typeof this.#state?.sessionKeysFor === 'function') {
|
||||
for (const key of this.#state.sessionKeysFor(sessionId)) {
|
||||
const chatId = typeof this.#state?.chatTargetFor === 'function'
|
||||
? this.#state.chatTargetFor(key)
|
||||
: null;
|
||||
if (chatId) targets.add(chatId);
|
||||
}
|
||||
}
|
||||
if (targets.size === 0) return;
|
||||
const title = await this.#sessionTitleFor(sessionId);
|
||||
for (const chatId of targets) {
|
||||
try {
|
||||
await this.#sendCard(chatId, completionCard(sessionId, title, reason));
|
||||
} catch (error) {
|
||||
this.#logger.warn?.('[dsh-feishu] completion push failed:', error.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async #sessionTitleFor(sessionId) {
|
||||
if (this.#sessionTitleCache.has(sessionId)) return this.#sessionTitleCache.get(sessionId);
|
||||
let title = '暂无标题';
|
||||
try {
|
||||
const currentPath = typeof this.#harness?.currentWorkspace === 'function'
|
||||
? this.#harness.currentWorkspace()
|
||||
: null;
|
||||
if (currentPath && typeof this.#harness?.listWorkspaceSessions === 'function') {
|
||||
const listed = await this.#harness.listWorkspaceSessions(currentPath);
|
||||
const session = (listed?.sessions ?? []).find((candidate) => candidate.sessionId === sessionId);
|
||||
if (session) title = String(session.title ?? '').replace(/\s+/gu, ' ').trim() || '暂无标题';
|
||||
}
|
||||
} catch {
|
||||
// Best-effort title; the card falls back to the id.
|
||||
}
|
||||
this.#sessionTitleCache.set(sessionId, title);
|
||||
return title;
|
||||
}
|
||||
|
||||
/**
|
||||
* Replay turn/end events missed while the mux was disconnected. Reads
|
||||
* each watched session's recent history and feeds unseen events through
|
||||
* the normal handler, whose `sessionId:turnId` dedup makes the replay
|
||||
* overlap-safe against live frames.
|
||||
*/
|
||||
async #compensateMissedEvents() {
|
||||
if (typeof this.#harness?.rpc !== 'function') return;
|
||||
const sessionIds = typeof this.#state?.watchedSessionIds === 'function'
|
||||
? this.#state.watchedSessionIds()
|
||||
: [];
|
||||
for (const sessionId of sessionIds) {
|
||||
try {
|
||||
const history = await this.#harness.rpc('session.history', { sessionId, maxMessages: 20 });
|
||||
const events = history?.events ?? [];
|
||||
for (const event of events) {
|
||||
if (event?.type !== 'turn/end' || typeof event.seq !== 'number') continue;
|
||||
const keys = typeof this.#state?.keysWatching === 'function'
|
||||
? this.#state.keysWatching(sessionId)
|
||||
: [];
|
||||
const allSeen = keys.length > 0 && keys.every((key) => {
|
||||
const entry = this.#state.watchEntry?.(key, sessionId);
|
||||
return entry && typeof entry.lastSeq === 'number' && entry.lastSeq >= event.seq;
|
||||
});
|
||||
if (allSeen) continue;
|
||||
this.#onHarnessEvent({ sessionId, event });
|
||||
}
|
||||
} catch (error) {
|
||||
this.#logger.warn?.(`[dsh-feishu] watch compensation failed for ${sessionId}:`, error.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#interactionAskOptions(event, key) {
|
||||
return {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
|
|
|
|||
|
|
@ -153,3 +153,32 @@ export function menuHelpText() {
|
|||
'/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 === 'stopped' ? '已停止' : reason === 'aborted' ? '已中止' : '已完成';
|
||||
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,18 @@
|
|||
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: {}, chatTargets: {} });
|
||||
|
||||
/** 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 +30,8 @@ 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 : {},
|
||||
chatTargets: parsed.chatTargets && typeof parsed.chatTargets === 'object' ? parsed.chatTargets : {},
|
||||
};
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
|
|
@ -63,6 +76,78 @@ 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];
|
||||
}
|
||||
|
||||
/** Conversation keys whose bound session is the given id. */
|
||||
sessionKeysFor(sessionId) {
|
||||
return Object.entries(this.#state.sessions)
|
||||
.filter(([, bound]) => bound === sessionId)
|
||||
.map(([key]) => key);
|
||||
}
|
||||
|
||||
// ── Persistent chat delivery targets for bound sessions ─────────────────
|
||||
|
||||
chatTargetFor(key) {
|
||||
const target = this.#state.chatTargets[key];
|
||||
return typeof target === 'string' && target.length > 0 ? target : null;
|
||||
}
|
||||
|
||||
async setChatTarget(key, chatId) {
|
||||
if (this.#state.chatTargets[key] === chatId) return;
|
||||
this.#state.chatTargets[key] = chatId;
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
async #persist() {
|
||||
const snapshot = JSON.stringify(this.#state, null, 2) + '\n';
|
||||
this.#writeQueue = this.#writeQueue.then(async () => {
|
||||
|
|
|
|||
|
|
@ -1223,6 +1223,85 @@ 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.
|
||||
*/
|
||||
watchHarnessEvents({ signal, onSessionEvent, onReconnect }) {
|
||||
if (typeof onSessionEvent !== 'function') {
|
||||
return Promise.reject(new TypeError('watchHarnessEvents requires onSessionEvent'));
|
||||
}
|
||||
const url = new URL('/api/events.mux', this.#baseUrl);
|
||||
url.protocol = url.protocol === 'https:' ? 'wss:' : 'ws:';
|
||||
let stopped = false;
|
||||
const connect = () => {
|
||||
if (stopped || signal.aborted) return Promise.resolve();
|
||||
return new Promise((resolve) => {
|
||||
const socket = this.#createWebSocket(url.toString());
|
||||
const settle = () => {
|
||||
socket.removeEventListener('open', handleOpen);
|
||||
socket.removeEventListener('message', handleMessage);
|
||||
socket.removeEventListener('close', handleClose);
|
||||
socket.removeEventListener('error', handleError);
|
||||
signal.removeEventListener('abort', handleAbort);
|
||||
resolve();
|
||||
};
|
||||
const handleOpen = () => {
|
||||
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') return;
|
||||
if (payload.type !== 'session/event') return;
|
||||
if (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 = () => {
|
||||
settle();
|
||||
if (!stopped && !signal.aborted) {
|
||||
setTimeout(() => { void connect(); }, 2000);
|
||||
}
|
||||
};
|
||||
const handleError = () => {
|
||||
try {
|
||||
socket.close();
|
||||
} catch {
|
||||
// The close event drives reconnection; nothing to do here.
|
||||
}
|
||||
};
|
||||
const handleAbort = () => {
|
||||
stopped = true;
|
||||
try {
|
||||
socket.close();
|
||||
} catch {
|
||||
// Already closed.
|
||||
}
|
||||
};
|
||||
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();
|
||||
});
|
||||
};
|
||||
return connect();
|
||||
}
|
||||
|
||||
stopManagedProcess() {
|
||||
if (this.#managedProcess?.exitCode === null) this.#managedProcess.kill('SIGTERM');
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2404,3 +2404,142 @@ 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 = [];
|
||||
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: history } : null),
|
||||
watchHarnessEvents: ({ onSessionEvent, onReconnect }) => {
|
||||
listeners.push({ onSessionEvent, onReconnect });
|
||||
return Promise.resolve();
|
||||
},
|
||||
_listeners: listeners,
|
||||
};
|
||||
}
|
||||
|
||||
async function watchStoreFixture(seedSessions = []) {
|
||||
const store = new StateStore(join(tmpdir(), `dsh-im-watch-test-${Math.random().toString(36).slice(2)}.json`));
|
||||
await store.load();
|
||||
for (const [key, sessionId] of seedSessions) await store.setSession(key, sessionId);
|
||||
return { 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 { state } = await watchStoreFixture();
|
||||
await state.setWatch('p2p:ou_owner', { sessionId: 'kept-session', title: 'Kept', chatId: 'oc_chat', lastSeq: 3 });
|
||||
const harness = watchHarness();
|
||||
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: textClient(async () => {}),
|
||||
channel: {},
|
||||
harness,
|
||||
state,
|
||||
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: null });
|
||||
const harness = watchHarness({
|
||||
sessionsByWorkspace: { 'C:/work': [{ sessionId: 'watched-session', title: 'Watched' }] },
|
||||
history: [
|
||||
{ type: 'turn/end', seq: 10, data: { turn: 't1', reason: { kind: 'completed' } } },
|
||||
{ type: 'turn/end', seq: 11, data: { turn: 't2', 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);
|
||||
|
||||
// Live event: t1 arrives before any replay.
|
||||
harness._listeners[0].onSessionEvent({ sessionId: 'watched-session', event: { type: 'turn/end', seq: 10, data: { turn: 't1', reason: { kind: 'completed' } } } });
|
||||
await bridge.waitForIdle();
|
||||
assert.equal(cards.length, 1, 'one completion for the live event');
|
||||
|
||||
// Reconnect: history replays t1 (dedup) and t2 (new).
|
||||
await harness._listeners[0].onReconnect();
|
||||
await new Promise((resolve) => setTimeout(resolve, 20));
|
||||
assert.equal(cards.length, 2, 'exactly one new completion from compensation (t1 deduped)');
|
||||
assert.match(JSON.stringify(cards[1]), /watched-session/);
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue