记忆系统:从 Embedding 到 Memory Flush
源码分布在
memory/目录下 24 个文件中:embedding.py、flush.py、store.py、retrieval.py...
一、架构总览:子系统的完整分工
OpenSquilla 的记忆系统由多个专门的文件各司其职:
memory/
├── embedding.py ← 向量嵌入引擎
├── embedding_resolver.py ← 自动选择嵌入方案
├── flush.py ← Memory Flush 主逻辑
├── flush_config.py ← Flush 触发配置
├── flush_status.py ← Flush 状态管理
├── store.py ← 向量存储与索引
├── retrieval.py ← 记忆检索(语义 + 关键词)
├── archive.py ← 记忆归档
├── checkpoint.py ← 检查点/恢复
├── dream_factory.py ← Dream 构建工厂
├── manager.py ← 记忆管理器(总入口)
├── meta.py ← 记忆元数据
├── retention.py ← 记忆保留策略
├── redaction.py ← 敏感信息脱敏
├── sync_manager.py ← 多会话同步
├── turn_capture.py ← Turn 级记忆捕获
├── session_flush.py ← 会话级 Flush
├── session_source.py ← 会话数据源
└── source_paths.py ← 数据文件路径二、本地嵌入引擎
2.1 为什么本地?
云端嵌入方案:
用户消息 → 发到 OpenAI API → 返回嵌入向量
问题:隐私泄露、网络延迟、API 费用
OpenSquilla 方案:
用户消息 → 本地 Ollama → 本地 ONNX 模型 → 嵌入向量
优势:零网络调用、零 API 费用、隐私完全保护2.2 embedding.py:嵌入引擎接口
# embedding.py 简化还原
class EmbeddingEngine:
"""本地嵌入引擎,支持多种后端"""
def __init__(self, provider: str = "ollama"):
self.provider = provider
self.model = self._load_model()
async def embed(self, text: str) -> list[float]:
"""将文本转为向量嵌入"""
if self.provider == "ollama":
return await self._embed_ollama(text)
elif self.provider == "onnx":
return await self._embed_onnx(text)
else:
return await self._embed_fallback(text)
async def _embed_ollama(self, text: str) -> list[float]:
"""通过本地 Ollama 服务嵌入"""
response = await ollama.embeddings(
model="nomic-embed-text",
prompt=text,
)
return response["embedding"] # 768 维向量2.3 embedding_resolver.py:自动选择最佳方案
# embedding_resolver.py 的逻辑
def resolve_embedding_provider(config: dict) -> str:
"""
自动选择可用的嵌入方案:
1. 如果配置了 embedding provider → 使用配置的方案
2. 如果检测到本地 Ollama → 使用 Ollama
3. 如果 ONNX 模型已下载 → 使用 ONNX(最快)
4. 回退到关键词检索(无需嵌入)
"""
if config.get("embedding_provider"):
return config["embedding_provider"]
if detect_ollama():
return "ollama"
if has_onnx_model():
return "onnx"
return "keyword" # 回退方案设计哲学: 永远有回退方案。即使没有嵌入模型,记忆系统也能通过关键词检索工作。
三、Memory Store:向量存储与检索
3.1 store.py:双路检索
# store.py 核心逻辑还原
class MemoryStore:
"""记忆存储,支持向量检索 + 关键词回退"""
def __init__(self, db_path: str):
self.db = sqlite3.connect(db_path)
self.embedder = EmbeddingEngine()
self._init_tables()
def _init_tables(self):
self.db.execute("""
CREATE TABLE IF NOT EXISTS memories (
id TEXT PRIMARY KEY,
content TEXT NOT NULL,
embedding BLOB, -- 向量嵌入
category TEXT,
importance REAL DEFAULT 0.5,
created_at TEXT,
last_accessed TEXT
)
""")
self.db.execute("""
CREATE VIRTUAL TABLE IF NOT EXISTS memories_fts
USING fts5(content) -- 全文搜索索引
""")
async def search(self, query: str, k: int = 5) -> list[Memory]:
"""
双路检索:
1. 向量语义检索(如果有嵌入模型)
2. 关键词全文检索(FTS5)
3. 合并结果,去重排序
"""
results = []
# 路径 1: 语义检索
if self.embedder.available:
query_vec = await self.embedder.embed(query)
semantic_results = self._cosine_search(query_vec, k * 2)
results.extend(semantic_results)
# 路径 2: 关键词检索(FTS5)
keyword_results = self._keyword_search(query, k * 2)
results.extend(keyword_results)
# 合并 + 去重 + 按相关性排序
return self._merge_and_rank(results, k)
def _cosine_search(self, vec: list, k: int) -> list:
"""余弦相似度检索"""
# 加载所有嵌入向量(实践中会使用 ANN 索引优化)
rows = self.db.execute("SELECT id, content, embedding FROM memories").fetchall()
scores = []
for row in rows:
stored_vec = pickle.loads(row[2]) if row[2] else None
if stored_vec:
sim = cosine_similarity(vec, stored_vec)
scores.append((row[0], row[1], sim))
return sorted(scores, key=lambda x: x[2], reverse=True)[:k]
def _keyword_search(self, query: str, k: int) -> list:
"""FTS5 全文检索"""
rows = self.db.execute(
"SELECT id, content, rank FROM memories_fts WHERE content MATCH ? ORDER BY rank LIMIT ?",
(query, k)
).fetchall()
return [(r[0], r[1], r[2]) for r in rows]3.2 部署细节:无服务器,单文件
# 所有记忆存储在一个 SQLite 文件中
# 路径:~/.opensquilla/memory.db
#
# 一个文件包含:
# - 记忆数据表
# - FTS5 全文索引
# - 嵌入向量
#
# 无需安装向量数据库(Milvus/Qdrant/Pinecone)四、Memory Flush:触发与写入
4.1 flush.py:周期性刷新
# flush.py 核心逻辑还原
class MemoryFlushEngine:
"""记忆刷新引擎"""
def __init__(self, store: MemoryStore, config: FlushConfig):
self.store = store
self.config = config
self.last_flush_time = 0
self.turn_count_since_flush = 0
async def should_flush(self, turn_context: TurnContext) -> bool:
"""
判断是否应该执行 Flush:
触发条件(任何一条满足即触发):
1. 距离上次 Flush 超过 N 分钟(默认 30)
2. 对话轮数超过 N 轮(默认 50)
3. 会话即将结束
4. 上下文压缩事件后
5. 手动触发
"""
if self.config.manual_trigger:
return True
time_since_flush = time.time() - self.last_flush_time
if time_since_flush > self.config.interval_minutes * 60:
return True
if self.turn_count_since_flush > self.config.max_turns:
return True
if turn_context.is_session_end:
return True
if turn_context.compaction_just_occurred:
return True
return False
async def flush(self, session_messages: list):
"""
执行 Memory Flush:
1. 从对话中提取关键信息
2. 去重(与已有记忆比较)
3. 合并(更新已有记忆)
4. 写入存储
5. 更新嵌入索引
"""
# 提取
memories = await self._extract_memories(session_messages)
# 去重 + 合并
for mem in memories:
existing = self.store.find_similar(mem.content)
if existing:
self._merge(existing, mem) # 更新访问时间 + 重要度
else:
self.store.add(mem) # 新增
# 更新索引
self.store.rebuild_fts()
self.store.update_embeddings()
self.last_flush_time = time.time()
self.turn_count_since_flush = 04.2 flush_config.py:可配置的 Flush 策略
class FlushConfig:
interval_minutes: int = 30 # 刷新间隔(分钟)
max_turns: int = 50 # 最大对话轮数
auto_archive_days: int = 90 # 自动归档天数
max_memory_count: int = 1000 # 最大记忆条数
importance_threshold: float = 0.3 # 重要度阈值五、记忆检索与归档
5.1 retrieval.py:记忆检索总入口
# retrieval.py 检索流程
async def retrieve_memories(query: str, store: MemoryStore, k: int = 5):
"""
三步检索流程:
1. 语义检索(如果有嵌入)
2. 关键词检索(FTS5,总是可用)
3. 混合排序:相关性分 + 重要度分 + 新鲜度分
"""
def score_memory(mem, query_relevance):
"""综合评分"""
relevance = query_relevance # 0-1
importance = mem.importance # 0-1
freshness = decay(mem.last_accessed) # 0-1,越近越高
return relevance * 0.6 + importance * 0.25 + freshness * 0.155.2 archive.py:记忆归档
# archive.py 归档逻辑
async def archive_old_memories(store: MemoryStore, days: int = 90):
"""
自动归档策略:
1. 标记 90 天未访问的记忆为 'archive'
2. 重要度 < 0.2 且 30 天未访问 → 归档
3. 归档记忆移到 archive 表(保留嵌入向量)
4. FTS 索引中移除(节省空间)
5. 归档记忆仍然可检索(通过 archive API)
"""
cutoff = datetime.now() - timedelta(days=days)
store.db.execute("""
UPDATE memories SET status = 'archived'
WHERE last_accessed < ? AND importance < 0.2
""", (cutoff,))六、记忆系统的完整数据流
┌──────────────────────────────────────────────────────────┐
│ Memory 数据流 │
├──────────────────────────────────────────────────────────┤
│ │
│ 会话进行中... │
│ ↓ │
│ turn_capture.py: 每个 turn 捕获关键信息 │
│ ↓ │
│ flush.py: 触发条件满足? │
│ ↓ YES │
│ 提取记忆 → 去重 → 合并 → 写入 store │
│ ↓ │
│ embedding.py: 为新记忆生成嵌入向量(本地) │
│ ↓ │
│ store.py: SQLite + FTS5 索引更新 │
│ ↓ │
│ 用户下次对话 → retrieval.py 检索 → 注入 prompt │
│ ↓ │
│ 90 天后 → archive.py 自动归档 │
│ │
└──────────────────────────────────────────────────────────┘七、记忆搜索:在 Agent 中的实际使用
# 在 prompt_assembler_stage.py 中
# 每次对话开始前,注入相关记忆
async def inject_memories(context: TurnContext):
store = MemoryStore("~/.opensquilla/memory.db")
# 搜索与当前对话相关的记忆
memories = await retrieve_memories(
query=context.user_input,
store=store,
k=5
)
# 构建记忆提示
memory_text = "## Your Memories of This User ##\n"
for mem in memories:
memory_text += f"- [{mem.category}] {mem.content}\n"
context.memory_snapshot = memory_text
context.memory_fingerprint = hash(memory_text) # 用于后续去重八、总结
OpenSquilla 记忆系统的设计原则:
| 原则 | 实现 |
|---|---|
| 本地优先 | Ollama/ONNX 嵌入,零云端依赖 |
| 双路检索 | 向量语义 + 关键词全文,互为备份 |
| 自动触发 | 时间/轮数/事件驱动的 Flush |
| 智能归档 | 按重要度和访问频率自动清理 |
| 无服务器 | 一个 SQLite 文件搞定一切 |
下期预告
OpenSquilla 源码解析 ④:Meta Skills — Agent 的自演进机制
- • Meta Skill 的生命周期源码
- • 从 Auto-Propose 到 E2E Test 的完整链路
- • Skill 过滤器的依赖注入机制
- • 与 SquillaRouter 的协同工作
科技洞察助手 | 源码深度解析系列
夜雨聆风