/** * @license * Copyright 2026 Vybestack LLC * SPDX-License-Identifier: Apache-2.0 */ import * as acp from '@agentclientprotocol/sdk'; import { type Agent, type AgentEvent } from '@vybestack/llxprt-code-agents'; import { EmojiFilter, type ContentBlock, type DebugLogger, type FilterConfiguration, getErrorStatus, } from '@vybestack/llxprt-code-core'; import { handleZedAgentEvent } from './zed-agent-event-handler.js'; import { StreamBatcher } from './zed-stream-batcher.js'; import type { ZedPathResolver } from './zed-path-resolver.js'; import { TerminalManager } from './zed-terminal-manager.js'; type SendUpdateFn = (update: acp.SessionUpdate) => Promise; type SendUsageFn = ( usage: Extract['usage'], ) => Promise; type HandleConfirmationFn = ( event: Extract, ) => Promise; export interface SessionStreamDeps { readonly agent: Agent; readonly terminals: TerminalManager | null; readonly sendUpdate: SendUpdateFn; readonly sendUsage: SendUsageFn; readonly handleConfirmation: HandleConfirmationFn; readonly isPromptStale: ( promptGeneration: number, pendingSend: AbortController, ) => boolean; readonly maxTurns: number; readonly logger: Pick; } /** * Consumes the agent event stream, forwarding each event to the Zed handler * and managing terminal lifecycle for shell tool calls. */ export async function consumeAgentStream( deps: SessionStreamDeps, parts: readonly ContentBlock[], pendingSend: AbortController, promptId: string, promptGeneration: number, batcher: StreamBatcher, ): Promise { const eventStream = deps.agent.stream(parts, { signal: pendingSend.signal, promptId, maxTurns: deps.maxTurns, }); let terminalStopReason: acp.StopReason | null = null; try { for await (const event of eventStream) { const result = await processStreamEvent( event, deps, batcher, promptGeneration, pendingSend, ); if (result === 'cancelled') return 'cancelled'; if (result !== null) terminalStopReason = result; } return terminalStopReason; } finally { if (deps.terminals !== null) { try { await deps.terminals.settleAll(); } catch (error) { deps.logger.debug(() => `Terminal cleanup failed: ${String(error)}`); } } } } async function processStreamEvent( event: AgentEvent, deps: SessionStreamDeps, batcher: StreamBatcher, promptGeneration: number, pendingSend: AbortController, ): Promise { if (deps.isPromptStale(promptGeneration, pendingSend)) { return 'cancelled'; } const stopReason = await handleZedAgentEvent(event, batcher, { sendUpdate: deps.sendUpdate, sendUsage: deps.sendUsage, handleConfirmation: deps.handleConfirmation, resolveToolKind: (toolName) => deps.agent.tools.get(toolName)?.kind, }); await trackTerminalEvent(event, deps); return stopReason; } /** * Best-effort terminal bookkeeping after each event. Extracted so terminal * tracking errors never mask the stop reason from `handleZedAgentEvent`. * The single `deps.terminals !== null` check addresses the repeated null guard. */ async function trackTerminalEvent( event: AgentEvent, deps: SessionStreamDeps, ): Promise { if (deps.terminals === null) return; try { if ( event.type === 'tool-call' && TerminalManager.isShellToolCall( event.call, deps.agent.tools.get(event.call.name)?.kind, ) ) { await deps.terminals.observeToolCall(event.call); } else if (event.type === 'tool-result') { deps.terminals.completeToolCall(event.result.id); } } catch (error) { deps.logger.debug(() => `Terminal tracking failed: ${String(error)}`); } } export interface PromptTurnDeps { readonly pathResolver: ZedPathResolver; readonly emojiFilterMode: FilterConfiguration['mode']; readonly streamDeps: SessionStreamDeps; } export async function runPromptTurn( deps: PromptTurnDeps, params: acp.PromptRequest, pendingSend: AbortController, promptId: string, promptGeneration: number, ): Promise { let parts: ContentBlock[]; try { parts = await deps.pathResolver.resolvePrompt( params.prompt, pendingSend.signal, ); } catch (error) { if ( pendingSend.signal.aborted || (error instanceof Error && error.name === 'AbortError') ) { return { stopReason: 'cancelled' }; } throw error; } const batcher = new StreamBatcher( new EmojiFilter({ mode: deps.emojiFilterMode }), (u) => deps.streamDeps.sendUpdate(u), ); let terminalStopReason: acp.StopReason | null = null; let thrownError: unknown = null; try { terminalStopReason = await consumeAgentStream( deps.streamDeps, parts, pendingSend, promptId, promptGeneration, batcher, ); } catch (error) { thrownError = error; } await safeFlush(batcher, deps.streamDeps.logger); batcher.dispose(); if (thrownError !== null) { if ( pendingSend.signal.aborted || (thrownError instanceof Error && thrownError.name === 'AbortError') ) { return { stopReason: 'cancelled' }; } if (getErrorStatus(thrownError) === 429) { const reason = thrownError instanceof Error ? thrownError.message : String(thrownError); throw new acp.RequestError(429, 'Rate limit exceeded. Try again later.', { reason, }); } throw thrownError; } if (terminalStopReason !== null) { return { stopReason: terminalStopReason }; } if (pendingSend.signal.aborted) { return { stopReason: 'cancelled' }; } return { stopReason: 'end_turn' }; } async function safeFlush( batcher: StreamBatcher, logger: Pick, ): Promise { try { await batcher.flush(); } catch (error) { logger.debug(() => `Stream flush failed: ${String(error)}`); } }