Skip to content

@helix-agents/store-redis ​

Redis implementations of store interfaces for production use. Provides durable state storage and stream management across processes.

Installation ​

bash
npm install @helix-agents/store-redis ioredis

RedisStateStore ​

Redis-backed state storage.

typescript
import { RedisStateStore } from '@helix-agents/store-redis';

const stateStore = new RedisStateStore({
  host: 'localhost',
  port: 6379,
  password: 'optional',
  db: 0,
  keyPrefix: 'helix:', // Optional prefix for keys
  ttl: 86400 * 7, // Optional TTL in seconds (default: 7 days)
});

Configuration Options ​

typescript
interface RedisStateStoreOptions {
  // Connection options (ioredis compatible)
  host?: string;
  port?: number;
  password?: string;
  db?: number;

  // Or provide existing client
  client?: Redis;

  // Key prefix for namespacing
  keyPrefix?: string;

  // TTL for state entries (seconds)
  ttl?: number;

  // Logger
  logger?: Logger;
}

Methods ​

Same interface as InMemoryStateStore:

typescript
await stateStore.saveState(sessionId, state);
const state = await stateStore.loadState(sessionId);
await stateStore.deleteSession(sessionId);
await stateStore.updateStatus(sessionId, status);
await stateStore.appendMessages(sessionId, messages);
const { messages, hasMore } = await stateStore.getMessages(sessionId, options);

Checkpoint Methods ​

typescript
// Get a specific checkpoint
const checkpoint = await stateStore.getCheckpoint('session-123', 'cpv1-session-123-s5-...');

// Get most recent checkpoint
const latest = await stateStore.getLatestCheckpoint('session-123');

// List checkpoints with pagination
const result = await stateStore.listCheckpoints('session-123', {
  limit: 10,
  cursor: 'next-page-cursor',
});

Known gap (RM-38)

After a saveState, getLatestCheckpoint may not return the most recent checkpoint. saveState points the session at a checkpoint id that was never written, so the lookup falls back to stepCount ordering (RM-38). See the conformance findings.

createSession is one atomic Lua script. A concurrent loadState sees either no session or the whole one, and a recreated session inherits no stale messages, runs, sub-session refs or checkpoint index (RM-50, fixed). appendMessages throws Session not found: <id> for a session that does not exist (see the v0.50 migration notes).

Staging Methods ​

typescript
// Stage changes for a tool call within a step
await stateStore.stageChanges('session-123', 'step-1', {
  toolCallId: 'call-456',
  writes: {
    ops: [{ kind: 'append', key: 'notes', items: ['New finding'] }],
    warnings: [],
  },
  timestamp: Date.now(),
});

// Commit staged changes
await stateStore.promoteStaging('session-123', 'step-1');

// Rollback staged changes
await stateStore.discardStaging('session-123', 'step-1');

Distributed Coordination ​

typescript
// Atomic compare-and-set for status
const result = await stateStore.compareAndSetStatus(
  'session-123',
  ['running'], // Expected statuses (array)
  'interrupted' // New status
);

if (result.ok) {
  // result.newVersion is the post-increment version
} else {
  // result.currentStatus / result.currentVersion show actual stored values
}

// Increment resume count (for tracking retries)
const count = await stateStore.incrementResumeCount('session-123');

Data Structure ​

State is stored as Redis hashes:

helix:session:{sessionId} -> {
  sessionId: string,
  agentType: string,
  status: string,
  stepCount: number,
  customState: JSON string,
  messages: JSON string,
  output: JSON string,
  error: string,
  ...
}

RedisStreamManager ​

Redis-backed stream management using Redis Streams.

typescript
import { RedisStreamManager } from '@helix-agents/store-redis';

const streamManager = new RedisStreamManager({
  host: 'localhost',
  port: 6379,
  keyPrefix: 'helix:',
  maxStreamLength: 10000, // Max chunks per stream
  blockTimeout: 5000, // Read timeout (ms)
});

Configuration Options ​

typescript
interface RedisStreamManagerOptions {
  // Connection options
  host?: string;
  port?: number;
  password?: string;
  db?: number;
  client?: Redis;

  // Stream options
  keyPrefix?: string;
  maxStreamLength?: number; // Trim streams to this length
  blockTimeout?: number; // XREAD BLOCK timeout
  logger?: Logger;
}

Methods ​

Same interface as InMemoryStreamManager:

typescript
const writer = await streamManager.createWriter(streamId, agentId, agentType);
const reader = await streamManager.createReader(streamId);
await streamManager.endStream(streamId);
await streamManager.failStream(streamId, error);
const info = await streamManager.getStreamInfo(streamId);

Resumable Reader ​

typescript
const reader = await streamManager.createResumableReader(streamId, {
  fromSequence: 100,
});

if (reader) {
  for await (const { chunk, sequence } of reader) {
    // Process chunk
    // sequence can be used for Last-Event-ID
  }
}

Data Structure ​

Uses Redis Streams:

