# Data Flow and State Management ## Overview Understanding how data flows through the agent system is crucial for debugging and extending functionality. --- ## Message Flow ### 1. Input Messages ```typescript // User input await harness.prompt("Build a web app"); // Internal messages await harness.steer("Wait, use React"); await harness.followUp("Now add tests"); await harness.nextTurn("Also deploy to production"); ``` **Normalization**: ```typescript function normalizePromptInput(input: string | AgentMessage | AgentMessage[]): AgentMessage[] { if (Array.isArray(input)) return input; if (typeof input !== "string") { return [input]; // Already a message } // String → user message return [{ role: "user", content: [{ type: "text", text: input }], timestamp: Date.now() }]; } ``` ### 2. AgentMessage Types ```typescript type AgentMessage = Message | CustomAgentMessages[keyof CustomAgentMessages] interface Message { role: "user" | "assistant" | "toolResult"; content: (TextContent | ImageContent)[]; api?: string; provider?: string; model?: string; usage?: Usage; stopReason?: StopReason; errorMessage?: string; timestamp: number; } interface TextContent { type: "text"; text: string; } interface ImageContent { type: "image"; mediaType: string; data: string; // Base64 } ``` ### 3. Message Lifecycle ``` User Input │ ▼ normalizePromptInput() → AgentMessage[] │ ▼ runPromptMessages() → runWithLifecycle() │ ├─► Set isStreaming=true ├─► Create abort controller └─► runAgentLoop() │ ▼ runLoop() │ ├─► message_start (user prompt) ├─► message_end ├─► streamAssistantResponse() │ ├─► message_start (assistant) │ ├─► message_update (chunks) │ └─► message_end ├─► executeToolCalls() │ └─► message_start/end (toolResults) └─► turn_end │ ▼ handleAgentEvent() (harness) │ ├─► session.appendMessage() │ └─► Storage: write entry └─► Emit: message_end (forwarded) ``` --- ## State Management ### Agent State ```typescript interface AgentState { systemPrompt: string; model: Model; thinkingLevel: ThinkingLevel; tools: AgentTool[]; messages: AgentMessage[]; isStreaming: boolean; streamingMessage?: AgentMessage; pendingToolCalls: Set; errorMessage?: string; } ``` **State changes**: | Event | State Changed | |-------|--------------| | `message_start` | `streamingMessage` = message | | `message_update` | `streamingMessage` = message | | `message_end` | `messages.push(message)`, `streamingMessage` = undefined | | `tool_execution_start` | `pendingToolCalls.add(toolCallId)` | | `tool_execution_end` | `pendingToolCalls.delete(toolCallId)` | | `turn_end` | `errorMessage` (if error) | | `agent_end` | `streamingMessage` = undefined | ### State Mutation Example ```typescript // In Agent.processEvents() private async processEvents(event: AgentEvent): Promise { switch (event.type) { case "message_start": this._state.streamingMessage = event.message; break; case "message_end": this._state.streamingMessage = undefined; this._state.messages.push(event.message); break; case "tool_execution_start": { const pending = new Set(this._state.pendingToolCalls); pending.add(event.toolCallId); this._state.pendingToolCalls = pending; break; } case "tool_execution_end": { const pending = new Set(this._state.pendingToolCalls); pending.delete(event.toolCallId); this._state.pendingToolCalls = pending; break; } } // Emit to listeners for (const listener of this.listeners) { await listener(event, signal); } } ``` --- ## Context Flow ### Context Snapshot ```typescript interface AgentContext { systemPrompt: string; messages: AgentMessage[]; tools?: AgentTool[]; } ``` **When created**: 1. `Agent.createContextSnapshot()` - before each LLM call 2. `AgentHarness.createContext()` - in turn state ### Context Transformation ```typescript // 1. transformContext() hook (AgentMessage[]) let messages = context.messages; if (config.transformContext) { messages = await config.transformContext(messages, signal); } // 2. convertToLlm() hook (AgentMessage[] → Message[]) const llmMessages = await config.convertToLlm(messages); // 3. Build LLM context (Message[]) const llmContext: Context = { systemPrompt: context.systemPrompt, messages: llmMessages, tools: context.tools }; ``` ### Context Transformations **Example: Prune old messages** ```typescript transformContext: async (messages) => { if (estimateTokens(messages) > MAX_TOKENS) { // Find cut point (preserve recent turns) const cutIndex = findCutPoint(messages, MAX_TOKENS * 0.7); return messages.slice(cutIndex); } return messages; } ``` **Example: Inject external context** ```typescript transformContext: async (messages) => { const externalData = await fetchExternalData(); const contextMessage: AgentMessage = { role: "user", content: [{ type: "text", text: externalData }], timestamp: Date.now() }; return [contextMessage, ...messages]; } ``` --- ## Hook Context Flow ### Hook Parameter Flow ``` Agent.prompt() │ ├─► transformContext(messages) [AgentLoopConfig] │ └─► Messages before LLM call │ ├─► convertToLlm(messages) │ └─► Messages to send to LLM │ ├─► beforeToolCall(context) [AgentLoopConfig] │ ├─► assistantMessage │ ├─► toolCall │ ├─► args (validated) │ └─► context (AgentContext) │ ├─► afterToolCall(context) [AgentLoopConfig] │ ├─► assistantMessage │ ├─► toolCall │ ├─► args │ ├─► result (executed) │ ├─► isError │ └─► context (AgentContext) │ ├─► shouldStopAfterTurn(context) [AgentLoopConfig] │ ├─► message (assistant) │ ├─► toolResults │ ├─► context (AgentContext) │ └─► newMessages │ ├─► prepareNextTurn(context) [AgentLoopConfig] │ └─► Return: context/model/thinkingLevel │ ├─► getSteeringMessages() [AgentLoopConfig] │ └─► Messages to inject now │ └─► getFollowUpMessages() [AgentLoopConfig] └─► Messages for after agent stops ``` ### Hook Return Value Flow ``` beforeToolCall() │ ├─► { block: true, reason } → Error tool result └─► undefined → Allow execution │ ▼ tool.execute() │ ▼ afterToolCall() │ ├─► Override: content, details, isError, usage, terminate └─► undefined → Use executed result │ ▼ Emit: tool_execution_end │ ▼ Create: ToolResultMessage │ ▼ Emit: message_start/end (toolResult) ``` --- ## Queue Flow ### Steering Queue **Purpose**: Interrupt agent while working. **Flow**: ``` steer("New instruction") │ ▼ steeringQueue.enqueue(message) │ ▼ After turn ends: │ ├─► getSteeringMessages() called │ ├─► Drain queue (mode: "all" or "one-at-a-time") │ └─► Return messages │ ▼ Inject messages into context │ ▼ Next LLM call includes steering messages ``` **Example**: ```typescript // User types while agent is working agent.steer("Wait, check this file first"); // Agent finishes current work // → Steering messages injected // → LLM sees: [original, ..., new user message] ``` ### Follow-up Queue **Purpose**: Queue messages for after agent stops naturally. **Flow**: ``` followUp("Next task") │ ▼ followUpQueue.enqueue(message) │ ▼ Agent would stop (no more tool calls) │ ├─► getFollowUpMessages() called │ ├─► Drain queue │ └─► Return messages │ ▼ Set as pendingMessages │ ▼ Inner loop continues ``` **Example**: ```typescript agent.followUp("Now create a README"); // Agent finishes current task // → Follow-up messages injected // → Agent continues with new task ``` ### Queue Modes **"all" Mode**: ``` Queued: [msg1, msg2, msg3] │ ▼ Drain: [msg1, msg2, msg3] │ ▼ All injected together ``` **"one-at-a-time" Mode**: ``` Queued: [msg1, msg2, msg3] │ ▼ Drain: [msg1] │ ▼ msg1 injected, msg2, msg3 remain │ ▼ After next turn: │ ▼ Drain: [msg2] │ ▼ ... and so on ``` --- ## Session Flow ### Session Tree Structure ``` root (parentId: null) ├─► message [id: 1, parentId: null] │ └─► message [id: 2, parentId: 1] │ └─► tool_result [id: 3, parentId: 2] │ └─► message [id: 4, parentId: 3] │ └─► compaction [id: 5, parentId: 4] │ ├─► retained: [msg6, msg7] │ └─► message [id: 8, parentId: 5] │ └─► leaf [id: 9, parentId: 8] ``` ### Context Building ```typescript async function buildContext(session: Session): Promise { // 1. Get path from leaf to root const pathEntries = await session.getBranch(); // [root, msg1, msg2, toolResult, msg4, compaction, msg8, leaf] // 2. Apply default transform (compaction logic) const contextEntries = defaultContextEntryTransform(pathEntries); // [compaction, retainedTail..., msg8] // 3. Project entries to messages const messages = contextEntries.flatMap(sessionEntryToContextMessages); // [compactionSummary, retainedMsgs..., msg8] // 4. Derive state const state = deriveSessionContextState(pathEntries); // { model, thinkingLevel, activeToolNames } return { ...state, messages }; } ``` ### Session Entry Types | Type | Stored When | |------|-------------| | `message` | Every user/assistant/toolResult | | `model_change` | `setModel()` called | | `thinking_level_change` | `setThinkingLevel()` called | | `active_tools_change` | `setActiveTools()` called | | `compaction` | `compact()` called | | `branch_summary` | Branching with summary | | `custom` | `appendCustomEntry()` | | `custom_message` | `appendCustomMessageEntry()` | | `label` | `appendLabel()` | | `leaf` | `setLeafId()` | | `session_info` | `appendSessionName()` | ### Pending Writes During active turns, writes are buffered: ```typescript async function appendMessage(message: AgentMessage): Promise { if (phase === "idle") { // Direct write await session.appendMessage(message); } else { // Buffer for later pendingSessionWrites.push({ type: "message", message }); } } async function flushPendingSessionWrites(): Promise { while (pendingSessionWrites.length > 0) { const write = pendingSessionWrites.shift(); if (write.type === "message") { await session.appendMessage(write.message); } else if (write.type === "model_change") { await session.appendModelChange(...); } // ... other types } } ``` --- ## Tool Execution State Flow ### Tool Call State ```typescript interface BeforeToolCallContext { assistantMessage: AssistantMessage; toolCall: AgentToolCall; args: unknown; // Validated context: AgentContext; // Snapshot } interface AfterToolCallContext { assistantMessage: AssistantMessage; toolCall: AgentToolCall; args: unknown; result: AgentToolResult; // Executed isError: boolean; context: AgentContext; } ``` ### Tool Result State ```typescript interface AgentToolResult { content: (TextContent | ImageContent)[]; // To model details: T; // For logs/UI usage?: Usage; // Tool-specific addedToolNames?: string[]; // New tools terminate?: boolean; // Early stop hint } ``` ### State Transition ``` Tool Call from LLM │ ▼ prepareToolCall() ├─► Find tool ├─► Validate args └─► beforeToolCall() ├─► block: true → Error └─► block: undefined → Continue │ ▼ tool.execute() ├─► onUpdate(partialResult) └─► Return final result │ ▼ afterToolCall() ├─► Override result └─► Use executed result │ ▼ createToolResultMessage() │ ▼ Emit: tool_execution_end │ ▼ Emit: message_start/end (toolResult) │ ▼ Push to context.messages ``` --- ## Abort Flow ### Abort Signal Propagation ```typescript // 1. Create abort controller const abortController = new AbortController(); // 2. Pass to all async operations await runAgentLoop(..., abortController.signal, ...); // 3. Check signal in long operations execute: async (id, params, signal, onUpdate) => { for await (const item of longProcess()) { if (signal?.aborted) { throw new Error("Aborted"); } } } // 4. Abort abortController.abort(); ``` ### Abort in Hooks ```typescript // Check signal at start beforeToolCall: async ({ toolCall }, signal) => { if (signal?.aborted) { return { block: true, reason: "Operation aborted" }; } return undefined; } // Check signal in async operations transformContext: async (messages, signal) => { if (signal?.aborted) { return messages; // Return safe fallback } // Long operation const result = await expensiveTransform(messages, signal); return result; } ``` --- ## Event Flow Diagram ``` ┌─────────────────────────────────────────────────────────────────────┐ │ AGENT LIFECYCLE │ ├─────────────────────────────────────────────────────────────────────┤ │ │ │ Agent.prompt("Hello") │ │ │ │ │ ├─► agent_start (event) │ │ ├─► turn_start (event) │ │ ├─► message_start (user) (event) │ │ ├─► message_end (user) (event) │ │ │ │ │ ├─► streamAssistantResponse() │ │ │ ├─► message_start (assistant) (event) │ │ │ ├─► message_update (text chunk 1) (event) │ │ │ ├─► message_update (text chunk 2) (event) │ │ │ ├─► message_update (toolCall) (event) │ │ │ └─► message_end (assistant) (event) │ │ │ │ │ ├─► executeToolCalls() │ │ │ ├─► tool_execution_start (event) │ │ │ ├─► tool_execute() │ │ │ │ └─► onUpdate(partial) (event) │ │ │ ├─► tool_execution_end (event) │ │ │ └─► message_start/end (toolResult) (events) │ │ │ │ │ ├─► turn_end (event) │ │ │ ├─► Should stop? → agent_end │ │ │ └─► Drain queues → another turn │ │ │ │ │ └─► agent_end (event) │ │ │ └─────────────────────────────────────────────────────────────────────┘ ``` --- ## Summary **Key data flows**: 1. Input messages → Normalized → Agent messages 2. Agent messages → Context transform → LLM messages 3. LLM response → Streamed → Agent messages 4. Tool calls → Executed → Tool results → Agent messages 5. All messages → Session storage → Tree structure **State management**: - Agent: In-memory state with mutation on events - Session: Persistent tree with entries - Hooks: Transform data at key points **Queue system**: - Steering: Interrupt current work - Follow-up: Queue for after agent stops - Modes: "all" or "one-at-a-time"