688 lines
17 KiB
Markdown
688 lines
17 KiB
Markdown
# 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<any>;
|
|
thinkingLevel: ThinkingLevel;
|
|
tools: AgentTool<any>[];
|
|
messages: AgentMessage[];
|
|
isStreaming: boolean;
|
|
streamingMessage?: AgentMessage;
|
|
pendingToolCalls: Set<string>;
|
|
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<void> {
|
|
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<any>[];
|
|
}
|
|
```
|
|
|
|
**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<SessionContext> {
|
|
// 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<void> {
|
|
if (phase === "idle") {
|
|
// Direct write
|
|
await session.appendMessage(message);
|
|
} else {
|
|
// Buffer for later
|
|
pendingSessionWrites.push({ type: "message", message });
|
|
}
|
|
}
|
|
|
|
async function flushPendingSessionWrites(): Promise<void> {
|
|
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<any>; // Executed
|
|
isError: boolean;
|
|
context: AgentContext;
|
|
}
|
|
```
|
|
|
|
### Tool Result State
|
|
|
|
```typescript
|
|
interface AgentToolResult<T> {
|
|
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"
|