How Nemori Handles Concurrent User Operations: Architecture and Implementation

Nemori uses per-user RLock instances to isolate concurrent operations, ThreadPoolExecutors for background processing, and sharded caching to enable high-throughput multi-user workloads without data races.

Nemori is an open-source memory system designed for AI applications that require persistent, searchable conversation history. When building production deployments of nemori-ai/nemori, handling concurrent user operations safely becomes critical to prevent data corruption and ensure responsive performance under load.

Per-User Isolation with Reentrant Locks

The RLock Architecture in MemorySystem

In src/core/memory_system.py, Nemori maintains a dictionary called _user_processing_locks that maps each owner_id to a dedicated threading.RLock instance. This design ensures that operations targeting different users proceed in parallel, while operations for the same user execute sequentially.

The helper method _get_user_processing_lock(owner_id) retrieves or creates the appropriate lock for a given user, as implemented in lines 17-36 of the memory system core.

Lock Acquisition in Public Methods

Every mutating public method—such as add_messages, delete_episode, and delete_semantic_memory—acquires the user-specific lock at entry using a with user_lock: context manager. This pattern appears in lines 71-74 of src/core/memory_system.py, ensuring that even if a single user sends rapid-fire requests, their data mutations remain atomic and race-condition free.

Background Processing with Thread Pools

ThreadPoolExecutor Configuration

To prevent blocking the main request thread during CPU-intensive or I/O-bound operations, Nemori delegates work to dedicated ThreadPoolExecutor instances. In src/core/memory_system.py, the system initializes a global executor specifically for semantic generation tasks (_semantic_generation_executor) at lines 64-70.

Additional executors handle bulk indexing operations with max_workers=2 (lines 84-106) and parallel episode creation in _batch_create_episodes (lines 51-66).

Semantic Generation Offloading

When add_messages triggers episode creation, the heavy lifting of generating semantic memories and vector embeddings occurs asynchronously. The semantic_task_manager.submit() method queues these tasks into the thread pool, allowing the API to return immediately while background workers process the LLM calls and index updates.

Resilience and Performance Optimizations

Retry Logic with SemanticTaskManager

Transient failures during LLM API calls could otherwise lose critical memory generation work. The SemanticTaskManager class in src/services/task_manager.py wraps the executor with simple retry logic (max_retries=1), implemented in the _run_with_retry method (lines 24-35). Tasks submitted via semantic_task_manager.submit(...) (lines 13-25) automatically retry once on failure before surfacing errors.

Sharded Caching Architecture

High-frequency cache access from concurrent threads can create contention bottlenecks. Nemori's PerformanceOptimizer (referenced in src/utils/performance.py) implements a sharded cache with num_cache_shards=40, as configured in lines 35-41 of the memory system. By distributing cached data across 40 independent shards, threads operating on different cache segments avoid blocking each other, significantly improving read/write throughput under load.

Event-Driven Decoupling

Synchronous processing of side effects—such as updating semantic indices after episode creation—would delay API responses. The EventBus in src/services/event_bus.py enables asynchronous event processing. When add_messages creates an episode, it publishes an episode_created event via event_bus.publish() (lines 122-127 in src/core/memory_system.py). Listeners, such as the semantic generation subscriber, consume these events on separate threads, completely decoupling the request lifecycle from background processing.

Practical Code Examples

The following examples demonstrate how Nemori's concurrency mechanisms operate in practice:


# Example: Adding messages from multiple threads concurrently

from threading import Thread
from nemori import MemorySystem, MemoryConfig

mem = MemorySystem(config=MemoryConfig())

def send_messages(user_id, msgs):
    result = mem.add_messages(user_id, msgs)
    print(f"{user_id}: {result['episodes_created']} episodes created")

# Two users send messages simultaneously

Thread(target=send_messages, args=("userA", [{"role":"user","content":"Hi"}]*5)).start()
Thread(target=send_messages, args=("userB", [{"role":"user","content":"Hello"}]*5)).start()

Each call acquires a distinct per-user lock, so userA and userB operations proceed in parallel without blocking each other.


# Example: Deleting an episode with pending background tasks

mem.delete_episode(owner_id="userA", episode_id="ep123")

# The method cancels any pending semantic-generation futures for that episode

# (see cancellation loop in delete_episode at lines 83-92 of memory_system.py)

# Example: Direct task submission with automatic retry

future = mem.semantic_task_manager.submit(
    mem.semantic_generator.generate_semantic_memory,
    owner_id="userA",
    episode=some_episode,
)

# The manager retries once on failure before raising (see _run_with_retry 

# in src/services/task_manager.py lines 24-35)

Summary

Nemori handles concurrent user operations through a multi-layered architecture that prioritizes data isolation and responsive performance:

  • Per-user RLock instances in src/core/memory_system.py ensure that operations for the same user execute sequentially while different users process in parallel.
  • ThreadPoolExecutors offload CPU-bound and I/O-bound work (semantic generation, bulk indexing) to background threads, preventing API blocking.
  • Automatic retry logic via SemanticTaskManager in src/services/task_manager.py protects against transient LLM failures.
  • Sharded caching with 40 independent shards reduces contention for high-frequency cache access.
  • Event-driven architecture using EventBus in src/services/event_bus.py decouples request handling from asynchronous side effects like index updates.

Frequently Asked Questions

How does Nemori prevent race conditions when multiple users access the system simultaneously?

Nemori prevents race conditions through per-user locking. Each user receives a dedicated threading.RLock instance stored in _user_processing_locks. When any mutating method like add_messages or delete_episode is called, it acquires the specific lock for that owner_id via _get_user_processing_lock. This ensures that while operations for different users proceed in parallel, operations for the same user are serialized, eliminating race conditions on user data.

What happens if two operations target the same user at the exact same time?

When two threads attempt to modify the same user's data simultaneously, the per-user lock acquired via _get_user_processing_lock(owner_id) serializes access. The second thread blocks until the first completes its operation and releases the lock. This blocking behavior occurs within the with user_lock: context manager used in methods like add_messages, ensuring that even rapid-fire requests from a single user maintain data consistency and prevent corruption.

How does Nemori prevent background tasks from blocking API responses?

Nemori uses ThreadPoolExecutor instances to offload heavy work such as semantic memory generation and vector indexing. When add_messages triggers episode creation, the system submits generation tasks via semantic_task_manager.submit() to a background thread pool. This allows the API to return immediately while worker threads handle LLM calls and index updates asynchronously, ensuring that expensive computational work never blocks the main request thread.

What mechanisms protect against failures in background processing?

The SemanticTaskManager class in src/services/task_manager.py implements retry logic with max_retries=1 in its _run_with_retry method. When tasks submitted via semantic_task_manager.submit() encounter transient errors—such as network timeouts during LLM API calls—the manager automatically retries the operation once before surfacing the error to the caller. This ensures robustness against temporary service interruptions without requiring manual intervention.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →