Retry Logic and Stale Message Recovery in Claude-Mem's PendingMessageStore
Claude-Mem implements a crash-safe SQLite message queue with configurable retry limits (default 3 attempts) and automatic stale message detection (default 5-minute threshold) to ensure reliable processing of SDK observations and summaries.
Claude-Mem persists every SDK-generated observation and summary in a robust SQLite-based queue managed by the PendingMessageStore class. Understanding the retry logic and stale message recovery mechanisms in Claude-Mem's PendingMessageStore is essential for building reliable worker processes that handle transient failures and crash scenarios gracefully.
Understanding the PendingMessageStore Architecture
The PendingMessageStore in src/services/sqlite/PendingMessageStore.ts functions as a persistent FIFO queue that guarantees at-most-once processing through atomic state transitions. Each message traverses a strict lifecycle: pending → processing → (confirmed | failed).
The Message Lifecycle
- Enqueue: New observations enter the queue with
status = 'pending'via theenqueuemethod. - Claim: Workers atomically flip messages to
status = 'processing'usingclaimAndDelete(which updates rather than deletes to enable recovery). - Success:
confirmProcessedpermanently deletes the row upon successful handling. - Failure:
markFailedimplements the retry logic by either incrementingretry_countor marking the message permanently failed.
Retry Logic Implementation
The retry logic prevents infinite loops while ensuring transient errors (network timeouts, rate limits) don't permanently drop critical observations.
Configurable Retry Limits
The constructor accepts a maxRetries parameter defaulting to 3 attempts (PendingMessageStore.ts, lines 43-46). This value persists for the store instance and applies to all messages in the queue.
The markFailed Method
When processing fails, the store invokes markFailed (PendingMessageStore.ts, lines 8-33), which implements the core retry decision:
- If
retry_count < maxRetries: Increment the counter, resetstatustopending, and queue for immediate retry. - If
retry_count >= maxRetries: Updatestatustofailedand stop further attempts.
This ensures messages receive exactly maxRetries opportunities before entering a terminal failed state.
Manual Retry Controls
Developers can force immediate reprocessing via retryMessage (PendingMessageStore.ts, lines 2-8), which resets a specific message's retry_count to 0 and restores pending status regardless of previous failures.
import { PendingMessageStore } from './src/services/sqlite/PendingMessageStore.js';
// Initialize with default 3 retries
const pendingStore = new PendingMessageStore(databaseConnection, 3);
// Process messages with automatic retry handling
const message = pendingStore.claimAndDelete(sessionId);
if (message) {
try {
await processObservation(message);
pendingStore.confirmProcessed(message.id);
} catch (error) {
// Automatically increments retry_count or marks failed
pendingStore.markFailed(message.id);
}
}
// Manually rescue a permanently failed message
pendingStore.retryMessage(stuckMessageId);
Stale Message Recovery Mechanisms
When workers crash or network partitions occur, messages may remain stranded in the processing state indefinitely. The store implements multiple recovery strategies to detect and resurrect these stale messages.
Automatic Detection on Startup
The resetStaleProcessingMessages method (PendingMessageStore.ts, lines 31-48) runs automatically when workers initialize. It queries for messages with status = 'processing' where started_processing_at_epoch exceeds a configurable threshold (default 5 minutes), updating them back to pending status.
This ensures that even if a worker dies mid-processing, another worker will reclaim the message after the timeout expires.
Administrative Recovery Methods
For operational intervention, the store exposes two additional recovery functions:
resetStuckMessages(PendingMessageStore.ts, lines 35-50): Accepts a custom epoch threshold (or0to reset all processing messages regardless of age), useful for recovering from extended outages.resetProcessingToPending(PendingMessageStore.ts, lines 12-20): Targets a specificsessionDbId, resetting only that session's stranded messages without affecting other concurrent workflows.
// Automatic recovery on worker startup (default 5 min threshold)
pendingStore.resetStaleProcessingMessages();
// Emergency recovery after extended downtime
pendingStore.resetStuckMessages(0); // Reset all processing messages immediately
// Recover only a specific session after a partial crash
pendingStore.resetProcessingToPending(specificSessionId);
Complete Implementation Example
The following example demonstrates the full lifecycle including enqueue, claim, retry handling, and crash recovery:
import { PendingMessageStore } from './src/services/sqlite/PendingMessageStore.js';
import type { PendingMessage } from './src/services/worker-types.js';
// Initialize store with 3 retry attempts
const db = /* obtain Database instance */;
const pendingStore = new PendingMessageStore(db, 3);
// 1. Enqueue a new observation
const pendingId = pendingStore.enqueue(
sessionDbId,
contentSessionId,
{
type: 'observation',
tool_name: 'search',
tool_input: { query: 'Claude-Mem architecture' },
prompt_number: 1,
} as PendingMessage,
);
// 2. Worker startup: recover stale messages from previous crashes
pendingStore.resetStaleProcessingMessages(); // 5 minute threshold
// 3. Processing loop
const msg = pendingStore.claimAndDelete(sessionDbId);
if (msg) {
try {
// Process the observation...
await handleObservation(msg);
// Confirm success
pendingStore.confirmProcessed(msg.id);
} catch (error) {
// Let the store handle retry logic
pendingStore.markFailed(msg.id);
// If retry_count < 3, message returns to pending queue
// If retry_count >= 3, status becomes 'failed'
}
}
// 4. Manual intervention if needed
pendingStore.retryMessage(failedMessageId); // Force retry regardless of count
Summary
- Persistent SQLite Queue:
PendingMessageStoreinsrc/services/sqlite/PendingMessageStore.tsprovides durable storage for SDK observations with atomic state transitions. - Configurable Retry Logic: The constructor accepts a
maxRetriesparameter (default 3) that controls how many timesmarkFailedwill re-queue a message before marking it permanently failed. - Automatic Stale Recovery:
resetStaleProcessingMessages(default 5-minute threshold) automatically returns crashed or hung messages to the pending queue on worker startup. - Administrative Controls:
resetStuckMessagesandresetProcessingToPendingprovide manual intervention capabilities for specific sessions or emergency recovery scenarios. - At-Most-Once Guarantee: The combination of
claimAndDelete(atomic claim),confirmProcessed(deletion on success), and recovery mechanisms ensures messages are processed exactly once despite worker crashes or transient failures.
Frequently Asked Questions
How does Claude-Mem handle message processing failures?
When a worker encounters an error processing a message, it calls markFailed in PendingMessageStore.ts (lines 8-33). This method checks the current retry_count against the configured maxRetries (default 3). If the message hasn't exceeded its limit, the method increments the counter and resets the status to pending for immediate reprocessing. Once the limit is reached, the status changes permanently to failed and the message requires manual intervention via retryMessage.
What happens to messages when a worker crashes?
Messages stranded in the processing state due to worker crashes are automatically recovered through resetStaleProcessingMessages (PendingMessageStore.ts, lines 31-48). This method runs on worker startup and identifies messages where started_processing_at_epoch exceeds a configurable threshold (default 5 minutes). It atomically updates these records back to pending status, ensuring another worker can claim them without data loss or duplicate processing.
Can I customize the retry behavior for specific messages?
While the maxRetries parameter applies globally to all messages in a PendingMessageStore instance, you can force immediate retry of specific failed messages using retryMessage (PendingMessageStore.ts, lines 2-8). This method resets both the retry_count to 0 and the status to pending, effectively giving the message a fresh set of retry attempts regardless of its previous failure history.
Where is the pending message queue stored?
The pending message queue persists in a SQLite database managed by the PendingMessageStore class in src/services/sqlite/PendingMessageStore.ts. The store uses atomic SQL transactions to ensure that state transitions—such as claiming messages with claimAndDelete or confirming completion with confirmProcessed—remain consistent even during application crashes or power failures.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →