How Nemori's Message Buffering and Segmentation Works: A Technical Deep Dive
Nemori buffers incoming chat messages per user in thread-safe MessageBuffer instances and automatically segments them into coherent episodes using an LLM-driven BatchSegmenter when configurable thresholds are reached.
Nemori's message buffering and segmentation system forms the backbone of its episodic memory architecture. This pipeline efficiently manages high-throughput chat data by storing messages in per-user buffers and intelligently grouping them into meaningful conversation episodes. The implementation spans several core modules in the nemori-ai/nemori repository, ensuring thread-safe operations and LLM-guided coherence.
Core Components of Nemori's Buffering Architecture
Message and MessageBuffer Data Models
At the foundation lies the Message dataclass in src/models/message.py, which represents a single chat turn containing role, content, timestamp, and metadata. The MessageBuffer class maintains an ordered list of these Message objects for individual users, providing methods like add_message, add_messages, clear, size, and is_empty.
Notably, the is_timeout method currently returns False (timeout disabled), meaning segmentation relies purely on message count thresholds rather than temporal gaps.
Thread-Safe Buffer Management
The MessageBufferManager in src/core/message_buffer.py orchestrates per-user buffer lifecycle. It maintains a mapping of owner_id → MessageBuffer protected by user-level locks for thread-safe concurrent access.
Key methods include:
get_or_create_buffer: Retrieves existing or initializes new buffer for a useradd_message/add_messages: Appends messages to specific user's bufferis_buffer_ready/is_buffer_full: Checks buffer state against thresholds
The Segmentation Pipeline: From Buffer to Episode
Threshold-Based Triggering
Segmentation initiates when the BatchSegmenter in src/generation/batch_segmenter.py detects threshold breaches. The should_create_episode method compares current buffer size against config.batch_threshold:
# src/generation/batch_segmenter.py
def should_create_episode(self, buffer_size: int) -> tuple[bool, str]:
if buffer_size >= self.threshold:
return True, f"Batch threshold reached: {buffer_size}/{self.threshold}"
return False, ""
In src/core/memory_system.py, the add_messages method checks this trigger after each insertion when config.enable_batch_segmentation is enabled:
# src/core/memory_system.py (excerpt)
if self.config.enable_batch_segmentation and self._batch_segmenter:
buffer.add_messages(message_objects)
should_process, reason = self._batch_segmenter.should_create_episode(buffer.size())
if should_process:
created_episodes = self._process_batch_segmentation(owner_id, buffer, reason)
LLM-Driven Batch Segmentation
When triggered, BatchSegmenter.segment_batch sends all buffered messages to an LLM with a segmentation prompt. The LLM returns a JSON structure grouping messages into episodes by 1-based indices:
{
"episodes": [
{"indices": [1,2,3], "topic": "greeting"},
{"indices": [4,5,6,7], "topic": "product discussion"},
{"indices": [8,10,11], "topic": "pricing"},
{"indices": [9,12], "topic": "closing"}
]
}
If the LLM fails to return valid groups, the fallback mechanism places all messages into a single episode. The implementation uses:
# src/generation/batch_segmenter.py (excerpt)
response = self.llm_client.generate_json_response(
prompt=prompt,
temperature=0.2,
max_tokens=8192,
category="batch_segmentation"
)
episode_groups = [ep["indices"] for ep in response.get("episodes", [])]
Episode Generation and Persistence
For each group of indices, MemorySystem._process_batch_segmentation extracts corresponding Message objects and invokes EpisodeGenerator.generate_episode from src/generation/episode_generator.py:
# src/core/memory_system.py (excerpt)
group_messages = [buffer_messages[idx-1] for idx in indices]
episode = self.episode_generator.generate_episode(
user_id=owner_id,
messages=group_messages,
boundary_reason=f"Batch segmentation (group {i+1}/{len(episode_groups)})"
)
The EpisodeGenerator constructs a conversational transcript, prompts the LLM for a title and content summary, and determines the episode timestamp (preferring LLM-provided timestamps, otherwise using the earliest message timestamp). The resulting Episode dataclass is then:
- Persisted via
_episode_repository.save - Indexed in both lexical (
_lexical_index) and vector (_vector_index) stores - Cached in
self.episode_cache
If episode merging is enabled, the first episode of a batch may be merged with existing episodes via EpisodeMerger before indexing.
Finally, after processing all groups, the buffer is cleared:
buffer.clear()
Code Examples: Working with the Buffer
Configuring Auto-Segmentation
Enable automatic episode creation when the buffer reaches 10 messages:
from nemori.main import MemorySystem, MemoryConfig
cfg = MemoryConfig(
enable_batch_segmentation=True,
batch_threshold=10, # trigger after 10 messages
buffer_size_max=50,
buffer_size_min=5,
)
mem = MemorySystem(config=cfg)
# Simulated chat turn payloads
msgs = [
{"role": "user", "content": "Hi!"},
{"role": "assistant", "content": "Hello! How can I help?"},
# … add more messages …
]
result = mem.add_messages(owner_id="user123", messages=msgs)
print(result["episodes_created"]) # list of episodes produced when threshold hit
Manual Buffer Inspection
Check current buffer state for a specific user:
buf = mem.buffer_manager.get_buffer("user123")
print(buf.size()) # current number of stored messages
print(buf.get_conversation_text()) # raw conversation view
Forcing Batch Segmentation
Bypass the threshold and immediately segment current buffer contents:
buffer = mem.buffer_manager.get_or_create_buffer("user123")
# … add messages …
episode_groups = mem._batch_segmenter.segment_batch(buffer.get_messages())
print(episode_groups) # e.g. [[1,2,3], [4,5,6,7], …]
Summary
- Per-user isolation: Nemori maintains separate MessageBuffer instances for each user, managed by MessageBufferManager with fine-grained locking for thread safety.
- Threshold-driven processing: The BatchSegmenter monitors buffer size against
batch_thresholdand triggers segmentation when exceeded. - LLM intelligence: Segmentation uses an LLM to analyze message semantics and return JSON groupings of 1-based indices, creating coherent episodes by topic rather than arbitrary splits.
- Complete lifecycle: EpisodeGenerator transforms message groups into searchable Episode objects with titles, content summaries, and timestamps, then persists them to repositories and indexes (lexical and vector) before clearing the buffer.
Frequently Asked Questions
How does Nemori handle concurrent message writes from the same user?
Nemori uses user-level locks within the MessageBufferManager class in src/core/message_buffer.py. Each owner_id has its own lock, ensuring that only one thread can modify a user's buffer at a time while allowing concurrent operations across different users.
What happens if the LLM fails to segment messages properly?
The BatchSegmenter in src/generation/batch_segmenter.py includes a fallback mechanism. If the LLM returns invalid JSON, empty episode groups, or fails entirely, the system defaults to placing all buffered messages into a single episode, ensuring no data loss occurs during the segmentation process.
Can I disable automatic segmentation and manually control when episodes are created?
Yes. Set enable_batch_segmentation=False in your MemoryConfig. When disabled, MemorySystem.add_messages will accumulate messages in the buffer without triggering automatic segmentation. You can then manually invoke MemorySystem._process_batch_segmentation or use BatchSegmenter.segment_batch directly when you determine the buffer is ready for processing.
How does Nemori determine the timestamp for a generated episode?
The EpisodeGenerator in src/generation/episode_generator.py prefers timestamps provided by the LLM in its response. If the LLM does not supply a timestamp, the generator defaults to the earliest message timestamp within the episode's message group, ensuring chronological consistency in the memory store.
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 →