How PrimeAgent Handles Asynchronous Operations: A Deep Dive into Its Event-Stream Architecture

PrimeAgent handles asynchronous operations through a unified event-stream abstraction centered on AssistantMessageEventStream, which exposes cancellable, retry-aware streams from every LLM provider.

PrimeAgent's core AI package in packages/ai implements a sophisticated asynchronous pipeline. This design lets developers consume LLM responses as incremental events or await a final message—without managing provider-specific streaming quirks. Here's how PrimeAgent handles asynchronous operations from request to completion.

The Central Stream Interface

PrimeAgent's entry points for asynchronous LLM calls live in [packages/ai/src/stream.ts](https://github.com/PrimeIntellect-ai/prime-agent/blob/main/packages/ai/src/stream.ts). The two primary functions are:

  • stream() — Returns an AssistantMessageEventStream for incremental consumption
  • streamSimple() — A lighter variant with streamlined options

Both resolve the appropriate provider via resolveApiProvider() and delegate to provider-specific implementations. The returned AssistantMessageEventStream acts as a universal handle regardless of which LLM backend serves the request.

For callers who only need the final result, complete() and completeSimple() wrap stream() and await result(), returning a single AssistantMessage promise.

The Event-Stream Engine: AssistantMessageEventStream

The concrete implementation resides in [packages/ai/src/utils/event-stream.ts](https://github.com/PrimeIntellect-ai/prime-agent/blob/main/packages/ai/src/utils/event-stream.ts). This class provides:

  • Typed event emission: start, text_delta, thinking_delta, toolcall_start, toolcall_delta, toolcall_end, done, error
  • Internal buffering: Partial AssistantMessage objects accumulate as chunks arrive
  • Final resolution: The result() method returns a promise that resolves with the complete message
import { stream } from "packages/ai/src/stream.js";

const s = stream(model, context, options);

// Option 1: Iterate over events
for await (const ev of s) {
  if (ev.type === "text_delta") {
    process.stdout.write(ev.delta);
  }
}

// Option 2: Await final result directly
const finalMessage = await s.result();

StreamOptions: Fine-Grained Async Control

PrimeAgent handles asynchronous operation configuration through StreamOptions defined in [packages/ai/src/types.ts](https://github.com/PrimeIntellect-ai/prime-agent/blob/main/packages/ai/src/types.ts#L78-L125). Key parameters include:

Parameter Purpose
signal AbortSignal for request cancellation
timeoutMs Client-side timeout before abort
maxRetries Number of retry attempts on transient failures
maxRetryDelayMs Cap on exponential backoff delay
sessionId Enables prompt caching across requests
onPayload / onResponse Callbacks for debugging and header inspection
const stream = stream(model, ctx, {
  signal: AbortSignal.timeout(30_000),  // Auto-cancel after 30s
  maxRetries: 3,
  maxRetryDelayMs: 5_000,
  sessionId: "user-123-session",        // Cache reuse identifier
});

Provider Implementations and Async I/O Isolation

Each LLM provider in packages/ai/src/providers/ (e.g., [openai-completions.ts](https://github.com/PrimeIntellect-ai/prime-agent/blob/main/packages/ai/src/providers/openai-completions.ts), [anthropic.ts](https://github.com/PrimeIntellect-ai/prime-agent/blob/main/packages/ai/src/providers/anthropic.ts), [amazon-bedrock.ts](https://github.com/PrimeIntellect-ai/prime-agent/blob/main/packages/ai/src/providers/amazon-bedrock.ts)) follows the same pattern:

  1. Create HTTP request or SDK call with StreamOptions
  2. Instantiate AssistantMessageEventStream
  3. Parse raw response chunks and push typed events into the stream
  4. Handle signal.aborted checks before each chunk
  5. Implement retry logic respecting maxRetries and maxRetryDelayMs

This isolation ensures PrimeAgent handles asynchronous operations consistently whether streaming from OpenAI, Anthropic, or Amazon Bedrock.

Cancellation and Error Handling

PrimeAgent's async pipeline respects cancellation at multiple layers:

  • Pre-flight: signal.aborted checked before initiating request
  • Per-chunk: Providers check abort status before yielding each chunk
  • Retry exhaustion: When maxRetries is exceeded, stream emits error event

When aborted, the stream emits:

{
  type: "error",
  stopReason: "aborted",  // or "error" for failures
  errorMessage: string
}

Complete Usage Examples

Streaming with Live UI Updates

import { stream } from "packages/ai/src/stream.js";
import { getModel } from "packages/ai/src/models.js";

async function handleStreamingAsync() {
  const model = await getModel("gpt-4o-mini");
  const ctx = {
    messages: [{
      role: "user",
      content: "Write a haiku about async programming.",
      timestamp: Date.now()
    }]
  };

  const s = stream(model, ctx, {
    temperature: 0.7,
    signal: AbortSignal.timeout(30_000),
  });

  for await (const ev of s) {
    switch (ev.type) {
      case "text_delta":
        process.stdout.write(ev.delta);
        break;
      case "thinking_delta":
        // Reasoning model intermediate steps
        break;
      case "toolcall_start":
        console.log(`\n[Tool: ${ev.toolName}]`);
        break;
    }
  }

  const final = await s.result();
  console.log("\n\nComplete message:", final);
}

One-Shot Completion

import { complete } from "packages/ai/src/stream.js";

async function handleSimpleAsync() {
  const model = await getModel("claude-3-5-sonnet");
  const ctx = {
    messages: [{
      role: "user",
      content: "Explain async/await in two sentences.",
      timestamp: Date.now()
    }]
  };

  const assistantMsg = await complete(model, ctx, {
    maxTokens: 200,
    onResponse: (resp) => {
      console.log("Rate limit remaining:", resp.headers.get("x-ratelimit-remaining"));
    },
  });

  return assistantMsg.content;
}

Key Architectural Files

File Role in Async Handling
packages/ai/src/stream.ts Central façade; resolves providers, returns AssistantMessageEventStream
packages/ai/src/utils/event-stream.ts Core async iterator implementation; event buffering and result() resolution
packages/ai/src/types.ts StreamOptions, AssistantMessageEvent type definitions
packages/ai/src/providers/*.ts Provider-specific HTTP/SDK streaming implementations
packages/ai/src/models.ts Model definitions with per-model async capabilities

Summary

  • Unified abstraction: AssistantMessageEventStream normalizes async behavior across all LLM providers
  • Flexible consumption: Iterate events for live UI, or await result() for one-shot usage
  • Robust control: StreamOptions provides cancellation, timeouts, retries, and session caching
  • Clean isolation: Provider implementations handle raw I/O without leaking complexity upstream

Frequently Asked Questions

How does PrimeAgent cancel an in-flight LLM request?

PrimeAgent checks options.signal?.aborted before initiating requests and between chunks during streaming. Pass an AbortSignal via StreamOptions.signal—including AbortSignal.timeout() for automatic time limits. When aborted, the stream emits an error event with stopReason: "aborted".

What's the difference between stream() and complete()?

stream() returns an AssistantMessageEventStream that yields incremental events (text_delta, toolcall_*, etc.) as the LLM generates tokens. complete() calls stream() internally but immediately awaits result(), returning a promise that resolves with the final AssistantMessage—hiding the event-stream details when you only need the completed response.

How does PrimeAgent handle retry logic for failed requests?

The StreamOptions interface includes maxRetries and maxRetryDelayMs parameters. Provider implementations use these to retry transient failures with exponential backoff, capped at maxRetryDelayMs. After exhausting retries, the stream emits an error event rather than throwing, allowing graceful degradation.

Can PrimeAgent reuse cached prompts across async requests?

Yes—set StreamOptions.sessionId to enable prompt caching with providers that support it (OpenAI, Anthropic). The session identifier attaches to requests, letting the remote service cache system prompts and reduce latency for subsequent calls with the same sessionId.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →