diff --git a/sdk/src/__test-utils__/pi-fixtures.ts b/sdk/src/__test-utils__/pi-fixtures.ts new file mode 100644 index 0000000..17c909b --- /dev/null +++ b/sdk/src/__test-utils__/pi-fixtures.ts @@ -0,0 +1,63 @@ +/** + * Wire-correct Pi frame builders, transcribed from + * `java/rust/broker/clients/pi-codes/tests/protocol_tests.rs`. + */ +import type { + AssistantContent, + AssistantMessage, + AssistantMessageEvent, + MessageUpdateEvent, + StopReason, + ToolResultMessage, +} from "../pi/protocol/index.js"; + +export function assistantMessage( + content: AssistantContent[], + stopReason: StopReason, +): AssistantMessage { + return { + role: "assistant", + content, + api: "anthropic-messages", + provider: "anthropic", + model: "claude-sonnet", + usage: { + input: 10, + output: 4, + cacheRead: 2, + cacheWrite: 1, + totalTokens: 17, + cost: { input: 0.1, output: 0.2, cacheRead: 0.01, cacheWrite: 0.02, total: 0.33 }, + }, + stopReason, + timestamp: 1_700_000_000_000, + }; +} + +/** + * An assistant message for a request that failed before any token accounting + * was recorded: `stopReason: "error"`, an `errorMessage`, and no `usage`. + */ +export function assistantErrorMessageWithoutUsage(): AssistantMessage { + const { usage: _usage, ...message } = assistantMessage([], "error"); + return { ...message, errorMessage: "overloaded" }; +} + +export function toolResultMessage(): ToolResultMessage { + return { + role: "toolResult", + toolCallId: "call-1", + toolName: "bash", + content: [{ type: "text", text: "/work" }], + isError: false, + timestamp: 1_700_000_000_001, + }; +} + +export function messageUpdate(delta: AssistantMessageEvent): MessageUpdateEvent { + return { + type: "message_update", + message: assistantMessage([{ type: "text", text: "Hi" }], "stop"), + assistantMessageEvent: delta, + }; +} diff --git a/sdk/src/pi/classify-pi-axon-event.test.ts b/sdk/src/pi/classify-pi-axon-event.test.ts new file mode 100644 index 0000000..b512f7b --- /dev/null +++ b/sdk/src/pi/classify-pi-axon-event.test.ts @@ -0,0 +1,207 @@ +import type { AxonEventView } from "@runloop/api-client/resources/axons"; +import { describe, expect, it } from "vitest"; +import { assistantMessage, toolResultMessage } from "../__test-utils__/pi-fixtures.js"; +import { classifyPiAxonEvent, isPiProtocolEventType } from "./classify-pi-axon-event.js"; +import { PI_EVENT_TYPES } from "./protocol/index.js"; +import { + isPiAgentEndEvent, + isPiAgentSettledEvent, + isPiMessageUpdateEvent, + isPiResponseEvent, + isPiToolExecutionEndEvent, + isTurnCompletedEvent, + isUnknownTimelineEvent, +} from "./timeline-event-guards.js"; + +function frame( + eventType: string, + payload: unknown, + origin: AxonEventView["origin"] = "AGENT_EVENT", +): AxonEventView { + return { + axon_id: "axn_pi", + event_type: eventType, + origin, + payload: JSON.stringify(payload), + sequence: 7, + source: "pi", + timestamp_ms: 1_752_192_000_000, + }; +} + +/** One wire-correct payload per modelled event type, keyed by `event_type`. */ +const EVENT_PAYLOADS: Record = { + agent_start: { type: "agent_start" }, + message_start: { + type: "message_start", + message: { role: "user", content: "hello", timestamp: 1_700_000_000_000 }, + }, + message_update: { + type: "message_update", + message: assistantMessage([{ type: "text", text: "Hi" }], "stop"), + assistantMessageEvent: { + type: "text_delta", + contentIndex: 0, + delta: "Hi", + partial: assistantMessage([{ type: "text", text: "Hi" }], "stop"), + }, + }, + message_end: { + type: "message_end", + message: assistantMessage( + [ + { type: "thinking", thinking: "checking" }, + { type: "toolCall", id: "call-1", name: "bash", arguments: { command: "pwd" } }, + ], + "toolUse", + ), + }, + tool_execution_start: { + type: "tool_execution_start", + toolCallId: "call-1", + toolName: "bash", + args: { command: "pwd" }, + }, + tool_execution_update: { + type: "tool_execution_update", + toolCallId: "call-1", + toolName: "bash", + args: { command: "pwd" }, + partialResult: { content: [{ type: "text", text: "/work" }] }, + }, + tool_execution_end: { + type: "tool_execution_end", + toolCallId: "call-1", + toolName: "bash", + result: { content: [{ type: "text", text: "/work" }] }, + isError: false, + }, + turn_start: { type: "turn_start" }, + turn_end: { + type: "turn_end", + message: assistantMessage([{ type: "text", text: "done" }], "stop"), + toolResults: [toolResultMessage()], + }, + agent_end: { + type: "agent_end", + messages: [assistantMessage([{ type: "text", text: "retrying" }], "error")], + willRetry: true, + }, + agent_settled: { type: "agent_settled" }, +}; + +describe("classifyPiAxonEvent", () => { + it.each(PI_EVENT_TYPES)("classifies the %s event as pi_protocol", (eventType) => { + const event = classifyPiAxonEvent(frame(eventType, EVENT_PAYLOADS[eventType])); + expect(event.kind).toBe("pi_protocol"); + if (event.kind === "pi_protocol") expect(event.eventType).toBe(eventType); + }); + + it("classifies a message_update carrying its streaming delta", () => { + const event = classifyPiAxonEvent(frame("message_update", EVENT_PAYLOADS.message_update)); + expect(isPiMessageUpdateEvent(event)).toBe(true); + if (isPiMessageUpdateEvent(event)) { + const delta = event.data.assistantMessageEvent; + expect(delta.type).toBe("text_delta"); + if (delta.type === "text_delta") expect(delta.delta).toBe("Hi"); + } + }); + + it("classifies a tool_execution_end result and its error flag", () => { + const event = classifyPiAxonEvent( + frame("tool_execution_end", EVENT_PAYLOADS.tool_execution_end), + ); + expect(isPiToolExecutionEndEvent(event)).toBe(true); + if (isPiToolExecutionEndEvent(event)) { + expect(event.data.toolCallId).toBe("call-1"); + expect(event.data.isError).toBe(false); + } + }); + + it("classifies agent_end with its retry flag separately from agent_settled", () => { + const end = classifyPiAxonEvent(frame("agent_end", EVENT_PAYLOADS.agent_end)); + expect(isPiAgentEndEvent(end)).toBe(true); + if (isPiAgentEndEvent(end)) expect(end.data.willRetry).toBe(true); + expect(isPiAgentSettledEvent(end)).toBe(false); + expect( + isPiAgentSettledEvent(classifyPiAxonEvent(frame("agent_settled", { type: "agent_settled" }))), + ).toBe(true); + }); + + it("classifies a get_state acknowledgement including its session state", () => { + const event = classifyPiAxonEvent( + frame("response", { + type: "response", + id: "req-1", + command: "get_state", + success: true, + data: { + model: null, + thinkingLevel: "medium", + isStreaming: false, + isCompacting: false, + steeringMode: "all", + followUpMode: "one-at-a-time", + sessionFile: "/sessions/current.jsonl", + sessionId: "session-1", + autoCompactionEnabled: true, + messageCount: 5, + pendingMessageCount: 0, + }, + }), + ); + expect(isPiResponseEvent(event)).toBe(true); + if (isPiResponseEvent(event)) { + expect(event.data.command).toBe("get_state"); + expect(event.data.success).toBe(true); + } + }); + + it("classifies a rejected prompt acknowledgement with its error", () => { + const event = classifyPiAxonEvent( + frame("response", { + type: "response", + command: "prompt", + success: false, + error: "agent is streaming", + }), + ); + expect(isPiResponseEvent(event)).toBe(true); + if (isPiResponseEvent(event)) expect(event.data.error).toBe("agent is streaming"); + }); + + it("prefers shared SYSTEM_EVENT classification", () => { + const event = classifyPiAxonEvent( + frame("turn.completed", { turn_id: "turn-1", stop_reason: "end_turn" }, "SYSTEM_EVENT"), + ); + expect(isTurnCompletedEvent(event)).toBe(true); + }); + + it("leaves unmodelled Pi frames unknown with the payload reachable", () => { + const event = classifyPiAxonEvent(frame("queue_update", { type: "queue_update", queued: 2 })); + expect(isUnknownTimelineEvent(event)).toBe(true); + expect(event.data).toEqual({ type: "queue_update", queued: 2 }); + expect(event.axonEvent.payload).toBe(JSON.stringify({ type: "queue_update", queued: 2 })); + }); + + it("falls through to unknown when a protocol payload is not an object", () => { + expect(isUnknownTimelineEvent(classifyPiAxonEvent(frame("agent_settled", "not-a-frame")))).toBe( + true, + ); + }); +}); + +describe("isPiProtocolEventType", () => { + it.each(PI_EVENT_TYPES)("recognizes %s", (eventType) => { + expect(isPiProtocolEventType(eventType)).toBe(true); + }); + + it("recognizes response and rejects unmodelled and outbound event types", () => { + expect(isPiProtocolEventType("response")).toBe(true); + expect(isPiProtocolEventType("queue_update")).toBe(false); + expect(isPiProtocolEventType("auto_retry_start")).toBe(false); + // Outbound broker control names never arrive as agent events. + expect(isPiProtocolEventType("turn/start")).toBe(false); + expect(isPiProtocolEventType("cancel")).toBe(false); + }); +}); diff --git a/sdk/src/pi/classify-pi-axon-event.ts b/sdk/src/pi/classify-pi-axon-event.ts new file mode 100644 index 0000000..10a4961 --- /dev/null +++ b/sdk/src/pi/classify-pi-axon-event.ts @@ -0,0 +1,23 @@ +import { createClassifier } from "../shared/timeline.js"; +import { PI_EVENT_TYPE_SET, PI_RESPONSE_EVENT_TYPE } from "./protocol/index.js"; +import type { PiProtocolTimelineEvent } from "./types.js"; + +/** Returns whether an Axon event type is one of the classified Pi wire frames. */ +export function isPiProtocolEventType(eventType: string): boolean { + return PI_EVENT_TYPE_SET.has(eventType) || eventType === PI_RESPONSE_EVENT_TYPE; +} + +/** Classifies a raw Axon event as Pi protocol, shared system, or unknown. */ +export const classifyPiAxonEvent = createClassifier({ + label: "classifyPiAxonEvent", + isProtocolEventType: isPiProtocolEventType, + toProtocolEvent: (data, ev) => { + if (typeof data !== "object" || data === null) return null; + return { + kind: "pi_protocol", + eventType: ev.event_type, + data, + axonEvent: ev, + } as PiProtocolTimelineEvent; + }, +}); diff --git a/sdk/src/pi/protocol/index.ts b/sdk/src/pi/protocol/index.ts new file mode 100644 index 0000000..64576b3 --- /dev/null +++ b/sdk/src/pi/protocol/index.ts @@ -0,0 +1,460 @@ +/** + * Pi JSONL RPC wire types and broker event-type constants. + * + * Hand-written from Pi `0.82.1` + * (`@earendil-works/pi-coding-agent/dist/modes/rpc/rpc-types.d.ts`) and kept + * field-for-field aligned with the broker's Rust bindings in + * `java/rust/broker/clients/pi-codes/src/{commands,events,response,types}.rs`, + * so the two can be diffed directly. + * + * Pi exchanges one JSON object per line — commands on stdin, responses and + * events on stdout — correlated by an optional `id`. There is no JSON-RPC + * envelope. Fields Pi itself declares `any` are `unknown` here, matching the + * crate's `serde_json::Value`. + */ + +// --------------------------------------------------------------------------- +// Scalars +// --------------------------------------------------------------------------- + +export type ThinkingLevel = "off" | "minimal" | "low" | "medium" | "high" | "xhigh" | "max"; + +/** How queued prompts are released: all at once, or one per turn. */ +export type QueueMode = "all" | "one-at-a-time"; + +/** + * Why an assistant message stopped. `toolUse` continues the turn; the broker + * maps `stop` to `EndTurn`, `length` to `MaxTokens`, `aborted` to `Cancelled` + * and `error` to `Error`. + */ +export type StopReason = "stop" | "length" | "toolUse" | "error" | "aborted"; + +/** Input modality a model accepts. */ +export type InputKind = "text" | "image"; + +/** How a prompt submitted while the agent is streaming is queued. */ +export type StreamingBehavior = "steer" | "followUp"; + +// --------------------------------------------------------------------------- +// Session state +// --------------------------------------------------------------------------- + +export interface ModelCost { + input: number; + output: number; + cacheRead: number; + cacheWrite: number; + tiers?: unknown[]; +} + +export interface PiModel { + id: string; + name: string; + api: string; + provider: string; + baseUrl: string; + reasoning: boolean; + input: InputKind[]; + cost: ModelCost; + contextWindow: number; + maxTokens: number; + thinkingLevelMap?: unknown; + headers?: unknown; + compat?: unknown; +} + +/** + * Snapshot returned by the `get_state` command. `sessionFile` is the path the + * broker persists as its resume state and replays with `switch_session`. + */ +export interface PiSessionState { + model: PiModel | null; + thinkingLevel: ThinkingLevel; + isStreaming: boolean; + isCompacting: boolean; + steeringMode: QueueMode; + followUpMode: QueueMode; + sessionFile?: string; + sessionId: string; + sessionName?: string; + autoCompactionEnabled: boolean; + messageCount: number; + pendingMessageCount: number; +} + +/** Result of `new_session` and `switch_session`. */ +export interface SessionChange { + cancelled: boolean; +} + +// --------------------------------------------------------------------------- +// Content blocks +// --------------------------------------------------------------------------- + +export interface TextContent { + type: "text"; + text: string; + textSignature?: string; +} + +export interface ImageContent { + type: "image"; + data: string; + mimeType: string; +} + +export interface ThinkingContent { + type: "thinking"; + thinking: string; + thinkingSignature?: string; + redacted?: boolean; +} + +export interface ToolCall { + type: "toolCall"; + id: string; + name: string; + arguments: unknown; + thoughtSignature?: string; +} + +export type UserContentBlock = + | { type: "text"; text: string; textSignature?: string } + | { type: "image"; data: string; mimeType: string }; + +/** A user message body: plain text, or a list of content blocks. */ +export type UserContent = string | UserContentBlock[]; + +export type AssistantContent = + | { type: "text"; text: string; textSignature?: string } + | { type: "thinking"; thinking: string; thinkingSignature?: string; redacted?: boolean } + | { type: "toolCall"; id: string; name: string; arguments: unknown; thoughtSignature?: string }; + +export type ToolResultContent = + | { type: "text"; text: string; textSignature?: string } + | { type: "image"; data: string; mimeType: string }; + +// --------------------------------------------------------------------------- +// Usage +// --------------------------------------------------------------------------- + +export interface Cost { + input: number; + output: number; + cacheRead: number; + cacheWrite: number; + total: number; +} + +export interface Usage { + input: number; + output: number; + cacheRead: number; + cacheWrite: number; + cacheWrite1h?: number; + reasoning?: number; + totalTokens: number; + cost: Cost; +} + +// --------------------------------------------------------------------------- +// Messages +// --------------------------------------------------------------------------- + +export interface UserMessage { + role: "user"; + content: UserContent; + timestamp: number; +} + +export interface AssistantMessage { + role: "assistant"; + content: AssistantContent[]; + api: string; + provider: string; + model: string; + responseModel?: string; + responseId?: string; + diagnostics?: unknown[]; + /** + * Token accounting, absent when the request failed before any was recorded. + * The Rust crate materialises a zeroed `Usage` in that case. + */ + usage?: Usage; + stopReason: StopReason; + errorMessage?: string; + timestamp: number; +} + +export interface ToolResultMessage { + role: "toolResult"; + toolCallId: string; + toolName: string; + content: ToolResultContent[]; + details?: unknown; + addedToolNames?: string[]; + isError: boolean; + timestamp: number; +} + +export interface BashExecutionMessage { + role: "bashExecution"; + command: string; + output: string; + exitCode: number | null; + cancelled: boolean; + truncated: boolean; + fullOutputPath?: string; + timestamp: number; + excludeFromContext?: boolean; +} + +export interface CustomMessage { + role: "custom"; + customType: string; + content: UserContent; + display: boolean; + details?: unknown; + timestamp: number; +} + +export interface BranchSummaryMessage { + role: "branchSummary"; + summary: string; + fromId: string; + timestamp: number; +} + +export interface CompactionSummaryMessage { + role: "compactionSummary"; + summary: string; + tokensBefore: number; + timestamp: number; +} + +/** Any message in a Pi session transcript. Discriminate on `role`. */ +export type AgentMessage = + | UserMessage + | AssistantMessage + | ToolResultMessage + | BashExecutionMessage + | CustomMessage + | BranchSummaryMessage + | CompactionSummaryMessage; + +// --------------------------------------------------------------------------- +// Assistant streaming deltas +// --------------------------------------------------------------------------- + +export type AssistantMessageEventSuccessReason = "stop" | "length" | "toolUse"; +export type AssistantMessageEventErrorReason = "aborted" | "error"; + +/** + * The streaming delta carried by a `message_update` event. Every variant + * except the two terminal ones carries `partial`, the assistant message + * accumulated so far. + */ +export type AssistantMessageEvent = + | { type: "start"; partial: AssistantMessage } + | { type: "text_start"; contentIndex: number; partial: AssistantMessage } + | { type: "text_delta"; contentIndex: number; delta: string; partial: AssistantMessage } + | { type: "text_end"; contentIndex: number; content: string; partial: AssistantMessage } + | { type: "thinking_start"; contentIndex: number; partial: AssistantMessage } + | { type: "thinking_delta"; contentIndex: number; delta: string; partial: AssistantMessage } + | { type: "thinking_end"; contentIndex: number; content: string; partial: AssistantMessage } + | { type: "toolcall_start"; contentIndex: number; partial: AssistantMessage } + | { type: "toolcall_delta"; contentIndex: number; delta: string; partial: AssistantMessage } + | { + type: "toolcall_end"; + contentIndex: number; + toolCall: ToolCall; + partial: AssistantMessage; + } + | { type: "done"; reason: AssistantMessageEventSuccessReason; message: AssistantMessage } + | { type: "error"; reason: AssistantMessageEventErrorReason; error: AssistantMessage }; + +// --------------------------------------------------------------------------- +// Events +// --------------------------------------------------------------------------- + +export interface AgentStartEvent { + type: "agent_start"; +} + +export interface MessageStartEvent { + type: "message_start"; + message: AgentMessage; +} + +export interface MessageUpdateEvent { + type: "message_update"; + message: AgentMessage; + assistantMessageEvent: AssistantMessageEvent; +} + +export interface MessageEndEvent { + type: "message_end"; + message: AgentMessage; +} + +export interface ToolExecutionStartEvent { + type: "tool_execution_start"; + toolCallId: string; + toolName: string; + args: unknown; +} + +export interface ToolExecutionUpdateEvent { + type: "tool_execution_update"; + toolCallId: string; + toolName: string; + args: unknown; + partialResult: unknown; +} + +export interface ToolExecutionEndEvent { + type: "tool_execution_end"; + toolCallId: string; + toolName: string; + result: unknown; + isError: boolean; +} + +export interface TurnStartEvent { + type: "turn_start"; +} + +export interface TurnEndEvent { + type: "turn_end"; + message: AgentMessage; + toolResults: ToolResultMessage[]; +} + +/** + * The agent stopped producing output. Not the end of the turn: `willRetry` + * marks an auto-retry, after which more events follow. Only `agent_settled` + * ends a turn. + */ +export interface AgentEndEvent { + type: "agent_end"; + messages: AgentMessage[]; + willRetry: boolean; +} + +/** The agent is idle with nothing queued. Terminates a turn. */ +export interface AgentSettledEvent { + type: "agent_settled"; +} + +/** Any event Pi emits while the agent runs. Discriminate on `type`. */ +export type PiEvent = + | AgentStartEvent + | MessageStartEvent + | MessageUpdateEvent + | MessageEndEvent + | ToolExecutionStartEvent + | ToolExecutionUpdateEvent + | ToolExecutionEndEvent + | TurnStartEvent + | TurnEndEvent + | AgentEndEvent + | AgentSettledEvent; + +// --------------------------------------------------------------------------- +// Responses +// --------------------------------------------------------------------------- + +/** + * Pi's acknowledgement of a command. `success: false` carries `error`; a + * `prompt` acknowledgement means the prompt was accepted, not that the turn + * finished. `data` is the command's payload — {@link PiSessionState} for + * `get_state`, {@link SessionChange} for `new_session` / `switch_session`. + */ +export interface PiResponse { + type: "response"; + id?: string; + command: string; + success: boolean; + error?: string; + data?: unknown; +} + +/** Any frame Pi writes to stdout. */ +export type PiOutput = PiResponse | PiEvent; + +// --------------------------------------------------------------------------- +// Commands +// --------------------------------------------------------------------------- + +/** + * The commands this SDK wraps. `steer` and `follow_up` reach Pi verbatim as + * broker Control frames and so have no counterpart in `pi-codes`, which types + * only the commands the broker itself constructs. + */ +export type PiCommand = + | { + type: "prompt"; + id?: string; + message: string; + images?: ImageContent[]; + streamingBehavior?: StreamingBehavior; + } + | { type: "steer"; id?: string; message: string; images?: ImageContent[] } + | { type: "follow_up"; id?: string; message: string; images?: ImageContent[] } + | { type: "abort"; id?: string } + | { type: "get_state"; id?: string } + | { type: "new_session"; id?: string; parentSession?: string } + | { type: "switch_session"; id?: string; sessionPath: string }; + +/** + * Any Pi command frame, including the ~25 commands this SDK does not wrap + * (`set_model`, `compact`, `bash`, `get_messages`, `export_html`, …). + */ +export interface PiCommandFrame { + type: string; + id?: string; + [key: string]: unknown; +} + +// --------------------------------------------------------------------------- +// Broker event-type constants +// --------------------------------------------------------------------------- + +/** + * Outbound broker event type that starts a turn. Slash-style, unlike the + * inbound `turn_start` Pi event — the broker's `classify_input` matches this + * exact string and silently ignores anything else it does not recognize. + */ +export const PI_TURN_START_EVENT_TYPE = "turn/start"; + +/** Outbound broker event type that aborts the in-flight turn. */ +export const PI_CANCEL_EVENT_TYPE = "cancel"; + +/** Event type the broker assigns to every Pi acknowledgement frame. */ +export const PI_RESPONSE_EVENT_TYPE = "response"; + +/** Event types the broker assigns to Pi's modelled events, in `pi-codes` order. */ +export const PI_EVENT_TYPES = [ + "agent_start", + "message_start", + "message_update", + "message_end", + "tool_execution_start", + "tool_execution_update", + "tool_execution_end", + "turn_start", + "turn_end", + "agent_end", + "agent_settled", +] as const; + +export type PiEventType = (typeof PI_EVENT_TYPES)[number]; + +export const PI_EVENT_TYPE_SET: ReadonlySet = new Set(PI_EVENT_TYPES); + +/** + * Id prefix the Pi adapter uses for the commands it issues itself. Client + * frames must not use it, or their acknowledgements would collide with the + * broker's. + */ +export const RESERVED_REQUEST_ID_PREFIX = "broker-"; diff --git a/sdk/src/pi/timeline-event-guards.test.ts b/sdk/src/pi/timeline-event-guards.test.ts new file mode 100644 index 0000000..5e9d9c0 --- /dev/null +++ b/sdk/src/pi/timeline-event-guards.test.ts @@ -0,0 +1,169 @@ +import type { AxonEventView } from "@runloop/api-client/resources/axons"; +import { describe, expect, it } from "vitest"; +import { + assistantErrorMessageWithoutUsage, + assistantMessage, + messageUpdate, +} from "../__test-utils__/pi-fixtures.js"; +import type { PiEvent, PiResponse } from "./protocol/index.js"; +import { + isPiAgentEndEvent, + isPiAgentSettledEvent, + isPiAgentStartEvent, + isPiAssistantTextDeltaEvent, + isPiAssistantThinkingDeltaEvent, + isPiMessageEndEvent, + isPiMessageStartEvent, + isPiMessageUpdateEvent, + isPiProtocolEvent, + isPiResponseEvent, + isPiToolExecutionEndEvent, + isPiToolExecutionStartEvent, + isPiToolExecutionUpdateEvent, + isPiTurnEndEvent, + isPiTurnStartEvent, + isSystemTimelineEvent, + isTurnStartedEvent, + isUnknownTimelineEvent, +} from "./timeline-event-guards.js"; +import type { PiProtocolTimelineEvent, PiTimelineEvent } from "./types.js"; + +function axonEvent(eventType: string, payload: unknown): AxonEventView { + return { + axon_id: "axn_pi", + event_type: eventType, + origin: "AGENT_EVENT", + payload: JSON.stringify(payload), + sequence: 3, + source: "pi", + timestamp_ms: 1_752_192_000_000, + }; +} + +function protocolEvent(frame: PiEvent | PiResponse): PiTimelineEvent { + return { + kind: "pi_protocol", + eventType: frame.type, + data: frame, + axonEvent: axonEvent(frame.type, frame), + } as PiProtocolTimelineEvent; +} + +const systemEvent: PiTimelineEvent = { + kind: "system", + data: { type: "turn.started", turnId: "turn-1" }, + axonEvent: axonEvent("turn.started", { turn_id: "turn-1" }), +}; + +const unknownEvent: PiTimelineEvent = { + kind: "unknown", + data: { type: "queue_update" }, + axonEvent: axonEvent("queue_update", { type: "queue_update" }), +}; + +describe("Pi protocol guards", () => { + it.each([ + [isPiAgentStartEvent, { type: "agent_start" } as PiEvent], + [ + isPiMessageStartEvent, + { type: "message_start", message: assistantMessage([], "stop") } as PiEvent, + ], + [ + isPiMessageUpdateEvent, + messageUpdate({ type: "start", partial: assistantMessage([], "stop") }) as PiEvent, + ], + [ + isPiMessageEndEvent, + { type: "message_end", message: assistantMessage([], "stop") } as PiEvent, + ], + [ + isPiToolExecutionStartEvent, + { type: "tool_execution_start", toolCallId: "c", toolName: "bash", args: {} } as PiEvent, + ], + [ + isPiToolExecutionUpdateEvent, + { + type: "tool_execution_update", + toolCallId: "c", + toolName: "bash", + args: {}, + partialResult: {}, + } as PiEvent, + ], + [ + isPiToolExecutionEndEvent, + { + type: "tool_execution_end", + toolCallId: "c", + toolName: "bash", + result: {}, + isError: false, + } as PiEvent, + ], + [isPiTurnStartEvent, { type: "turn_start" } as PiEvent], + [ + isPiTurnEndEvent, + { type: "turn_end", message: assistantMessage([], "stop"), toolResults: [] } as PiEvent, + ], + [isPiAgentEndEvent, { type: "agent_end", messages: [], willRetry: false } as PiEvent], + [isPiAgentSettledEvent, { type: "agent_settled" } as PiEvent], + ])("matches only its own event type", (guard, frame) => { + const event = protocolEvent(frame); + expect(guard(event)).toBe(true); + expect(isPiProtocolEvent(event)).toBe(true); + expect(guard(systemEvent)).toBe(false); + expect(guard(unknownEvent)).toBe(false); + // Every guard is exclusive: no other modelled frame satisfies it. + expect(guard(protocolEvent({ type: "response", command: "prompt", success: true }))).toBe( + false, + ); + }); + + it("classifies a message_end whose assistant message carries no usage", () => { + const message = assistantErrorMessageWithoutUsage(); + const event = protocolEvent({ type: "message_end", message }); + expect(isPiMessageEndEvent(event)).toBe(true); + if (isPiMessageEndEvent(event) && event.data.message.role === "assistant") { + expect(event.data.message.usage).toBeUndefined(); + expect(event.data.message.stopReason).toBe("error"); + expect(event.data.message.errorMessage).toBe("overloaded"); + } + }); + + it("matches acknowledgement frames", () => { + const event = protocolEvent({ type: "response", command: "prompt", success: true }); + expect(isPiResponseEvent(event)).toBe(true); + if (isPiResponseEvent(event)) expect(event.data.command).toBe("prompt"); + expect(isPiResponseEvent(protocolEvent({ type: "agent_settled" }))).toBe(false); + }); + + it("narrows message_update to text and thinking deltas", () => { + const partial = assistantMessage([{ type: "text", text: "Hi" }], "stop"); + const text = protocolEvent( + messageUpdate({ type: "text_delta", contentIndex: 0, delta: "Hi", partial }), + ); + const thinking = protocolEvent( + messageUpdate({ type: "thinking_delta", contentIndex: 0, delta: "hmm", partial }), + ); + const start = protocolEvent(messageUpdate({ type: "start", partial })); + + expect(isPiAssistantTextDeltaEvent(text)).toBe(true); + if (isPiAssistantTextDeltaEvent(text)) expect(text.data.assistantMessageEvent.delta).toBe("Hi"); + expect(isPiAssistantThinkingDeltaEvent(thinking)).toBe(true); + if (isPiAssistantThinkingDeltaEvent(thinking)) + expect(thinking.data.assistantMessageEvent.delta).toBe("hmm"); + + expect(isPiAssistantTextDeltaEvent(thinking)).toBe(false); + expect(isPiAssistantThinkingDeltaEvent(text)).toBe(false); + expect(isPiAssistantTextDeltaEvent(start)).toBe(false); + expect(isPiAssistantThinkingDeltaEvent(start)).toBe(false); + expect(isPiAssistantTextDeltaEvent(unknownEvent)).toBe(false); + }); + + it("re-exports the shared system and unknown guards", () => { + expect(isSystemTimelineEvent(systemEvent)).toBe(true); + expect(isTurnStartedEvent(systemEvent)).toBe(true); + expect(isUnknownTimelineEvent(unknownEvent)).toBe(true); + expect(isPiProtocolEvent(systemEvent)).toBe(false); + }); +}); diff --git a/sdk/src/pi/timeline-event-guards.ts b/sdk/src/pi/timeline-event-guards.ts new file mode 100644 index 0000000..1586d53 --- /dev/null +++ b/sdk/src/pi/timeline-event-guards.ts @@ -0,0 +1,99 @@ +/** Type guards for Pi timeline events. */ +import type { + PiAgentEndTimelineEvent, + PiAgentSettledTimelineEvent, + PiAgentStartTimelineEvent, + PiAssistantTextDeltaTimelineEvent, + PiAssistantThinkingDeltaTimelineEvent, + PiMessageEndTimelineEvent, + PiMessageStartTimelineEvent, + PiMessageUpdateTimelineEvent, + PiProtocolTimelineEvent, + PiResponseTimelineEvent, + PiTimelineEvent, + PiToolExecutionEndTimelineEvent, + PiToolExecutionStartTimelineEvent, + PiToolExecutionUpdateTimelineEvent, + PiTurnEndTimelineEvent, + PiTurnStartTimelineEvent, +} from "./types.js"; + +export type { + AgentErrorTimelineEvent, + AgentLogTimelineEvent, + BrokerErrorTimelineEvent, + DevboxLifecycleTimelineEvent, + TurnCompletedTimelineEvent, + TurnFailedTimelineEvent, + TurnStartedTimelineEvent, +} from "../shared/timeline-event-guards.js"; +export { + createCustomEventGuard, + isAgentErrorEvent, + isAgentLogEvent, + isBrokerErrorEvent, + isDevboxLifecycleEvent, + isSystemTimelineEvent, + isTurnCompletedEvent, + isTurnFailedEvent, + isTurnStartedEvent, + isUnknownTimelineEvent, +} from "../shared/timeline-event-guards.js"; + +export function isPiProtocolEvent(event: PiTimelineEvent): event is PiProtocolTimelineEvent { + return event.kind === "pi_protocol"; +} + +function hasEventType( + event: PiTimelineEvent, + eventType: M, +): event is Extract { + return event.kind === "pi_protocol" && event.eventType === eventType; +} + +export const isPiAgentStartEvent = (event: PiTimelineEvent): event is PiAgentStartTimelineEvent => + hasEventType(event, "agent_start"); +export const isPiMessageStartEvent = ( + event: PiTimelineEvent, +): event is PiMessageStartTimelineEvent => hasEventType(event, "message_start"); +export const isPiMessageUpdateEvent = ( + event: PiTimelineEvent, +): event is PiMessageUpdateTimelineEvent => hasEventType(event, "message_update"); +export const isPiMessageEndEvent = (event: PiTimelineEvent): event is PiMessageEndTimelineEvent => + hasEventType(event, "message_end"); +export const isPiToolExecutionStartEvent = ( + event: PiTimelineEvent, +): event is PiToolExecutionStartTimelineEvent => hasEventType(event, "tool_execution_start"); +export const isPiToolExecutionUpdateEvent = ( + event: PiTimelineEvent, +): event is PiToolExecutionUpdateTimelineEvent => hasEventType(event, "tool_execution_update"); +export const isPiToolExecutionEndEvent = ( + event: PiTimelineEvent, +): event is PiToolExecutionEndTimelineEvent => hasEventType(event, "tool_execution_end"); +export const isPiTurnStartEvent = (event: PiTimelineEvent): event is PiTurnStartTimelineEvent => + hasEventType(event, "turn_start"); +export const isPiTurnEndEvent = (event: PiTimelineEvent): event is PiTurnEndTimelineEvent => + hasEventType(event, "turn_end"); +export const isPiAgentEndEvent = (event: PiTimelineEvent): event is PiAgentEndTimelineEvent => + hasEventType(event, "agent_end"); +export const isPiAgentSettledEvent = ( + event: PiTimelineEvent, +): event is PiAgentSettledTimelineEvent => hasEventType(event, "agent_settled"); +export const isPiResponseEvent = (event: PiTimelineEvent): event is PiResponseTimelineEvent => + hasEventType(event, "response"); + +/** Narrows a `message_update` to an assistant text delta. */ +export function isPiAssistantTextDeltaEvent( + event: PiTimelineEvent, +): event is PiAssistantTextDeltaTimelineEvent { + return isPiMessageUpdateEvent(event) && event.data.assistantMessageEvent?.type === "text_delta"; +} + +/** Narrows a `message_update` to an assistant thinking delta. */ +export function isPiAssistantThinkingDeltaEvent( + event: PiTimelineEvent, +): event is PiAssistantThinkingDeltaTimelineEvent { + return ( + isPiMessageUpdateEvent(event) && event.data.assistantMessageEvent?.type === "thinking_delta" + ); +} diff --git a/sdk/src/pi/types.ts b/sdk/src/pi/types.ts new file mode 100644 index 0000000..3e35c61 --- /dev/null +++ b/sdk/src/pi/types.ts @@ -0,0 +1,87 @@ +/** Types for Pi timeline classification. */ +import type { + BaseTimelineEvent, + SystemTimelineEvent, + UnknownTimelineEvent, +} from "../shared/types.js"; +import type { + AssistantMessageEvent, + MessageUpdateEvent, + PiEvent, + PiResponse, +} from "./protocol/index.js"; + +type EventFrame = Extract; + +type ProtocolTimelineEvent = BaseTimelineEvent & { + kind: "pi_protocol"; + eventType: M; + data: D; +}; + +export type PiAgentStartTimelineEvent = ProtocolTimelineEvent< + "agent_start", + EventFrame<"agent_start"> +>; +export type PiMessageStartTimelineEvent = ProtocolTimelineEvent< + "message_start", + EventFrame<"message_start"> +>; +export type PiMessageUpdateTimelineEvent = ProtocolTimelineEvent< + "message_update", + EventFrame<"message_update"> +>; +export type PiMessageEndTimelineEvent = ProtocolTimelineEvent< + "message_end", + EventFrame<"message_end"> +>; +export type PiToolExecutionStartTimelineEvent = ProtocolTimelineEvent< + "tool_execution_start", + EventFrame<"tool_execution_start"> +>; +export type PiToolExecutionUpdateTimelineEvent = ProtocolTimelineEvent< + "tool_execution_update", + EventFrame<"tool_execution_update"> +>; +export type PiToolExecutionEndTimelineEvent = ProtocolTimelineEvent< + "tool_execution_end", + EventFrame<"tool_execution_end"> +>; +export type PiTurnStartTimelineEvent = ProtocolTimelineEvent< + "turn_start", + EventFrame<"turn_start"> +>; +export type PiTurnEndTimelineEvent = ProtocolTimelineEvent<"turn_end", EventFrame<"turn_end">>; +export type PiAgentEndTimelineEvent = ProtocolTimelineEvent<"agent_end", EventFrame<"agent_end">>; +export type PiAgentSettledTimelineEvent = ProtocolTimelineEvent< + "agent_settled", + EventFrame<"agent_settled"> +>; +export type PiResponseTimelineEvent = ProtocolTimelineEvent<"response", PiResponse>; + +/** A `message_update` narrowed to one variant of its nested streaming delta. */ +type AssistantDeltaTimelineEvent = + PiMessageUpdateTimelineEvent & { + data: MessageUpdateEvent & { + assistantMessageEvent: Extract; + }; + }; + +export type PiAssistantTextDeltaTimelineEvent = AssistantDeltaTimelineEvent<"text_delta">; +export type PiAssistantThinkingDeltaTimelineEvent = AssistantDeltaTimelineEvent<"thinking_delta">; + +export type PiProtocolTimelineEvent = + | PiAgentStartTimelineEvent + | PiMessageStartTimelineEvent + | PiMessageUpdateTimelineEvent + | PiMessageEndTimelineEvent + | PiToolExecutionStartTimelineEvent + | PiToolExecutionUpdateTimelineEvent + | PiToolExecutionEndTimelineEvent + | PiTurnStartTimelineEvent + | PiTurnEndTimelineEvent + | PiAgentEndTimelineEvent + | PiAgentSettledTimelineEvent + | PiResponseTimelineEvent; + +export type PiTimelineEvent = PiProtocolTimelineEvent | SystemTimelineEvent | UnknownTimelineEvent;