How Nemori Implements Episodic Memory: A Deep Dive into the Open-Source Architecture
Nemori implements episodic memory through a multi-layer pipeline that buffers user messages into thread-safe per-user queues, generates narrative Episode objects via LLM prompts, persists them as JSONL with in-memory indexing, and enables hybrid BM25 and vector search retrieval.
Nemori is an open-source memory system for AI agents developed by nemori-ai. Understanding how Nemori implements episodic memory reveals a sophisticated architecture designed for concurrency, durability, and multi-modal retrieval. This article examines the source code to explain the end-to-end flow from message ingestion to searchable episodic memories.
Message Ingestion and Buffering
The episodic memory pipeline begins with the MessageBufferManager in src/core/message_buffer.py. This component manages per-user message queues using fine-grained locks to ensure thread safety during concurrent access.
When MemorySystem.process_messages in src/core/memory_system.py receives incoming messages, it pushes them into the buffer. The buffer triggers episode creation when it reaches config.buffer_size_max or when a custom batch segmentation condition detects a natural conversation boundary. This design allows Nemori to balance between memory granularity and storage efficiency while preventing data races through per-user serialization.
Episode Generation via LLM
Once the buffer triggers, MemorySystem._batch_create_episodes calls _create_episode_from_messages, which delegates to the EpisodeGenerator in src/generation/episode_generator.py.
The generator follows a three-step process:
- Prompt Formatting: It formats the conversation history using a prompt template designed for narrative extraction.
- LLM Invocation: It calls
LLMClient.generate_json_responseto obtain structured JSON containing at leasttitle,content, and optionallytimestamp. - Fallback Handling: If the LLM fails or returns invalid JSON, the generator falls back to a deterministic template that extracts key information from the message list without external API calls.
This hybrid approach ensures robustness while maintaining high-quality narrative summaries when possible.
The Episode Data Model
The core data structure is the Episode dataclass defined in src/models/episode.py. This immutable representation stores:
- Metadata:
episode_id,title,content,timestamp - Provenance:
original_messages(the raw message list),message_count,boundary_reason(why this episode was created) - Extensibility: Optional
metadatadictionary andtagsfor categorization
By preserving the original messages alongside the generated narrative, Nemori maintains full auditability while providing compressed, searchable content.
Persistence and Storage
Durability is handled by EpisodeStorage in src/storage/episode_storage.py. This component implements:
- JSONL Format: Each episode serializes as a single JSON line in
<user_id>_episodes.jsonlfiles, enabling append-only writes and easy streaming. - Concurrent Safety: Per-user file locks prevent write conflicts during parallel episode creation.
- In-Memory Indexing: The storage maintains
_episode_indexand_user_indexdictionaries for O(1) lookup of episode file paths, eliminating the need to scan disk during retrieval.
The write buffer flushes immediately (size=1) to ensure durability at the cost of slightly higher I/O, prioritizing data safety over performance.
Hybrid Search Indexing
Retrieval is implemented through the UnifiedSearchEngine in src/search/unified_search.py, which combines lexical and semantic search:
- BM25Search (
src/search/bm25_search.py): Indexes episode titles and content for keyword-based retrieval. - ChromaSearchEngine (
src/search/chroma_search.py): Maintains vector embeddings for semantic similarity search using FAISS or Chroma backends.
When MemorySystem._handle_episode_created_event detects a new episode, it publishes an episode_created event. The search engine subscribes to these events and immediately indexes the episode in both backends, enabling real-time hybrid queries.
Event-Driven Semantic Extraction
The architecture supports downstream processing through the EventBus in src/services/event_bus.py. When semantic memory generation is enabled (config.enable_semantic_memory), the SemanticTaskManager in src/services/task_manager.py listens for episode_created events and triggers asynchronous extraction of semantic facts from the episode content.
This decouples episodic storage from knowledge graph construction, allowing each component to scale independently.
Configuration
Behavior is controlled by MemoryConfig in src/config.py, which exposes:
buffer_size_max: Threshold for automatic episode creation- Batch segmentation toggles for custom boundary detection
- Index rebuild policies and cache TTLs
- Feature flags for semantic memory and hybrid search
Code Examples
Automatically Creating Episodes from Message Streams
The most common pattern involves feeding messages to the MemorySystem and letting the buffer manager trigger episode creation:
from nemori.main import Nemori
from nemori.models import Message
nemori = Nemori()
user_id = "alice"
# Simulate conversation
messages = [
Message(role="user", content="Hey, I booked a flight to Paris tomorrow."),
Message(role="assistant", content="Great! Do you need a hotel recommendation?"),
Message(role="user", content="Yes, please.")
]
# Process messages (buffering happens internally)
for msg in messages:
result = nemori.process_message(user_id, msg)
# Access created episodes when buffer triggers
if result["episodes_created"]:
episode = result["episodes_created"][0]["episode_object"]
print(f"Episode: {episode.title}")
Relevant source: MemorySystem.process_messages in src/core/memory_system.py handles the buffering logic and episode creation triggers.
Manually Generating Episodes
For custom workflows, use the EpisodeGenerator directly:
from nemori.generation.episode_generator import EpisodeGenerator
from nemori.utils import LLMClient
from nemori.config import MemoryConfig
from nemori.models import Message
llm = LLMClient(api_key="...", model="gpt-4o-mini")
generator = EpisodeGenerator(llm, MemoryConfig())
episode = generator.generate_episode(
user_id="bob",
messages=[
Message(role="user", content="I started a new Python project."),
Message(role="assistant", content="What's the goal?"),
Message(role="user", content="A CLI tool for file backup.")
],
boundary_reason="Buffer size reached"
)
print(episode.title) # Generated narrative title
print(episode.content) # LLM-generated summary
Relevant source: EpisodeGenerator.generate_episode in src/generation/episode_generator.py implements the LLM prompt and fallback logic.
Retrieving Episodes with Hybrid Search
Search across episodic memories using both keyword and semantic similarity:
from nemori.search.unified_search import UnifiedSearchEngine
from nemori.utils import EmbeddingClient
from nemori.config import MemoryConfig
search_engine = UnifiedSearchEngine(
embedding_client=EmbeddingClient(api_key="...", model="text-embedding-3-large"),
config=MemoryConfig(),
language="en"
)
results = search_engine.search_episodes(
user_id="alice",
query="flight to Paris tomorrow",
top_k=5,
search_method="hybrid" # Combines BM25 and vector scores
)
for result in results:
print(f"{result['title']} (Score: {result['score']})")
Relevant source: UnifiedSearchEngine.search_episodes in src/search/unified_search.py orchestrates the hybrid retrieval logic.
Key Files in the Episodic Memory System
| File | Role | Direct Link |
|---|---|---|
src/models/episode.py |
Defines the Episode dataclass with metadata, provenance, and content fields |
src/models/episode.py |
src/generation/episode_generator.py |
LLM-driven narrative generation with JSON parsing and fallback templates | src/generation/episode_generator.py |
src/storage/episode_storage.py |
Thread-safe JSONL persistence and in-memory indexes for O(1) lookups | src/storage/episode_storage.py |
src/core/message_buffer.py |
Per-user message buffering with fine-grained locks and segmentation logic | src/core/message_buffer.py |
src/core/memory_system.py |
Central orchestrator coordinating buffering, generation, storage, and events | src/core/memory_system.py |
src/search/unified_search.py |
Hybrid search engine combining BM25 and vector similarity | src/search/unified_search.py |
src/search/bm25_search.py |
Lexical indexing implementation for keyword-based retrieval | src/search/bm25_search.py |
src/search/chroma_search.py |
Vector indexing using FAISS or Chroma backends | src/search/chroma_search.py |
src/services/event_bus.py |
Event publication system for episode_created events |
src/services/event_bus.py |
src/config.py |
Configuration management for buffer sizes and feature flags | src/config.py |
Summary
- Nemori implements episodic memory through a pipeline that buffers messages per-user, generates narrative summaries via LLM, and persists them as immutable
Episodeobjects with full provenance tracking. - Thread-safe ingestion in
MessageBufferManageruses fine-grained locks to handle concurrent message streams without data loss. - Robust generation employs
EpisodeGeneratorwith LLM JSON parsing and deterministic fallbacks to ensure episodes are created even when APIs fail. - Durable storage uses JSONL files with in-memory indexes in
EpisodeStorage, providing O(1) lookup performance and immediate persistence via per-user file locks. - Hybrid retrieval combines BM25 lexical search and vector similarity through
UnifiedSearchEngine, enabling both keyword and semantic memory recall. - Event-driven architecture decouples storage from downstream processing, allowing asynchronous semantic memory extraction via
EventBusandSemanticTaskManager.
Frequently Asked Questions
How does Nemori decide when to create a new episode?
Nemori triggers episode creation when the per-user message buffer reaches the buffer_size_max threshold defined in MemoryConfig, or when a custom batch segmentation condition detects a natural conversation boundary. The MessageBufferManager in src/core/message_buffer.py manages this logic using thread-safe locks to prevent race conditions during concurrent message processing.
What happens if the LLM fails during episode generation?
The EpisodeGenerator in src/generation/episode_generator.py implements a robust fallback mechanism. If LLMClient.generate_json_response returns invalid JSON or the API call fails, the system falls back to a deterministic template that extracts key information from the message list without requiring external API access. This ensures that episodic memory creation never blocks on LLM availability.
How does Nemori ensure thread safety when storing episodes?
Thread safety is enforced through multiple mechanisms in the storage layer. EpisodeStorage in src/storage/episode_storage.py uses per-user file locks to prevent concurrent write conflicts on JSONL files. Additionally, the MessageBufferManager uses fine-grained locks per user ID to serialize buffer operations, and the MemorySystem coordinates these components using a thread-pool executor for parallel processing without data races.
Can Nemori retrieve episodic memories using both keywords and semantic meaning?
Yes, Nemori implements hybrid search through the UnifiedSearchEngine in src/search/unified_search.py. This engine combines BM25Search for lexical keyword matching and ChromaSearchEngine for vector similarity search using embeddings. When querying, users can specify search_method="hybrid" to combine BM25 and vector scores, enabling retrieval of episodes that match specific keywords while also capturing semantically related but lexically distinct memories.
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 →