Real-time streaming from LLM providers with automatic reconnection and error handling.
The streaming module provides framework-agnostic tools for handling real-time AI responses:
- AIStream: Client-side streaming with SSE (Server-Sent Events)
- StreamingResponse: Server-side SSE stream formatting
- Provider Adapters: OpenAI and Anthropic streaming adapters
- Token Counting: Real-time token usage tracking
npm install @ainative/ai-kit-coreimport { AIStream, StreamingResponse } from '@ainative/ai-kit-core/streaming';Client-side AI streaming with automatic reconnection and error handling.
new AIStream(config: StreamConfig, options?: StreamOptions)Parameters:
-
config: StreamConfig- Stream configurationendpoint: string- API endpoint URLmodel: string- Model identifier (e.g., 'gpt-4', 'claude-3-opus')headers?: Record<string, string>- Custom headerssystemPrompt?: string- System prompt for the conversationonToken?: (token: string) => void- Callback for each tokenonCost?: (usage: Usage) => void- Callback for usage updatesonError?: (error: Error) => void- Error callbackretry?: RetryConfig- Retry configuration
-
options?: StreamOptions- Additional options (reserved for future use)
Example:
import { AIStream } from '@ainative/ai-kit-core/streaming';
const stream = new AIStream({
endpoint: '/api/chat',
model: 'gpt-4',
systemPrompt: 'You are a helpful assistant.',
onToken: (token) => console.log('Token:', token),
onCost: (usage) => console.log('Usage:', usage),
retry: {
maxRetries: 3,
backoff: 'exponential',
initialDelay: 1000,
maxDelay: 10000
}
});Send a message and start streaming the response.
await stream.send('Hello, how are you?');Parameters:
content: string- User message content
Returns: Promise<void>
Events Emitted:
message- New message added (user or assistant)streaming-start- Streaming startedtoken- New token receivedusage- Usage information updatedstreaming-end- Streaming completed
Reset the conversation, clearing all messages and usage statistics.
stream.reset();Events Emitted:
reset- Conversation reset
Retry the last message. Removes the last assistant response and re-streams.
await stream.retry();Returns: Promise<void>
Stop the current stream immediately.
stream.stop();Get all messages in the conversation.
const messages = stream.getMessages();
console.log(messages);
// [
// { id: '...', role: 'user', content: 'Hello', timestamp: 1234567890 },
// { id: '...', role: 'assistant', content: 'Hi there!', timestamp: 1234567891 }
// ]Returns: Message[] - Copy of messages array
Check if currently streaming a response.
if (stream.getIsStreaming()) {
console.log('Currently streaming...');
}Returns: boolean
Get current token usage statistics.
const usage = stream.getUsage();
console.log('Prompt tokens:', usage.promptTokens);
console.log('Completion tokens:', usage.completionTokens);
console.log('Total tokens:', usage.totalTokens);Returns: Usage - Copy of usage object
AIStream extends EventEmitter and emits the following events:
// Message added (user or assistant)
stream.on('message', (message: Message) => {
console.log('New message:', message);
});
// Streaming started
stream.on('streaming-start', () => {
console.log('Streaming started');
});
// New token received
stream.on('token', (token: string) => {
process.stdout.write(token);
});
// Usage information updated
stream.on('usage', (usage: Usage) => {
console.log('Usage:', usage);
});
// Streaming completed
stream.on('streaming-end', () => {
console.log('Streaming ended');
});
// Error occurred
stream.on('error', (error: Error) => {
console.error('Error:', error);
});
// Retry attempted
stream.on('retry', ({ attempt, delay }) => {
console.log(`Retry attempt ${attempt} in ${delay}ms`);
});
// Conversation reset
stream.on('reset', () => {
console.log('Conversation reset');
});import { AIStream } from '@ainative/ai-kit-core/streaming';
const stream = new AIStream({
endpoint: '/api/chat',
model: 'gpt-4',
systemPrompt: 'You are a helpful coding assistant.',
retry: {
maxRetries: 3,
backoff: 'exponential'
}
});
// Listen to events
stream.on('token', (token) => {
process.stdout.write(token);
});
stream.on('usage', (usage) => {
console.log(`\nTokens used: ${usage.totalTokens}`);
});
stream.on('error', (error) => {
console.error('Stream error:', error.message);
});
// Send messages
await stream.send('Write a function to calculate factorial');
// Get conversation history
const messages = stream.getMessages();
console.log('Conversation:', messages);
// Reset when done
stream.reset();Server-side SSE (Server-Sent Events) stream formatting for Node.js/Express.
new StreamingResponse(response: ResponseLike, options?: StreamingOptions)Parameters:
response: ResponseLike- Node.js ServerResponse or Express Response objectoptions?: StreamingOptions- Configuration optionsenableHeartbeat?: boolean- Enable heartbeat (default: false)heartbeatInterval?: number- Heartbeat interval in ms (default: 30000)compressionEnabled?: boolean- Enable compression (default: false)customHeaders?: Record<string, string>- Custom HTTP headers
Example:
import { StreamingResponse } from '@ainative/ai-kit-core/streaming';
import express from 'express';
const app = express();
app.post('/api/chat', async (req, res) => {
const stream = new StreamingResponse(res, {
enableHeartbeat: true,
heartbeatInterval: 30000
});
stream.start();
// ... send tokens ...
stream.end();
});Initialize the SSE stream with proper headers.
const stream = new StreamingResponse(res);
stream.start();Returns: this - For method chaining
Headers Set:
Content-Type: text/event-streamCache-Control: no-cache, no-transformConnection: keep-aliveX-Accel-Buffering: no(disables nginx buffering)
Send a text token to the client.
stream.sendToken('Hello');
stream.sendToken(' world!');Parameters:
token: string- Text token to sendindex?: number- Optional token index
Returns: this - For method chaining
Send token usage metadata to the client.
stream.sendUsage({
promptTokens: 100,
completionTokens: 50,
totalTokens: 150
});Parameters:
usage: UsageEvent- Token usage informationpromptTokens: numbercompletionTokens: numbertotalTokens: number
Returns: this - For method chaining
Send an error event to the client.
stream.sendError('Rate limit exceeded');
// or
stream.sendError({
error: 'Rate limit exceeded',
code: 'RATE_LIMIT'
});Parameters:
error: string | ErrorEvent- Error message or error object
Returns: this - For method chaining
Send custom metadata to the client.
stream.sendMetadata({
model: 'gpt-4',
temperature: 0.7,
customField: 'value'
});Parameters:
metadata: MetadataEvent- Metadata object (any key-value pairs)
Returns: this - For method chaining
End the SSE stream gracefully.
stream.end();Sends a DONE event and closes the connection.
Abort the stream with an error.
stream.abort('Internal server error');Parameters:
error: string | ErrorEvent- Error message or error object
Sends an error event and then ends the stream.
Check if the stream is currently active.
if (stream.isStreamActive()) {
console.log('Stream is active');
}Returns: boolean
Get the number of messages sent.
const count = stream.getMessageCount();
console.log(`Sent ${count} messages`);Returns: number
The StreamingResponse sends different event types:
START- Stream initializedTOKEN- Text token receivedUSAGE- Usage metadataERROR- Error occurredMETADATA- Custom metadataDONE- Stream completed
import { StreamingResponse } from '@ainative/ai-kit-core/streaming';
import express from 'express';
import OpenAI from 'openai';
const app = express();
const openai = new OpenAI({ apiKey: process.env.OPENAI_API_KEY });
app.post('/api/chat', async (req, res) => {
const { messages } = req.body;
const stream = new StreamingResponse(res, {
enableHeartbeat: true
});
try {
stream.start();
const completion = await openai.chat.completions.create({
model: 'gpt-4',
messages,
stream: true
});
for await (const chunk of completion) {
const content = chunk.choices[0]?.delta?.content;
if (content) {
stream.sendToken(content);
}
// Send usage on last chunk
if (chunk.usage) {
stream.sendUsage({
promptTokens: chunk.usage.prompt_tokens,
completionTokens: chunk.usage.completion_tokens,
totalTokens: chunk.usage.total_tokens
});
}
}
stream.end();
} catch (error) {
stream.abort(error.message);
}
});
app.listen(3000);Adapters for streaming from different LLM providers.
import { OpenAIAdapter } from '@ainative/ai-kit-core/streaming/adapters';
import OpenAI from 'openai';
const openai = new OpenAI({ apiKey: process.env.OPENAI_API_KEY });
const response = await openai.chat.completions.create({
model: 'gpt-4',
messages: [{ role: 'user', content: 'Hello!' }],
stream: true
});
// Convert to standard stream
const stream = OpenAIAdapter(response);
// Use in framework
return new Response(stream);import { AnthropicAdapter } from '@ainative/ai-kit-core/streaming/adapters';
import Anthropic from '@anthropic-ai/sdk';
const anthropic = new Anthropic({ apiKey: process.env.ANTHROPIC_API_KEY });
const response = await anthropic.messages.create({
model: 'claude-3-opus-20240229',
messages: [{ role: 'user', content: 'Hello!' }],
stream: true
});
// Convert to standard stream
const stream = AnthropicAdapter(response);
// Use in framework
return new Response(stream);interface StreamConfig {
endpoint: string;
model: string;
headers?: Record<string, string>;
systemPrompt?: string;
onToken?: (token: string) => void;
onCost?: (usage: Usage) => void;
onError?: (error: Error) => void;
retry?: RetryConfig;
}interface RetryConfig {
maxRetries?: number; // Default: 3
backoff?: 'exponential' | 'linear'; // Default: 'exponential'
initialDelay?: number; // Default: 1000ms
maxDelay?: number; // Default: 10000ms
}interface Message {
id: string;
role: 'user' | 'assistant' | 'system';
content: string;
timestamp: number;
}interface Usage {
promptTokens: number;
completionTokens: number;
totalTokens: number;
}interface StreamingOptions {
enableHeartbeat?: boolean;
heartbeatInterval?: number;
compressionEnabled?: boolean;
customHeaders?: Record<string, string>;
}enum SSEEventType {
START = 'start',
TOKEN = 'token',
USAGE = 'usage',
ERROR = 'error',
METADATA = 'metadata',
DONE = 'done'
}// app/api/chat/route.ts
import { StreamingResponse } from '@ainative/ai-kit-core/streaming';
import { OpenAI } from 'openai';
export async function POST(req: Request) {
const { messages } = await req.json();
const openai = new OpenAI();
const response = await openai.chat.completions.create({
model: 'gpt-4',
messages,
stream: true
});
const stream = new ReadableStream({
async start(controller) {
for await (const chunk of response) {
const content = chunk.choices[0]?.delta?.content;
if (content) {
controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify({ token: content })}\n\n`));
}
}
controller.close();
}
});
return new Response(stream, {
headers: {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache'
}
});
}import express from 'express';
import { StreamingResponse } from '@ainative/ai-kit-core/streaming';
import { OpenAI } from 'openai';
const app = express();
const openai = new OpenAI();
app.post('/api/chat', async (req, res) => {
const stream = new StreamingResponse(res);
stream.start();
try {
const completion = await openai.chat.completions.create({
model: 'gpt-4',
messages: req.body.messages,
stream: true
});
for await (const chunk of completion) {
const content = chunk.choices[0]?.delta?.content;
if (content) {
stream.sendToken(content);
}
}
stream.end();
} catch (error) {
stream.abort(error.message);
}
});stream.on('error', (error) => {
console.error('Stream error:', error);
// Implement fallback logic
});const stream = new AIStream({
endpoint: '/api/chat',
model: 'gpt-4',
retry: {
maxRetries: 3,
backoff: 'exponential',
initialDelay: 1000
}
});const stream = new StreamingResponse(res, {
enableHeartbeat: true,
heartbeatInterval: 30000 // 30 seconds
});// Stop streaming when component unmounts
useEffect(() => {
return () => {
stream.stop();
};
}, []);stream.on('usage', (usage) => {
if (usage.totalTokens > BUDGET_LIMIT) {
console.warn('Token budget exceeded');
stream.stop();
}
});const stream = new StreamingResponse(res);
stream.start();
res.on('close', () => {
console.log('Client disconnected');
// Clean up resources
});// Network error - will retry automatically
stream.on('error', (error) => {
if (error.message.includes('fetch failed')) {
console.log('Network error, retrying...');
}
});
// Rate limit error
stream.on('error', (error) => {
if (error.message.includes('429')) {
console.log('Rate limited');
}
});
// Authentication error
stream.on('error', (error) => {
if (error.message.includes('401')) {
console.log('Invalid API key');
}
});const stream = new AIStream({
endpoint: '/api/chat',
model: 'gpt-4',
retry: {
maxRetries: 5,
backoff: 'exponential',
initialDelay: 1000,
maxDelay: 30000
}
});
stream.on('retry', ({ attempt, delay }) => {
console.log(`Retry ${attempt} in ${delay}ms`);
});- Use Token Streaming: Stream tokens as they arrive for better UX
- Implement Caching: Cache responses for repeated queries
- Monitor Token Usage: Track tokens to optimize costs
- Use Compression: Enable compression for large responses
- Implement Timeouts: Set reasonable timeouts for requests
- Pool Connections: Reuse connections when possible