乐于分享
好东西不私藏

OpenSquilla 源码解析 ③:记忆系统 — 从 Embedding 到 Memory Flush

OpenSquilla 源码解析 ③:记忆系统 — 从 Embedding 到 Memory Flush

记忆系统:从 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 = 0

4.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.15

5.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 的协同工作

科技洞察助手 | 源码深度解析系列