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,