# How Nemori Handles Concurrent User Operations: Architecture and Implementation

> Discover how Nemori tackles concurrent user operations with per-user RLock, ThreadPoolExecutors, and sharded caching for high-throughput, race-free workloads. Learn its architecture and implementation.

- Repository: [Nemori AI/nemori](https://github.com/nemori-ai/nemori)
- Tags: architecture
- Published: 2026-03-08

---

**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`](https://github.com/nemori-ai/nemori/blob/main/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`](https://github.com/nemori-ai/nemori/blob/main/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`](https://github.com/nemori-ai/nemori/blob/main/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`](https://github.com/nemori-ai/nemori/blob/main/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`](https://github.com/nemori-ai/nemori/blob/main/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`](https://github.com/nemori-ai/nemori/blob/main/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`](https://github.com/nemori-ai/nemori/blob/main/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:

```python

# 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.

```python

# 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)

```

```python

# 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`](https://github.com/nemori-ai/nemori/blob/main/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`](https://github.com/nemori-ai/nemori/blob/main/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`](https://github.com/nemori-ai/nemori/blob/main/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`](https://github.com/nemori-ai/nemori/blob/main/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.