@cognikit/core
v0.1.3
Published
Core runtime for CogniKit.
Maintainers
Readme
@cognikit/core
Framework-agnostic runtime for AI streaming. Types, streaming pipeline, state machines, and utilities - everything you need to talk to LLM APIs from any JavaScript environment.
npm install @cognikit/core
# or
pnpm add @cognikit/coreFeatures
- Typed streaming protocol - discriminated union (
StreamChunk) for text, reasoning, tool calls, results, errors, and metadata - Multi-provider auto-detection - OpenAI, Anthropic (including extended thinking), and Gemini wire formats decoded automatically, including parallel tool calls and multiple completion choices/candidates, plus a simpler CogniKit-native format for custom backends (see Connecting Your Own Backend)
- Multi-modal content -
Message.contentaccepts plain text or an array of text/image/file parts - Non-streaming mode - set
stream: falsefor a plain JSON response, decoded into the sameStreamChunkshapes as a streaming one - SSE parser - Incremental, handles partial chunks (including split multi-byte UTF-8 characters), CR/LF/CRLF, comments, and missing terminators
- Streaming pipeline - Full request lifecycle: fetch → parse → retry → middleware → dedup
- Automatic retry - Exponential backoff with full jitter,
Retry-Afterheader support, configurable predicates - Middleware system - Intercept and transform requests, chunks, and errors
- Request deduplication - In-flight caching prevents duplicate API calls
- State machines - XState v5 machines for chat and stream lifecycles
- Zero UI dependencies - Runs in Node.js, Deno, Cloudflare Workers, browsers, anywhere
Quick Start
import { runChat, authMiddleware } from '@cognikit/core';
runChat({
api: 'https://api.openai.com/v1/chat/completions',
model: 'gpt-4o',
messages: [{ id: '1', role: 'user', content: 'Tell me a joke' }],
middleware: [authMiddleware(`Bearer ${process.env.OPENAI_API_KEY}`)],
onChunk(chunk) {
if (chunk.type === 'text-delta') process.stdout.write(chunk.textDelta);
if (chunk.type === 'error') console.error(chunk.error);
},
onFinish(response) {
console.log(`\nDone. Tokens: ${response.usage?.totalTokens}`);
},
onError(error) {
console.error('Request failed:', error.message);
},
});Core Concepts
Messages
import type { Message } from '@cognikit/core';
const messages: Message[] = [
{ id: '1', role: 'system', content: 'You are a helpful assistant.' },
{ id: '2', role: 'user', content: 'What is the weather in London?' },
{ id: '3', role: 'assistant', content: '', toolCalls: [{ ... }] },
{ id: '4', role: 'tool', content: '', toolResults: [{ ... }] },
];content can also be an array of multi-modal parts instead of a plain string:
const message: Message = {
id: '5',
role: 'user',
content: [
{ type: 'text', text: 'What is in this image?' },
{ type: 'image', url: 'https://example.com/cat.png' },
],
};Use getTextContent(message) to extract just the text regardless of which form content is in.
Stream Chunks
Every chunk is a member of the StreamChunk discriminated union. Switch on type for full exhaustiveness:
import type { StreamChunk } from '@cognikit/core';
function handleChunk(chunk: StreamChunk) {
switch (chunk.type) {
case 'text-delta': return chunk.textDelta;
case 'tool-call-delta': return chunk.delta;
case 'tool-call': return chunk.toolCall;
case 'tool-result': return chunk.toolResult;
case 'error': return chunk.error;
case 'finish': return chunk.finishReason;
case 'metadata': return chunk.metadata;
}
}Streaming Pipeline
runChat is the primary API - execution starts immediately and callbacks handle the stream:
import { runChat, authMiddleware, debugMiddleware, createDedupCache } from '@cognikit/core';
const { abort, done } = runChat({
api: '/api/chat',
model: 'gpt-4o',
messages,
temperature: 0.7,
maxTokens: 1000,
headers: { 'X-Request-ID': crypto.randomUUID() },
middleware: [authMiddleware(token), debugMiddleware()],
retry: { maxRetries: 3, initialDelay: 1000 },
dedupCache: createDedupCache(),
onChunk(chunk) {
console.log('Chunk:', chunk.type);
},
onFinish(response) {
console.log('Complete:', response.message.content);
console.log('Tokens:', response.usage?.totalTokens);
},
onError(error) {
console.error('Failed:', error.message);
},
});
// Cancel mid-stream
abort();
// Wait for completion
await done;For low-level control without callbacks, use streamChat as a pull-based async iterable:
import { streamChat } from '@cognikit/core';
for await (const chunk of streamChat({ api: '/api/chat', messages })) {
if (chunk.type === 'text-delta') appendToUI(chunk.textDelta);
if (chunk.type === 'finish') console.log('Done:', chunk.finishReason);
}Middleware
Plug into the request/response lifecycle without modifying the pipeline:
import type { StreamMiddleware } from '@cognikit/core';
const loggingMiddleware: StreamMiddleware = {
name: 'logger',
transformRequest(ctx) {
console.log('Sending to:', ctx.url);
},
transformChunk(chunk) {
if (chunk.type === 'metadata') return null; // filter out metadata
},
onError(error) {
return { type: 'error', error: `Recovered: ${error.message}` };
},
};Built-in middleware: authMiddleware(token) and debugMiddleware(prefix).
Retry
Transient failures are retried automatically with exponential backoff + full jitter:
import { runChat, isRetryableError } from '@cognikit/core';
runChat({
api: '/api/chat',
messages,
retry: {
maxRetries: 5,
initialDelay: 500, // ms
maxDelay: 10_000, // cap at 10s
backoffFactor: 2,
retryOn(error) {
return isRetryableError(error) || error.status === 403; // custom logic
},
},
onChunk(chunk) { /* ... */ },
});Generation Parameters
runChat({
api: '/api/chat',
messages,
topP: 0.9,
stop: ['\n\n'],
presencePenalty: 0.5,
frequencyPenalty: 0.2,
seed: 42,
responseFormat: { type: 'json_object' }, // or { type: 'json_schema', schema }
toolChoice: 'required', // 'auto' | 'none' | 'required' | { type: 'function', name }
onChunk(chunk) { /* ... */ },
});Non-Streaming Mode
Set stream: false to get a plain JSON response instead of SSE - runChat/streamChat yield the exact same StreamChunk shapes either way, just delivered in one batch instead of incrementally:
const { done } = runChat({ api: '/api/chat', messages, stream: false });
const response = await done; // same ChatResponse shape as streamingRequest Deduplication
Prevent duplicate API calls when two parts of your app send identical requests simultaneously:
import { runChat, createDedupCache } from '@cognikit/core';
const dedup = createDedupCache();
// Both calls with the same key share a single network request
await Promise.all([
runChat({ api: '/chat', messages, dedupCache: dedup }).done,
runChat({ api: '/chat', messages, dedupCache: dedup }).done,
]);State Machines
XState v5 machines for predictable state management:
import { chatMachine, streamMachine } from '@cognikit/core';
import { createActor } from 'xstate';
const actor = createActor(chatMachine);
actor.start();
actor.send({ type: 'SEND', messages: [{ id: '1', role: 'user', content: 'Hi' }] });
// state: 'connecting'
actor.send({ type: 'CHUNK', chunk: { type: 'text-delta', textDelta: 'Hello' } });
// state: 'streaming', context.streamBuffer = 'Hello'
actor.send({ type: 'DONE' });
// state: 'done'Connecting Your Own Backend
@cognikit/core never calls a provider directly - runChat/streamChat always send one fixed request shape to whatever api URL you configure, and your backend decides what to do with it. This section documents that contract for anyone building a custom or proxy backend rather than pointing straight at a named provider.
The outgoing request
Every request is a POST with this JSON body (fields are only present when set):
{
"messages": [
{ "role": "user", "content": "Hello" },
{ "role": "assistant", "content": "", "tool_calls": [ /* ToolCall[] */ ] },
{
"role": "tool",
"content": "",
"tool_results": [
{ "tool_call_id": "call_1", "tool_name": "get_weather", "result": "72°F", "is_error": false }
]
}
],
"stream": true,
"model": "gpt-4o",
"system": "You are a helpful assistant.",
"max_tokens": 1000,
"temperature": 0.7,
"top_p": 0.9,
"stop": ["\n\n"],
"presence_penalty": 0.5,
"frequency_penalty": 0.2,
"seed": 42,
"response_format": { "type": "json_object" },
"tool_choice": "auto",
"tools": [{ "name": "get_weather", "description": "...", "parameters": { "type": "object" } }]
}This already matches OpenAI's Chat Completions request shape (that's why the Quick Start example above can point straight at api.openai.com) - but your backend is free to translate this body into any upstream API's shape, or none at all. @cognikit/core doesn't know or care what's behind api.
Each message can also carry metadata (forwarded verbatim) and reasoning (an assistant message's prior [{text, signature?}] extended-thinking blocks, forwarded verbatim so your backend can replay them to Anthropic exactly as required for multi-turn tool use with extended thinking).
The response: auto-detected, 4 formats
readStreamChunks() (streaming) and the internal non-streaming decoder try each registered adapter in order and use the first one that recognizes the payload: OpenAI → Anthropic → Gemini → CogniKit-native. Your backend has two options:
- Relay a real provider's response untouched (proxy pattern) - stream back exactly what OpenAI/Anthropic/Gemini sent, byte for byte; it's auto-detected with zero client-side configuration.
- Emit CogniKit's own native format - simpler than any provider's actual wire format, since it's just
StreamChunkserialized directly, with no translation layer to write or maintain server-side.
Native streaming format
Each SSE event's data: is one JSON-encoded StreamChunk - literally data: ${JSON.stringify(chunk)}\n\n:
data: {"type":"text-delta","textDelta":"Hel"}
data: {"type":"text-delta","textDelta":"lo!"}
data: {"type":"finish","finishReason":"stop","usage":{"totalTokens":12}}
Recognized type values are StreamChunk's own discriminants: text-delta, tool-call-delta, tool-call, tool-result, error, finish, metadata, reasoning-delta. An unrecognized type is silently ignored rather than treated as an error, so older client versions stay forward-compatible if your backend starts sending new chunk types. You can decode this format directly with the exported decodeCogniKitChunk(event) (alongside decodeOpenAIChunk, decodeAnthropicChunk, decodeGeminiChunk - see Chunk Decoder).
Native non-streaming format (stream: false)
A single JSON body:
{
"message": { "content": "Hello!", "toolCalls": [] },
"finishReason": "stop",
"usage": { "totalTokens": 12 }
}This response-side shape isn't exposed as a standalone exported function today - it's decoded internally by runChat/streamChat when stream: false (see Non-Streaming Mode).
Low-Level APIs
SSE Parser
Parse raw ReadableStream bytes into SSE events:
import { parseSSEStream } from '@cognikit/core';
const response = await fetch('https://api.example.com/stream');
for await (const event of parseSSEStream(response.body!)) {
console.log(event.event, event.data);
}Chunk Decoder
Map SSE events to typed StreamChunk objects:
import { readStreamChunks, collectStreamChunks } from '@cognikit/core';
for await (const chunk of readStreamChunks(response.body!)) {
// typed StreamChunk - switch on chunk.type
}
const allChunks = await collectStreamChunks(response.body!);API Reference
Types
| Export | Description |
|--------|-------------|
| Message | Conversation message (id, role, content, toolCalls, toolResults, reasoning) |
| ContentPart | A multi-modal content part: text, image, or file |
| getTextContent(message) | Extract just the text from content, whether it's a string or ContentPart[] |
| ToolCall | Model-requested function invocation |
| ToolResult | Result of a tool execution (including error, forwarded to the model on failure) |
| StreamChunk | Discriminated union of all stream event types, including reasoning-delta |
| ChatOptions | Configuration for chat requests (generation params, stream, etc.) |
| Provider | Protocol-agnostic API endpoint config |
| ChatResponse | Final result of a completed chat request |
Streaming
| Export | Description |
|--------|-------------|
| runChat(options) | Push-based streaming: onChunk/onFinish/onError callbacks, returns { abort, done } |
| streamChat(options) | Pull-based streaming: for await...of async iterable, no callbacks |
| parseSSEStream(stream) | Parse ReadableStream<Uint8Array> into SSEEvent objects |
| readStreamChunks(stream) | Parse + decode a stream into typed StreamChunk objects |
| collectStreamChunks(stream) | Collect entire stream into StreamChunk[] |
Middleware
| Export | Description |
|--------|-------------|
| StreamMiddleware | Interface for request/chunk/error interception |
| authMiddleware(token) | Adds Bearer authorization header |
| debugMiddleware(prefix?) | Logs every chunk to console |
Retry
| Export | Description |
|--------|-------------|
| RetryConfig | Configuration (maxRetries, delays, backoff factor, predicate) |
| isRetryableError(error) | Default predicate: network errors, 429, 5xx |
| calculateBackoff(attempt, config) | Exponential backoff with full jitter |
| DEFAULT_RETRY_CONFIG | Sensible defaults (3 retries, 1s initial, 30s cap) |
Deduplication
| Export | Description |
|--------|-------------|
| createDedupCache() | Creates a dedup cache for preventing duplicate requests |
| defaultDedupKey(options) | Generates a cache key from request options |
State Machines
| Export | Description |
|--------|-------------|
| chatMachine | XState v5 machine: idle → connecting → streaming → done/error/aborted |
| streamMachine | XState v5 machine for raw streaming lifecycle |
Known Limitations
- Retries restart the whole request, not just the missing tail. If a connection drops mid-stream, the retry loop re-runs
fetchfrom scratch rather than resuming from where the previous attempt left off - there's no standardized resume-cursor/replay contract across OpenAI/Anthropic/Gemini today, and none of them honor SSE's ownLast-Event-Idreconnection mechanism. If your use case needs this, de-duplicate already-received text on your end (e.g. by trackingchunk.indexand appended length) before appending further deltas from a retried attempt. runChat()andchatMachinearen't safe for multiple completion choices (n > 1). Both accumulatetext-deltachunks into one buffer regardless ofchunk.index, so text from different candidates will interleave. UsestreamChat()directly and group itsindex-tagged chunks yourself if you need multiple candidates.defaultDedupKeydoesn't account fortoolsormaxTokens. Two concurrent requests differing only in those fields will incorrectly dedupe to the same in-flight stream.
Requirements
- Runtime: ES2022+ (Node.js >= 18, modern browsers, Deno, Cloudflare Workers)
- TypeScript: 5.0+
- Dependencies:
xstate(^5.19.0)
License
MIT © CogniKit contributors