helix:stream:{streamId} -> XADD entries with:
  - chunk: JSON-encoded StreamChunk
  - sequence: Auto-incrementing ID

helix:stream:{streamId}:meta -> {
  status: 'active' | 'ended' | 'failed',
  error?: string,
  createdAt: ISO timestamp,
}

RedisResumableStorage ​

Lower-level resumable stream storage.

typescript
import { RedisResumableStorage } from '@helix-agents/store-redis';

const storage = new RedisResumableStorage({
  client: redisClient,
  keyPrefix: 'helix:',
});

// Create stream
await storage.create(streamId);

// Append chunk with sequence
await storage.append(streamId, chunk, sequence);

// Read from sequence
const chunks = await storage.readFrom(streamId, fromSequence, limit);

// Subscribe to new chunks
const unsubscribe = storage.subscribe(streamId, (chunk, sequence) => {
  console.log(`New chunk [${sequence}]:`, chunk);
});

// Get metadata
const metadata = await storage.getMetadata(streamId);

// End stream
await storage.end(streamId);

Error Classes ​

typescript
import {
  SequenceConflictError, // Sequence number conflict
  MalformedChunkError, // Invalid chunk data
} from '@helix-agents/store-redis';

try {
  await storage.append(streamId, chunk, sequence);
} catch (error) {
  if (error instanceof SequenceConflictError) {
    // Handle sequence conflict
  }
}

Usage Example ​

typescript
import { RedisStateStore, RedisStreamManager } from '@helix-agents/store-redis';
import { JSAgentExecutor } from '@helix-agents/runtime-js';
import { VercelAIAdapter } from '@helix-agents/llm-vercel';

// Create Redis-backed stores
const stateStore = new RedisStateStore({
  host: process.env.REDIS_HOST,
  port: parseInt(process.env.REDIS_PORT || '6379'),
  password: process.env.REDIS_PASSWORD,
  keyPrefix: 'myapp:agents:',
});

const streamManager = new RedisStreamManager({
  host: process.env.REDIS_HOST,
  port: parseInt(process.env.REDIS_PORT || '6379'),
  password: process.env.REDIS_PASSWORD,
  keyPrefix: 'myapp:agents:',
});

// Create executor
const executor = new JSAgentExecutor(stateStore, streamManager, new VercelAIAdapter());

// Execute agent
const handle = await executor.execute(MyAgent, 'Hello', { sessionId: 'my-session-1' });
const result = await handle.result();

Production Considerations ​

Connection Pooling ​

Use a shared Redis client:

typescript
import Redis from 'ioredis';

const redisClient = new Redis({
  host: process.env.REDIS_HOST,
  maxRetriesPerRequest: 3,
  enableReadyCheck: true,
});

const stateStore = new RedisStateStore({ client: redisClient });
const streamManager = new RedisStreamManager({ client: redisClient });

Cluster Support ​

Redis Cluster is not supported. Use a single Redis primary: standalone, or with replicas and Sentinel failover. RedisStateStore runs multi-key Lua scripts, and those keys are not co-located in one hash slot:

  • The global secondary-index keys (<prefix>:idx:sessions:*) hash to different slots from a session's own keys.
  • Several scripts build per-field custom-state list keys inside the script instead of declaring them in KEYS.

On a cluster these scripts fail with CROSSSLOT errors or touch undeclared keys. RedisStreamManager has not been verified on Cluster either. The key layout would need hash tags to support Cluster, which is not implemented.

Error Handling ​

typescript
stateStore.on('error', (error) => {
  console.error('Redis state store error:', error);
});

streamManager.on('error', (error) => {
  console.error('Redis stream manager error:', error);
});

TTL Management ​

Set appropriate TTLs based on your retention needs:

typescript
const stateStore = new RedisStateStore({
  ttl: 86400 * 30, // 30 days for state
});

RedisUsageStore ​

Redis-backed usage tracking storage. Track LLM tokens, tool executions, sub-agent calls, and custom metrics with persistence.

typescript
import Redis from 'ioredis';
import { RedisUsageStore } from '@helix-agents/store-redis';

const redis = new Redis(process.env.REDIS_URL);
const usageStore = new RedisUsageStore(redis, {
  keyPrefix: 'myapp', // Key prefix (default: 'helix')
  ttlSeconds: 86400 * 7, // Retention period (default: 24h)
});

Configuration Options ​

typescript
interface RedisUsageStoreOptions {
  keyPrefix?: string; // Prefix for Redis keys
  ttlSeconds?: number; // TTL for usage entries
}

Basic Usage ​

Pass to executor to enable usage tracking:

typescript
const handle = await executor.execute(agent, 'Do the task', { usageStore });
await handle.result();

// Get aggregated usage
const rollup = await handle.getUsageRollup();
console.log(`Total tokens: ${rollup?.tokens.total}`);

Methods ​

recordEntry ​

Record a usage entry (called internally by the framework).

typescript
await usageStore.recordEntry({
  kind: 'tokens',
  sessionId: 'session-123',
  stepCount: 1,
  timestamp: Date.now(),
  source: { type: 'agent', name: 'my-agent' },
  model: 'gpt-4o',
  tokens: { prompt: 100, completion: 50, total: 150 },
});

