Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
63 changes: 63 additions & 0 deletions packages/runtime/src/__tests__/plugin-executor-backend.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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');
Expand Down Expand Up @@ -457,6 +488,38 @@ function fixture(
return { root, binding: service.bind('session-a', 'remote'), dispose };
}

function recordingFixture(execute: Parameters<PluginExecutorService['register']>[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}`;
Expand Down
29 changes: 21 additions & 8 deletions packages/runtime/src/plugin-executor-backend.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,12 @@ interface ActiveExecution {
readonly settled: Promise<void>;
}

/**
* 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;
Expand Down Expand Up @@ -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;
Expand All @@ -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
Expand Down Expand Up @@ -153,14 +162,15 @@ export class PluginExecutorBackend implements AgentBackend {
messageId: string,
signal: AbortSignal,
queue: AsyncEventQueue<SessionEvent>,
): Promise<boolean> {
): Promise<ExecutionSettlement> {
const turnId = input.turnId;
let thinkingText = '';
const toolUseIds = new Map<string, string>();
const toolOutputSequences = new Map<string, number>();
let result: PluginExecutorResult | undefined;
let failure: unknown;
let failed = false;
let settled = false;
try {
result = await this.#binding.execute(
{
Expand All @@ -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(
Expand Down Expand Up @@ -212,7 +225,7 @@ export class PluginExecutorBackend implements AgentBackend {
false,
queue,
);
return false;
return settled ? 'rejected' : 'unsettled';
}
if (result === undefined) {
this.#publishFailure(
Expand All @@ -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(
Expand Down
7 changes: 7 additions & 0 deletions packages/runtime/src/plugin-executor-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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);
Expand Down