import type { AgentSessionEvent } from "@mariozechner/pi-coding-agent"; import type { PiMessage, SessionBackendEvent } from "./pi-events.js"; import { extractAssistantText, extractToolFullOutputPath, translatePiEvent, } from "./session-protocol.js"; import type { EventProcessorSessionState, ExtensionUIRequest, SessionEventProcessor, } from "./session-events.js"; import type { SessionStopCoordinator, StopSessionState } from "./session-stop.js"; import type { SessionTurnCoordinator, TurnSessionState } from "./session-turns.js"; import type { ServerMessage } from "./types.js"; export interface SessionAgentEventState extends EventProcessorSessionState, TurnSessionState, StopSessionState { subscribers: Set<(msg: ServerMessage) => void>; toolFullOutputPaths: Map; } export interface SessionAgentEventCoordinatorDeps { getActiveSession: (key: string) => SessionAgentEventState | undefined; eventProcessor: SessionEventProcessor; stopCoordinator: SessionStopCoordinator; turnCoordinator: SessionTurnCoordinator; broadcast: (key: string, message: ServerMessage) => void; resetIdleTimer: (key: string) => void; markQueuedMessageStarted?: (key: string, message: PiMessage) => void; } export class SessionAgentEventCoordinator { private static readonly LOGGED_EVENT_TYPES = new Set([ "agent_start", "agent_end", "turn_start", "turn_end", "message_end", "tool_execution_start", "tool_execution_end", "compaction_start", "compaction_end", "auto_retry_start", "auto_retry_end", ]); private static readonly STATUS_BROADCAST_TYPES = new Set([ "agent_start", "agent_end", "message_end", "tool_execution_start", ]); constructor(private readonly deps: SessionAgentEventCoordinatorDeps) {} handlePiEvent(key: string, data: SessionBackendEvent): void { const active = this.deps.getActiveSession(key); if (!active) { return; } if (data.type === "extension_ui_request") { this.handleExtensionUIRequest(key, data); return; } if (data.type === "extension_error") { console.error("[pi] extension error", { sessionId: active.session.id, extensionPath: data.extensionPath, error: data.error, }); this.deps.resetIdleTimer(key); return; } if (data.type === "prompt_error") { console.error("[pi] prompt error", { sessionId: active.session.id, error: data.error, }); this.deps.broadcast(key, { type: "error", error: data.error }); this.deps.resetIdleTimer(key); return; } const event = data; if (SessionAgentEventCoordinator.LOGGED_EVENT_TYPES.has(event.type)) { const toolName = "toolName" in event && typeof event.toolName === "string" ? event.toolName : undefined; console.log("[pi] event", { sessionId: active.session.id, eventType: event.type, toolName, subscriberCount: active.subscribers.size, }); } if (event.type === "message_start" && event.message.role === "user") { this.deps.markQueuedMessageStarted?.(key, event.message); } // pi may not always emit user message_start for queued steer/follow-up. // Reconcile from SDK queue on each turn_start so queue chips cannot linger // when dequeues happen without a matching message_start payload. if (event.type === "turn_start") { this.deps.markQueuedMessageStarted?.(key, {}); } const ctx = this.deps.eventProcessor.translationContext(active); const messages = translatePiEvent(event, ctx); active.streamedAssistantText = ctx.streamedAssistantText; active.hasStreamedThinking = ctx.hasStreamedThinking; if (event.type === "tool_execution_end") { const toolCallId = typeof event.toolCallId === "string" && event.toolCallId.length > 0 ? event.toolCallId : null; if (toolCallId) { const fullOutputPath = extractToolFullOutputPath(event.result?.details); if (fullOutputPath) { active.toolFullOutputPaths.set(toolCallId, fullOutputPath); } } } for (const message of messages) { this.deps.broadcast(key, message); } if (event.type === "agent_start") { this.deps.turnCoordinator.markNextTurnStarted(key, active); } this.deps.eventProcessor.updateSessionFromEvent(key, active, event); if (event.type === "agent_end") { this.deps.stopCoordinator.finishPendingStopOnAgentEnd(key, active); } if (event.type === "message_end") { const role = event.message.role; if (role === "assistant" || role === "user") { this.deps.broadcast(key, { type: "message_end", role, content: extractAssistantText(event.message), }); } } if (SessionAgentEventCoordinator.STATUS_BROADCAST_TYPES.has(event.type)) { console.log("[pi] status update", { sessionId: active.session.id, status: active.session.status, }); this.deps.broadcast(key, { type: "state", session: active.session }); } this.deps.resetIdleTimer(key); } handleExtensionUIRequest(key: string, req: ExtensionUIRequest): void { const active = this.deps.getActiveSession(key); if (!active) { return; } this.deps.eventProcessor.handleExtensionUIRequest(key, active, req); } }