getEntries ​

Get usage entries for a run with optional filtering.

typescript
// All entries
const entries = await usageStore.getEntries('session-123');

// Filter by kind
const tokenEntries = await usageStore.getEntries('session-123', {
  kinds: ['tokens'],
});

// Filter by step range
const midRunEntries = await usageStore.getEntries('session-123', {
  stepRange: { min: 5, max: 10 },
});

// Filter by time range
const recentEntries = await usageStore.getEntries('session-123', {
  timeRange: { start: Date.now() - 3600000, end: Date.now() },
});

// Pagination
const page = await usageStore.getEntries('session-123', {
  limit: 10,
  offset: 20,
});

getRollup ​

Get aggregated usage rollup.

typescript
// This agent's usage only
const rollup = await usageStore.getRollup('session-123');

// Include sub-agent usage (lazy aggregation)
const totalRollup = await usageStore.getRollup('session-123', {
  includeSubAgents: true,
});

exists ​

Check if usage data exists by inspecting getRollup() — it returns null when no entries exist for the session.

typescript
const rollup = await usageStore.getRollup('session-123');
const hasUsage = rollup !== null;

delete ​

Delete usage data for a run.

typescript
await usageStore.delete('session-123');

getEntryCount ​

Get entry count without fetching entries.

typescript
const count = await usageStore.getEntryCount('session-123');

findSessionIds ​

Find all tracked session IDs.

typescript
const sessionIds = await usageStore.findSessionIds();

Data Structure ​

Usage entries are stored as JSON in Redis lists:

{keyPrefix}:usage:entries:{sessionId} -> LIST of JSON entries

Each entry includes:

  • id - Unique entry identifier
  • kind - Entry type (tokens, tool, subagent, custom)
  • runId - Associated run ID
  • stepCount - Step number when recorded
  • timestamp - Recording timestamp
  • source - Source information
  • Kind-specific fields (model, tokens, toolName, etc.)

TTL Management ​

Entries automatically expire based on ttlSeconds:

typescript
const usageStore = new RedisUsageStore(redis, {
  ttlSeconds: 86400 * 30, // 30 days retention
});

RedisLockManager ​

Distributed lock manager using Redis. Prevents concurrent execution of the same agent across multiple processes.

typescript
import Redis from 'ioredis';
import { RedisLockManager } from '@helix-agents/store-redis';

const redis = new Redis(process.env.REDIS_URL);
const lockManager = new RedisLockManager({
  redis,
  keyPrefix: 'myapp:locks:', // Optional
  retryCount: 0, // Default: fail-fast on contention
});

Configuration Options ​

typescript
interface RedisLockManagerOptions {
  keyPrefix?: string; // Prefix for lock keys
  defaultTTLMs?: number; // Default TTL in milliseconds
}

Methods ​

acquire ​

Acquire a lock with fencing token.

typescript
const lock = await lockManager.acquire('session-123', {
  ttlMs: 30000, // Lock TTL (optional, uses default)
  owner: 'executor-1', // Owner identifier
});

if (lock) {
  console.log('Lock acquired');
  console.log('Fencing token:', lock.fencingToken);

  // Do work...

  // Refresh to extend TTL
  await lock.refresh();

  // Release when done
  await lock.release();
} else {
  console.log('Lock held by another process');
}

isLocked ​

Check if a lock is currently held.

typescript
const locked = await lockManager.isLocked('session-123');

Fencing Tokens ​

Fencing tokens prevent split-brain scenarios:

typescript
const lock = await lockManager.acquire('session-123');

// Pass fencing token to operations
await stateStore.saveState(state, { fencingToken: lock.fencingToken });

// Operations with stale tokens are rejected

Usage with Executor ​

typescript
import { JSAgentExecutor } from '@helix-agents/runtime-js';
import { RedisLockManager } from '@helix-agents/store-redis';

const lockManager = new RedisLockManager({ redis });
const executor = new JSAgentExecutor(stateStore, streamManager, llmAdapter);

// Acquire lock before executing to prevent concurrent runs
await lockManager.withLock('my-session-1', async () => {
  await executor.execute(MyAgent, 'Hello', { sessionId: 'my-session-1' });
});

Lock Data Structure ​

Locks are stored as Redis strings with TTL:

{keyPrefix}lock:{sessionId} -> JSON { owner, fencingToken, acquiredAt }
TTL: configured ttlMs

Best Practices ​

  1. Set appropriate TTL - Long enough for operations, short enough for quick recovery
  2. Refresh locks - For long-running operations, periodically refresh
  3. Handle lock failures - Retry with backoff when lock acquisition fails
  4. Use fencing tokens - Pass to operations to prevent split-brain

Schemas ​

Zod schema for stored state:

typescript
import { AgentStateSchema } from '@helix-agents/store-redis';

// Validate state
const result = AgentStateSchema.safeParse(data);

See Also ​

Released under the MIT License.