From 5425513b2c3789f5c50abe16a50a2a382926abb6 Mon Sep 17 00:00:00 2001 From: Sun-GLiang <81428141+Sun-GLiang@users.noreply.github.com> Date: Tue, 29 Sep 2026 17:13:35 +0800 Subject: [PATCH 01/30] feat(acp): add mode selection and scoped catalog lifecycle Generated-by: Codex --- .../main/__tests__/executor-selection.test.ts | 17 +- .../runtime-host-session-catalog-ipc-main.ts | 5 +- apps/desktop/src/preload/bridge-contract.d.ts | 2 +- apps/desktop/src/preload/preload.ts | 4 +- .../controller/use-executor-selection.ts | 56 ++- .../conversation/model/executor-composer.ts | 2 +- .../conversation/model/executor-submission.ts | 23 +- .../renderer/features/conversation/ports.ts | 2 +- docs/antigravity-acp-plugin-rebuild.md | 14 +- .../archive/antigravity-acp-pr4-acceptance.md | 47 +++ .../src/__tests__/acp-executor-plugin.test.ts | 302 ++++++++++++--- packages/acp-executor-plugin/src/index.ts | 363 +++++++++++++----- .../src/__tests__/executor-catalog.test.ts | 31 +- packages/core/src/executor-catalog.ts | 37 +- .../external-agent-setup-coordinator.test.ts | 25 ++ .../session-catalog-coordinator.test.ts | 52 ++- .../session-catalog-protocol.test.ts | 4 +- packages/runtime-host/src/protocol/index.ts | 4 +- .../src/protocol/plugin-platform.ts | 12 +- .../src/server/execution-composition.ts | 19 +- .../external-agent-setup-coordinator.ts | 3 + .../src/server/session-catalog-coordinator.ts | 47 ++- .../__tests__/plugin-executor-service.test.ts | 6 + .../runtime/src/plugin-executor-service.ts | 33 +- packages/runtime/src/session-manager.ts | 1 + .../__tests__/executor-model-picker.test.tsx | 23 ++ packages/ui/src/executor-model-picker.tsx | 93 ++++- packages/ui/src/model-picker-panel.tsx | 8 +- 28 files changed, 1017 insertions(+), 218 deletions(-) create mode 100644 docs/archive/antigravity-acp-pr4-acceptance.md diff --git a/apps/desktop/src/main/__tests__/executor-selection.test.ts b/apps/desktop/src/main/__tests__/executor-selection.test.ts index 1616973d24e..6a166a0dad5 100644 --- a/apps/desktop/src/main/__tests__/executor-selection.test.ts +++ b/apps/desktop/src/main/__tests__/executor-selection.test.ts @@ -24,7 +24,7 @@ import { act, createElement } from 'react'; import { createRoot } from 'react-dom/client'; import type { ExecutorCatalogEntry } from '@maka/core/executor-catalog'; import type { SessionSummary } from '@maka/core/session'; -import { useExecutorSelection, newTaskConfiguration, ConversationServicesProvider, type ConversationServices } from '../../renderer/features/conversation/index.js'; +import { useExecutorSelection, newTaskConfiguration, executorSubmissionError, ConversationServicesProvider, type ConversationServices } from '../../renderer/features/conversation/index.js'; const entry: ExecutorCatalogEntry = { id: 'external', displayName: 'External', readiness: 'ready', models: [{ id: 'selected', name: 'Selected' }], supportsAttachments: false, supportsModelChange: true }; @@ -246,3 +246,18 @@ test('an executor choice uses its exact model without inheriting native thinking assert.equal(configuration.collaborationMode, 'agent'); assert.equal(configuration.orchestrationMode, 'default'); }); + +test('draft submission keeps the selected mode and blocks a removed catalog choice', () => { + const executorSelection = { executorId: 'external', configuration: { model: 'selected', mode: 'auto' } }; + const configuration = newTaskConfiguration({ + executorSelection, + newChatModel: null, pendingNewChatThinkingLevel: undefined, + newChatPermissionChoice: undefined, newChatCollaborationMode: 'agent', + newChatOrchestrationMode: 'default', + }); + assert.deepEqual(configuration.executorConfig, { model: 'selected', mode: 'auto' }); + assert.match(executorSubmissionError({ + executorSelection, + executorEntry: { ...entry, modes: [{ id: 'ask', name: 'Ask' }] }, + }, 0, 'en') ?? '', /no longer available/u); +}); diff --git a/apps/desktop/src/main/runtime-host-session-catalog-ipc-main.ts b/apps/desktop/src/main/runtime-host-session-catalog-ipc-main.ts index efbaa949484..86b30c7aa97 100644 --- a/apps/desktop/src/main/runtime-host-session-catalog-ipc-main.ts +++ b/apps/desktop/src/main/runtime-host-session-catalog-ipc-main.ts @@ -116,9 +116,10 @@ export function registerRuntimeHostSessionCatalogIpc( const actionIds = (sessionId: string, options: unknown) => resolveSessionActionIds(() => listSessions(), sessionId, options); - handleReconnectableRead(ipcMain, 'sessions:executorCatalog', async (_event, cwd: string) => { + handleReconnectableRead(ipcMain, 'sessions:executorCatalog', async (_event, cwd: string, refresh?: boolean) => { if (typeof cwd !== 'string' || !cwd) throw new Error('Executor discovery requires a workspace'); - return (await deps.queryExecutors?.({ kind: 'catalog', cwd }))?.items ?? []; + if (refresh !== undefined && typeof refresh !== 'boolean') throw new Error('Invalid executor refresh flag'); + return (await deps.queryExecutors?.({ kind: 'catalog', cwd, ...(refresh ? { refresh: true } : {}) }))?.items ?? []; }); handleReconnectableRead(ipcMain, 'sessions:executorState', async (_event, sessionId: string) => (await deps.queryExecutors?.({ kind: 'conversation', sessionId }))?.items ?? [], diff --git a/apps/desktop/src/preload/bridge-contract.d.ts b/apps/desktop/src/preload/bridge-contract.d.ts index d3db91690bd..0ee6008b55e 100644 --- a/apps/desktop/src/preload/bridge-contract.d.ts +++ b/apps/desktop/src/preload/bridge-contract.d.ts @@ -1019,7 +1019,7 @@ export interface MakaBridge { }; newTasks: { - getExecutors(target: DesktopNewTaskTarget, cwd: string): Promise; + getExecutors(target: DesktopNewTaskTarget, cwd: string, refresh?: boolean): Promise; getCatalog(): Promise; subscribeChanges(handler: () => void): () => void; addProject(host: DesktopNewTaskHostRef, name?: string): Promise< diff --git a/apps/desktop/src/preload/preload.ts b/apps/desktop/src/preload/preload.ts index 58a87563bef..c73cc2ef767 100644 --- a/apps/desktop/src/preload/preload.ts +++ b/apps/desktop/src/preload/preload.ts @@ -1925,8 +1925,8 @@ const makaBridge = { }, }, newTasks: { - async getExecutors(target, cwd) { - return ipcRenderer.invoke('sessions:executorCatalog', await runtimeHostScope(target), cwd); + async getExecutors(target, cwd, refresh) { + return ipcRenderer.invoke('sessions:executorCatalog', await runtimeHostScope(target), cwd, refresh); }, getCatalog(): Promise { return loadNewTaskCatalog(); diff --git a/apps/desktop/src/renderer/features/conversation/controller/use-executor-selection.ts b/apps/desktop/src/renderer/features/conversation/controller/use-executor-selection.ts index 587594f3d88..04f9825615f 100644 --- a/apps/desktop/src/renderer/features/conversation/controller/use-executor-selection.ts +++ b/apps/desktop/src/renderer/features/conversation/controller/use-executor-selection.ts @@ -41,15 +41,28 @@ export function useExecutorSelection(input: { const inFlight = useRef(undefined); const sessionId = input.session?.id; const executorId = input.session?.executorId; - const key = sessionId ?? input.key; + const key = sessionId ?? JSON.stringify([ + input.key, + input.target?.hostId, + input.target?.profileId, + input.target?.projectId, + input.cwd, + ]); const current = useRef(key); current.current = key; const revision = useRef(0); const refreshes = useRef(new Map>()); const pendingInvalidations = useRef(new Set()); - const refresh = useCallback((): Promise => { + const pendingForce = useRef(new Set()); + const refresh = useCallback((force = false): Promise => { const existing = refreshes.current.get(key); - if (existing) return existing; + if (existing) { + if (force) { + pendingInvalidations.current.add(key); + pendingForce.current.add(key); + } + return existing; + } const run = (async () => { const attempt = ++revision.current; if (sessionId && !executorId) { @@ -65,7 +78,7 @@ export function useExecutorSelection(input: { try { const catalog = sessionId ? ((await services.sessions.getExecutorState?.(sessionId)) ?? []) - : ((await services.newTasks.getExecutors?.(input.target!, input.cwd!)) ?? []); + : ((await services.newTasks.getExecutors?.(input.target!, input.cwd!, force)) ?? []); if (current.current === key && revision.current === attempt) setSnapshot({ key, catalog, loading: false }); } catch (error) { @@ -82,7 +95,8 @@ export function useExecutorSelection(input: { tracked = run.finally(() => { if (refreshes.current.get(key) !== tracked) return; refreshes.current.delete(key); - if (pendingInvalidations.current.delete(key) && current.current === key) void refresh(); + const forced = pendingForce.current.delete(key); + if (pendingInvalidations.current.delete(key) && current.current === key) void refresh(forced); }); refreshes.current.set(key, tracked); return tracked; @@ -91,6 +105,7 @@ export function useExecutorSelection(input: { sessionId, executorId, input.session?.executorConfig?.model, + input.session?.executorConfig?.mode, input.target?.hostId, input.target?.profileId, input.target?.projectId, @@ -122,6 +137,7 @@ export function useExecutorSelection(input: { revision.current++; refreshes.current.delete(key); pendingInvalidations.current.delete(key); + pendingForce.current.delete(key); unsubscribe(); unSession(); clearTimeout(timer); @@ -133,19 +149,24 @@ export function useExecutorSelection(input: { const catalog = snapshot?.key === key ? snapshot.catalog : []; const inspected = catalog.find(candidate => candidate.id === executorId); const selection = executorId - ? { executorId, configuration: inspected?.readiness === 'ready' && inspected.currentModel - ? { model: inspected.currentModel } : input.session?.executorConfig ?? {} } + ? { executorId, configuration: inspected?.readiness === 'ready' + ? { ...input.session?.executorConfig, + ...(inspected.currentModel ? { model: inspected.currentModel } : {}), + ...(inspected.currentMode ? { mode: inspected.currentMode } : {}) } + : input.session?.executorConfig ?? {} } : sessionId ? undefined - : draft?.key === input.key + : draft?.key === key ? draft.selection : undefined; const select = async (next: ExecutorSelection | undefined) => { if (inFlight.current === key) throw new Error('Executor configuration is pending'); if (!sessionId) { if (next && !catalog.some(entry => entry.id === next.executorId && entry.readiness === 'ready' && - entry.models.some(model => model.id === next.configuration.model))) throw new Error('Executor model is unavailable'); - setDraft({ key: input.key, selection: next }); + (!next.configuration.model || entry.models.some(model => model.id === next.configuration.model)) && + (!next.configuration.mode || entry.modes?.some(mode => mode.id === next.configuration.mode)))) + throw new Error('Executor configuration is unavailable'); + setDraft({ key, selection: next }); return; } if (!next || next.executorId !== executorId || !services.sessions.setExecutorModelConfiguration) @@ -158,14 +179,15 @@ export function useExecutorSelection(input: { next.configuration, ); if (!result.ok) throw new Error(result.code); - if (result.session.executorConfig?.model !== next.configuration.model) - throw new Error('Executor model change was not confirmed'); + if ((next.configuration.model && result.session.executorConfig?.model !== next.configuration.model) || + (next.configuration.mode && result.session.executorConfig?.mode !== next.configuration.mode)) + throw new Error('Executor configuration change was not confirmed'); if (current.current === key) { revision.current++; setSnapshot(previous => ({ key, loading: false, catalog: (previous?.key === key ? previous.catalog : []).map(entry => entry.id === executorId - ? { ...entry, currentModel: result.session.executorConfig!.model } : entry), + ? { ...entry, currentModel: result.session.executorConfig?.model, currentMode: result.session.executorConfig?.mode } : entry), })); await refresh(); } @@ -188,8 +210,10 @@ export function useExecutorSelection(input: { const restore = async () => { if (!sessionId || !executorId) throw new Error('Executor Session is unavailable'); if (inFlight.current === key) throw new Error('Executor configuration is pending'); - const model = input.session?.executorConfig?.model ?? inspected?.currentModel; - if (!model) throw new Error('Executor model is unavailable'); + const configuration = input.session?.executorConfig ?? { + ...(inspected?.currentModel ? { model: inspected.currentModel } : {}), + ...(inspected?.currentMode ? { mode: inspected.currentMode } : {}), + }; setSnapshot((previous) => previous?.key === key ? { @@ -200,7 +224,7 @@ export function useExecutorSelection(input: { } : previous, ); - await select({ executorId, configuration: { model } }); + await select({ executorId, configuration }); }; return { selection, diff --git a/apps/desktop/src/renderer/features/conversation/model/executor-composer.ts b/apps/desktop/src/renderer/features/conversation/model/executor-composer.ts index 9be292d6207..1a30e9b0593 100644 --- a/apps/desktop/src/renderer/features/conversation/model/executor-composer.ts +++ b/apps/desktop/src/renderer/features/conversation/model/executor-composer.ts @@ -45,7 +45,7 @@ export function executorComposerProps( onSelect: (selection) => executor.select(selection), onRestore: () => executor.restore(), onRetry: () => { - void executor.refresh(); + void executor.refresh(true); }, onSetup: input.onSetup, onNewTask: input.onNewTask, diff --git a/apps/desktop/src/renderer/features/conversation/model/executor-submission.ts b/apps/desktop/src/renderer/features/conversation/model/executor-submission.ts index ce9c474e1ab..58404f0e023 100644 --- a/apps/desktop/src/renderer/features/conversation/model/executor-submission.ts +++ b/apps/desktop/src/renderer/features/conversation/model/executor-submission.ts @@ -27,36 +27,41 @@ import type { NewChatModel } from './shell-chat-model-selection.js'; export interface ExecutorSubmission { executorSelection?: ExecutorSelection; - executorEntry?: Pick; + executorEntry?: Pick & + Partial>; } const SUBMISSION_COPY = { en: { attachments: 'Remove unsupported attachments or select Maka. Your draft is preserved.', unavailable: 'Check External Agents settings or start a new task.', + invalid: 'The selected Agent configuration is no longer available. Choose again.', }, 'zh-CN': { attachments: '请移除不支持的附件或选择 Maka,草稿会保留。', unavailable: '请检查外部 Agent 设置,或新建任务。', + invalid: '所选 Agent 配置已失效,请重新选择。', }, 'zh-TW': { attachments: '請移除不支援的附件或選擇 Maka,草稿會保留。', unavailable: '請檢查外部 Agent 設定,或建立新任務。', + invalid: '所選 Agent 設定已失效,請重新選擇。', }, -} satisfies UiCatalog<{ attachments: string; unavailable: string }>; +} satisfies UiCatalog<{ attachments: string; unavailable: string; invalid: string }>; export function executorSubmissionError( input: ExecutorSubmission, attachments: number, locale: UiLocale, ): string | undefined { - if ( - !input.executorSelection || - (input.executorEntry?.readiness === 'ready' && - (!attachments || input.executorEntry.supportsAttachments)) - ) - return; - return SUBMISSION_COPY[locale][attachments ? 'attachments' : 'unavailable']; + if (!input.executorSelection) return; + const entry = input.executorEntry; + if (!entry || entry.readiness !== 'ready') return SUBMISSION_COPY[locale].unavailable; + if (attachments && !entry.supportsAttachments) return SUBMISSION_COPY[locale].attachments; + const { model, mode } = input.executorSelection.configuration; + if ((model && entry.models && !entry.models.some(candidate => candidate.id === model)) || + (mode && !entry.modes?.some(candidate => candidate.id === mode))) + return SUBMISSION_COPY[locale].invalid; } export function newTaskConfiguration( diff --git a/apps/desktop/src/renderer/features/conversation/ports.ts b/apps/desktop/src/renderer/features/conversation/ports.ts index 83ffa98476b..8e9f96da01c 100644 --- a/apps/desktop/src/renderer/features/conversation/ports.ts +++ b/apps/desktop/src/renderer/features/conversation/ports.ts @@ -88,7 +88,7 @@ export interface ConversationServices extends Pick< ): Promise; }; readonly newTasks: { - getExecutors?(target: ConversationNewTaskTarget, cwd: string): Promise; + getExecutors?(target: ConversationNewTaskTarget, cwd: string, refresh?: boolean): Promise; subscribeChanges(handler: () => void): () => void; listInvocableSkills( target: ConversationNewTaskTarget, diff --git a/docs/antigravity-acp-plugin-rebuild.md b/docs/antigravity-acp-plugin-rebuild.md index edc33aa862e..6a6b1f08936 100644 --- a/docs/antigravity-acp-plugin-rebuild.md +++ b/docs/antigravity-acp-plugin-rebuild.md @@ -157,11 +157,13 @@ This does not introduce a new application architecture or another backend. ## Scope and acceptance -PR 2 covers local macOS arm64 Desktop and local Runtime Host. PR 3 alone will restore an external -Session after process loss. PR 4 owns modes, account/directory invalidation and the expanded catalog -lifecycle. Remote execution, OAuth forwarding, external child orchestration, steering, rollback and +PR 2 covers local macOS arm64 Desktop and local Runtime Host. PR 3 added explicit restoration of the +same external Session after process loss, with a visible history gap when replay cannot be aligned. +PR 4 extends the generic executor configuration with opaque Agent modes and scopes catalog caches +by workspace. A successful setup/login or configuration change invalidates provider discovery; +manual refresh requests a new probe for the selected workspace. Remote execution, OAuth forwarding, external child orchestration, steering and cross-Agent continuation are not added here. -See [PR 2 acceptance evidence](archive/antigravity-acp-pr2-acceptance.md) for controlled-process coverage, -official Agent verification, Desktop verification and the remaining merge gate. The issue's PR 2 -checkbox stays unchecked until the PR is reviewed and merged. +See [PR 2 acceptance evidence](archive/antigravity-acp-pr2-acceptance.md), +[PR 3 acceptance evidence](archive/antigravity-acp-pr3-acceptance.md), and +[PR 4 acceptance evidence](archive/antigravity-acp-pr4-acceptance.md) for the checks and their limits. diff --git a/docs/archive/antigravity-acp-pr4-acceptance.md b/docs/archive/antigravity-acp-pr4-acceptance.md new file mode 100644 index 00000000000..bc5451c9527 --- /dev/null +++ b/docs/archive/antigravity-acp-pr4-acceptance.md @@ -0,0 +1,47 @@ + + +# Antigravity ACP PR 4 acceptance + +PR 4 follows [issue #5103](https://github.com/apache/maka/issues/5103). +This record separates controlled protocol fixtures from official Agent results. + +## Official Agent capability gate, 2026-09-29 + +- Platform: macOS arm64. Client: ACP SDK 1.4.0, Node.js 24.19.0. +- No official Antigravity executable or authenticated Agent home was present in this development environment. The official Google macOS arm64 1.1.1 archive was downloaded from the URL in `docs/antigravity-acp-settings.md`. Its server and helper SHA-256 values matched the already recorded distribution hashes there. +- An isolated, disposable ACP client used a temporary toy directory. `initialize` returned protocol version 1, Agent version `agy_acp_server_1.1.1`, and resume/load capabilities. `session/new` returned JSON-RPC `-32000 Authentication required`. The probe sent no prompt, created no Maka task, did not log an external Session ID, and terminated its process group and toy directory. +- The [official ACP registry](https://github.com/agentclientprotocol/registry/blob/main/antigravity-acp/agent.json) listed version 1.2.1. Its Google macOS arm64 archive contained the matching server and helper (SHA-256 `c93c86c0f505fcdf8b13c695bed26d306141ef5446189d591397074d324db34e` and `1b8a2b712ca312c9769e425b800bfbcceec4770f19736404474d1e8e50d65456`). `initialize` returned Agent version `1.2.1`, protocol version 1 and resume/load capabilities; `session/new` again returned `-32000 Authentication required`. Authentication remains required before a real `configOptions` list can be observed. + +The gate is **blocked by authentication**, not proven protocol incompatibility. Real mode IDs, mode/model interaction, confirmation responses, different-directory catalogs, restoration of mode, and the full Desktop acceptance path are **not verified** in this environment. Controlled fixtures below cannot satisfy those real-Agent checklist items. No account or proxy configuration was changed by this work. + +## Implementation and controlled checks + +The generic executor configuration and catalog now carry optional opaque mode IDs. ACP maps only real `select` mode options from the Agent; omitted mode preserves the Agent default. The same Host query and Desktop picker carry models and modes. Catalogs are keyed by resolved directory, share one bounded probe per directory, and can be invalidated by setup/login, policy changes, expiration, or explicit refresh. The retained task's configuration is inspected independently of draft discovery. + +Configuration updates validate the complete target before applying it, use Agent `setConfigOption` confirmations, and attempt a complete rollback on failure. The saved continuity record keeps the confirmed mode while retaining compatibility with PR 3 records containing only `confirmedModel`. Restore compares the Agent's returned configuration against that record and refuses a mismatch. + +Controlled tests cover the contract, prompt application, idle mode update, combined-option rollback, same-Session restore, directory cache isolation, refresh, login invalidation, and late notification handling. + +## Final validation + +- `npm run build`, `npm run typecheck`, `npm run lint`, `npm run format:check`, `npm run check:renderer-architecture`, `npm run check:locale-hygiene`, and `npm run check:asf-headers`: passed. +- `node scripts/run-workspace-tests-parallel.mjs --concurrency=1` with the bundled Node.js 24.19.0 and Python 3.12.14: all workspaces passed. A separate three-workspace concurrent run was stopped after unrelated timing-sensitive Runtime Host integration cases failed under load; it is not counted as passing validation. +- The protocol epoch changed from 198 to 199 for the additive mode and refresh wire fields. `node scripts/protocol-epoch-check.mjs --staged` passed against the complete staged diff. +- Authenticated official-Agent and Desktop end-to-end checks remain blocked by the `session/new` authentication failure above. In particular, this run cannot claim a real mode list, confirmed mode switch, or restart continuation for an authenticated task. Those issue #5103 acceptance items require a signed-in official Agent and a fresh Desktop run. diff --git a/packages/acp-executor-plugin/src/__tests__/acp-executor-plugin.test.ts b/packages/acp-executor-plugin/src/__tests__/acp-executor-plugin.test.ts index 7b50b0f4d95..17af2ff03bd 100644 --- a/packages/acp-executor-plugin/src/__tests__/acp-executor-plugin.test.ts +++ b/packages/acp-executor-plugin/src/__tests__/acp-executor-plugin.test.ts @@ -1022,6 +1022,204 @@ test('discovery shares a disposable probe, does not mark a task, and first promp } }); +test('ACP modes remain distinct from models across discovery, prompt, idle change and restoration', async () => { + const fixture = await executableFixture(); + const protocol = fakeProtocol(); + protocol.hasMode = true; + protocol.supportsRestore = true; + const storage = durableState(); + const make = () => + new AcpExecutor( + adapter, + { executable: fixture.executable }, + { + createConnection: protocol.factory, + state: storage.state, + }, + ); + const first = make(); + try { + const catalog = await first.discover({ + cwd: fixture.root, + signal: new AbortController().signal, + }); + assert.deepEqual( + catalog.modes?.map((mode) => mode.id), + ['ask', 'auto'], + ); + assert.equal(catalog.currentMode, 'ask'); + assert.equal(catalog.supportsModeChange, true); + const requestWithMode = { ...request('first'), configuration: { model: 'fast', mode: 'auto' } }; + assert.equal((await first.execute(requestWithMode, executorContext([]))).status, 'completed'); + assert.deepEqual(protocol.promptModels, ['fast']); + assert.deepEqual(protocol.promptModes, ['auto']); + assert.equal(storage.record().confirmedMode, 'auto'); + await first.acknowledgeExecution('session-a', 'turn-first'); + await first.configureConversation( + { conversationKey: 'session-a', cwd: process.cwd(), configuration: { mode: 'ask' } }, + new AbortController().signal, + ); + assert.equal(protocol.selectedModel, 'fast'); + assert.equal(protocol.selectedMode, 'ask'); + assert.equal(storage.record().confirmedMode, 'ask'); + await assert.rejects( + first.configureConversation( + { conversationKey: 'session-a', cwd: process.cwd(), configuration: { mode: 'invented' } }, + new AbortController().signal, + ), + ); + assert.equal(protocol.selectedMode, 'ask'); + await first.dispose(); + const restored = make(); + try { + const saved = await restored.inspectConversation({ + conversationKey: 'session-a', + cwd: process.cwd(), + configuration: { model: 'fast', mode: 'ask' }, + }); + assert.equal(saved.readiness, 'restorable'); + assert.deepEqual( + saved.models.map((model) => model.id), + ['fast'], + ); + assert.deepEqual( + saved.modes?.map((mode) => mode.id), + ['ask'], + ); + await restored.configureConversation( + { + conversationKey: 'session-a', + cwd: process.cwd(), + configuration: { model: 'fast', mode: 'ask' }, + }, + new AbortController().signal, + ); + assert.equal( + (await restored.inspectConversation({ conversationKey: 'session-a', cwd: process.cwd() })) + .currentMode, + 'ask', + ); + assert.equal( + protocol.sessions, + 2, + 'one catalog probe and one retained Session; restore does not create another', + ); + } finally { + await restored.dispose(); + } + } finally { + await first.dispose(); + await rm(fixture.root, { recursive: true, force: true }); + } +}); + +test('a failed second configuration option restores both model and mode', async () => { + const fixture = await executableFixture(); + const protocol = fakeProtocol(); + protocol.hasMode = true; + const executor = new AcpExecutor( + adapter, + { executable: fixture.executable }, + { createConnection: protocol.factory }, + ); + try { + assert.equal( + (await executor.execute(request('first'), executorContext([]))).status, + 'completed', + ); + protocol.configurationFailure = 'mode_once'; + await assert.rejects( + executor.configureConversation( + { + conversationKey: 'session-a', + cwd: process.cwd(), + configuration: { model: 'fast', mode: 'auto' }, + }, + new AbortController().signal, + ), + ); + assert.equal(protocol.selectedModel, 'default'); + assert.equal(protocol.selectedMode, 'ask'); + const state = await executor.inspectConversation({ + conversationKey: 'session-a', + cwd: process.cwd(), + }); + assert.equal(state.readiness, 'ready'); + assert.equal(state.currentModel, 'default'); + assert.equal(state.currentMode, 'ask'); + } finally { + await executor.dispose(); + await rm(fixture.root, { recursive: true, force: true }); + } +}); + +test('late configuration notifications cannot overwrite an idle confirmed change', async () => { + const fixture = await executableFixture(); + const protocol = fakeProtocol(); + const executor = new AcpExecutor( + adapter, + { executable: fixture.executable }, + { + createConnection: protocol.factory, + }, + ); + try { + assert.equal( + (await executor.execute(request('first'), executorContext([]))).status, + 'completed', + ); + await executor.configureConversation( + { conversationKey: 'session-a', cwd: process.cwd(), configuration: { model: 'fast' } }, + new AbortController().signal, + ); + protocol.notifyConfiguration('default'); + assert.equal( + (await executor.inspectConversation({ conversationKey: 'session-a', cwd: process.cwd() })) + .currentModel, + 'fast', + ); + } finally { + await executor.dispose(); + await rm(fixture.root, { recursive: true, force: true }); + } +}); + +test('catalog reuse is scoped to cwd and refresh only replaces that scope', async () => { + const fixture = await executableFixture(); + const other = await mkdtemp(join(tmpdir(), 'maka-acp-catalog-other-')); + const protocol = fakeProtocol(); + const executor = new AcpExecutor( + adapter, + { executable: fixture.executable }, + { + createConnection: protocol.factory, + }, + ); + const signal = new AbortController().signal; + try { + await Promise.all([ + executor.discover({ cwd: fixture.root, signal }), + executor.discover({ cwd: fixture.root, signal }), + ]); + assert.equal(protocol.connections, 1); + await executor.discover({ cwd: other, signal }); + assert.equal(protocol.connections, 2); + await executor.discover({ cwd: fixture.root, signal }); + assert.equal(protocol.connections, 2); + await executor.discover({ cwd: fixture.root, signal, refresh: true }); + assert.equal(protocol.connections, 3); + await executor.discover({ cwd: other, signal }); + assert.equal(protocol.connections, 3, 'refreshing one cwd must retain another cwd cache'); + executor.invalidateCatalog(); + await executor.discover({ cwd: other, signal }); + assert.equal(protocol.connections, 4); + } finally { + await executor.dispose(); + await rm(fixture.root, { recursive: true, force: true }); + await rm(other, { recursive: true, force: true }); + } +}); + for (const failure of ['response_lost', 'timeout'] as const) { test(`an applied model change with ${failure} prevents prompts using stale configuration`, async () => { const fixture = await executableFixture(); @@ -1160,9 +1358,12 @@ function fakeProtocol(): { prompts: number; disposals: number; selectedModel?: string; - configurationFailure?: 'response_lost' | 'unconfirmed' | 'timeout' | 'once'; + selectedMode?: string; + hasMode: boolean; + configurationFailure?: 'response_lost' | 'unconfirmed' | 'timeout' | 'once' | 'mode_once'; sessionCreationFailure?: 'once'; promptModels: string[]; + promptModes: string[]; supportsRestore: boolean; restoreFailure?: boolean; resumes: number; @@ -1179,14 +1380,18 @@ function fakeProtocol(): { prompts: 0, disposals: 0, selectedModel: undefined as string | undefined, + selectedMode: 'ask', + hasMode: false, configurationFailure: undefined as | 'response_lost' | 'unconfirmed' | 'timeout' | 'once' + | 'mode_once' | undefined, sessionCreationFailure: undefined as 'once' | undefined, promptModels: [] as string[], + promptModes: [] as string[], supportsRestore: false, restoreFailure: false, resumes: 0, @@ -1198,6 +1403,33 @@ function fakeProtocol(): { notifyConfiguration: (_model: string): void => {}, factory: undefined as unknown as AcpConnectionFactory, }; + const configOptions = () => [ + { + type: 'select', + id: 'model', + name: 'Model', + currentValue: fixture.selectedModel ?? 'default', + options: [ + { value: 'default', name: 'Default' }, + { value: 'fast', name: 'Fast' }, + ], + }, + ...(fixture.hasMode + ? [ + { + type: 'select', + id: 'mode', + category: 'mode', + name: 'Mode', + currentValue: fixture.selectedMode, + options: [ + { value: 'ask', name: 'Ask before edits' }, + { value: 'auto', name: 'Autonomous' }, + ], + }, + ] + : []), + ]; fixture.factory = (input) => { fixture.connections += 1; const notifications = new Map unknown>(); @@ -1286,20 +1518,7 @@ function fakeProtocol(): { params: { sessionId: 'acp-session', update } as never, }); } - return { - configOptions: [ - { - type: 'select', - id: 'model', - name: 'Model', - currentValue: fixture.selectedModel ?? 'default', - options: [ - { value: 'default', name: 'Default' }, - { value: 'fast', name: 'Fast' }, - ], - }, - ], - }; + return { configOptions: configOptions() }; } if (method === methods.agent.session.new) { fixture.sessions += 1; @@ -1307,25 +1526,15 @@ function fakeProtocol(): { fixture.sessionCreationFailure = undefined; throw new Error('Session creation response was lost'); } - return { - sessionId: 'acp-session', - configOptions: [ - { - type: 'select', - id: 'model', - name: 'Model', - currentValue: 'default', - options: [ - { value: 'default', name: 'Default' }, - { value: 'fast', name: 'Fast' }, - ], - }, - ], - }; + return { sessionId: 'acp-session', configOptions: configOptions() }; } if (method === methods.agent.session.setConfigOption) { - fixture.selectedModel = String(params.value); - if (fixture.configurationFailure === 'once') { + if (params.configId === 'mode') fixture.selectedMode = String(params.value); + else fixture.selectedModel = String(params.value); + if ( + fixture.configurationFailure === 'once' || + (fixture.configurationFailure === 'mode_once' && params.configId === 'mode') + ) { fixture.configurationFailure = undefined; throw new Error('Transient rejection after mutation'); } @@ -1339,24 +1548,20 @@ function fakeProtocol(): { ); } return { - configOptions: [ - { - type: 'select', - id: 'model', - name: 'Model', - currentValue: - fixture.configurationFailure === 'unconfirmed' ? 'default' : params.value, - options: [ - { value: 'default', name: 'Default' }, - { value: 'fast', name: 'Fast' }, - ], - }, - ], + configOptions: + fixture.configurationFailure === 'unconfirmed' + ? configOptions().map((option) => + option.id === params.configId + ? { ...option, currentValue: option.id === 'mode' ? 'ask' : 'default' } + : option, + ) + : configOptions(), }; } if (method === methods.agent.session.prompt) { fixture.prompts += 1; fixture.promptModels.push(fixture.selectedModel ?? 'default'); + fixture.promptModes.push(fixture.selectedMode); const text = (params.prompt as Array<{ text: string }>)[0]!.text; await requests.get(methods.client.session.requestPermission)?.({ params: { @@ -1489,22 +1694,27 @@ for (const failure of ['once', 'unconfirmed'] as const) test('idle Agent configuration notifications update the inspected model without another prompt', async () => { const fixture = await executableFixture(); const protocol = fakeProtocol(); + const storage = durableState(); const executor = new AcpExecutor( adapter, { executable: fixture.executable }, - { createConnection: protocol.factory }, + { createConnection: protocol.factory, state: storage.state }, ); try { await executor.execute( { ...request('first'), configuration: { model: 'default' } }, executorContext([]), ); + await executor.acknowledgeExecution('session-a', 'turn-first'); protocol.notifyConfiguration('fast'); assert.equal( (await executor.inspectConversation({ conversationKey: 'session-a', cwd: process.cwd() })) .currentModel, 'fast', ); + for (let attempt = 0; attempt < 20 && storage.record().confirmedModel !== 'fast'; attempt++) + await new Promise((resolve) => setTimeout(resolve, 5)); + assert.equal(storage.record().confirmedModel, 'fast'); assert.equal(protocol.prompts, 1); } finally { await executor.dispose(); diff --git a/packages/acp-executor-plugin/src/index.ts b/packages/acp-executor-plugin/src/index.ts index af24887252a..0b5bc834207 100644 --- a/packages/acp-executor-plugin/src/index.ts +++ b/packages/acp-executor-plugin/src/index.ts @@ -19,8 +19,7 @@ import { createHash } from 'node:crypto'; import { constants, createReadStream } from 'node:fs'; -import { access, mkdtemp, rm, realpath, stat } from 'node:fs/promises'; -import { tmpdir } from 'node:os'; +import { access, realpath, stat } from 'node:fs/promises'; import type { ExecutorCatalogEntry, ExecutorConfiguration } from '@maka/core/executor-catalog'; import { dirname, isAbsolute, resolve } from 'node:path'; import { @@ -120,6 +119,7 @@ export interface AcpContinuityRecord { readonly sessionId?: string; readonly pendingTurnId?: string; readonly confirmedModel?: string; + readonly confirmedMode?: string; readonly gapEvidence?: { readonly replayedUpdates: number; readonly replayedUserChunks: number; @@ -165,6 +165,8 @@ interface RetainedSession { initialization?: Promise; active?: ActivePrompt; configuring?: boolean; + suppressConflictingConfigUpdates?: boolean; + configPersistence?: Promise; lost: boolean; loss?: Promise; } @@ -179,8 +181,11 @@ export class AcpExecutor implements PluginExecutorProvider { readonly #state?: AcpConversationStateStore; readonly #sessions = new Map(); #disposed = false; - #catalog?: ExecutorCatalogEntry; - #discovery?: Promise; + readonly #catalog = new Map(); + readonly #discovery = new Map>(); + readonly #probes = new Map(); + readonly #scopeTokens = new Map(); + #catalogEpoch = 0; constructor( adapter: AcpAgentAdapter, @@ -198,13 +203,32 @@ export class AcpExecutor implements PluginExecutorProvider { this.#state = options.state; } - async discover(input: { cwd: string; signal: AbortSignal }): Promise { + async discover(input: { + cwd: string; + signal: AbortSignal; + refresh?: boolean; + }): Promise { if (this.#disposed) return this.#catalogEntry('unavailable'); - if (this.#catalog) return this.#catalog; - if (this.#discovery) return this.#discovery; + input.signal.throwIfAborted(); + const cwd = await realpath(resolve(input.cwd)).catch(() => undefined); + if (!cwd) return this.#catalogEntry('unavailable'); + if (input.refresh) { + this.#catalog.delete(cwd); + this.#scopeTokens.set(cwd, {}); + this.#probes.get(cwd)?.abort(new Error('ACP catalog scope refreshed')); + this.#discovery.delete(cwd); + } + const cached = this.#catalog.get(cwd); + if (cached && cached.expires > Date.now()) return cached.entry; + const existing = this.#discovery.get(cwd); + if (existing) return existing; + const epoch = this.#catalogEpoch; + const scopeToken = this.#scopeTokens.get(cwd) ?? {}; + this.#scopeTokens.set(cwd, scopeToken); + const controller = new AbortController(); + this.#probes.set(cwd, controller); const probe = async () => { // A bounded disposable ACP probe never creates a Maka task or joins the retained-session map. - const cwd = await mkdtemp(resolve(tmpdir(), 'maka-acp-catalog-')); const session: RetainedSession = { conversationKey: 'catalog-probe', cwd, @@ -212,9 +236,20 @@ export class AcpExecutor implements PluginExecutorProvider { lost: false, }; try { - await this.#initialize(session, input.signal, true); + await this.#initialize(session, controller.signal, true); const result = this.#catalogEntry('ready', session.configOptions); - this.#catalog = result; + if ( + epoch !== this.#catalogEpoch || + this.#scopeTokens.get(cwd) !== scopeToken || + this.#disposed + ) + return this.#catalogEntry('unavailable'); + this.#catalog.set(cwd, { entry: result, expires: Date.now() + 60_000 }); + while (this.#catalog.size > 16) { + const oldest = this.#catalog.keys().next().value!; + this.#catalog.delete(oldest); + if (!this.#probes.has(oldest)) this.#scopeTokens.delete(oldest); + } return result; } catch (error) { return this.#catalogEntry( @@ -224,13 +259,24 @@ export class AcpExecutor implements PluginExecutorProvider { ); } finally { await this.#disposeSession(session); - await rm(cwd, { recursive: true, force: true }); } }; - this.#discovery = probe().finally(() => { - this.#discovery = undefined; + const discovery = probe().finally(() => { + if (this.#probes.get(cwd) === controller) this.#probes.delete(cwd); + if (this.#discovery.get(cwd) === discovery) this.#discovery.delete(cwd); + if (!this.#catalog.has(cwd) && this.#scopeTokens.get(cwd) === scopeToken) + this.#scopeTokens.delete(cwd); }); - return this.#discovery; + this.#discovery.set(cwd, discovery); + return discovery; + } + + invalidateCatalog(): void { + this.#catalogEpoch++; + this.#catalog.clear(); + this.#scopeTokens.clear(); + for (const probe of this.#probes.values()) probe.abort(new Error('ACP catalog scope changed')); + this.#discovery.clear(); } async inspectConversation(input: { @@ -241,53 +287,64 @@ export class AcpExecutor implements PluginExecutorProvider { const session = this.#sessions.get(input.conversationKey); if (this.#disposed) return this.#catalogEntry('unavailable'); if (session && session.cwd !== resolve(input.cwd)) - return this.#catalogEntry('history_only', [], input.configuration?.model); + return this.#catalogEntry('history_only', [], input.configuration); if (session?.restoring) - return this.#catalogEntry('restoring', session.configOptions, input.configuration?.model); + return this.#catalogEntry('restoring', session.configOptions, input.configuration); if (session?.restoreFailed) - return this.#catalogEntry( - 'restore_failed', - session.configOptions, - input.configuration?.model, - ); + return this.#catalogEntry('restore_failed', session.configOptions, input.configuration); if (session?.historyGap) - return this.#catalogEntry('history_gap', session.configOptions, input.configuration?.model); + return this.#catalogEntry('history_gap', session.configOptions, input.configuration); if (session && !session.lost) - return this.#catalogEntry('ready', session.configOptions, input.configuration?.model); + return this.#catalogEntry('ready', session.configOptions, input.configuration); if (session?.lost && !session.record) - return this.#catalogEntry('history_only', session.configOptions, input.configuration?.model); + return this.#catalogEntry('history_only', session.configOptions, input.configuration); const stored = this.#state?.read ? decodeContinuity(await this.#state.read(input.conversationKey)) : (await this.#state?.has(input.conversationKey, input.cwd)) ? 'legacy' : undefined; - if (!stored) return this.#catalogEntry('ready', [], input.configuration?.model); + if (!stored) return this.#catalogEntry('ready', [], input.configuration); if (stored === 'invalid' || stored === 'legacy' || stored.cwd !== resolve(input.cwd)) - return this.#catalogEntry('history_only', [], input.configuration?.model); + return this.#catalogEntry('history_only', [], input.configuration); if (stored.phase === 'reserved') - return this.#catalogEntry('history_only', [], input.configuration?.model); + return this.#catalogEntry('history_only', [], input.configuration); return this.#catalogEntry( stored.phase === 'prompt_pending' || stored.phase === 'history_gap' ? 'history_gap' : 'restorable', [], - input.configuration?.model ?? stored.confirmedModel, + { + model: input.configuration?.model ?? stored.confirmedModel, + mode: input.configuration?.mode ?? stored.confirmedMode, + }, ); } async configureConversation( input: { conversationKey: string; cwd: string; configuration?: ExecutorConfiguration }, signal: AbortSignal, - ): Promise { - if (!input.configuration?.model) throw new Error('ACP model change is unavailable'); + ): Promise { + if (!input.configuration?.model && !input.configuration?.mode) + throw new Error('ACP configuration change is unavailable'); const session = await this.#session(input); - if (session.active || session.configuring) throw new Error('ACP model change is unavailable'); + if (session.active || session.configuring || session.restoring || session.awaitingAck) + throw new Error('ACP configuration change is unavailable while the Session is busy'); session.configuring = true; const wasConnected = !!session.connection; try { await this.#ensureInitialized(session, signal); - await this.#applyInitialConfig(session, { model: input.configuration.model }, signal, true); + await session.configPersistence; + await this.#applyInitialConfig(session, input.configuration, signal, true); + return { + ...(currentAcpModel(session.configOptions) + ? { model: currentAcpModel(session.configOptions) } + : {}), + ...(acpOption(session.configOptions, 'mode')?.currentValue + ? { mode: acpOption(session.configOptions, 'mode')!.currentValue } + : {}), + }; } catch (error) { + if (isAuthenticationFailure(error)) this.invalidateCatalog(); if (!wasConnected || session.lost) { await this.#lose(session); if (session.record?.sessionId && !session.historyGap) session.restoreFailed = true; @@ -304,7 +361,7 @@ export class AcpExecutor implements PluginExecutorProvider { #catalogEntry( readiness: ExecutorCatalogEntry['readiness'], options: readonly SessionConfigOption[] = [], - selected?: string, + selected?: ExecutorConfiguration, ): ExecutorCatalogEntry { const model = options.find( (option) => @@ -315,7 +372,20 @@ export class AcpExecutor implements PluginExecutorProvider { ? model.options .flatMap((entry) => ('options' in entry ? entry.options : [entry])) .map((entry) => ({ id: entry.value, name: entry.name })) - : (this.#catalog?.models ?? []); + : selected?.model + ? [{ id: selected.model, name: selected.model }] + : []; + const mode = options.find( + (option) => option.type === 'select' && (option.category === 'mode' || option.id === 'mode'), + ); + const modes = + mode?.type === 'select' + ? mode.options + .flatMap((entry) => ('options' in entry ? entry.options : [entry])) + .map((entry) => ({ id: entry.value, name: entry.name })) + : selected?.mode + ? [{ id: selected.mode, name: selected.mode }] + : []; return { id: this.id, displayName: this.displayName, @@ -324,13 +394,19 @@ export class AcpExecutor implements PluginExecutorProvider { ? (this.#adapter.describeModels?.(models) ?? { models }) : { models, - ...(this.#catalog?.modelGroups ? { modelGroups: this.#catalog.modelGroups } : {}), }), ...(model?.type === 'select' ? { currentModel: model.currentValue } - : selected - ? { currentModel: selected } + : selected?.model + ? { currentModel: selected.model } : {}), + ...(mode?.type === 'select' + ? { modes, currentMode: mode.currentValue, supportsModeChange: true } + : { + modes, + ...(selected?.mode ? { currentMode: selected.mode } : {}), + supportsModeChange: false, + }), supportsAttachments: false, supportsModelChange: model?.type === 'select', }; @@ -356,7 +432,7 @@ export class AcpExecutor implements PluginExecutorProvider { return failure('Conflicting executor models', 'acp_config_invalid'); const model = request.configuration?.model ?? (request.model === this.id ? undefined : request.model); - if (model) request = { ...request, configuration: { model } }; + if (model) request = { ...request, configuration: { ...request.configuration, model } }; let session: RetainedSession; try { session = await this.#session(request); @@ -369,13 +445,10 @@ export class AcpExecutor implements PluginExecutorProvider { session.active = active; try { await this.#ensureInitialized(session, context.signal); + await session.configPersistence; context.signal.throwIfAborted(); - if (request.configuration?.model) - await this.#applyInitialConfig( - session, - { model: request.configuration.model }, - context.signal, - ); + if (request.configuration?.model || request.configuration?.mode) + await this.#applyInitialConfig(session, request.configuration, context.signal); await this.#beginPrompt(session, request.turnId); const prompt = session.connection!.agent.request(methods.agent.session.prompt, { sessionId: session.acpSessionId!, @@ -401,6 +474,7 @@ export class AcpExecutor implements PluginExecutorProvider { } return { status: 'completed', text: active.text }; } catch (error) { + if (isAuthenticationFailure(error)) this.invalidateCatalog(); if (errorCode(error) === 'acp_history_gap') session.historyGap = true; await this.#lose(session); if (session.record?.sessionId && !session.historyGap) session.restoreFailed = true; @@ -435,6 +509,9 @@ export class AcpExecutor implements PluginExecutorProvider { async dispose(): Promise { if (this.#disposed) return; this.#disposed = true; + const probes = [...this.#discovery.values()]; + this.invalidateCatalog(); + await Promise.allSettled(probes); const sessions = [...this.#sessions.values()]; this.#sessions.clear(); const settlements = await Promise.allSettled( @@ -585,7 +662,16 @@ export class AcpExecutor implements PluginExecutorProvider { ]); session.configOptions = restored.configOptions ?? []; const restoredModel = currentAcpModel(session.configOptions); - if (restoredModel) await this.#persistConfirmedModel(session, restoredModel); + const restoredMode = acpOption(session.configOptions, 'mode')?.currentValue; + if ( + (stored.confirmedModel && restoredModel !== stored.confirmedModel) || + (stored.confirmedMode && restoredMode !== stored.confirmedMode) + ) + throw new AcpRuntimeError( + 'Restored Agent configuration differs from the saved task', + 'acp_config_unconfirmed', + ); + await this.#persistConfirmedConfig(session); session.restoring = false; if (pending) { // The Agent may have progressed beyond Maka's last durable event. Replay @@ -648,14 +734,13 @@ export class AcpExecutor implements PluginExecutorProvider { }; await this.#state?.write?.(session.conversationKey, established); session.record = established; - const createdModel = currentAcpModel(session.configOptions); - if (createdModel) await this.#persistConfirmedModel(session, createdModel); + await this.#persistConfirmedConfig(session); } if (!probe) { await this.#applyInitialConfig( session, - session.configuration?.model - ? { model: session.configuration.model } + session.configuration?.model || session.configuration?.mode + ? { ...launch.initialConfig, ...session.configuration } : (launch.initialConfig ?? {}), startupSignal, ); @@ -663,6 +748,8 @@ export class AcpExecutor implements PluginExecutorProvider { } async #beginPrompt(session: RetainedSession, turnId: string): Promise { + await session.configPersistence; + session.suppressConflictingConfigUpdates = false; if (!session.record) return; const pending: AcpContinuityRecord = { ...session.record, @@ -688,6 +775,12 @@ export class AcpExecutor implements PluginExecutorProvider { phase: 'committed', committedPrompts: session.record.committedPrompts + 1, pendingTurnId: undefined, + ...(currentAcpModel(session.configOptions) + ? { confirmedModel: currentAcpModel(session.configOptions) } + : {}), + ...(acpOption(session.configOptions, 'mode')?.currentValue + ? { confirmedMode: acpOption(session.configOptions, 'mode')!.currentValue } + : {}), }; try { await this.#state?.write?.(conversationKey, committed); @@ -698,6 +791,12 @@ export class AcpExecutor implements PluginExecutorProvider { } session.record = committed; session.awaitingAck = undefined; + try { + await this.#persistConfirmedConfig(session); + } catch (error) { + await this.#lose(session); + throw error; + } } /** The external turn settled, but Maka did not durably consume its terminal event. */ @@ -726,83 +825,116 @@ export class AcpExecutor implements PluginExecutorProvider { async #applyInitialConfig( session: RetainedSession, - values: Readonly>, + values: ExecutorConfiguration | Readonly>, signal: AbortSignal, restoreOnFailure = false, ): Promise { - for (const [key, value] of Object.entries(values)) { - const option = session.configOptions.find( - (candidate) => - candidate.type === 'select' && (candidate.id === key || candidate.category === key), - ); - if (!option || option.type !== 'select') + const target = Object.entries(values).filter( + (entry): entry is [string, string] => entry[1] !== undefined, + ); + const before = session.configOptions; + // Validate the entire target before the first Agent mutation. + for (const [key, value] of target) { + const option = acpOption(before, key); + if (!option) throw new AcpRuntimeError( `ACP configuration is unavailable: ${key}`, 'acp_config_unavailable', ); - const options = option.options.flatMap((entry) => + const candidates = option.options.flatMap((entry) => 'options' in entry ? entry.options : [entry], ); - if (!options.some((entry) => entry.value === value)) + if (!candidates.some((entry) => entry.value === value)) throw new AcpRuntimeError( `ACP configuration value is unavailable: ${key}`, 'acp_config_invalid', ); - if (option.currentValue === value) { - if (key === 'model') await this.#persistConfirmedModel(session, value); - continue; - } - signal.throwIfAborted(); - try { + } + let mutationAttempted = false; + try { + for (const [key, value] of target) { + const option = acpOption(session.configOptions, key); + if (!option) + throw new AcpRuntimeError( + `ACP configuration is unavailable: ${key}`, + 'acp_config_unavailable', + ); + if (option.currentValue === value) continue; + signal.throwIfAborted(); + mutationAttempted = true; const updated = await session.connection!.agent.request( methods.agent.session.setConfigOption, { sessionId: session.acpSessionId!, configId: option.id, value }, { cancellationSignal: signal }, ); - const confirmed = updated.configOptions.find((candidate) => candidate.id === option.id); - if (confirmed?.type !== 'select' || confirmed.currentValue !== value) + session.configOptions = updated.configOptions; + if (acpOption(updated.configOptions, key)?.currentValue !== value) throw new AcpRuntimeError( 'Agent did not confirm the selected configuration', 'acp_config_unconfirmed', ); - session.configOptions = updated.configOptions; - if (key === 'model') await this.#persistConfirmedModel(session, value); - } catch (error) { - // A rejected/unconfirmed mutation can have reached the Agent. Restore the - // previous real ID and require an acknowledgement before permitting retry. - let restored = false; - if (restoreOnFailure && !session.lost && !signal.aborted) { - try { + } + if ( + target.some(([key, value]) => acpOption(session.configOptions, key)?.currentValue !== value) + ) + throw new AcpRuntimeError( + 'Agent changed another selected configuration', + 'acp_config_unconfirmed', + ); + await this.#persistConfirmedConfig(session); + session.suppressConflictingConfigUpdates = true; + } catch (error) { + let restored = false; + if (restoreOnFailure && !session.lost && !signal.aborted) { + try { + // A rejected response can still have changed the Agent. Restore both + // fields, then verify the full original snapshot. + for (const key of ['mode', 'model']) { + const original = acpOption(before, key); + const current = acpOption(session.configOptions, key); + if ( + !original || + !current || + (!mutationAttempted && current.currentValue === original.currentValue) + ) + continue; const rollback = await session.connection!.agent.request( methods.agent.session.setConfigOption, - { sessionId: session.acpSessionId!, configId: option.id, value: option.currentValue }, + { + sessionId: session.acpSessionId!, + configId: current.id, + value: original.currentValue, + }, { cancellationSignal: AbortSignal.timeout(5_000) }, ); - const confirmed = rollback.configOptions.find( - (candidate) => candidate.id === option.id, - ); - if ( - confirmed?.type === 'select' && - confirmed.currentValue === option.currentValue && - !session.lost - ) { - session.configOptions = rollback.configOptions; - if (key === 'model') await this.#persistConfirmedModel(session, option.currentValue); - restored = true; - } - } catch { - /* Uncertain configuration remains history-only. */ + session.configOptions = rollback.configOptions; } + restored = ['model', 'mode'].every((key) => { + const original = acpOption(before, key); + return ( + !original || + acpOption(session.configOptions, key)?.currentValue === original.currentValue + ); + }); + } catch { + /* Uncertain configuration remains unavailable. */ } - if (!restored) await this.#lose(session); - throw error; } + if (!restored) await this.#lose(session); + throw error; } } - async #persistConfirmedModel(session: RetainedSession, model: string): Promise { - if (!session.record || session.record.confirmedModel === model) return; - const record = { ...session.record, confirmedModel: model }; + async #persistConfirmedConfig(session: RetainedSession): Promise { + if (!session.record) return; + const model = acpOption(session.configOptions, 'model')?.currentValue; + const mode = acpOption(session.configOptions, 'mode')?.currentValue; + if (session.record.confirmedModel === model && session.record.confirmedMode === mode) return; + const record = { + ...session.record, + ...(model ? { confirmedModel: model } : {}), + ...(mode ? { confirmedMode: mode } : {}), + }; await this.#state?.write?.(session.conversationKey, record); session.record = record; } @@ -854,7 +986,24 @@ export class AcpExecutor implements PluginExecutorProvider { if (session.historyGap || session.lost) return; if (update.sessionUpdate === 'config_option_update') { // During our mutation, its response (or rollback response) is authoritative. - if (!session.configuring && !session.lost) session.configOptions = update.configOptions; + if (!session.configuring && !session.lost) { + if ( + session.suppressConflictingConfigUpdates && + ['model', 'mode'].some((key) => { + const current = acpOption(session.configOptions, key)?.currentValue; + const notified = acpOption(update.configOptions, key)?.currentValue; + return current !== undefined && notified !== undefined && current !== notified; + }) + ) + return; + session.configOptions = update.configOptions; + if (!session.active && !session.awaitingAck && session.record?.phase !== 'prompt_pending') { + const previous = session.configPersistence ?? Promise.resolve(); + const persistence = previous.then(() => this.#persistConfirmedConfig(session)); + session.configPersistence = persistence; + void persistence.catch(() => this.#lose(session)).catch(() => undefined); + } + } return; } const active = session.active; @@ -995,11 +1144,25 @@ export class AcpExecutor implements PluginExecutorProvider { } function currentAcpModel(options: readonly SessionConfigOption[]): string | undefined { - const option = options.find( - (candidate) => - candidate.type === 'select' && (candidate.id === 'model' || candidate.category === 'model'), + return acpOption(options, 'model')?.currentValue; +} + +function isAuthenticationFailure(error: unknown): boolean { + return ( + (error as { code?: unknown })?.code === -32000 || + (error instanceof Error && + /authentication required|unauthorized/u.test(error.message.toLowerCase())) + ); +} + +function acpOption( + options: readonly SessionConfigOption[], + key: string, +): Extract | undefined { + return options.find( + (candidate): candidate is Extract => + candidate.type === 'select' && (candidate.id === key || candidate.category === key), ); - return option?.type === 'select' ? option.currentValue : undefined; } function toolTextContent(content: readonly ToolCallContent[]): string { @@ -1087,6 +1250,8 @@ function decodeContinuity(value: unknown): StoredContinuity { (typeof record.pendingTurnId !== 'string' || !record.pendingTurnId)) || (record.confirmedModel !== undefined && (typeof record.confirmedModel !== 'string' || !record.confirmedModel)) || + (record.confirmedMode !== undefined && + (typeof record.confirmedMode !== 'string' || !record.confirmedMode)) || (record.phase === 'reserved' && record.sessionId !== undefined) || (record.phase !== 'reserved' && !record.sessionId) || (record.phase === 'prompt_pending' && !record.pendingTurnId) || diff --git a/packages/core/src/__tests__/executor-catalog.test.ts b/packages/core/src/__tests__/executor-catalog.test.ts index 374de122851..3b65d2326d3 100644 --- a/packages/core/src/__tests__/executor-catalog.test.ts +++ b/packages/core/src/__tests__/executor-catalog.test.ts @@ -19,7 +19,11 @@ import assert from 'node:assert/strict'; import test from 'node:test'; -import { normalizeCatalogEntry, type ExecutorCatalogEntry } from '../executor-catalog.js'; +import { + isExecutorConfiguration, + normalizeCatalogEntry, + type ExecutorCatalogEntry, +} from '../executor-catalog.js'; const catalog: ExecutorCatalogEntry = { id: 'external', @@ -48,6 +52,31 @@ test('structured capabilities survive wire normalization with immutable exact re assert.ok(Object.isFrozen(output.modelGroups![0]!.variants[0])); assert.notEqual(output.modelGroups, catalog.modelGroups); }); +test('opaque modes round-trip through generic configuration and catalog validation', () => { + assert.equal(isExecutorConfiguration({ model: 'opaque-high', mode: 'ask' }), true); + assert.equal(isExecutorConfiguration({ mode: 'ask', unknown: true }), false); + assert.equal(isExecutorConfiguration({ mode: 'bad\nvalue' }), false); + const entry = normalizeCatalogEntry( + { + ...catalog, + modes: [ + { id: 'ask', name: 'Ask' }, + { id: 'auto', name: 'Automatic' }, + ], + currentMode: 'ask', + supportsModeChange: true, + }, + 'external', + ); + assert.deepEqual( + entry.modes?.map((mode) => mode.id), + ['ask', 'auto'], + ); + assert.ok(Object.isFrozen(entry.modes?.[0])); + assert.throws(() => + normalizeCatalogEntry({ ...entry, modes: [entry.modes![0]!, entry.modes![0]!] }, 'external'), + ); +}); for (const readiness of ['restorable', 'restoring', 'restore_failed', 'history_gap'] as const) test(`restoration readiness survives executor catalog normalization: ${readiness}`, () => { assert.equal( diff --git a/packages/core/src/executor-catalog.ts b/packages/core/src/executor-catalog.ts index 5d821d7083e..a05af8cd55e 100644 --- a/packages/core/src/executor-catalog.ts +++ b/packages/core/src/executor-catalog.ts @@ -24,6 +24,8 @@ import { isThinkingLevel, type ThinkingLevel } from './model-thinking.js'; /** Provider-owned choices; no model Connection or external protocol identity crosses this seam. */ export interface ExecutorConfiguration { readonly model?: string; + /** Opaque provider mode ID. Omission leaves the Agent's default unchanged. */ + readonly mode?: string; } export interface ExecutorSelection { @@ -47,6 +49,11 @@ export interface ExecutorModelChoice { readonly providerType?: ProviderType; } +export interface ExecutorModeChoice { + readonly id: string; + readonly name: string; +} + /** Presentation capability only: every variant references a real catalog model ID. */ export interface ExecutorModelGroup { readonly id: string; @@ -61,20 +68,28 @@ export interface ExecutorCatalogEntry { readonly models: readonly ExecutorModelChoice[]; readonly modelGroups?: readonly ExecutorModelGroup[]; readonly currentModel?: string; + readonly modes?: readonly ExecutorModeChoice[]; + readonly currentMode?: string; readonly supportsAttachments: boolean; readonly supportsModelChange: boolean; + readonly supportsModeChange?: boolean; } export function isExecutorConfiguration(value: unknown): value is ExecutorConfiguration { if (!value || typeof value !== 'object' || Array.isArray(value)) return false; const record = value as Record; return ( - Object.keys(record).every((key) => key === 'model') && + Object.keys(record).every((key) => key === 'model' || key === 'mode') && (record.model === undefined || (typeof record.model === 'string' && record.model.length > 0 && record.model.length <= 1024 && - !/[\0\r\n]/u.test(record.model))) + !/[\0\r\n]/u.test(record.model))) && + (record.mode === undefined || + (typeof record.mode === 'string' && + record.mode.length > 0 && + record.mode.length <= 1024 && + !/[\0\r\n]/u.test(record.mode))) ); } @@ -109,8 +124,17 @@ export function normalizeCatalogEntry( ) || new Set(value.models.map((model) => model.id)).size !== value.models.length || !isExecutorConfiguration({ model: value.currentModel }) || + (value.modes !== undefined && + (!Array.isArray(value.modes) || + value.modes.length > 64 || + !value.modes.every( + (mode) => mode && isExecutorConfiguration({ mode: mode.id }) && isCatalogText(mode.name), + ) || + new Set(value.modes.map((mode) => mode.id)).size !== value.modes.length)) || + !isExecutorConfiguration({ mode: value.currentMode }) || typeof value.supportsAttachments !== 'boolean' || - typeof value.supportsModelChange !== 'boolean' + typeof value.supportsModelChange !== 'boolean' || + (value.supportsModeChange !== undefined && typeof value.supportsModeChange !== 'boolean') ) throw new TypeError('Executor catalog is invalid'); const usedModels = new Set(); @@ -180,8 +204,15 @@ export function normalizeCatalogEntry( } : {}), ...(value.currentModel !== undefined ? { currentModel: value.currentModel } : {}), + ...(value.modes !== undefined + ? { modes: Object.freeze(value.modes.map(({ id, name }) => Object.freeze({ id, name }))) } + : {}), + ...(value.currentMode !== undefined ? { currentMode: value.currentMode } : {}), supportsAttachments: value.supportsAttachments, supportsModelChange: value.supportsModelChange, + ...(value.supportsModeChange !== undefined + ? { supportsModeChange: value.supportsModeChange } + : {}), }); } diff --git a/packages/runtime-host/src/__tests__/external-agent-setup-coordinator.test.ts b/packages/runtime-host/src/__tests__/external-agent-setup-coordinator.test.ts index 39cad24e6c3..b1300a6d59f 100644 --- a/packages/runtime-host/src/__tests__/external-agent-setup-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/external-agent-setup-coordinator.test.ts @@ -276,6 +276,31 @@ test('install admits an empty saved path, reports progress and checks without au await coordinator.close(); } }); + +test('successful login invalidates catalog even when the executable setting is unchanged', async () => { + let invalidations = 0; + const coordinator = new HostExternalAgentSetupCoordinator({ + readPolicy: async () => policy, + platform: 'darwin', + arch: 'arm64', + acquireResidency: () => ({ release() {} }), + onCleanupFailure() {}, + onSucceeded: () => { + invalidations++; + }, + capabilities: { callService: async () => ({ kind: 'presented' }) }, + run: async () => {}, + }); + try { + await coordinator.handlers['external_agents.setup.start'](input, context); + for (let i = 0; i < 30 && (await projection(coordinator)).phase !== 'succeeded'; i++) + await new Promise((resolve) => setTimeout(resolve, 5)); + assert.equal((await projection(coordinator)).phase, 'succeeded'); + assert.equal(invalidations, 1); + } finally { + await coordinator.close(); + } +}); test('installation projection rejects invented progress, phase and output paths', () => { const base = { ...input, diff --git a/packages/runtime-host/src/__tests__/session-catalog-coordinator.test.ts b/packages/runtime-host/src/__tests__/session-catalog-coordinator.test.ts index 48516f4ea23..be49cd62a2e 100644 --- a/packages/runtime-host/src/__tests__/session-catalog-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/session-catalog-coordinator.test.ts @@ -918,14 +918,14 @@ test('plugin executor creation bypasses model resolution and persists the execut executorId: 'codex', executorModel: 'gpt-codex', thinkingLevel: 'high', - executorConfig: { model: 'gpt-codex' }, + executorConfig: { model: 'gpt-codex', mode: 'auto' }, }, context, ); assert.equal(outcome.ok, true, JSON.stringify(outcome)); assert.equal(persistedInput?.executorId, 'codex'); - assert.deepEqual(persistedInput?.executorConfig, { model: 'gpt-codex' }); + assert.deepEqual(persistedInput?.executorConfig, { model: 'gpt-codex', mode: 'auto' }); assert.equal(persistedInput?.llmConnectionId, undefined); assert.equal(persistedInput?.llmConnectionSlug, 'executor:codex'); assert.equal(persistedInput?.model, 'gpt-codex'); @@ -2484,6 +2484,54 @@ test('executor model changes commit only after idle agent confirmation', async ( assert.equal(fixture.drainRequests(), 0); }); +test('mode-only changes retain the saved model and wait for Agent confirmation', async () => { + const fixture = createFixture({ + header: { + backend: 'plugin-executor', + executorId: 'remote', + executorConfig: { model: 'before', mode: 'ask' }, + }, + configureExecutor: async (_header, config) => { + assert.deepEqual(config, { mode: 'auto' }); + assert.deepEqual(fixture.header().executorConfig, { model: 'before', mode: 'ask' }); + }, + }); + const outcome = await fixture.coordinator.handlers['session.configuration.update']( + { + sessionId: fixture.sessionId, + expectedRevision: fixture.revision(), + patch: { executorConfig: { mode: 'auto' } }, + }, + context, + ); + assert.equal(outcome.ok, true, JSON.stringify(outcome)); + assert.deepEqual(fixture.header().executorConfig, { model: 'before', mode: 'auto' }); +}); + +test('Agent-confirmed mode side effects are persisted with a model change', async () => { + const fixture = createFixture({ + header: { + backend: 'plugin-executor', + executorId: 'remote', + executorConfig: { model: 'before', mode: 'ask' }, + }, + configureExecutor: async (_header, config) => { + assert.deepEqual(config, { model: 'after' }); + return { model: 'after', mode: 'auto' }; + }, + }); + const outcome = await fixture.coordinator.handlers['session.configuration.update']( + { + sessionId: fixture.sessionId, + expectedRevision: fixture.revision(), + patch: { executorConfig: { model: 'after' } }, + }, + context, + ); + assert.equal(outcome.ok, true, JSON.stringify(outcome)); + assert.deepEqual(fixture.header().executorConfig, { model: 'after', mode: 'auto' }); +}); + test('failed Session commit restores the confirmed executor model', async () => { const confirmed: string[] = []; const fixture = createFixture({ diff --git a/packages/runtime-host/src/__tests__/session-catalog-protocol.test.ts b/packages/runtime-host/src/__tests__/session-catalog-protocol.test.ts index 1df27371174..b33f24cb3a3 100644 --- a/packages/runtime-host/src/__tests__/session-catalog-protocol.test.ts +++ b/packages/runtime-host/src/__tests__/session-catalog-protocol.test.ts @@ -814,7 +814,9 @@ test('executor configuration rejects ambiguous routes and malformed values', () null, { model: '' }, { model: 'bad\nvalue' }, - { mode: 'yolo' }, + { mode: '' }, + { mode: 'bad\nvalue' }, + { mode: 4 }, { model: 4 }, ]) { assert.throws( diff --git a/packages/runtime-host/src/protocol/index.ts b/packages/runtime-host/src/protocol/index.ts index 2f41d5595c1..82e1c1e2675 100644 --- a/packages/runtime-host/src/protocol/index.ts +++ b/packages/runtime-host/src/protocol/index.ts @@ -104,7 +104,9 @@ export const RUNTIME_HOST_REGISTRATION_SCHEMA_VERSION = 1 as const; export const RUNTIME_HOST_PROTOCOL_VERSION = 0 as const; // Increment when the same protocol version no longer guarantees safe Client-Host // interoperability. Mismatches are rejected before domain commands are admitted. -export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 198 as const; +export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 199 as const; +// 199: Executor catalogs and Session configuration carry opaque mode IDs; +// catalog queries may request a provider refresh. Older peers reject these fields. // 198: Executor readiness exposes explicit restore, restore-failed and history-gap states. // 197: Queue reorder requests carry the expected queue revision. Epoch-196 // peers reject the required field or send an unfenced reorder request. diff --git a/packages/runtime-host/src/protocol/plugin-platform.ts b/packages/runtime-host/src/protocol/plugin-platform.ts index 518ed9a208f..02fe55e5841 100644 --- a/packages/runtime-host/src/protocol/plugin-platform.ts +++ b/packages/runtime-host/src/protocol/plugin-platform.ts @@ -102,7 +102,7 @@ export interface PluginPackageProjection { } export type PluginExecutorQueryInput = - | { readonly kind: 'catalog'; readonly cwd: string } + | { readonly kind: 'catalog'; readonly cwd: string; readonly refresh?: boolean } | { readonly kind: 'conversation'; readonly sessionId: string }; export interface PluginExecutorQueryResult { readonly items: readonly ExecutorCatalogEntry[]; @@ -342,8 +342,14 @@ export const PLUGIN_PLATFORM_OPERATION_SPECS = { decodeInput: (value) => { const record = requireRecord(value, 'Executor query'); if (record.kind === 'catalog') { - const input = requireExactRecord(record, 'Executor catalog', ['kind', 'cwd']); - return { kind: 'catalog', cwd: requireString(input.cwd, 'cwd', 4096) }; + const input = requireExactRecord(record, 'Executor catalog', ['kind', 'cwd', 'refresh']); + if (input.refresh !== undefined && typeof input.refresh !== 'boolean') + throw invalidProtocolFrame('Invalid executor refresh flag'); + return { + kind: 'catalog', + cwd: requireString(input.cwd, 'cwd', 4096), + ...(input.refresh ? { refresh: true } : {}), + }; } const input = requireExactRecord(record, 'Executor conversation', ['kind', 'sessionId']); if (input.kind !== 'conversation') throw invalidProtocolFrame('Invalid executor query'); diff --git a/packages/runtime-host/src/server/execution-composition.ts b/packages/runtime-host/src/server/execution-composition.ts index 19e2f649e05..ef3f7d093fe 100644 --- a/packages/runtime-host/src/server/execution-composition.ts +++ b/packages/runtime-host/src/server/execution-composition.ts @@ -422,7 +422,7 @@ export async function createExecutionRuntimeHostComposition( ...(header.executorConfig ? { configuration: header.executorConfig } : {}), }); } - const catalog = await pluginExecutors.catalog({ cwd: input.cwd }); + const catalog = await pluginExecutors.catalog({ cwd: input.cwd, refresh: input.refresh }); return withBuiltinExternalAgentCatalog(catalog); }, ); @@ -1692,6 +1692,10 @@ export async function createExecutionRuntimeHostComposition( context.retainUntilProcessExit(); context.requestDrain(); }, + onSucceeded: () => { + pluginExecutors.invalidateCatalog(); + hostChanges.publishConfiguration(); + }, capabilities: clientCapabilities, }); oauth = new HostOAuthCoordinator({ @@ -2197,6 +2201,7 @@ export async function createExecutionRuntimeHostComposition( }); async function applyRuntimePolicyMutationEffects(): Promise { try { + pluginExecutors.invalidateCatalog(); await builtinExternalAgentPlugins.reconcile(); await requireMemory(memory).refreshAfterPolicyMutation(); } catch (error) { @@ -2220,9 +2225,10 @@ export async function createExecutionRuntimeHostComposition( continuity: continuityCoordinator, workspaceResolver, requestDrain: context.requestDrain, + isTurnBusy: (sessionId) => rootCoordinator?.hasActiveOrPendingTurn(sessionId) ?? false, configureExecutor: async (header, configuration) => { if (!header.executorId) throw new Error('Session has no executor'); - await pluginExecutors.configureConversation(header.id, header.executorId, { + return await pluginExecutors.configureConversation(header.id, header.executorId, { conversationKey: header.id, cwd: header.cwd, configuration, @@ -2239,12 +2245,19 @@ export async function createExecutionRuntimeHostComposition( if (!entry || entry.readiness !== 'ready') throw new Error('Executor is not ready'); // Catalog-managed executors pin their confirmed configuration on every create path. // Providers without model discovery retain main's executor-specific model contract. - if (entry.supportsModelChange || entry.models.length > 0) { + if ( + entry.supportsModelChange || + entry.models.length > 0 || + entry.supportsModeChange || + entry.modes?.length + ) { if ( configuration?.model && !entry.models.some((model) => model.id === configuration.model) ) throw new Error('Executor model is unavailable'); + if (configuration?.mode && !entry.modes?.some((mode) => mode.id === configuration.mode)) + throw new Error('Executor mode is unavailable'); return configuration ?? {}; } }, diff --git a/packages/runtime-host/src/server/external-agent-setup-coordinator.ts b/packages/runtime-host/src/server/external-agent-setup-coordinator.ts index db9f5eb91c9..5104f47f29a 100644 --- a/packages/runtime-host/src/server/external-agent-setup-coordinator.ts +++ b/packages/runtime-host/src/server/external-agent-setup-coordinator.ts @@ -61,6 +61,7 @@ export class HostExternalAgentSetupCoordinator { private readonly deps: { readPolicy(): Promise; onCleanupFailure(): void; + onSucceeded?(): void; acquireResidency(): OperationResidency; capabilities: Pick; install?(input: { @@ -172,6 +173,7 @@ export class HostExternalAgentSetupCoordinator { }); attempt.abort.signal.throwIfAborted(); attempt.projection = { ...attempt.projection, phase: 'succeeded', installedExecutable }; + this.deps.onSucceeded?.(); return; } await (this.deps.run ?? runAntigravitySetup)({ @@ -196,6 +198,7 @@ export class HostExternalAgentSetupCoordinator { } }, }); + if (!attempt.abort.signal.aborted) this.deps.onSucceeded?.(); attempt.projection = { ...attempt.projection, phase: attempt.abort.signal.aborted ? 'cancelled' : 'succeeded', diff --git a/packages/runtime-host/src/server/session-catalog-coordinator.ts b/packages/runtime-host/src/server/session-catalog-coordinator.ts index e66a6bb6f4b..4c7dfc345a3 100644 --- a/packages/runtime-host/src/server/session-catalog-coordinator.ts +++ b/packages/runtime-host/src/server/session-catalog-coordinator.ts @@ -207,10 +207,11 @@ export interface HostSessionCatalogCoordinatorOptions { readonly continuity: SessionContinuity; readonly workspaceResolver: HostWorkspaceResolver; readonly requestDrain: () => void; + readonly isTurnBusy?: (sessionId: string) => boolean; readonly configureExecutor?: ( header: SessionHeader, config: import('@maka/core/executor-catalog').ExecutorConfiguration, - ) => Promise; + ) => Promise; readonly retireExecutor?: (sessionId: string) => Promise; readonly assertExecutorAvailable?: ( sessionId: string, @@ -329,6 +330,7 @@ export class HostSessionCatalogCoordinator { readonly #turnIndex: SessionTurnIndexReader; readonly #runtimePolicy: SessionRuntimePolicyStores; readonly #manager: SessionConfigurationAuthority; + readonly #isTurnBusy?: (sessionId: string) => boolean; readonly #admission: SessionAdmissionGate; readonly #continuity: SessionContinuity; readonly #workspaceResolver: HostWorkspaceResolver; @@ -345,6 +347,7 @@ export class HostSessionCatalogCoordinator { this.#turnIndex = options.turnIndex; this.#runtimePolicy = options.runtimePolicy; this.#manager = options.manager; + this.#isTurnBusy = options.isTurnBusy; this.#admission = options.admission; this.#continuity = options.continuity; this.#workspaceResolver = options.workspaceResolver; @@ -828,7 +831,7 @@ export class HostSessionCatalogCoordinator { ); } - const configuration = await this.#mergeConfigurationPatch(current.header, input.patch); + let configuration = await this.#mergeConfigurationPatch(current.header, input.patch); const clearsConnectionBlock = input.patch.modelTarget !== undefined && current.header.blockedReason === 'NO_REAL_CONNECTION'; @@ -850,7 +853,8 @@ export class HostSessionCatalogCoordinator { if (input.patch.executorConfig) { if ( this.#manager.runningTurnIds(input.sessionId).length || - current.header.status === 'waiting_for_user' + this.#isTurnBusy?.(input.sessionId) || + current.header.status !== 'active' ) throw new SessionOperationFailure( 'operation_conflict', @@ -862,7 +866,24 @@ export class HostSessionCatalogCoordinator { 'Executor configuration is unavailable', ); try { - await this.#configureExecutor(current.header, input.patch.executorConfig); + const confirmed = await this.#configureExecutor( + current.header, + input.patch.executorConfig, + ); + if (confirmed) { + if ( + (input.patch.executorConfig.model && + confirmed.model !== input.patch.executorConfig.model) || + (input.patch.executorConfig.mode && + confirmed.mode !== input.patch.executorConfig.mode) + ) + throw new Error('Executor confirmation differs from the requested configuration'); + configuration = { + ...configuration, + executorConfig: { ...configuration.executorConfig, ...confirmed }, + model: confirmed.model ?? configuration.model, + }; + } confirmedExecutor = current.header; } catch { throw new SessionOperationFailure( @@ -916,7 +937,7 @@ export class HostSessionCatalogCoordinator { const actual = (await this.#stores.readHeaderRecordSnapshot(previous.id)).header; if ( actual.executorId !== previous.executorId || - !actual.executorConfig?.model || + !actual.executorConfig || !this.#configureExecutor ) throw new Error('Confirmed executor model is unavailable'); @@ -1443,7 +1464,9 @@ export class HostSessionCatalogCoordinator { return { backend: patch.modelTarget || current.backend === 'ai-sdk' ? 'ai-sdk' : 'plugin-executor', executorId: patch.modelTarget ? undefined : current.executorId, - executorConfig: patch.executorConfig ?? current.executorConfig, + executorConfig: patch.executorConfig + ? { ...current.executorConfig, ...patch.executorConfig } + : current.executorConfig, llmConnectionId: model.connectionId, llmConnectionSlug: model.connectionSlug, model: patch.executorConfig?.model ?? model.model, @@ -1473,7 +1496,9 @@ export class HostSessionCatalogCoordinator { executorConfig = await this.#assertExecutorAvailable?.( input.sessionId, input.executorId, - requestedModel ? { model: requestedModel } : input.executorConfig, + requestedModel + ? { ...input.executorConfig, model: requestedModel } + : input.executorConfig, cwd, ); } catch { @@ -1524,6 +1549,7 @@ function sessionConfigurationMatches( header.backend === configuration.backend && header.executorId === configuration.executorId && header.executorConfig?.model === configuration.executorConfig?.model && + header.executorConfig?.mode === configuration.executorConfig?.mode && header.llmConnectionId === configuration.llmConnectionId && header.llmConnectionSlug === configuration.llmConnectionSlug && header.model === configuration.model && @@ -1614,7 +1640,12 @@ function createRequestFingerprint( prepared.name, prepared.labels, input.executorId - ? ['executor', input.executorId, input.executorConfig?.model ?? input.executorModel ?? null] + ? [ + 'executor', + input.executorId, + input.executorConfig?.model ?? input.executorModel ?? null, + input.executorConfig?.mode ?? null, + ] : input.modelTarget?.kind === 'default' ? ['default'] : [ diff --git a/packages/runtime/src/__tests__/plugin-executor-service.test.ts b/packages/runtime/src/__tests__/plugin-executor-service.test.ts index aa37df2d481..e18930aa85c 100644 --- a/packages/runtime/src/__tests__/plugin-executor-service.test.ts +++ b/packages/runtime/src/__tests__/plugin-executor-service.test.ts @@ -337,6 +337,7 @@ test('catalog inspection stays process-free and model configuration is isolated let discoveries = 0, inspections = 0, configured = 0, + invalidations = 0, retired = ''; let release!: () => void; const catalog = { @@ -353,6 +354,9 @@ test('catalog inspection stays process-free and model configuration is isolated discoveries++; return catalog; }, + invalidateCatalog: () => { + invalidations++; + }, inspectConversation: async () => { inspections++; return { ...catalog, readiness: 'history_only' }; @@ -379,6 +383,8 @@ test('catalog inspection stays process-free and model configuration is isolated assert.equal(state?.readiness, 'history_only'); assert.equal(discoveries, 1); assert.equal(inspections, 1); + service.invalidateCatalog(); + assert.equal(invalidations, 1); const execution = service.execute('remote', request('session-a')); await service.configureConversation('session-b', 'remote', { conversationKey: 'session-b', diff --git a/packages/runtime/src/plugin-executor-service.ts b/packages/runtime/src/plugin-executor-service.ts index a8636d3658b..b88134790e1 100644 --- a/packages/runtime/src/plugin-executor-service.ts +++ b/packages/runtime/src/plugin-executor-service.ts @@ -157,6 +157,7 @@ export interface PluginExecutorContext { export interface PluginExecutorDiscoveryInput { readonly cwd: string; readonly signal: AbortSignal; + readonly refresh?: boolean; } export interface PluginExecutorConversationInput { readonly conversationKey: string; @@ -170,10 +171,12 @@ export interface PluginExecutorProvider { readonly capabilities?: PluginExecutorCapabilities; disposeConversation?(conversationKey: string): Promise; discover?(input: PluginExecutorDiscoveryInput): Promise; + /** Clear provider-owned discovery data after an observable account or setup change. */ + invalidateCatalog?(): void; configureConversation?( input: PluginExecutorConversationInput, signal: AbortSignal, - ): Promise; + ): Promise; inspectConversation?(input: PluginExecutorConversationInput): Promise; /** Confirm terminal consumption for any result; the provider decides whether it can checkpoint. */ acknowledgeExecution?(conversationKey: string, turnId: string): Promise; @@ -276,6 +279,9 @@ export class PluginExecutorService extends Service { ? { configureConversation: provider.configureConversation.bind(provider) } : {}), ...(provider.discover ? { discover: provider.discover.bind(provider) } : {}), + ...(provider.invalidateCatalog + ? { invalidateCatalog: provider.invalidateCatalog.bind(provider) } + : {}), ...(provider.inspectConversation ? { inspectConversation: provider.inspectConversation.bind(provider) } : {}), @@ -342,6 +348,7 @@ export class PluginExecutorService extends Service { async catalog( input: { cwd: string; + refresh?: boolean; sessionId?: string; /** Discover within this Session's scope without inspecting a retained conversation. */ discoverySessionId?: string; @@ -387,7 +394,11 @@ export class PluginExecutorService extends Service { cwd: input.cwd, ...(input.configuration ? { configuration: input.configuration } : {}), }) - : await entry.provider.discover?.({ cwd: input.cwd, signal: combined }); + : await entry.provider.discover?.({ + cwd: input.cwd, + signal: combined, + ...(input.refresh ? { refresh: true } : {}), + }); combined.throwIfAborted(); if (!result) return fallback; return normalizeCatalogEntry(result, entry.provider.id); @@ -399,11 +410,17 @@ export class PluginExecutorService extends Service { ); } + invalidateCatalog(): void { + for (const entry of this.registry.entries('profile')) { + if (!entry.retired) entry.provider.invalidateCatalog?.(); + } + } + async configureConversation( sessionId: string, executorId: string, input: PluginExecutorConversationInput, - ): Promise { + ): Promise { const entry = this.entry(sessionId, executorId); const configure = entry.provider.configureConversation; if ( @@ -413,11 +430,17 @@ export class PluginExecutorService extends Service { throw new Error('Executor configuration is unavailable or busy'); if (!isExecutorConfiguration(input.configuration)) throw new TypeError('Invalid executor configuration'); - await this.withActiveOperation( + return await this.withActiveOperation( entry, { conversationKey: input.conversationKey }, async (signal) => { - await configure(input, AbortSignal.any([signal, AbortSignal.timeout(30_000)])); + const confirmed = await configure( + input, + AbortSignal.any([signal, AbortSignal.timeout(30_000)]), + ); + if (confirmed !== undefined && !isExecutorConfiguration(confirmed)) + throw new TypeError('Invalid confirmed executor configuration'); + return confirmed; }, ); } diff --git a/packages/runtime/src/session-manager.ts b/packages/runtime/src/session-manager.ts index 313950124be..747de436a4e 100644 --- a/packages/runtime/src/session-manager.ts +++ b/packages/runtime/src/session-manager.ts @@ -5400,6 +5400,7 @@ function sessionConfigurationMatchesExceptPermissionMode( header.backend === configuration.backend && header.executorId === configuration.executorId && header.executorConfig?.model === configuration.executorConfig?.model && + header.executorConfig?.mode === configuration.executorConfig?.mode && header.llmConnectionId === configuration.llmConnectionId && header.llmConnectionSlug === configuration.llmConnectionSlug && header.connectionLocked === configuration.connectionLocked && diff --git a/packages/ui/src/__tests__/executor-model-picker.test.tsx b/packages/ui/src/__tests__/executor-model-picker.test.tsx index c143f8cd23f..39c6bf6e7eb 100644 --- a/packages/ui/src/__tests__/executor-model-picker.test.tsx +++ b/packages/ui/src/__tests__/executor-model-picker.test.tsx @@ -24,6 +24,7 @@ import type { ChatModelChoice } from '@maka/core/chat-model-choice'; import { ExecutorModelPicker, ExecutorThinkingLevelSelector, + ExecutorModeSelector, type ExecutorSelection, type ExecutorModelPickerProps, } from '../executor-model-picker.js'; @@ -76,6 +77,28 @@ const choices: ChatModelChoice[] = [ }, ]; +test('provider mode selector preserves the model and commits only a real mode ID', async () => { + const dom = installTranscriptDom(); + const selections: ExecutorSelection[] = []; + try { + await dom.render( { if (selection) selections.push(selection); }} + onSetup={() => {}} onRetry={() => {}} onNewTask={() => {}} + />); + const trigger = dom.document.querySelector('.maka-executor-mode-selector'); + assert.ok(trigger?.textContent?.includes('模式: 询问')); + await act(async () => { trigger!.dispatchEvent(new dom.window.Event('click', { bubbles: true })); }); + const auto = [...dom.document.querySelectorAll('[role="option"]')].find(row => row.textContent?.includes('自动')); + assert.ok(auto); + await act(async () => { auto.dispatchEvent(new dom.window.Event('click', { bubbles: true })); }); + assert.deepEqual(selections, [{ executorId: 'antigravity', configuration: { model: 'model-0', mode: 'auto' } }]); + } finally { await dom.cleanup(); } +}); + test('the picker shows a generic loading state until the executor catalog arrives', async () => { const dom = installTranscriptDom(); const selections: unknown[] = []; diff --git a/packages/ui/src/executor-model-picker.tsx b/packages/ui/src/executor-model-picker.tsx index 7b66869718e..f48ceac5903 100644 --- a/packages/ui/src/executor-model-picker.tsx +++ b/packages/ui/src/executor-model-picker.tsx @@ -60,6 +60,7 @@ interface ExecutorCopy { title: string; nativeOperations: string; search: string; + searchModes: string; manage: string; default: string; loading: string; @@ -76,6 +77,9 @@ interface ExecutorCopy { retry: string; newTask: string; selectionFailed: string; + invalid: string; + mode: string; + modeFailed: string; } const EXECUTOR_COPY = { @@ -83,6 +87,7 @@ const EXECUTOR_COPY = { title: 'Executor', nativeOperations: 'This operation requires Maka. Start a new Maka task.', search: 'Search models', + searchModes: 'Search modes', manage: 'Manage external agents', default: 'Agent default', loading: 'Loading agents and models…', @@ -101,11 +106,15 @@ const EXECUTOR_COPY = { retry: 'Retry', newTask: 'New task', selectionFailed: 'Model change failed. Try again.', + invalid: 'The selected Agent configuration is no longer available. Choose again.', + mode: 'Mode', + modeFailed: 'Mode change failed. Try again.', }, 'zh-CN': { title: '执行者', nativeOperations: '此操作仅支持 Maka。请新建 Maka 任务。', search: '搜索模型', + searchModes: '搜索模式', manage: '管理外部 Agent', default: 'Agent 默认', loading: '正在读取执行者与模型…', @@ -122,11 +131,15 @@ const EXECUTOR_COPY = { retry: '重试', newTask: '新建任务', selectionFailed: '模型切换失败,请重试。', + invalid: '所选 Agent 配置已失效,请重新选择。', + mode: '模式', + modeFailed: '模式切换失败,请重试。', }, 'zh-TW': { title: '執行者', nativeOperations: '此操作僅支援 Maka。請建立 Maka 任務。', search: '搜尋模型', + searchModes: '搜尋模式', manage: '管理外部 Agent', default: 'Agent 預設', loading: '正在讀取執行者與模型…', @@ -143,6 +156,9 @@ const EXECUTOR_COPY = { retry: '重試', newTask: '建立新任務', selectionFailed: '模型切換失敗,請重試。', + invalid: '所選 Agent 設定已失效,請重新選擇。', + mode: '模式', + modeFailed: '模式切換失敗,請重試。', }, } satisfies UiCatalog; @@ -175,6 +191,10 @@ export function ExecutorModelPicker(props: ExecutorModelPickerProps) { ? (props.selection?.configuration.model ?? browsed?.currentModel) : browsed?.currentModel; const selectedUnavailable = !!props.selection && selected?.readiness !== 'ready'; + const selectedInvalid = !!props.selection && !!selected && selected.readiness === 'ready' && ( + (!!props.selection.configuration.model && !selected.models.some(model => model.id === props.selection!.configuration.model)) || + (!!props.selection.configuration.mode && !selected.modes?.some(mode => mode.id === props.selection!.configuration.mode)) + ); const selectedModel = props.selection?.configuration.model ?? selected?.currentModel; const selectedGroup = selected && executorModelGroup(selected, selectedModel); const currentGroup = browsed && executorModelGroup(browsed, currentModel); @@ -270,6 +290,14 @@ export function ExecutorModelPicker(props: ExecutorModelPickerProps) { className="maka-executor-picker-entry maka-executor-picker-manage" onClick={openSetup} /> + {props.selection ? : props.nativeThinkingControl} - {(selectedUnavailable || props.error) && ( + + {(selectedUnavailable || selectedInvalid || props.error) && ( - {selected && selected.readiness !== 'ready' ? copy[selected.readiness] : copy.unavailable} + {selected && selected.readiness !== 'ready' ? copy[selected.readiness] : selectedInvalid ? copy.invalid : copy.selectionFailed} {selected?.readiness === 'restorable' || selected?.readiness === 'restore_failed' ? ( - {props.selection ? : props.nativeThinkingControl} - + {props.selection ? : props.nativeThinkingControl} + {(selectedUnavailable || selectedInvalid || props.error) && ( {selected && selected.readiness !== 'ready' ? copy[selected.readiness] : selectedInvalid ? copy.invalid : copy.selectionFailed} {selected?.readiness === 'restorable' || selected?.readiness === 'restore_failed' ? ( -