Files
pi_harness/packages/agent/learning/07-DATA-FLOW-STATE.md
T
2026-07-29 10:59:18 +07:00

17 KiB

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

// 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:

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

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

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

// 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

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

// 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

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

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:

// 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:

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

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:

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

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

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

// 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

// 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"