diff --git a/packages/runtime/src/__tests__/plugin-executor-backend.test.ts b/packages/runtime/src/__tests__/plugin-executor-backend.test.ts index f521703eff..97d8a2c5df 100644 --- a/packages/runtime/src/__tests__/plugin-executor-backend.test.ts +++ b/packages/runtime/src/__tests__/plugin-executor-backend.test.ts @@ -118,6 +118,37 @@ test('failed Plugin acknowledgement abandons the settled external execution', as } }); +test('a settled Plugin result that the Runtime rejects is abandoned once', async () => { + const { root, backend, acknowledged, abandoned } = recordingFixture(async () => ({ + status: 'completed', + text: 'x'.repeat(256 * 1024 + 1), + })); + try { + const events = await collect(backend.send({ turnId: 'turn-a', text: 'task' })); + assert.match(events[0]?.type === 'error' ? events[0].message : '', /completion text exceeds/u); + assert.deepEqual(acknowledged, []); + assert.deepEqual(abandoned, ['session-a/turn-a']); + } finally { + await backend.dispose(); + await root.fiber.dispose(); + } +}); + +test('a Plugin execution that fails before settling is neither acknowledged nor abandoned', async () => { + const { root, backend, acknowledged, abandoned } = recordingFixture(async () => { + throw new Error('transient transport failure'); + }); + try { + const events = await collect(backend.send({ turnId: 'turn-a', text: 'task' })); + assert.match(events[0]?.type === 'error' ? events[0].message : '', /transient transport/u); + assert.deepEqual(acknowledged, []); + assert.deepEqual(abandoned, []); + } finally { + await backend.dispose(); + await root.fiber.dispose(); + } +}); + test('executor backend converts plugin output and result to ordinary Session events', async () => { const { root, binding } = fixture(async (request, context) => { assert.equal(request.instructions, 'child instructions'); @@ -457,6 +488,38 @@ function fixture( return { root, binding: service.bind('session-a', 'remote'), dispose }; } +function recordingFixture(execute: Parameters[0]['execute']): { + root: Context; + backend: PluginExecutorBackend; + acknowledged: string[]; + abandoned: string[]; +} { + const root = new Context(); + const service = new PluginExecutorService(root); + const acknowledged: string[] = []; + const abandoned: string[] = []; + root + .extend({ + maka: { rootId: 'profile', packageId: 'fixture', entryId: 'provider', generation: 1 }, + }) + .executors.register({ + id: 'remote', + execute, + acknowledgeExecution: async (conversationKey, turnId) => { + acknowledged.push(`${conversationKey}/${turnId}`); + }, + abandonExecution: async (conversationKey, turnId) => { + abandoned.push(`${conversationKey}/${turnId}`); + }, + }); + const backend = new PluginExecutorBackend({ + sessionId: 'session-a', + cwd: '/workspace', + binding: service.bind('session-a', 'remote'), + }); + return { root, backend, acknowledged, abandoned }; +} + function ids(): () => string { let value = 0; return () => `id-${++value}`; diff --git a/packages/runtime/src/plugin-executor-backend.ts b/packages/runtime/src/plugin-executor-backend.ts index 81bbed4021..10e490aaaa 100644 --- a/packages/runtime/src/plugin-executor-backend.ts +++ b/packages/runtime/src/plugin-executor-backend.ts @@ -44,6 +44,12 @@ interface ActiveExecution { readonly settled: Promise; } +/** + * Whether the provider settled a terminal result, and whether it was published. + * Only a settled result may be acknowledged or abandoned. + */ +type ExecutionSettlement = 'published' | 'rejected' | 'unsettled'; + export interface PluginExecutorBackendInput { readonly sessionId: string; readonly cwd: string; @@ -99,8 +105,8 @@ export class PluginExecutorBackend implements AgentBackend { yield event; queue.ackConsumed(); } - const returnedResult = await producer; - if (returnedResult) { + const settlement = await producer; + if (settlement === 'published') { // The Runtime Kernel requests the next item only after onSessionEvent // resolves. Reaching this point means its terminal event was accepted. // The Plugin decides whether its external execution actually settled; @@ -116,9 +122,12 @@ export class PluginExecutorBackend implements AgentBackend { } finally { queue.noteConsumerDetached(); abort.abort(new Error('Plugin executor event consumer detached')); - const returnedResult = await producer.catch(() => false); + // A settled result that was not acknowledged is abandoned, including one + // the Runtime rejected before publishing it. A provider failure before + // settlement is neither acknowledged nor abandoned. + const settlement = await producer.catch((): ExecutionSettlement => 'unsettled'); if ( - returnedResult && + settlement !== 'unsettled' && this.#binding.acknowledgeExecution && !acknowledged && this.#binding.abandonExecution @@ -153,7 +162,7 @@ export class PluginExecutorBackend implements AgentBackend { messageId: string, signal: AbortSignal, queue: AsyncEventQueue, - ): Promise { + ): Promise { const turnId = input.turnId; let thinkingText = ''; const toolUseIds = new Map(); @@ -161,6 +170,7 @@ export class PluginExecutorBackend implements AgentBackend { let result: PluginExecutorResult | undefined; let failure: unknown; let failed = false; + let settled = false; try { result = await this.#binding.execute( { @@ -180,6 +190,9 @@ export class PluginExecutorBackend implements AgentBackend { }, { signal, + onSettled: () => { + settled = true; + }, onEvent: (event) => { if (event.type === 'thinking_delta') thinkingText += event.text; this.#publishOutputEvent( @@ -212,7 +225,7 @@ export class PluginExecutorBackend implements AgentBackend { false, queue, ); - return false; + return settled ? 'rejected' : 'unsettled'; } if (result === undefined) { this.#publishFailure( @@ -222,10 +235,10 @@ export class PluginExecutorBackend implements AgentBackend { false, queue, ); - return false; + return settled ? 'rejected' : 'unsettled'; } this.#publishResult(turnId, messageId, result, queue); - return true; + return 'published'; } async #requestPermission( diff --git a/packages/runtime/src/plugin-executor-service.ts b/packages/runtime/src/plugin-executor-service.ts index a8636d3658..a900a580c9 100644 --- a/packages/runtime/src/plugin-executor-service.ts +++ b/packages/runtime/src/plugin-executor-service.ts @@ -187,6 +187,8 @@ export interface PluginExecutorProvider { export interface PluginExecutorExecutionOptions { readonly signal?: AbortSignal; + /** Called when the provider returns a terminal result, before the Runtime validates it. */ + readonly onSettled?: () => void; readonly onEvent?: (event: PluginExecutorOutputEvent) => void; readonly onPermissionRequest?: ( request: PluginExecutorPermissionRequest, @@ -527,6 +529,11 @@ export class PluginExecutorService extends Service { return normalizePermissionResult(result, normalized); }, }); + try { + options.onSettled?.(); + } catch { + // A settlement observer must not change external execution. + } if (signal.aborted) { const cancelled = cancelledResult(signal.reason); const normalized = normalizeResult(result);