Why Standard Async Patterns Fail for AI
AI agent workloads break traditional asynchronous patterns. Unlike uniform HTTP requests, AI operations consume tokens at variable rates and unpredictable costs. A single prompt might use 500 tokens or 50,000 depending on context size and model usage.
Rate limits compound the problem. Most LLM APIs enforce both requests-per-minute and tokens-per-minute limits. Hit either one and you'll see 429 errors that ripple through the rest of the system.
Context loss is the most expensive failure mode. When an AI agent drops its conversation history mid-operation, rebuilding that context costs both tokens and time. Worse, the reconstructed context may not be identical, leading to inconsistent decisions.
A task queue becomes your source of truth. Each task stores the complete context needed to execute independently: the conversation history, the user's original request, intermediate results from prior operations, and metadata about what has been attempted. When a failure occurs, the retry pulls this context directly from the queue rather than trying to reconstruct it.
Building a Minimal AI Agent Task Queue
We'll build a queue that supports the patterns AI agents actually need: priority levels, adaptive rate limiting, dead letter handling, and context preservation. The implementation uses in-memory storage for simplicity, but the same structure maps cleanly to Redis or PostgreSQL when you need persistence and multiple workers.
class TaskQueue { constructor(options = {}) { this.tasks = new Map(); this.processing = new Set(); this.deadLetter = new Map(); this.maxRetries = options.maxRetries || 3; this.rateLimitPerMinute = options.rateLimitPerMinute || 60; this.tokenLimitPerMinute = options.tokenLimitPerMinute || 90000; this.recentRequests = []; this.recentTokens = []; this.priorities = { high: [], normal: [], low: [] }; } async add(task) { const taskId = task.id || `task_${Date.now()}_${Math.random()}`; const taskData = { id: taskId, priority: task.priority || 'normal', context: task.context, operation: task.operation, payload: task.payload, retries: 0, createdAt: Date.now(), status: 'pending' }; const contextHash = this._hashContext(task.context); const duplicate = this._findDuplicate(contextHash, task.operation); if (duplicate) { return { taskId: duplicate.id, isDuplicate: true }; } this.tasks.set(taskId, taskData); this.priorities[taskData.priority].push(taskId); return { taskId, isDuplicate: false }; } async process(handler) { while (true) { await this._waitForRateLimit(); const taskId = this._getNextTask(); if (!taskId) { await new Promise(resolve => setTimeout(resolve, 100)); continue; } const task = this.tasks.get(taskId); if (!task || this.processing.has(taskId)) continue; this.processing.add(taskId); task.status = 'processing'; try { const result = await handler(task); this._recordUsage(result.tokensUsed || 1000); this.tasks.delete(taskId); this.processing.delete(taskId); } catch (error) { this._handleFailure(task, error); } } } _getNextTask() { for (const priority of ['high', 'normal', 'low']) { const queue = this.priorities[priority]; while (queue.length > 0) { const taskId = queue.shift(); const task = this.tasks.get(taskId); if (task && task.status === 'pending') { return taskId; } } } return null; } async _waitForRateLimit() { const now = Date.now(); const oneMinuteAgo = now - 60000; this.recentRequests = this.recentRequests.filter(t => t > oneMinuteAgo); this.recentTokens = this.recentTokens.filter(t => t.timestamp > oneMinuteAgo); const currentRequests = this.recentRequests.length; const currentTokens = this.recentTokens.reduce((sum, t) => sum + t.count, 0); if (currentRequests >= this.rateLimitPerMinute || currentTokens >= this.tokenLimitPerMinute) { const waitTime = Math.max( this.recentRequests[0] + 60000 - now, this.recentTokens[0]?.timestamp + 60000 - now, 0 ); await new Promise(resolve => setTimeout(resolve, waitTime + 100)); } this.recentRequests.push(now); } _recordUsage(tokenCount) { this.recentTokens.push({ timestamp: Date.now(), count: tokenCount }); } _handleFailure(task, error) { this.processing.delete(task.id); task.retries++; if (task.retries >= this.maxRetries) { task.status = 'failed'; this.deadLetter.set(task.id, { ...task, error: error.message }); this.tasks.delete(task.id); } else { task.status = 'pending'; const delay = Math.min(1000 * Math.pow(2, task.retries), 30000); setTimeout(() => { this.priorities[task.priority].push(task.id); }, delay); } } _hashContext(context) { return JSON.stringify(context).split('').reduce( (hash, char) => ((hash << 5) - hash) + char.charCodeAt(0), 0 ); } _findDuplicate(contextHash, operation) { for (const [id, task] of this.tasks) { const taskHash = this._hashContext(task.context); if (taskHash === contextHash && task.operation === operation) { return task; } } return null; } getDeadLetterTasks() { return Array.from(this.deadLetter.values()); } async retryDeadLetter(taskId) { const task = this.deadLetter.get(taskId); if (!task) return false; task.retries = 0; task.status = 'pending'; this.tasks.set(taskId, task); this.priorities[task.priority].push(taskId); this.deadLetter.delete(taskId); return true; }} export default TaskQueue;Context Hashing and Deduplication for AI Workloads
AI agents frequently encounter duplicate requests: a user retries a failed operation, or multiple API calls generate identical analysis requests. Without deduplication, these become redundant token costs and unnecessary work.
The queue uses context hashing to detect duplicates without expensive string comparisons. The hash function is intentionally simple for performance, but effective at detecting identical conversation states. In production, you'd likely replace this with content-addressed storage or embedding-based similarity.
import Anthropic from '@anthropic-ai/sdk';import TaskQueue from './task-queue'; class AIAgent { constructor(apiKey) { this.client = new Anthropic({ apiKey }); this.queue = new TaskQueue({ maxRetries: 3, rateLimitPerMinute: 50, tokenLimitPerMinute: 80000 }); this.queue.process(this.executeTask.bind(this)); } async analyzeDocument(documentText, userId) { const context = { userId, conversationHistory: [], documentId: this._hashString(documentText) }; const analysisTask = await this.queue.add({ priority: 'high', context, operation: 'analyze', payload: { prompt: `Analyze this document and extract the 5 most important points:\n\n${documentText}`, onComplete: (result) => { this._generateSummary(result, context, userId); } } }); return analysisTask.taskId; } private _hashString(str) { return str.split('').reduce( (hash, char) => ((hash << 5) - hash) + char.charCodeAt(0), 0 ); }}Priority Assignment Based on User Impact
AI agents mix user-facing work with background processing. Requests that directly impact