diff --git a/apps/desktop/src/main/agent-binaries/__tests__/pi-runtime-recovery.test.ts b/apps/desktop/src/main/agent-binaries/__tests__/pi-runtime-recovery.test.ts new file mode 100644 index 0000000000..950e894ace --- /dev/null +++ b/apps/desktop/src/main/agent-binaries/__tests__/pi-runtime-recovery.test.ts @@ -0,0 +1,107 @@ +import { describe, expect, it, vi } from 'vitest'; + +import { createPiRuntimeRecovery } from '../pi-runtime-recovery.js'; + +describe('Pi runtime recovery', () => { + it('retries after the network returns and registers Pi once', async () => { + let online = false; + let prepareCalls = 0; + let registered = false; + const onRegistered = vi.fn(); + const recovery = createPiRuntimeRecovery({ + isOnline: () => online, + prepare: async () => { + prepareCalls += 1; + return { ready: true, path: '/tmp/pi' }; + }, + register: () => { + if (registered) return false; + registered = true; + return true; + }, + onRegistered, + retryDelayMs: 60_000, + setTimeout: (() => 0) as unknown as typeof setTimeout, + clearTimeout: (() => undefined) as unknown as typeof clearTimeout, + }); + + recovery.markUnavailable('manifest_failed'); + expect(await recovery.retryNow('offline')).toBe(false); + online = true; + expect(await recovery.retryNow('online')).toBe(true); + expect(await recovery.retryNow('duplicate')).toBe(false); + expect(prepareCalls).toBe(1); + expect(onRegistered).toHaveBeenCalledOnce(); + expect(recovery.isDisabled()).toBe(false); + recovery.dispose(); + }); + + it('deduplicates concurrent recovery and keeps retryable failure disabled', async () => { + let resolvePrepare!: (value: { ready: boolean; path?: string; error?: string }) => void; + const prepare = vi.fn( + () => new Promise<{ ready: boolean; path?: string; error?: string }>((resolve) => { + resolvePrepare = resolve; + }), + ); + const recovery = createPiRuntimeRecovery({ + isOnline: () => true, + prepare, + register: () => true, + onRegistered: vi.fn(), + retryDelayMs: 60_000, + setTimeout: (() => 0) as unknown as typeof setTimeout, + clearTimeout: (() => undefined) as unknown as typeof clearTimeout, + }); + + recovery.markUnavailable('manifest_failed'); + const first = recovery.retryNow(); + const second = recovery.retryNow(); + expect(first).toBe(second); + expect(prepare).toHaveBeenCalledOnce(); + resolvePrepare({ ready: false, error: 'still_offline' }); + expect(await first).toBe(false); + expect(recovery.isDisabled()).toBe(true); + recovery.dispose(); + }); + + it('does not schedule retries for permanent prepare errors', async () => { + const prepare = vi.fn(async () => ({ ready: true, path: '/tmp/pi' })); + const schedule = vi.fn(() => 0); + const recovery = createPiRuntimeRecovery({ + isOnline: () => true, + prepare, + register: () => true, + onRegistered: vi.fn(), + setTimeout: schedule as unknown as typeof setTimeout, + clearTimeout: (() => undefined) as unknown as typeof clearTimeout, + }); + + recovery.markUnavailable('asset_missing'); + expect(schedule).not.toHaveBeenCalled(); + expect(await recovery.retryNow('permanent')).toBe(false); + expect(prepare).not.toHaveBeenCalled(); + recovery.dispose(); + }); + + it('stops an existing retry loop when a later prepare becomes permanent', async () => { + const prepare = vi.fn(async () => ({ ready: false, error: 'asset_missing' })); + const schedule = vi.fn(() => 0); + const cancel = vi.fn(); + const recovery = createPiRuntimeRecovery({ + isOnline: () => true, + prepare, + register: () => true, + onRegistered: vi.fn(), + setTimeout: schedule as unknown as typeof setTimeout, + clearTimeout: cancel as unknown as typeof clearTimeout, + }); + + recovery.markUnavailable('manifest_failed'); + expect(schedule).toHaveBeenCalledOnce(); + expect(await recovery.retryNow('permanent-after-transient')).toBe(false); + expect(cancel).toHaveBeenCalledOnce(); + expect(schedule).toHaveBeenCalledOnce(); + expect(recovery.isDisabled()).toBe(true); + recovery.dispose(); + }); +}); diff --git a/apps/desktop/src/main/agent-binaries/pi-runtime-recovery.ts b/apps/desktop/src/main/agent-binaries/pi-runtime-recovery.ts new file mode 100644 index 0000000000..956b6acf17 --- /dev/null +++ b/apps/desktop/src/main/agent-binaries/pi-runtime-recovery.ts @@ -0,0 +1,145 @@ +import type { PrepareResult } from './types.js'; + +/** Retry delay for an optional Pi runtime that missed the startup network. */ +export const PI_RUNTIME_RECOVERY_RETRY_MS = 30_000; + +/** Only errors that may change when connectivity returns are worth retrying. */ +export function isRetryablePiPrepareError(error?: string): boolean { + return error === 'manifest_failed' + || error === 'NETWORK' + || error === 'HTTP_5XX' + || error === 'ABORTED'; +} + +export interface PiRuntimeRecoveryOptions { + isOnline: () => boolean; + prepare: () => Promise; + register: () => boolean; + onRegistered: () => void; + logWarn?: (message: string, error?: unknown) => void; + retryDelayMs?: number; + setTimeout?: typeof globalThis.setTimeout; + clearTimeout?: typeof globalThis.clearTimeout; +} + +export interface PiRuntimeRecovery { + /** Mark the startup prepare as unavailable and begin background recovery. */ + markUnavailable(error?: string): void; + /** Try recovery immediately; returns true only when Pi was registered. */ + retryNow(reason?: string): Promise; + /** Stop future retries during app shutdown or test cleanup. */ + dispose(): void; + isDisabled(): boolean; +} + +/** + * Owns the small recovery state machine for an optional Pi runtime. + * + * The CDN policy remains unchanged: every retry calls the managed prepare path, + * so a local runtime is accepted only after the manifest and verification rules + * have succeeded. Concurrent focus/timer signals share one prepare promise. + */ +export function createPiRuntimeRecovery(options: PiRuntimeRecoveryOptions): PiRuntimeRecovery { + const retryDelayMs = options.retryDelayMs ?? PI_RUNTIME_RECOVERY_RETRY_MS; + const schedule = options.setTimeout ?? globalThis.setTimeout; + const cancel = options.clearTimeout ?? globalThis.clearTimeout; + let disabled = false; + let disposed = false; + let retryable = false; + let retryTimer: ReturnType | null = null; + let inFlight: Promise | null = null; + + const logWarn = (message: string, error?: unknown): void => { + options.logWarn?.(message, error); + }; + + const scheduleRetry = (): void => { + if (disposed || !disabled || retryTimer !== null) return; + retryTimer = schedule(() => { + retryTimer = null; + void recovery.retryNow('timer'); + }, retryDelayMs); + const unref = (retryTimer as unknown as { unref?: () => void }).unref; + unref?.call(retryTimer); + }; + + const recovery: PiRuntimeRecovery = { + markUnavailable(error) { + disabled = true; + retryable = isRetryablePiPrepareError(error); + if (!retryable && retryTimer !== null) { + cancel(retryTimer); + retryTimer = null; + } + if (error) { + logWarn( + retryable + ? 'Pi runtime unavailable; scheduling recovery' + : 'Pi runtime unavailable; recovery not scheduled for permanent prepare error', + error, + ); + } + if (retryable) scheduleRetry(); + }, + + retryNow(reason = 'manual') { + if (disposed || !disabled || !retryable) return Promise.resolve(false); + if (inFlight) return inFlight; + + let online = false; + try { + online = options.isOnline(); + } catch (error) { + logWarn('Pi runtime network state probe failed', error); + } + if (!online) { + scheduleRetry(); + return Promise.resolve(false); + } + + const attempt = (async (): Promise => { + try { + const result = await options.prepare(); + if (!result.ready || !result.path) { + logWarn(`Pi runtime recovery prepare failed (${reason})`, result.error); + recovery.markUnavailable(result.error); + return false; + } + if (!options.register()) { + scheduleRetry(); + return false; + } + disabled = false; + retryable = false; + options.onRegistered(); + return true; + } catch (error) { + logWarn(`Pi runtime recovery threw (${reason})`, error); + recovery.markUnavailable(error instanceof Error ? error.message : String(error)); + return false; + } + })(); + inFlight = attempt; + void attempt.then(() => { + if (inFlight === attempt) inFlight = null; + }, () => { + if (inFlight === attempt) inFlight = null; + }); + return attempt; + }, + + dispose() { + disposed = true; + if (retryTimer !== null) { + cancel(retryTimer); + retryTimer = null; + } + }, + + isDisabled() { + return disabled; + }, + }; + + return recovery; +} diff --git a/apps/desktop/src/main/bootstrap-electron.ts b/apps/desktop/src/main/bootstrap-electron.ts index 00d7db5f5e..703c10a1a1 100644 --- a/apps/desktop/src/main/bootstrap-electron.ts +++ b/apps/desktop/src/main/bootstrap-electron.ts @@ -61,6 +61,8 @@ import { waitForTurnChangeSetPersistence, } from './turn-change-set/store.js'; +let retryPiRuntimeAfterNetworkRecovery: (() => void) | null = null; +let disposePiRuntimeRecovery: (() => void) | null = null; // Official Linux binaries total hundreds of MB. Keep one shared deadline for // both downloads, but allow normal consumer connections to finish while the // splash displays real byte progress. @@ -192,7 +194,6 @@ import { prepare as binaryPrepare, peekNeedsDownload as binaryPeekNeedsDownload, broadcastResetForStep as binaryBroadcastResetForStep, - getCachedBinaryStatus, type AgentBinaryKind, type PrepareResult, } from './agent-binaries'; @@ -489,7 +490,9 @@ import { setProviderAccessRuntimeRefreshListener, restartCodexAfterAuthModeChange, waitForInitialCustomMcpRefresh, + registerPiAgentIfAvailable, } from './maker-host/index.js'; +import { createPiRuntimeRecovery } from './agent-binaries/pi-runtime-recovery.js'; import { createDynamicMaker } from './maker-host/dynamic-maker.js'; import { ensureBundledRipgrepReady } from './maker-host/runtime-configs.js'; import { @@ -3146,6 +3149,9 @@ const windowsClosePromptFallback = createWindowsClosePromptFallbackController( app.on('before-quit', () => { isQuitting = true; + disposePiRuntimeRecovery?.(); + retryPiRuntimeAfterNetworkRecovery = null; + disposePiRuntimeRecovery = null; windowsClosePromptFallback.dispose(); destroyWindowsTray(); disposeUpdatePresentationRecovery(); @@ -3204,6 +3210,7 @@ function scheduleAppFocusSync(): void { } app.on('browser-window-focus', (_event, win) => { + retryPiRuntimeAfterNetworkRecovery?.(); if (win === mainWindowRef) updatePresentationRecovery?.onWindowFocused(); if (appFocusSyncTimer) { clearTimeout(appFocusSyncTimer); @@ -5388,9 +5395,28 @@ const registerIpcHandlers = () => { // getMaker() 在构造期就读 binary path, 早于 splash 调用会抛错; 第一次 splash 成功后置 true, // 后续 retry 走 check-environment 不重复注册 (重复 ipcMain.handle 会覆盖同名 handler)。 let makerIpcsRegistered = false; - // Pi 是“本次启动可选”的能力:一旦首次准备失败,就算后续清单/CDN恢复, - // 也不再把 Pi 动态塞回已经构造好的 Maker,避免返回状态与实际能力不一致。 - let piDisabledForLaunch = false; + // Pi 是“本次启动可选”的能力:准备失败不阻塞主界面,交给 recovery 在网络恢复后 + // 重新走 managed prepare,并在成功后动态注册到当前 Maker。 + const piRuntimeRecovery = createPiRuntimeRecovery({ + isOnline: () => net.isOnline(), + prepare: async () => { + const result = await binaryPrepare('pi', { + broadcastFailure: false, + broadcastProgress: false, + signal: AbortSignal.timeout(PI_AGENT_INSTALL_STARTUP_DEADLINE_MS), + }); + return result; + }, + register: () => registerPiAgentIfAvailable(), + onRegistered: () => { + console.info('[bootstrap-electron] Pi runtime recovered and agent registered'); + }, + logWarn: (message, error) => console.warn(`[bootstrap-electron] ${message}`, error ?? ''), + }); + retryPiRuntimeAfterNetworkRecovery = () => { + void piRuntimeRecovery.retryNow('window-focus'); + }; + disposePiRuntimeRecovery = () => piRuntimeRecovery.dispose(); const registerMakerIpcsAfterSplash = async (): Promise => { if (makerIpcsRegistered) return; // 模型供应商目录(providers.json)按「OSS 真源 / bundled 兜底」加载一次存内存:必须在第一次 @@ -5751,41 +5777,26 @@ const registerIpcHandlers = () => { resetBeforeSegment('pi', claudeRes.downloaded === true || codexRes.downloaded === true); let piInfo: { status: 'passed' | 'failed'; path?: string; error?: string }; - // 轮 27 LOW-4:首次准备失败后账号切换(同一进程),二进制可能已由后台 - // 下载/手动放置变得可用 —— 轻量重试:意外可用则清除标志继续准备。 - if (piDisabledForLaunch) { - const cached = getCachedBinaryStatus('pi'); - if (cached?.binaryPath && cached.binaryPath.length > 0) { - piDisabledForLaunch = false; - } - } - if (piDisabledForLaunch) { + try { + const piRes = await binaryPrepare('pi', { + ...stepOptsFor('pi'), + broadcastFailure: false, + signal: piInstallSignal, + }); + piInfo = + piRes.ready && piRes.path + ? { status: 'passed' as const, path: piRes.path } + : { status: 'failed' as const, error: piRes.error ?? 'pi binary not available' }; + } catch (err: unknown) { piInfo = { status: 'failed' as const, - error: 'pi disabled for this launch after an earlier prepare failure', + error: err instanceof Error ? err.message : String(err), }; - } else { - try { - const piRes = await binaryPrepare('pi', { - ...stepOptsFor('pi'), - broadcastFailure: false, - signal: piInstallSignal, - }); - piInfo = - piRes.ready && piRes.path - ? { status: 'passed' as const, path: piRes.path } - : { status: 'failed' as const, error: piRes.error ?? 'pi binary not available' }; - } catch (err: unknown) { - piInfo = { - status: 'failed' as const, - error: err instanceof Error ? err.message : String(err), - }; - } } if (piInfo.status === 'failed') { - piDisabledForLaunch = true; + piRuntimeRecovery.markUnavailable(piInfo.error); console.warn( - `[bootstrap-electron] pi binary prepare failed (non-fatal, pi disabled for this launch): ${piInfo.error}`, + `[bootstrap-electron] pi binary prepare failed (non-fatal, recovery scheduled): ${piInfo.error}`, ); } diff --git a/apps/desktop/src/main/maker-host/__tests__/piBinaryDistribution.test.ts b/apps/desktop/src/main/maker-host/__tests__/piBinaryDistribution.test.ts index 7a2906fa19..7c880bd385 100644 --- a/apps/desktop/src/main/maker-host/__tests__/piBinaryDistribution.test.ts +++ b/apps/desktop/src/main/maker-host/__tests__/piBinaryDistribution.test.ts @@ -52,14 +52,23 @@ describe('Pi binary distribution contract', () => { expect(binaryVersion).toContain("if (kind === 'pi') return null;"); }); - it('keeps a failed Pi disabled when check-environment is retried', () => { + it('schedules Pi recovery after a failed optional prepare', () => { const bootstrap = fs.readFileSync( path.join(desktopRoot, 'src/main/bootstrap-electron.ts'), 'utf8', ); - expect(bootstrap).toContain('let piDisabledForLaunch = false;'); - expect(bootstrap).toContain('if (piDisabledForLaunch)'); - expect(bootstrap).toContain('piDisabledForLaunch = true;'); + expect(bootstrap).toContain('createPiRuntimeRecovery'); + expect(bootstrap).toContain('piRuntimeRecovery.markUnavailable'); + expect(bootstrap).toContain('registerPiAgentIfAvailable'); + }); + + it('does not override Pi memory when the runtime recovers', () => { + const host = fs.readFileSync( + path.join(desktopRoot, 'src/main/maker-host/index.ts'), + 'utf8', + ); + + expect(host).not.toContain("syncNativeAgentsOff(['pi'])"); }); }); diff --git a/apps/desktop/src/main/maker-host/index.ts b/apps/desktop/src/main/maker-host/index.ts index 000d50d07e..bf968857d5 100644 --- a/apps/desktop/src/main/maker-host/index.ts +++ b/apps/desktop/src/main/maker-host/index.ts @@ -270,6 +270,7 @@ type RemoteCcQuery = Awaited< >; let _maker: Maker | null = null; +let _registerPiAgent: (() => boolean) | null = null; /** 视觉桥实例(层 A/B/C 共用),在 resetMaker 时释放缓存。 */ let _visionBridgeInstance: ReturnType | null = null; @@ -1680,10 +1681,12 @@ export function getMaker(): Maker { _codexAgent = codexAgent; // 装配第二步: 把 agents 引用挂回 manager (manager.enable() 时遍历 setMemory(false))。 - attachAgentsToMakerMemory(makerMemoryManager, { + // 保留同一份可变 map,供可选 Pi runtime 在 Maker 构造后动态注册。 + const makerAgents: ConstructorParameters[0]['agents'] = { 'claude-code': claudeAgent, codex: codexAgent, - }); + }; + attachAgentsToMakerMemory(makerMemoryManager, makerAgents); if (makerMemoryManager.isEnabled()) { void makerMemoryManager.enable(); } @@ -1823,7 +1826,7 @@ export function getMaker(): Maker { registerCustomMcpArrays(claudeMcpProviders, codexMcpProviders, piMcpProviders); _initialCustomMcpRefresh = refreshCustomMcpProviders(); - const piAgent = buildPiAgent({ + const buildPiAgentForDesktop = () => buildPiAgent({ logger: desktopMakerLogger, turnChangeCapture: { beforeKnownFileWrite: captureKnownFileBefore, @@ -2034,6 +2037,8 @@ export function getMaker(): Maker { return getRemoteAgentProxyEnv(remoteHost); }, }); + const piAgent = buildPiAgentForDesktop(); + if (piAgent) makerAgents.pi = piAgent; setVisionGatewayKeyReader(readClaudeApiKey); _visionBridgeInstance = createVisionBridge({ @@ -2058,11 +2063,7 @@ export function getMaker(): Maker { }); _maker = new Maker({ - agents: { - 'claude-code': claudeAgent, - codex: codexAgent, - ...(piAgent ? { pi: piAgent } : {}), - }, + agents: makerAgents, storage: desktopSessionStorage, logger: desktopMakerLogger, makerMemory: makerMemoryManager, @@ -2147,6 +2148,21 @@ export function getMaker(): Maker { }, }, }); + _registerPiAgent = () => { + if (!_maker || _maker.listAvailableAgents().includes('pi')) return false; + const next = buildPiAgentForDesktop(); + if (!next) return false; + const registered = _maker.registerAgent('pi', next); + if (!registered) { + void next.dispose().catch((error) => { + desktopMakerLogger.warn('discarding Pi agent after registration race failed', { + error: error instanceof Error ? error.message : String(error), + }); + }); + return false; + } + return true; + }; setVisionBridgeController({ shouldBridge: _visionBridgeInstance.isTargetModel, describeImage: _visionBridgeInstance.describeImage, @@ -2211,12 +2227,39 @@ export function getMakerIfReady(): Maker | null { return _maker; } +/** Register Pi after a managed runtime retry and notify local renderers. */ +export function registerPiAgentIfAvailable(): boolean { + const register = _registerPiAgent; + if (!register) return false; + try { + if (_maker?.listAvailableAgents().includes('pi')) return true; + const registered = register(); + if (!registered) return false; + for (const win of BrowserWindow.getAllWindows()) { + if (win.isDestroyed()) continue; + try { + win.webContents.send(MAKER_PUSH.AGENTS_CHANGED); + } catch { + // Window teardown may race the broadcast; other windows still receive it. + } + } + tapWindowBroadcast(MAKER_PUSH.AGENTS_CHANGED, {}); + return true; + } catch (error) { + desktopMakerLogger.warn('Pi agent registration after runtime recovery failed', { + error: error instanceof Error ? error.message : String(error), + }); + return false; + } +} + /** * 重置 Maker 单例(切账号 / 测试用)。 */ export function resetMaker(): void { cancelCodexAuthModeChange(); _maker = null; + _registerPiAgent = null; _codexAgent = null; // coordinator 闭包捕获了刚作废的那个 maker —— 不清掉的话,换账号窗口期内到达的 auth // 事件会拿旧实例去拉模型清单(串号)。下次 getMaker() 会带着干净记账重建它。 diff --git a/apps/desktop/src/main/maker-ipc/channels.ts b/apps/desktop/src/main/maker-ipc/channels.ts index 38291d8d27..44c1655f92 100644 --- a/apps/desktop/src/main/maker-ipc/channels.ts +++ b/apps/desktop/src/main/maker-ipc/channels.ts @@ -800,6 +800,8 @@ export const MAKER_SEND = { export const MAKER_PUSH = { EVENT: 'maker:event', + /** Runtime agent roster changed after an optional agent recovery. */ + AGENTS_CHANGED: 'maker:agents:changed', TURN_CHANGE_SET_UPDATED: 'maker:turn-change-set:updated', STATUS_CHANGED: 'maker:status-changed', /** 用户从独立 Computer Use 授权引导浮窗主动取消。 */ diff --git a/apps/desktop/src/preload/preload.ts b/apps/desktop/src/preload/preload.ts index dc9b46a6af..1f11955dda 100644 --- a/apps/desktop/src/preload/preload.ts +++ b/apps/desktop/src/preload/preload.ts @@ -651,6 +651,7 @@ const fanOutHookControlWorkspaceProviderSource = createIpcFanOut( // ─── Maker Core 一阶段重构(新链路)── 与 cc-agent:* / codex:* 双轨并行 ───── const fanOutMakerEvent = createIpcFanOut('maker:event'); +const fanOutMakerAgentsChanged = createIpcFanOut('maker:agents:changed'); const fanOutMakerTurnChangeSetUpdated = createIpcFanOut('maker:turn-change-set:updated'); const fanOutMakerStatusChanged = createIpcFanOut('maker:status-changed'); const fanOutMakerInputProjection = createIpcFanOut('maker:input:projection'); @@ -5212,6 +5213,7 @@ contextBridge.exposeInMainWorld('electronAPI', { maker: { listAvailableAgents: (): Promise> => ipcRenderer.invoke('maker:list-available-agents'), + onAgentsChanged: fanOutMakerAgentsChanged, getCapabilities: (agentKind: 'claude-code' | 'codex' | 'pi'): Promise => ipcRenderer.invoke('maker:get-capabilities', agentKind), listTurnChangeSets: ( diff --git a/apps/desktop/src/renderer/__tests__/newMakerProjectPicker.test.ts b/apps/desktop/src/renderer/__tests__/newMakerProjectPicker.test.ts index ea22ab3263..e696dd11c2 100644 --- a/apps/desktop/src/renderer/__tests__/newMakerProjectPicker.test.ts +++ b/apps/desktop/src/renderer/__tests__/newMakerProjectPicker.test.ts @@ -551,6 +551,9 @@ describe('Shared create project picker', () => { ); // claude-code → cc 归一,fail-open(未加载不隐藏)。 expect(availableAgentsHookSource).toContain("agent === 'claude-code' ? 'cc' : agent"); + expect(availableAgentsHookSource).toContain('refreshLocalCapabilities'); + expect(availableAgentsHookSource).toContain('evictDeviceCapabilities'); + expect(availableAgentsHookSource).toContain('prefetchDeviceCapabilities'); // 未加载完成时不隐藏任何入口(loaded 保持 false → 空 hidden)。 expect(availableAgentsHookSource).toMatch(/loaded/); diff --git a/apps/desktop/src/renderer/hooks/__tests__/useAvailableAgents.test.ts b/apps/desktop/src/renderer/hooks/__tests__/useAvailableAgents.test.ts new file mode 100644 index 0000000000..fe950c1502 --- /dev/null +++ b/apps/desktop/src/renderer/hooks/__tests__/useAvailableAgents.test.ts @@ -0,0 +1,188 @@ +// @vitest-environment jsdom +import { act, renderHook, waitFor } from '@testing-library/react'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; + +vi.mock('@/lib/logger', () => ({ + createLogger: () => ({ + warn: vi.fn(), + }), +})); + +vi.mock('../useAgentCapabilities', () => ({ + evictDeviceCapabilities: vi.fn(), + prefetchDeviceCapabilities: vi.fn(async () => {}), + refreshLocalCapabilities: vi.fn(async () => {}), +})); + +type RuntimeAgentKind = 'claude-code' | 'codex' | 'pi'; +type PresenceListener = (snapshot: { deviceId: string; online: boolean }) => void; +type StatusListener = (payload: { status: 'stopped' | 'connecting' | 'online' }) => void; + +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise((nextResolve) => { + resolve = nextResolve; + }); + return { promise, resolve }; +} + +function installMakerApi() { + const listeners = new Set<() => void>(); + const api = { + listAvailableAgents: vi.fn<() => Promise>(), + onAgentsChanged: vi.fn((listener: () => void) => { + listeners.add(listener); + return () => listeners.delete(listener); + }), + }; + (window as unknown as { electronAPI: { maker: typeof api } }).electronAPI = { maker: api }; + return { api, listeners }; +} + +function installDeviceLinkApi() { + const presenceListeners = new Set(); + const statusListeners = new Set(); + const api = { + invoke: vi.fn<(deviceId: string, channel: string, args: unknown[]) => Promise>(), + onPresenceChanged: vi.fn((listener: PresenceListener) => { + presenceListeners.add(listener); + return () => presenceListeners.delete(listener); + }), + onStatusChanged: vi.fn((listener: StatusListener) => { + statusListeners.add(listener); + return () => statusListeners.delete(listener); + }), + onRemotePush: vi.fn(() => () => {}), + }; + (window as unknown as { electronAPI: { deviceLink: typeof api } }).electronAPI = { + deviceLink: api, + }; + return { api, presenceListeners, statusListeners }; +} + +describe('useAvailableAgents roster cache', () => { + beforeEach(() => { + vi.resetModules(); + vi.clearAllMocks(); + delete (window as unknown as { electronAPI?: unknown }).electronAPI; + }); + + it('ignores a pre-change roster response that resolves after the change push', async () => { + const first = deferred(); + const second = deferred(); + const { api, listeners } = installMakerApi(); + api.listAvailableAgents.mockReturnValueOnce(first.promise).mockReturnValueOnce(second.promise); + + const { useAvailableAgents } = await import('../useAvailableAgents'); + const { result, unmount } = renderHook(() => useAvailableAgents()); + + await waitFor(() => expect(api.onAgentsChanged).toHaveBeenCalledTimes(1)); + act(() => { + for (const listener of listeners) listener(); + }); + expect(api.listAvailableAgents).toHaveBeenCalledTimes(2); + + await act(async () => { + first.resolve(['claude-code', 'codex']); + await first.promise; + }); + expect(result.current.availableVendors.has('pi')).toBe(false); + unmount(); + const remounted = renderHook(() => useAvailableAgents()); + expect(remounted.result.current.loaded).toBe(false); + + await act(async () => { + second.resolve(['claude-code', 'codex', 'pi']); + await second.promise; + }); + await waitFor(() => expect(remounted.result.current.availableVendors.has('pi')).toBe(true)); + }); + + it('refetches a remote roster after the selected device reconnects', async () => { + const first = deferred(); + const second = deferred(); + const { api, presenceListeners } = installDeviceLinkApi(); + api.invoke.mockReturnValueOnce(first.promise).mockReturnValueOnce(second.promise); + + const { useAvailableAgents } = await import('../useAvailableAgents'); + const { result } = renderHook(() => useAvailableAgents('device-1')); + + await waitFor(() => expect(api.invoke).toHaveBeenCalledTimes(1)); + await act(async () => { + first.resolve(['claude-code', 'codex']); + await first.promise; + }); + expect(result.current.availableVendors.has('pi')).toBe(false); + + act(() => { + for (const listener of presenceListeners) listener({ deviceId: 'device-1', online: false }); + for (const listener of presenceListeners) listener({ deviceId: 'device-1', online: true }); + }); + await waitFor(() => expect(api.invoke).toHaveBeenCalledTimes(2)); + + await act(async () => { + second.resolve(['claude-code', 'codex', 'pi']); + await second.promise; + }); + await waitFor(() => expect(result.current.availableVendors.has('pi')).toBe(true)); + }); + + it('shares one invalidation and result across mounted consumers', async () => { + const first = deferred(); + const second = deferred(); + const { api, listeners } = installMakerApi(); + api.listAvailableAgents.mockReturnValueOnce(first.promise).mockReturnValueOnce(second.promise); + + const { useAvailableAgents } = await import('../useAvailableAgents'); + const firstHook = renderHook(() => useAvailableAgents()); + const secondHook = renderHook(() => useAvailableAgents()); + await waitFor(() => expect(api.listAvailableAgents).toHaveBeenCalledTimes(1)); + + await act(async () => { + first.resolve(['claude-code', 'codex']); + await first.promise; + }); + act(() => { + for (const listener of listeners) listener(); + }); + expect(api.listAvailableAgents).toHaveBeenCalledTimes(2); + + await act(async () => { + second.resolve(['claude-code', 'codex', 'pi']); + await second.promise; + }); + await waitFor(() => { + expect(firstHook.result.current.availableVendors.has('pi')).toBe(true); + expect(secondHook.result.current.availableVendors.has('pi')).toBe(true); + }); + }); + + it('keeps the source subscription alive while no roster consumer is mounted', async () => { + const first = deferred(); + const second = deferred(); + const { api, listeners } = installMakerApi(); + api.listAvailableAgents.mockReturnValueOnce(first.promise).mockReturnValueOnce(second.promise); + + const { useAvailableAgents } = await import('../useAvailableAgents'); + const firstHook = renderHook(() => useAvailableAgents()); + await waitFor(() => expect(api.listAvailableAgents).toHaveBeenCalledTimes(1)); + await act(async () => { + first.resolve(['claude-code', 'codex']); + await first.promise; + }); + firstHook.unmount(); + + act(() => { + for (const listener of listeners) listener(); + }); + expect(api.listAvailableAgents).toHaveBeenCalledTimes(1); + + const remounted = renderHook(() => useAvailableAgents()); + await waitFor(() => expect(api.listAvailableAgents).toHaveBeenCalledTimes(2)); + await act(async () => { + second.resolve(['claude-code', 'codex', 'pi']); + await second.promise; + }); + await waitFor(() => expect(remounted.result.current.availableVendors.has('pi')).toBe(true)); + }); +}); diff --git a/apps/desktop/src/renderer/hooks/useAvailableAgents.ts b/apps/desktop/src/renderer/hooks/useAvailableAgents.ts index 27fa928ab1..a74970aa2b 100644 --- a/apps/desktop/src/renderer/hooks/useAvailableAgents.ts +++ b/apps/desktop/src/renderer/hooks/useAvailableAgents.ts @@ -17,9 +17,36 @@ import { useEffect, useState } from 'react'; import type { MakerVendor } from '@/lib/ccAgent.types'; import { createLogger } from '@/lib/logger'; +import { + evictDeviceCapabilities, + prefetchDeviceCapabilities, + refreshLocalCapabilities, +} from './useAgentCapabilities'; const log = createLogger('useAvailableAgents'); +let localCapabilitiesRefreshInFlight: Promise | null = null; +const remoteCapabilitiesRefreshInFlight = new Map>(); + +function refreshLocalCapabilitiesOnce(): void { + if (localCapabilitiesRefreshInFlight) return; + const pending = refreshLocalCapabilities().finally(() => { + if (localCapabilitiesRefreshInFlight === pending) localCapabilitiesRefreshInFlight = null; + }); + localCapabilitiesRefreshInFlight = pending; +} + +function refreshRemoteCapabilitiesOnce(deviceId: string): void { + if (remoteCapabilitiesRefreshInFlight.has(deviceId)) return; + evictDeviceCapabilities(deviceId); + const pending = prefetchDeviceCapabilities(deviceId).finally(() => { + if (remoteCapabilitiesRefreshInFlight.get(deviceId) === pending) { + remoteCapabilitiesRefreshInFlight.delete(deviceId); + } + }); + remoteCapabilitiesRefreshInFlight.set(deviceId, pending); +} + type RuntimeAgentKind = 'claude-code' | 'codex' | 'pi'; /** runtime agent id → NewMaker vendor(其余保持同名)。 */ @@ -29,9 +56,15 @@ function toVendor(agent: RuntimeAgentKind): MakerVendor { interface MakerApiShape { listAvailableAgents: () => Promise; + onAgentsChanged: (cb: () => void) => () => void; } interface DeviceLinkShape { invoke: (deviceId: string, channel: string, args: unknown[]) => Promise; + onPresenceChanged?: (cb: (snapshot: { deviceId: string; online: boolean }) => void) => () => void; + onStatusChanged?: (cb: (payload: { status: 'stopped' | 'connecting' | 'online' }) => void) => () => void; + onRemotePush?: ( + cb: (payload: { deviceId: string; channel: string; payload: unknown }) => void, + ) => () => void; } function getMakerApi(): MakerApiShape | null { @@ -73,6 +106,92 @@ interface AgentsCacheEntry { } const agentsCache = new Map(); const inFlight = new Map>>(); +/** + * 每个 roster 缓存 key 的失效代际。收到 roster push 后递增,使 push 之前发起的 + * 在途请求不能把旧的 agent 集合重新写回缓存或覆盖当前 hook 状态。 + */ +const agentsCacheGeneration = new Map(); +/** 同一 roster push 会同步通知多个 mounted consumer;同一 tick 只失效一次。 */ +const agentsCacheInvalidationScheduled = new Set(); +const rosterListeners = new Map void>>(); +let localRosterSourceInstalled = false; +let remoteRosterSourceInstalled = false; +let relayStatus: 'stopped' | 'connecting' | 'online' | null = null; +const deviceOnline = new Map(); + +function currentAgentsCacheGeneration(key: string): number { + return agentsCacheGeneration.get(key) ?? 0; +} + +function invalidateAgentsCache(key: string): boolean { + if (agentsCacheInvalidationScheduled.has(key)) return false; + agentsCacheInvalidationScheduled.add(key); + queueMicrotask(() => agentsCacheInvalidationScheduled.delete(key)); + const generation = currentAgentsCacheGeneration(key) + 1; + agentsCacheGeneration.set(key, generation); + agentsCache.delete(key); + inFlight.delete(key); + return true; +} + +function remoteRosterKeys(): Set { + const keys = new Set([ + ...agentsCache.keys(), + ...inFlight.keys(), + ...rosterListeners.keys(), + ]); + keys.delete(''); + return keys; +} + +function notifyRosterChanged(deviceId?: string): void { + const key = cacheKeyOf(deviceId); + if (!invalidateAgentsCache(key)) return; + if (deviceId) refreshRemoteCapabilitiesOnce(deviceId); + else refreshLocalCapabilitiesOnce(); + for (const listener of rosterListeners.get(key) ?? []) listener(); +} + +function ensureRosterSource(key: string): void { + if (key === '') { + if (localRosterSourceInstalled) return; + const api = getMakerApi(); + if (!api?.onAgentsChanged) return; + localRosterSourceInstalled = true; + api.onAgentsChanged(() => notifyRosterChanged()); + return; + } + if (remoteRosterSourceInstalled) return; + const dl = getDeviceLink(); + if (!dl) return; + remoteRosterSourceInstalled = true; + dl.onRemotePush?.((push) => { + if (push.channel === 'maker:agents:changed') notifyRosterChanged(push.deviceId); + }); + dl.onPresenceChanged?.((snapshot) => { + const wasOnline = deviceOnline.get(snapshot.deviceId); + deviceOnline.set(snapshot.deviceId, snapshot.online); + if (snapshot.online && wasOnline !== true) notifyRosterChanged(snapshot.deviceId); + }); + dl.onStatusChanged?.(({ status }) => { + const wasOnline = relayStatus; + relayStatus = status; + if (status === 'online' && wasOnline !== 'online') { + for (const deviceId of remoteRosterKeys()) notifyRosterChanged(deviceId); + } + }); +} + +function subscribeRosterChanges(key: string, listener: () => void): () => void { + const bucket = rosterListeners.get(key) ?? new Set<() => void>(); + bucket.add(listener); + rosterListeners.set(key, bucket); + ensureRosterSource(key); + return () => { + bucket.delete(listener); + if (bucket.size === 0) rosterListeners.delete(key); + }; +} function cacheKeyOf(deviceId?: string | null): string { return deviceId ?? ''; @@ -82,14 +201,19 @@ function loadAvailableAgents(deviceId?: string | null): Promise { const vendors: ReadonlySet = new Set(agents.map(toVendor)); - agentsCache.set(key, { vendors, fetchedAt: Date.now() }); + // A roster push can invalidate this request while the IPC call is pending. Do not + // let that pre-change response become the authoritative cache entry. + if (currentAgentsCacheGeneration(key) === requestGeneration) { + agentsCache.set(key, { vendors, fetchedAt: Date.now() }); + } return vendors; }) .finally(() => { - inFlight.delete(key); + if (inFlight.get(key) === promise) inFlight.delete(key); }); inFlight.set(key, promise); return promise; @@ -119,10 +243,12 @@ export function useAvailableAgents(deviceId?: string | null): UseAvailableAgents const hit = agentsCache.get(cacheKeyOf(deviceId)); setAvailableVendors(hit?.vendors ?? new Set()); setLoaded(hit !== undefined); + const key = cacheKeyOf(deviceId); const run = (): void => { + const requestGeneration = currentAgentsCacheGeneration(key); loadAvailableAgents(deviceId) .then((vendors) => { - if (cancelled) return; + if (cancelled || currentAgentsCacheGeneration(key) !== requestGeneration) return; setAvailableVendors(vendors); setLoaded(true); }) @@ -136,6 +262,7 @@ export function useAvailableAgents(deviceId?: string | null): UseAvailableAgents }); }; run(); + const offRosterChanged = subscribeRosterChanges(key, run); // 会话期间 Pi 二进制可能被按需下载补齐:窗口重新聚焦时再拉一次,让入口及时出现。 // 节流:补齐是分钟级的事,秒级来回切窗口不必反复打 IPC(远程还要过隧道)。 const onFocus = (): void => { @@ -147,6 +274,7 @@ export function useAvailableAgents(deviceId?: string | null): UseAvailableAgents return () => { cancelled = true; window.removeEventListener('focus', onFocus); + offRosterChanged(); }; }, [deviceId]); @@ -155,6 +283,14 @@ export function useAvailableAgents(deviceId?: string | null): UseAvailableAgents /** 测试用 —— 清进程内缓存与在途请求(其它代码不应调用)。 */ export function __resetAvailableAgentsCacheForTest(): void { + for (const key of new Set([ + ...agentsCache.keys(), + ...inFlight.keys(), + ...agentsCacheGeneration.keys(), + ])) { + invalidateAgentsCache(key); + } agentsCache.clear(); inFlight.clear(); + agentsCacheInvalidationScheduled.clear(); } diff --git a/apps/desktop/src/renderer/vite-env.d.ts b/apps/desktop/src/renderer/vite-env.d.ts index 1577820a41..85c5eebdcd 100644 --- a/apps/desktop/src/renderer/vite-env.d.ts +++ b/apps/desktop/src/renderer/vite-env.d.ts @@ -4734,6 +4734,7 @@ interface ElectronAPI { */ maker: { listAvailableAgents: () => Promise>; + onAgentsChanged: (cb: () => void) => () => void; getCapabilities: (agentKind: 'claude-code' | 'codex' | 'pi') => Promise; /** workflow 逐 agent 进度树(只读);读不到 / 解析失败返回 null → 回退 workflow 级卡片。 */ getWorkflowProgress: ( diff --git a/apps/mobile/app/sessions/new.tsx b/apps/mobile/app/sessions/new.tsx index 2812b30ef9..1e0e33a983 100644 --- a/apps/mobile/app/sessions/new.tsx +++ b/apps/mobile/app/sessions/new.tsx @@ -103,6 +103,7 @@ import { import { buildAgentCapabilitiesCacheKey, commitAgentCapabilities, + evictAgentCapabilitiesForDevice, getAgentCapabilitiesGeneration, getCachedAgentCapabilities, isAgentCapabilitiesGenerationCurrent, @@ -428,6 +429,8 @@ export default function NewRemoteSessionScreen() { invoke, openLink, subscribe, + unsubscribe, + onAgentsChanged, status: deviceLinkStatus, connectionEpoch, presenceVersion, @@ -515,6 +518,13 @@ export default function NewRemoteSessionScreen() { // 建出最终 requireAgent 报 not-registered 的会话(codex review P2)。 const [availableAgentKinds, setAvailableAgentKinds] = useState | null>(null); + const [availableAgentRefreshNonce, setAvailableAgentRefreshNonce] = useState(0); + const [availableAgentRosterRefreshNonce, setAvailableAgentRosterRefreshNonce] = useState(0); + const rosterRecoveryIdentityRef = useRef<{ + deviceId: string; + connectionEpoch: number; + presenceVersion: number; + } | null>(null); // worktree 开关(project 模式 + 已选目录时显示):勾选值存工作端(get-new-maker-defaults // 播种 / 显式点击写穿),资格由 worktree:detect-cwd 探测(目录变化即重探,seq 防竞态)。 const [worktreeProbe, setWorktreeProbe] = useState(null); @@ -1750,7 +1760,7 @@ export default function NewRemoteSessionScreen() { cancelled = true; unsubscribe(); }; - }, [selectedDeviceId, draft.agentKind, maker, openLink]); + }, [availableAgentRefreshNonce, selectedDeviceId, draft.agentKind, maker, openLink]); // Fast 记忆延迟恢复(codex review P1):切/恢复 agent 的瞬间,目标 agent 的能力表 // 尚未到达(或残留着切换前 agent 的),恢复点只能保守置 false;真正的恢复在这里—— @@ -1845,7 +1855,67 @@ export default function NewRemoteSessionScreen() { return () => { cancelled = true; }; - }, [selectedDeviceId, maker, openLink]); + }, [selectedDeviceId, maker, openLink, availableAgentRosterRefreshNonce]); + + useEffect(() => { + if (!selectedDeviceId) { + rosterRecoveryIdentityRef.current = null; + return; + } + const previous = rosterRecoveryIdentityRef.current; + rosterRecoveryIdentityRef.current = { deviceId: selectedDeviceId, connectionEpoch, presenceVersion }; + if (!previous || previous.deviceId !== selectedDeviceId) return; + if ( + previous.connectionEpoch === connectionEpoch + && previous.presenceVersion === presenceVersion + ) return; + // Roster pushes are edge-triggered and are not replayed by topic rehydration. + // Refresh both the runtime roster and the selected agent capabilities after a + // relay/target recovery, even when the screen stayed mounted and focused. + setAvailableAgentRefreshNonce((value) => value + 1); + setAvailableAgentRosterRefreshNonce((value) => value + 1); + }, [connectionEpoch, presenceVersion, selectedDeviceId]); + + useEffect(() => { + if (!selectedDeviceId) return; + let cancelled = false; + void withTransientRemoteRetry(async () => { + await openLink(selectedDeviceId); + if (cancelled) return; + await subscribe(`new-session:${selectedDeviceId}`, selectedDeviceId, ['sessions']); + }).catch(() => { + /* The existing roster fetch remains fail-open; the next focus/retry can resubscribe. */ + }); + return () => { + cancelled = true; + void unsubscribe(`new-session:${selectedDeviceId}`, selectedDeviceId, ['sessions']).catch(() => undefined); + }; + }, [openLink, selectedDeviceId, subscribe, unsubscribe]); + + useEffect(() => { + if (!selectedDeviceId) return; + let cancelled = false; + void withTransientRemoteRetry(async () => { + await openLink(selectedDeviceId); + if (cancelled) return; + await subscribe(`new-session:${selectedDeviceId}`, selectedDeviceId, ['sessions']); + }).catch(() => { + /* The mount-time owner remains held; the next connection/presence change retries. */ + }); + return () => { + cancelled = true; + }; + }, [connectionEpoch, openLink, presenceVersion, selectedDeviceId, subscribe]); + + useEffect(() => { + if (!selectedDeviceId) return; + return onAgentsChanged((deviceId) => { + if (deviceId !== selectedDeviceId) return; + evictAgentCapabilitiesForDevice(deviceId); + setAvailableAgentRefreshNonce((value) => value + 1); + setAvailableAgentRosterRefreshNonce((value) => value + 1); + }); + }, [onAgentsChanged, selectedDeviceId]); useEffect(() => { if (!selectedDeviceId || composerTrigger.kind !== 'slash') { diff --git a/apps/mobile/src/__tests__/newSession.test.ts b/apps/mobile/src/__tests__/newSession.test.ts index 9dc5d3c5db..6ead7e4d98 100644 --- a/apps/mobile/src/__tests__/newSession.test.ts +++ b/apps/mobile/src/__tests__/newSession.test.ts @@ -1154,6 +1154,12 @@ describe('new session model', () => { const newSource = readTextLf(resolve(process.cwd(), 'app/sessions/new.tsx'), 'utf8'); // 拉被控端 runtime 注册集合,渲染按可用集过滤,选中不可用时 coerce。 expect(newSource).toContain('maker.listAvailableAgents()'); + expect(newSource).toContain('onAgentsChanged'); + expect(newSource).toContain('evictAgentCapabilitiesForDevice(deviceId);'); + expect(newSource).toContain('const rosterRecoveryIdentityRef = useRef<'); + expect(newSource).toContain('setAvailableAgentRosterRefreshNonce((value) => value + 1);'); + expect(newSource).toContain('previous.connectionEpoch === connectionEpoch'); + expect(newSource.match(/subscribe\(`new-session:\$\{selectedDeviceId\}`/g)?.length ?? 0).toBeGreaterThanOrEqual(2); expect(newSource).toContain('availableNewSessionAgentOptions(availableAgentKinds).map'); expect(newSource).toMatch(/availableAgentKinds\.has\(draft\.agentKind\)/); // 传输层 passthrough 到 allowlisted channel。 diff --git a/apps/mobile/src/debug/visualMock.ts b/apps/mobile/src/debug/visualMock.ts index e02b2c42f7..37ee75ea3f 100644 --- a/apps/mobile/src/debug/visualMock.ts +++ b/apps/mobile/src/debug/visualMock.ts @@ -162,6 +162,7 @@ export function createVisualMockDeviceLinkContext(): DeviceLinkContextValue { invoke: visualMockInvoke, subscribe: async () => undefined, unsubscribe: async () => undefined, + onAgentsChanged: () => () => undefined, }; } diff --git a/apps/mobile/src/device-link/DeviceLinkContext.tsx b/apps/mobile/src/device-link/DeviceLinkContext.tsx index de9973d6b7..12c0df82ed 100644 --- a/apps/mobile/src/device-link/DeviceLinkContext.tsx +++ b/apps/mobile/src/device-link/DeviceLinkContext.tsx @@ -159,6 +159,8 @@ export interface DeviceLinkContextValue { // when its last owner unsubscribes. subscribe(owner: string, deviceId: string, topics: string[]): Promise; unsubscribe(owner: string, deviceId: string, topics: string[]): Promise; + /** 被控端 runtime Agent roster 发生变化时通知当前控制端页面。 */ + onAgentsChanged: (listener: (deviceId: string) => void) => () => void; } const DeviceLinkContext = createContext(null); @@ -178,6 +180,7 @@ const CONTROLLER_CAPABILITIES = [ // 响应性熔断;只用于判定并发返回的 unavailable 是否已被更晚目标应答推翻。 const remoteResponseEvidenceEpochs = createPresenceAvailabilityEpochs(); const remoteResponseEvidenceListeners = new Set<(deviceId: string) => void>(); +const remoteAgentRosterListeners = new Set<(deviceId: string) => void>(); // 永久 link-close 后被抑制后台重建的设备(见 updateRehydrateSuppressionOnLinkClose)。 // 模块级(与 remoteResponseEvidenceEpochs 同模式):sendOpenLink 等模块级函数也需要 @@ -197,6 +200,13 @@ function subscribeRemoteResponseEvidence( return () => remoteResponseEvidenceListeners.delete(listener); } +function subscribeRemoteAgentRoster( + listener: (deviceId: string) => void, +): () => void { + remoteAgentRosterListeners.add(listener); + return () => remoteAgentRosterListeners.delete(listener); +} + interface RehydrateState { inFlight: Promise | null; rerun: boolean; @@ -896,6 +906,9 @@ export function DeviceLinkProvider({ children }: { children: ReactNode }) { .catch(() => { /* 下次进入选择器或重连补齐时继续重试。 */ }); void refreshDeviceCapabilities(client, deviceId); }, + onAgentsChanged: (deviceId) => { + for (const listener of remoteAgentRosterListeners) listener(deviceId); + }, })); // 与 transport-timeout link-close 同族的链路死锁自救(互为兜底):对端还在按 // 可靠流给本机发帧,而本机侧 link 未就绪——典型成因是 link-accept 在弱网丢失 @@ -1160,6 +1173,7 @@ export function DeviceLinkProvider({ children }: { children: ReactNode }) { invoke, subscribe, unsubscribe, + onAgentsChanged: subscribeRemoteAgentRoster, }), [ closeLink, connectionEpoch, @@ -1173,6 +1187,7 @@ export function DeviceLinkProvider({ children }: { children: ReactNode }) { status, subscribe, unsubscribe, + subscribeRemoteAgentRoster, ]); return {children}; @@ -1191,6 +1206,7 @@ export function routeFrame(env: Envelope, handlers: { onAccessRevoked?: (deviceId: string) => void; onLinkClosed?: (deviceId: string, reason?: string) => void; onProviderChanged?: (deviceId: string) => void; + onAgentsChanged?: (deviceId: string) => void; } = {}): void { const peerLinkClosed = handlePeerLinkCloseFrame( env, @@ -1207,6 +1223,10 @@ export function routeFrame(env: Envelope, handlers: { handlers.onProviderChanged?.(env.src); return; } + if (push.channel === 'maker:agents:changed') { + handlers.onAgentsChanged?.(env.src); + return; + } if (push.channel === 'maker:schedule:event') { remoteScheduleEventStore.apply(env.src, push.payload); } diff --git a/packages/device-link/src/__tests__/allowlist.test.ts b/packages/device-link/src/__tests__/allowlist.test.ts index 2cb0edc120..2b30cbca3e 100644 --- a/packages/device-link/src/__tests__/allowlist.test.ts +++ b/packages/device-link/src/__tests__/allowlist.test.ts @@ -346,6 +346,7 @@ describe('PUSH_FORWARD_ALLOWLIST', () => { 'maker:interaction-dismissed', 'maker:auto-permission:fallback', 'maker:provider:changed', + 'maker:agents:changed', 'maker:schedule:event', 'maker:orca:worker-changed', 'usage:message-turn-cost', diff --git a/packages/device-link/src/__tests__/topics.test.ts b/packages/device-link/src/__tests__/topics.test.ts index 51aeed471e..ef3457c4d1 100644 --- a/packages/device-link/src/__tests__/topics.test.ts +++ b/packages/device-link/src/__tests__/topics.test.ts @@ -34,6 +34,7 @@ describe('topicForPush', () => { it('账号 / 全局级 channel → sessions(随列表订阅走)', () => { expect(topicForPush('maker:provider:changed', { revision: 42 })).toBe('sessions'); + expect(topicForPush('maker:agents:changed', {})).toBe('sessions'); expect(topicForPush('maker:schedule:event', { kind: 'x' })).toBe('sessions'); expect(topicForPush('maker:project-automation:event', {})).toBe('sessions'); // 被控端当前草稿全量变更(无 sessionId)→ 并入 sessions topic。 diff --git a/packages/device-link/src/allowlist.ts b/packages/device-link/src/allowlist.ts index 3b854c4d29..c06bc8bacd 100644 --- a/packages/device-link/src/allowlist.ts +++ b/packages/device-link/src/allowlist.ts @@ -552,6 +552,8 @@ export const REMOTE_INVOKE_ALLOWLIST: ReadonlySet = new Set([ export const PUSH_FORWARD_ALLOWLIST: ReadonlySet = new Set([ // maker-ipc MAKER_PUSH 'maker:event', + // Device-level runtime Agent roster changes; controllers refresh their local availability cache. + 'maker:agents:changed', 'maker:status-changed', 'maker:input:projection', 'maker:interaction-request', diff --git a/packages/device-link/src/topics.ts b/packages/device-link/src/topics.ts index e97ac8e1a1..21982460c7 100644 --- a/packages/device-link/src/topics.ts +++ b/packages/device-link/src/topics.ts @@ -150,6 +150,8 @@ const SESSION_LIST_CHANNELS: ReadonlySet = new Set([ const ACCOUNT_CHANNELS: ReadonlySet = new Set([ // provider 目录是设备级快照;控制端订阅 sessions 后按来源 deviceId 精确刷新。 'maker:provider:changed', + // runtime Agent roster is a device-level snapshot; controllers refresh after this push. + 'maker:agents:changed', 'maker:schedule:event', 'maker:project-automation:event', // 被控端「当前 New Maker 草稿」全量变更:账号 / 全局级(无 sessionId),并入 `sessions` topic diff --git a/packages/maker-core/src/agents/pi/windows-git-path-powershell.test.ts b/packages/maker-core/src/agents/pi/windows-git-path-powershell.test.ts index b24e645f86..9a81514656 100644 --- a/packages/maker-core/src/agents/pi/windows-git-path-powershell.test.ts +++ b/packages/maker-core/src/agents/pi/windows-git-path-powershell.test.ts @@ -221,7 +221,9 @@ describe('Windows Git PATH PowerShell probes', () => { '}', ].join('\n'), ], - { stdio: 'ignore', timeout: 5_000, windowsHide: true }, + // PowerShell startup can exceed five seconds on a busy hosted Windows runner; + // allow cleanup to observe the child process exit without changing the probe budget. + { stdio: 'ignore', timeout: 10_000, windowsHide: true }, ); } finally { if (coordinatorPid) { @@ -238,7 +240,7 @@ describe('Windows Git PATH PowerShell probes', () => { } } }, - 12_000, + 20_000, ); it('reports recoverable script failures only when a logger is supplied', () => { diff --git a/packages/maker-core/src/maker.test.ts b/packages/maker-core/src/maker.test.ts index 03ac3e5724..b4b03292b5 100644 --- a/packages/maker-core/src/maker.test.ts +++ b/packages/maker-core/src/maker.test.ts @@ -91,6 +91,19 @@ describe('Maker agent status', () => { authReady: false, }); }); + + it('registers an optional agent after construction idempotently', () => { + const maker = new Maker({ + agents: {}, + storage: createStorage(), + logger: createLogger(), + }); + const pi = createAgent(async () => undefined, 'pi'); + + expect(maker.registerAgent('pi', pi)).toBe(true); + expect(maker.registerAgent('pi', pi)).toBe(false); + expect(maker.listAvailableAgents()).toEqual(['pi']); + }); }); function createAgent( diff --git a/packages/maker-core/src/maker.ts b/packages/maker-core/src/maker.ts index e4956b0bb9..2881b5bb25 100644 --- a/packages/maker-core/src/maker.ts +++ b/packages/maker-core/src/maker.ts @@ -1101,6 +1101,21 @@ export class Maker { return Object.keys(this.agents) as AgentKind[]; } + /** + * Register an optional agent after Maker construction. + * + * Hosts may provision optional runtimes asynchronously (for example after a + * transient network failure during startup). Registration is intentionally + * additive and idempotent so existing sessions and agent instances remain + * untouched. + */ + registerAgent(kind: AgentKind, agent: BaseAgent): boolean { + if (this.shutdownStarted) return false; + if (this.agents[kind]) return false; + this.agents[kind] = agent; + return true; + } + /** * Agent 内置 command (palette 'agent-builtin' 类目) —— 同步硬编码白名单。 * 见 agents//commands.ts。