diff --git a/packages/agent-runtime/src/pi/event-translation.test.ts b/packages/agent-runtime/src/pi/event-translation.test.ts index 407f91e2e5..06ee8dd9a6 100644 --- a/packages/agent-runtime/src/pi/event-translation.test.ts +++ b/packages/agent-runtime/src/pi/event-translation.test.ts @@ -3,6 +3,7 @@ import { resolve, dirname } from "node:path"; import { fileURLToPath } from "node:url"; import { describe, expect, it } from "vitest"; import { threadScope, turnScope } from "@bb/domain"; +import { queueAcceptedUserMessage } from "@bb/provider-bridge-protocol/bridge-kit"; import type { AgentSessionEvent } from "@earendil-works/pi-coding-agent"; import { createPiEventTranslator } from "./event-translation.js"; @@ -77,6 +78,74 @@ interface PiBashStartEventArgs { toolCallId: string; } +interface PiProcessNotificationEventArgs { + attention: "context" | "ignore" | "turn"; + content: string; + eventType?: "message_end" | "message_start"; + kind: "log_match" | "success"; + processId: string; + threadId?: string; +} + +function createPiProcessNotificationEvent( + args: PiProcessNotificationEventArgs, +) { + return { + jsonrpc: "2.0" as const, + method: "sdk/message", + params: { + threadId: args.threadId ?? "pi-thread-1", + message: { + type: args.eventType ?? "message_start", + message: { + role: "custom", + customType: "ad-process:notification", + content: args.content, + display: true, + details: { + attention: args.attention, + kind: args.kind, + processId: args.processId, + }, + timestamp: 1_786_919_243_630, + }, + }, + }, + }; +} + +function createPiEmptyAgentEndEvent(): AgentSessionEvent { + return { + type: "agent_end", + willRetry: false, + messages: [ + { + role: "assistant", + content: [], + stopReason: "stop", + api: "openai-responses", + provider: "openai-codex", + model: "gpt-5.6-sol", + usage: { + input: 10, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 10, + cost: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + total: 0, + }, + }, + timestamp: 1_786_919_246_950, + }, + ], + }; +} + function createPiBashStartEvent(args: PiBashStartEventArgs): AgentSessionEvent { return { type: "tool_execution_start", @@ -221,6 +290,405 @@ describe("pi event translation", () => { expect(events.some((event) => event.type === "provider/error")).toBe(false); }); + it("completes extension-triggered turns when agent_end includes string custom content", () => { + const translator = createTranslator(); + const context = { threadId: "pi-thread-1" }; + const agentEndEvent = { + type: "agent_end", + messages: [ + { + role: "custom", + customType: "pi-processes", + content: "Process completed successfully", + display: true, + timestamp: 1777995780000, + }, + { + role: "assistant", + content: [{ type: "text", text: "The process finished." }], + api: "anthropic-messages", + provider: "anthropic", + model: "claude-haiku-4-5", + usage: { + input: 10, + output: 5, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 15, + cost: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + total: 0, + }, + }, + stopReason: "stop", + timestamp: 1777995781000, + }, + ], + willRetry: false, + } satisfies AgentSessionEvent; + + translator.translatePiEvent( + { + jsonrpc: "2.0", + method: "sdk/message", + params: { + threadId: context.threadId, + message: { type: "agent_start" }, + }, + }, + context, + ); + + const events = translator.translatePiEvent( + { + jsonrpc: "2.0", + method: "sdk/message", + params: { + threadId: context.threadId, + message: agentEndEvent, + }, + }, + context, + ); + + expect(events).toContainEqual( + expect.objectContaining({ + type: "turn/completed", + scope: turnScope("turn-1"), + status: "completed", + }), + ); + expect(events.some((event) => event.type === "provider/unhandled")).toBe( + false, + ); + }); + + it("projects one turn for a process notification and accepts it in agent_end", () => { + const translator = createTranslator(); + const context = { threadId: "pi-thread-1" }; + const content = ''; + + const notificationEvents = translator.translatePiEvent( + createPiProcessNotificationEvent({ + attention: "turn", + content, + kind: "success", + processId: "proc_sleep", + }), + context, + ); + const duplicateEndEvents = translator.translatePiEvent( + createPiProcessNotificationEvent({ + attention: "turn", + content, + eventType: "message_end", + kind: "success", + processId: "proc_sleep", + }), + context, + ); + const providerStartEvents = translator.translatePiEvent( + loadFixture("agent-start.json"), + context, + ); + const terminalEvents = translator.translatePiEvent( + { + type: "agent_end", + messages: [ + { + role: "custom", + customType: "ad-process:notification", + content, + details: { attention: "turn" }, + }, + { + role: "assistant", + content: [{ type: "text", text: "The process finished." }], + stopReason: "stop", + }, + ], + providerCheckpointId: "pi-entry-process", + willRetry: false, + }, + context, + ); + const events = [ + ...notificationEvents, + ...duplicateEndEvents, + ...providerStartEvents, + ...terminalEvents, + ]; + + expect(notificationEvents).toEqual([ + expect.objectContaining({ + type: "turn/started", + scope: turnScope("turn-1"), + }), + expect.objectContaining({ + type: "item/completed", + scope: turnScope("turn-1"), + item: expect.objectContaining({ + type: "userMessage", + content: [{ type: "text", text: content }], + }), + }), + ]); + expect(duplicateEndEvents).toEqual([]); + expect(providerStartEvents).toEqual([]); + expect( + events.filter((event) => event.type === "turn/started"), + ).toHaveLength(1); + expect(events).toContainEqual( + expect.objectContaining({ + type: "item/completed", + item: expect.objectContaining({ + type: "agentMessage", + text: "The process finished.", + }), + }), + ); + expect(events).toContainEqual( + expect.objectContaining({ + type: "turn/completed", + providerCheckpointId: "pi-entry-process", + }), + ); + expect(events.some((event) => event.type === "provider/unhandled")).toBe( + false, + ); + }); + + it.each(["context", "ignore"] as const)( + "does not open an idle turn for a process notification with %s attention", + (attention) => { + const translator = createTranslator(); + const context = { threadId: "pi-thread-1" }; + + const events = translator.translatePiEvent( + createPiProcessNotificationEvent({ + attention, + content: ``, + kind: "success", + processId: `proc_${attention}`, + }), + context, + ); + + expect(events).toEqual([]); + expect( + translator.translatePiEvent(loadFixture("agent-start.json"), context), + ).toContainEqual( + expect.objectContaining({ + type: "turn/started", + scope: turnScope("turn-1"), + }), + ); + }, + ); + + it("uses the translator item prefix for process notification IDs", () => { + const firstTranslator = createPiEventTranslator({ + providerId: "pi", + itemIdPrefix: "resume-a-", + }); + const secondTranslator = createPiEventTranslator({ + providerId: "pi", + itemIdPrefix: "resume-b-", + }); + const notification = createPiProcessNotificationEvent({ + attention: "turn", + content: '', + kind: "success", + processId: "proc_sleep", + }); + + const firstEvents = firstTranslator.translatePiEvent(notification, { + threadId: "pi-thread-1", + }); + const secondEvents = secondTranslator.translatePiEvent(notification, { + threadId: "pi-thread-1", + }); + + expect(firstEvents).toContainEqual( + expect.objectContaining({ + type: "item/completed", + item: expect.objectContaining({ + type: "userMessage", + id: "resume-a-process-notification-1", + }), + }), + ); + expect(secondEvents).toContainEqual( + expect.objectContaining({ + type: "item/completed", + item: expect.objectContaining({ + type: "userMessage", + id: "resume-b-process-notification-1", + }), + }), + ); + }); + + it("coalesces clustered process notifications into one active turn", () => { + const translator = createTranslator(); + const context = { threadId: "pi-thread-1" }; + + const lifecycleEvents = translator.translatePiEvent( + createPiProcessNotificationEvent({ + attention: "turn", + content: '', + kind: "success", + processId: "proc_sleep", + }), + context, + ); + const logMatchEvents = translator.translatePiEvent( + createPiProcessNotificationEvent({ + attention: "turn", + content: + 'done', + kind: "log_match", + processId: "proc_sleep", + }), + context, + ); + const events = [...lifecycleEvents, ...logMatchEvents]; + + expect( + events.filter((event) => event.type === "turn/started"), + ).toHaveLength(1); + expect( + events.filter( + (event) => + event.type === "item/completed" && event.item.type === "userMessage", + ), + ).toHaveLength(2); + }); + + it("warns when a process-triggered turn ends without assistant text", () => { + const translator = createTranslator(); + const context = { threadId: "pi-thread-1" }; + translator.translatePiEvent( + createPiProcessNotificationEvent({ + attention: "turn", + content: '', + kind: "success", + processId: "proc_sleep", + }), + context, + ); + + const events = translator.translatePiEvent( + createPiEmptyAgentEndEvent(), + context, + ); + + expect(events).toContainEqual( + expect.objectContaining({ + type: "provider/warning", + scope: turnScope("turn-1"), + summary: + "Pi completed the process notification turn without a text response", + }), + ); + expect(events).toContainEqual( + expect.objectContaining({ + type: "turn/completed", + scope: turnScope("turn-1"), + }), + ); + }); + + it("projects context notifications into an active user turn without reclassifying it", () => { + const translator = createTranslator(); + const context = { threadId: "pi-thread-1" }; + translator.translatePiEvent(loadFixture("agent-start.json"), context); + + const notificationEvents = translator.translatePiEvent( + createPiProcessNotificationEvent({ + attention: "context", + content: '', + kind: "success", + processId: "proc_sleep", + }), + context, + ); + const terminalEvents = translator.translatePiEvent( + createPiEmptyAgentEndEvent(), + context, + ); + + expect(notificationEvents).toEqual([ + expect.objectContaining({ + type: "item/completed", + item: expect.objectContaining({ type: "userMessage" }), + }), + ]); + expect( + terminalEvents.some((event) => event.type === "provider/warning"), + ).toBe(false); + }); + + it("ignores ignore-attention notifications during an active turn", () => { + const translator = createTranslator(); + const context = { threadId: "pi-thread-1" }; + translator.translatePiEvent(loadFixture("agent-start.json"), context); + + const events = translator.translatePiEvent( + createPiProcessNotificationEvent({ + attention: "ignore", + content: '', + kind: "success", + processId: "proc_sleep", + }), + context, + ); + + expect(events).toEqual([]); + }); + + it("keeps an accepted user prompt from becoming a process-triggered turn", () => { + const translator = createTranslator(); + const context = { threadId: "pi-thread-1" }; + queueAcceptedUserMessage({ + clientRequestId: "creq_abcdefghjk", + state: translator.resolveState(context), + }); + + const notificationEvents = translator.translatePiEvent( + createPiProcessNotificationEvent({ + attention: "turn", + content: '', + kind: "success", + processId: "proc_sleep", + }), + context, + ); + const providerStartEvents = translator.translatePiEvent( + loadFixture("agent-start.json"), + context, + ); + const terminalEvents = translator.translatePiEvent( + createPiEmptyAgentEndEvent(), + context, + ); + + expect(notificationEvents).toContainEqual( + expect.objectContaining({ + type: "turn/input/accepted", + clientRequestId: "creq_abcdefghjk", + scope: turnScope("turn-1"), + }), + ); + expect(providerStartEvents).toEqual([]); + expect( + terminalEvents.some((event) => event.type === "provider/warning"), + ).toBe(false); + }); + it("translateEvent agent_end surfaces Pi assistant stop errors as failed turns", () => { const translator = createTranslator(); const context = { threadId: "bb-thread-1" }; diff --git a/packages/agent-runtime/src/pi/event-translation.ts b/packages/agent-runtime/src/pi/event-translation.ts index 50ef647e56..ac6e719d5d 100644 --- a/packages/agent-runtime/src/pi/event-translation.ts +++ b/packages/agent-runtime/src/pi/event-translation.ts @@ -140,6 +140,24 @@ const piIgnoredEventSchema = z .passthrough() .refine((event) => PI_IGNORED_EVENT_TYPES.has(event.type)); +const piProcessNotificationEventSchema = z + .object({ + type: z.enum(["message_start", "message_end"]), + message: z + .object({ + role: z.literal("custom"), + customType: z.literal("ad-process:notification"), + content: z.string(), + details: z + .object({ + attention: z.enum(["turn", "context", "ignore"]), + }) + .passthrough(), + }) + .passthrough(), + }) + .passthrough(); + const piMessageContentBlockSchema = z .object({ type: z.string(), @@ -171,7 +189,9 @@ const piAssistantMessageSchema = z const piConversationMessageSchema = z .object({ role: z.string(), - content: z.array(piMessageContentBlockSchema).optional(), + content: z + .union([z.string(), z.array(piMessageContentBlockSchema)]) + .optional(), stopReason: z.string().optional(), errorMessage: z.string().optional(), provider: z.string().optional(), @@ -405,6 +425,8 @@ export interface PiTurnState { openAssistantMessageIdsByScope: Map; openReasoningItemIdsByScope: Map; pendingAcceptedUserMessages: AcceptedUserMessageState["pendingAcceptedUserMessages"]; + processNotificationCounter: number; + processNotificationTurnId: string | undefined; reasoningItemCounter: number; toolItemsByCallId: Map; } @@ -452,6 +474,10 @@ export function createPiEventTranslator( options.itemIdPrefix === undefined ? "pi-assistant" : `${options.itemIdPrefix}assistant`; + const processNotificationIdPrefix = + options.itemIdPrefix === undefined + ? "pi-process-notification" + : `${options.itemIdPrefix}process-notification`; const piReasoningItemIds = createScopedItemIdFactory({ prefix: options.itemIdPrefix === undefined @@ -477,9 +503,14 @@ export function createPiEventTranslator( openAssistantMessageIdsByScope: new Map(), openReasoningItemIdsByScope: new Map(), pendingAcceptedUserMessages: [], + processNotificationCounter: 0, + processNotificationTurnId: undefined, reasoningItemCounter: 0, toolItemsByCallId: new Map(), }), + onTurnFinish: ({ state }) => { + state.processNotificationTurnId = undefined; + }, onTurnStart: ({ events, state, threadId, turnId }) => { resetPiCommandOutputSnapshots(state); drainAcceptedUserMessages({ @@ -555,6 +586,9 @@ export function createPiEventTranslator( ) { return []; } + const isProcessNotification = piProcessNotificationEventSchema.safeParse( + sdkEnvelope.data.params.message, + ).success; const parentToolCallId = sdkEnvelope.data.params.parent_tool_use_id ?? context?.parentToolCallId; const translated = translatePiEvent(sdkEnvelope.data.params.message, { @@ -562,7 +596,7 @@ export function createPiEventTranslator( ...(parentToolCallId ? { parentToolCallId } : {}), }); const fallbackTurnId = resolvePiActiveTurnId(context); - return translated.length > 0 + return translated.length > 0 || isProcessNotification ? translated : buildUnhandledPiEvent({ rawEvent: { @@ -661,6 +695,55 @@ export function createPiEventTranslator( }); } + const processNotification = + piProcessNotificationEventSchema.safeParse(event); + if (processNotification.success) { + if (processNotification.data.type === "message_end") { + return []; + } + + const stateKey = context?.threadId ?? ""; + const state = turnState.getOrCreate({ threadId: stateKey }); + const attention = processNotification.data.message.details.attention; + if ( + attention === "ignore" || + (state.currentTurnId === undefined && attention !== "turn") + ) { + return []; + } + + const events: ThreadEvent[] = []; + const startsProcessTurn = + state.currentTurnId === undefined && + state.pendingAcceptedUserMessages.length === 0; + const turnId = turnState.ensureTurnStarted({ + events, + state, + threadId: UNSTAMPED_THREAD_ID, + }); + state.processNotificationCounter += 1; + if (startsProcessTurn) { + state.processNotificationTurnId = turnId; + } + events.push({ + type: "item/completed", + threadId: UNSTAMPED_THREAD_ID, + providerThreadId: "", + scope: turnScope(turnId), + item: { + type: "userMessage", + id: `${processNotificationIdPrefix}-${state.processNotificationCounter}`, + content: [ + { + type: "text", + text: processNotification.data.message.content, + }, + ], + }, + }); + return events; + } + const eventType = piEventTypeSchema.safeParse(event); if (!eventType.success) { return []; @@ -807,22 +890,32 @@ export function createPiEventTranslator( }), ]; } - if (lastAssistant) { - const text = extractAssistantText(lastAssistant); - if (text) { - const itemId = turnState.resolveCompletedAssistantMessageId({ - assistantIdPrefix, - parentToolCallId: context?.parentToolCallId, - state, - }); - events.push({ - type: "item/completed", - threadId, - providerThreadId: "", - scope: turnScope(currentTurnId), - item: { type: "agentMessage", id: itemId, text }, - }); - } + const assistantText = lastAssistant + ? extractAssistantText(lastAssistant) + : undefined; + if (assistantText) { + const itemId = turnState.resolveCompletedAssistantMessageId({ + assistantIdPrefix, + parentToolCallId: context?.parentToolCallId, + state, + }); + events.push({ + type: "item/completed", + threadId, + providerThreadId: "", + scope: turnScope(currentTurnId), + item: { type: "agentMessage", id: itemId, text: assistantText }, + }); + } else if (state.processNotificationTurnId === currentTurnId) { + events.push({ + type: "provider/warning", + threadId, + providerThreadId: "", + scope: turnScope(currentTurnId), + category: "general", + summary: + "Pi completed the process notification turn without a text response", + }); } const tokenUsage = extractPiTokenUsage( lastAssistant,