System theme

Building an AI Agent Task Queue

Why standard async patterns fail for AI and how task queues solve the cascade failures, context loss, and token waste problem

KK13Updated 2 min read

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.

src/task-queue.tsTypeScript
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.

src/ai-agent.tsTypeScript
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