Skip to content
Closed
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
426 changes: 426 additions & 0 deletions packages/agent-runtime/src/pi/event-translation.test.ts

Large diffs are not rendered by default.

126 changes: 110 additions & 16 deletions packages/agent-runtime/src/pi/event-translation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,8 @@ const piEventTypeSchema = z
"agent_start",
"compaction_end",
"compaction_start",
"message_end",
"message_start",
"message_update",
"tool_execution_end",
"tool_execution_start",
Expand Down Expand Up @@ -139,6 +141,25 @@ const piIgnoredEventSchema = z
.passthrough()
.refine((event) => PI_IGNORED_EVENT_TYPES.has(event.type));

const piCustomMessageBoundaryEventSchema = z
.object({
type: z.enum(["message_end", "message_start"]),
message: z
.object({
role: z.literal("custom"),
content: z.string(),
display: z.boolean().default(false),
details: z
.object({
attention: z.enum(["context", "ignore", "turn"]).optional(),
})
.passthrough()
.optional(),
})
.passthrough(),
})
.passthrough();

const piMessageContentBlockSchema = z
.object({
type: z.string(),
Expand Down Expand Up @@ -447,8 +468,11 @@ export interface PiTurnState {
openAssistantMessageIdsByScope: Map<string, string>;
openScopedItemIdsByScope: Map<string, string>;
pendingAcceptedUserMessages: AcceptedUserMessageState["pendingAcceptedUserMessages"];
providerInputCounter: number;
providerInputTurnId: string | undefined;
scopedItemCounter: number;
toolItemsByCallId: Map<string, ThreadEventItem>;
turnStartedWithAcceptedUserMessage: boolean;
}

const piCompactionItemIds = createScopedItemIdFactory({
Expand Down Expand Up @@ -494,6 +518,10 @@ export function createPiEventTranslator(
options.itemIdPrefix === undefined
? "pi-assistant"
: `${options.itemIdPrefix}assistant`;
const providerInputIdPrefix =
options.itemIdPrefix === undefined
? "pi-provider-input"
: `${options.itemIdPrefix}provider-input`;
const piReasoningItemIds = createScopedItemIdFactory({
prefix:
options.itemIdPrefix === undefined
Expand All @@ -519,11 +547,20 @@ export function createPiEventTranslator(
openAssistantMessageIdsByScope: new Map(),
openScopedItemIdsByScope: new Map(),
pendingAcceptedUserMessages: [],
providerInputCounter: 0,
providerInputTurnId: undefined,
scopedItemCounter: 0,
toolItemsByCallId: new Map(),
turnStartedWithAcceptedUserMessage: false,
}),
onTurnFinish: ({ state }) => {
state.providerInputTurnId = undefined;
state.turnStartedWithAcceptedUserMessage = false;
},
onTurnStart: ({ state }) => {
resetPiCommandOutputSnapshots(state);
state.turnStartedWithAcceptedUserMessage =
state.pendingAcceptedUserMessages.length > 0;
},
turnIdPrefix: options.turnIdPrefix,
});
Expand Down Expand Up @@ -715,6 +752,53 @@ export function createPiEventTranslator(
});

switch (eventType.data.type) {
case "message_end":
case "message_start": {
const piEvent = piCustomMessageBoundaryEventSchema.safeParse(event);
if (!piEvent.success) {
return buildUnexpectedEvent(event);
}
if (
piEvent.data.type === "message_end" ||
!piEvent.data.message.display ||
piEvent.data.message.details?.attention === "ignore"
) {
return [];
}
if (
state.currentTurnId === undefined &&
piEvent.data.message.details?.attention === "context"
) {
return [];
}
const turnId = turnState.ensureTurnStarted({
events,
state,
threadId,
});
state.providerInputCounter += 1;
if (!state.turnStartedWithAcceptedUserMessage) {
state.providerInputTurnId = turnId;
}
events.push({
type: "item/completed",
threadId,
providerThreadId: "",
scope: turnScope(turnId),
item: {
type: "userMessage",
id: `${providerInputIdPrefix}-${state.providerInputCounter}`,
content: [
{
type: "text",
text: piEvent.data.message.content,
},
],
},
});
break;
}

case "agent_start": {
const piEvent = piAgentStartEventSchema.safeParse(event);
if (!piEvent.success) {
Expand Down Expand Up @@ -859,22 +943,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.providerInputTurnId === currentTurnId) {
events.push({
type: "provider/warning",
threadId,
providerThreadId: "",
scope: turnScope(currentTurnId),
category: "general",
summary:
"Pi completed the provider-triggered turn without a text response",
});
}
const tokenUsage = extractPiTokenUsage(
lastAssistant,
Expand Down
18 changes: 17 additions & 1 deletion packages/agent-runtime/src/pi/visibility.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,12 @@ type PiAssistantEventType =
| "toolcall_start"
| "unknown";

type PiMessageBoundaryRole = "assistant" | "toolResult" | "user" | "unknown";
type PiMessageBoundaryRole =
| "assistant"
| "custom"
| "toolResult"
| "user"
| "unknown";

type PiSdkEventType =
| "agent_end"
Expand Down Expand Up @@ -75,6 +80,7 @@ interface PiSimpleSdkRawEvent {
}

interface PiMessageBoundaryRawEvent {
display: boolean;
kind: "sdk/message-boundary";
role: PiMessageBoundaryRole;
sdkType: "message_end" | "message_start";
Expand Down Expand Up @@ -134,6 +140,7 @@ function toPiMessageBoundaryRole(
): PiMessageBoundaryRole {
switch (role) {
case "assistant":
case "custom":
case "toolResult":
case "user":
return role;
Expand Down Expand Up @@ -214,6 +221,7 @@ function parsePiRawEvent(event: JsonRpcMessage): PiRawEvent {
case "message_end": {
const payload = getRecordProperty(message, "message");
return {
display: payload?.["display"] === true,
kind: "sdk/message-boundary",
sdkType,
role: toPiMessageBoundaryRole(
Expand Down Expand Up @@ -319,6 +327,14 @@ function describeParsedPiRawEvent(
switch (event.role) {
case "assistant":
return { kind, coverage: "noise" };
case "custom":
return {
kind,
coverage:
event.sdkType === "message_start" && event.display
? "normalized"
: "noise",
};
case "toolResult":
case "user":
return { kind, coverage: "noise" };
Expand Down
7 changes: 6 additions & 1 deletion packages/host-daemon-contract/src/protocol.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,8 @@
// Version 136 translates visible Pi custom-message boundaries into provider
// input items, preserves their attention semantics, and warns when a
// provider-triggered turn returns no assistant text. Older daemons still send
// those boundaries as unhandled events and omit the warning.
//
// Version 135 adds the `compaction-skipped` provider warning category. The Pi
// bridge now reports a refused manual compaction ("Nothing to compact") as
// that warning plus a completed turn instead of a failed turn. An older daemon
Expand Down Expand Up @@ -44,7 +49,7 @@
//
// The version mismatch is what triggers the enrolled daemon's automatic update
// instead of an `invalid-message` reconnect loop.
export const HOST_DAEMON_PROTOCOL_VERSION = 135 as const;
export const HOST_DAEMON_PROTOCOL_VERSION = 136 as const;

/**
* Absolute ceiling for any executable artifact delivered to a host daemon —
Expand Down
4 changes: 3 additions & 1 deletion packages/host-daemon-contract/test/contract.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1133,6 +1133,8 @@ describe("host-daemon command schemas", () => {
// Version 116 reports provider exits that happen while a turn start is
// pending. Older daemons can leave the server thread active until the live
// command timeout, so enrolled machines must update before handling turns.
// Version 136 projects visible Pi custom messages as provider input and
// reports empty provider-triggered turns instead of silently completing them.
// Version 115 settles zero-work provider prompts with a complete synthetic
// turn lifecycle. Older daemons can leave locally handled prompts active
// indefinitely, so enrolled machines must update for reliable completion.
Expand All @@ -1142,7 +1144,7 @@ describe("host-daemon command schemas", () => {
// mixed version. Version 113 carried the Devin Desktop open target rename
// and remains part of the protocol lineage.
it("uses the current host-daemon protocol version", () => {
expect(HOST_DAEMON_PROTOCOL_VERSION).toBe(135);
expect(HOST_DAEMON_PROTOCOL_VERSION).toBe(136);
expect(HOST_ARTIFACT_MAX_BYTES).toBe(256 * 1024 * 1024);
});

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import { describe, expect, it } from "vitest";
import { createUnhandledProviderEvent } from "./provider-unhandled-event.js";

describe("provider unhandled events", () => {
it("does not throw when raw event params are not JSON-serializable", () => {
it("preserves JSON-compatible fields when provider params contain undefined", () => {
const event = createUnhandledProviderEvent({
providerId: "test-provider",
rawType: "sdk/custom",
Expand All @@ -27,8 +27,8 @@ describe("provider unhandled events", () => {
jsonrpc: "2.0",
method: "sdk/message",
params: {
serializationError:
"Provider raw event params were not JSON-serializable.",
threadId: "thread-1",
nested: {},
},
},
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,14 @@ export interface BuildUnhandledProviderEventsArgs {
}

function toProviderRawEvent(rawEvent: JsonRpcMessage): ProviderRawEvent {
const parsed = providerRawEventSchema.safeParse(rawEvent);
if (parsed.success) {
return parsed.data;
try {
const normalized = JSON.parse(JSON.stringify(rawEvent));
const parsed = providerRawEventSchema.safeParse(normalized);
if (parsed.success) {
return parsed.data;
}
} catch {
// Fall through to the diagnostic envelope below.
}

return {
Expand Down
7 changes: 7 additions & 0 deletions packages/thread-view/src/build-event-projection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ import {
parseRejectedUsersFromClientRequest,
parseUsersFromClientRequest,
parseLegacyUserMessage,
parseProviderUserMessage,
} from "./user-message-parsing.js";
import { isTerminalBufferedTextFlushEvent } from "./assistant-buffering.js";
import {
Expand Down Expand Up @@ -824,6 +825,12 @@ function buildFlatProjectionData(
continue;
}

const providerUserMessage = parseProviderUserMessage(decoded, meta);
if (providerUserMessage) {
appendProjectedUserMessage(state, providerUserMessage);
continue;
}

const legacyUserMessage = parseLegacyUserMessage(decoded, meta);
if (legacyUserMessage) {
flushToolActivityBeforeNonToolMessage(state);
Expand Down
18 changes: 17 additions & 1 deletion packages/thread-view/src/format-timeline-text.ts
Original file line number Diff line number Diff line change
Expand Up @@ -443,14 +443,30 @@ function formatConversationRequestLabel(
return "steer";
}

function conversationRoleLabel(row: TimelineConversationViewRow): string {
if (row.role === "assistant") {
return "Assistant";
}
switch (row.initiator) {
case "agent":
return "Agent";
case "system":
return "System";
case "user":
return "User";
default:
return assertNever(row.initiator);
}
}

function formatRow(
row: ThreadTimelineViewRow,
context: TimelineTextFormatContext,
): string {
switch (row.kind) {
case "conversation":
return [
rowHeader(row.role === "user" ? "User" : "Assistant", context),
rowHeader(conversationRoleLabel(row), context),
maybeTruncateBodyLinesForAudit(row.text.split("\n"), context).join(
"\n",
),
Expand Down
Loading
Loading