dApp Docs/AI Agent RAG 知识库摄取指南
Development reference. Not independently verified for production.

AI Agent RAG 知识库摄取指南

链 ID: msg-chain-1 | 共识: DAR | 虚拟机: CosmWasm (WasmVM)
状态: 规划文档 — 主网裁决为 No-Go,所有数据均为主网预演
Gas: 1,000,000,000 attoMSG/gas
Gas 分配: 40% 验证者 / 30% 开发者 / 20% 燃烧 / 10% 基金会金库
地址格式: SHA3-512(前40位) + SHA-256 校验 | 签名: Dilithium-5 (公钥2592字节 / 私钥4864字节 / 签名4595字节)
AI Agent安全边界: 永不自主创建新合约,永不自主调整 Gas 参数,永不自主铸造/销毁代币


目录

  1. 引言 — 为什么需要 RAG 知识库摄取
  2. 白皮书机器层的 RAG 分层架构
  3. 6 步批处理流程详解
  4. 元数据策略(推荐字段与索引设计)
  5. 检索策略(Chunk 优先、Rerank、降级)
  6. 边界声明与服务策略
  7. 向量库选择与部署实践(Chroma/Milvus/Pinecone 集成)
  8. 增量更新策略
  9. 示例:构建一个 MSG Chain RAG 知识库
  10. 总结

1. 引言 — 为什么需要 RAG 知识库摄取

1.1 背景

随着 AI Agent 和大型语言模型(LLM)的广泛应用,链上技术与白皮书中积累的结构化知识需要被高效地检索和引用。MSG Chain 白皮书系统采用机器可读的模块化架构,通过标准化的 JSON 导出格式,使 AI Agent 能够以 RAG(Retrieval-Augmented Generation)方式摄取、索引和检索知识。

1.2 为什么需要专门的摄取流程

链上白皮书知识具有以下特点,决定了通用文档摄取方案无法胜任:

特点 说明 对 RAG 的要求
模块化结构 白皮书按模块组织,每个模块有独立的状态和元数据 需保留模块层级关系和状态标记
状态语义 每个模块标记为 implemented / partial / planned 检索需按状态优先级排序
边界声明 模块和段落带有边界条款和证据引用 检索结果必须附带边界声明,避免过度声称
版本演化 白皮书随时间更新,模块状态可能变化 支持增量更新和版本追踪
机器可读 所有内容以结构化 JSON 发布 无需解析非结构化文本,直接结构化摄取

1.3 目标读者

本指南面向:

1.4 本文档使用的约定


2. 白皮书机器层的 RAG 分层架构

2.1 架构概览

MSG Chain 白皮书系统为 RAG 摄取设计了明确的分层入口架构。AI Agent 通过以下四个入口点发现和访问白皮书知识:

agent_entry.json
    └── whitepaper_manifest.json
            ├── module_exports/index.json
            │       └── module_exports/*.json
            └── module_chunks/index.json
                    └── module_chunks/*.json

2.2 入口点 1:agent_entry.json

agent_entry.json 是 AI Agent 接入白皮书系统的顶级入口文件。Agent 首次发现白皮书时,应首先访问此文件。它包含:

获取方式:

GET https://msgchain.org/whitepaper/integration_examples/agent_entry.json

Agent 层面的建议逻辑:

def discover_whitepaper():
    entry = fetch_json("https://msgchain.org/whitepaper/integration_examples/agent_entry.json")
    if entry.get("schema_version") == "v1":
        manifest_url = resolve_relative(entry, "../whitepaper_manifest.json")
        return fetch_manifest(manifest_url)
    raise ValueError("不兼容的 schema 版本")

2.3 入口点 2:whitepaper_manifest.json

whitepaper_manifest.json 是整个白皮书系统的清单文件,记录:

示例结构:

{
  "schema_version": "v1",
  "generated_at_utc": "2025-06-30T12:00:00Z",
  "entry_points": [
    "module_exports/index.json",
    "module_chunks/index.json"
  ]
}

RAG 摄取管线中的用途:

  1. 版本校验:确认当前白皮书版本与索引兼容
  2. 入口枚举:发现所有需要摄取的子索引文件
  3. 时间追踪:记录生成时间,用于增量更新的判断

2.4 入口点 3:module_exports/index.json

module_exports/index.json 是模块导出索引,枚举所有白皮书模块。每个模块包含:

Agent 处理逻辑:

def load_module_index(manifest):
    index_url = resolve_relative(manifest, "module_exports/index.json")
    index = fetch_json(index_url)
    modules = []
    for entry in index["modules"]:
        module_url = resolve_relative(index, entry["module_filename"])
        module_data = fetch_json(module_url)
        modules.append(module_data)
    return modules

2.5 入口点 4:module_chunks/index.json

module_chunks/index.json 是分块索引,枚举所有文本块及其所属模块的映射关系。每个分块条目包含:

分块原则:

MSG Chain 白皮书的分块策略遵循语义边界,而非简单的字符截断。每个 chunk 对应于白皮书中的一个语义段落或子章节,确保:

2.6 分层设计的优势

层级 作用 访问频率 变更频率
agent_entry.json 发现入口 一次(首次接入) 极低
whitepaper_manifest.json 版本和入口清单 每次摄取 低
module_exports/index.json 模块索引 每次摄取 中
module_exports/*.json 模块详细数据 每次摄取 中
module_chunks/index.json 分块索引 每次摄取 高
module_chunks/*.json 分块详细数据 每次摄取 高

这种分层结构允许 Agent 和 RAG 管线按需加载,避免一次性加载全部数据。


3. 6 步批处理流程详解

3.1 流程总览

官方 RAG Ingest Flow 定义了 6 个串联的批处理步骤,从加载清单到写入向量库:

Step 1: load_manifest
    ↓
Step 2: load_module_index
    ↓
Step 3: load_chunk_index
    ↓
Step 4: ingest_chunk_documents
    ↓
Step 5: ingest_module_documents
    ↓
Step 6: attach_boundaries

3.2 Step 1:load_manifest

输入: whitepaper_manifest.json

目的: 获取当前白皮书的 schema 版本、生成时间和入口点列表。

具体操作:

def step_load_manifest():
    manifest = fetch_json(
        "https://msgchain.org/whitepaper/integration_examples/whitepaper_manifest.json"
    )
    assert manifest["schema_version"] == "v1", "不兼容的 schema 版本"
    return {
        "schema_version": manifest["schema_version"],
        "generated_at": manifest["generated_at_utc"],
        "entry_points": manifest["entry_points"]
    }

校验规则:

检查项 要求 失败处理
schema_version 必须是 v1 中止摄取,记录错误
generated_at_utc 必须存在且为合法 UTC 时间 跳过时间校验,记录警告
entry_points 必须包含 module_exports/index.json 和 module_chunks/index.json 中止摄取

3.3 Step 2:load_module_index

输入: module_exports/index.json

目的: 枚举所有模块,保留每个模块的元数据(状态、分组、标签、URL 等)。

具体操作:

def step_load_module_index(manifest):
    index = fetch_json(resolve_relative(
        manifest["base_url"], "module_exports/index.json"
    ))
    modules = []
    for m in index["modules"]:
        modules.append({
            "module_filename": m["module_filename"],
            "module_title": m["module_title"],
            "status": m["status"],
            "status_label": m["status_label"],
            "group": m["group"],
            "group_label": m["group_label"],
            "tags": m["tags"],
            "module_public_url": m["module_public_url"]
        })
    return modules

元数据保留策略:

模块状态语义:

状态 含义 检索权重
implemented 已实现 最高
partial 部分实现 中等
planned 计划中 最低(不作为主要证据)

3.4 Step 3:load_chunk_index

输入: module_chunks/index.json

目的: 枚举所有分块文件,建立 chunk 到 module 的映射关系。

具体操作:

def step_load_chunk_index(manifest):
    index = fetch_json(resolve_relative(
        manifest["base_url"], "module_chunks/index.json"
    ))
    chunks = []
    for c in index["chunks"]:
        chunks.append({
            "chunk_id": c["chunk_id"],
            "module_filename": c["module_filename"],
            "chunk_public_url": c["chunk_public_url"]
        })
    return chunks

映射关系:

每个 chunk → 一个 module(多对一)
每个 module → 多个 chunk(一对多)

3.5 Step 4:ingest_chunk_documents

输入: module_chunks/*.json

目的: 将分块文本作为主要可检索文档体写入向量数据库。这是 RAG 检索的主要数据源。

具体操作:

def step_ingest_chunk_documents(chunks, vector_store):
    for chunk_ref in chunks:
        chunk_data = fetch_json(resolve_relative(
            base_url, chunk_ref["chunk_id"]
        ))
        document = {
            "id": chunk_ref["chunk_id"],
            "text": chunk_data["text"],
            "metadata": {
                "chunk_id": chunk_ref["chunk_id"],
                "module_filename": chunk_ref["module_filename"],
                "module_public_url": chunk_data["module_public_url"],
                "chunk_public_url": chunk_ref["chunk_public_url"],
                "evidence_refs": chunk_data.get("evidence_refs", []),
                "boundary_clauses": chunk_data.get("boundary_clauses", [])
            }
        }
        vector_store.upsert(
            id=document["id"],
            vector=embedding_model.embed(document["text"]),
            metadata=document["metadata"]
        )

向量化策略:

3.6 Step 5:ingest_module_documents

输入: module_exports/*.json

目的: 存储模块级别的元数据,用于检索路由、rerank 排序和降级回退检索。

具体操作:

def step_ingest_module_documents(modules, vector_store):
    for module_ref in modules:
        module_data = fetch_json(resolve_relative(
            base_url, module_ref["module_filename"]
        ))
        document = {
            "id": f"module:{module_ref['module_filename']}",
            "text": module_data.get("abstract", module_data["module_title"]),
            "metadata": {
                "module_filename": module_ref["module_filename"],
                "module_title": module_ref["module_title"],
                "status": module_ref["status"],
                "status_label": module_ref["status_label"],
                "group": module_ref["group"],
                "group_label": module_ref["group_label"],
                "tags": module_ref["tags"],
                "module_public_url": module_ref["module_public_url"]
            }
        }
        vector_store.upsert(
            id=document["id"],
            vector=embedding_model.embed(document["text"]),
            metadata=document["metadata"]
        )

与 chunk 的关系:

维度 Chunk Document Module Document
检索角色 主要(primary) 次要(secondary)
内容粒度 语义段落 / 子章节 模块摘要或标题
用途 直接回答用户问题 路由、rerank、降级回退
检索权重 高 低

3.7 Step 6:attach_boundaries

输入: boundary_clauses + evidence_refs + code_refs

目的: 为每个文档附加边界声明和证据引用,确保检索结果不会过度声称或误导用户。

具体操作:

def step_attach_boundaries(chunks, modules, vector_store):
    for chunk_ref in chunks:
        chunk_data = fetch_json(resolve_relative(
            base_url, chunk_ref["chunk_id"]
        ))
        boundaries = {
            "evidence_refs": chunk_data.get("evidence_refs", []),
            "boundary_clauses": chunk_data.get("boundary_clauses", []),
            "code_refs": chunk_data.get("code_refs", [])
        }
        vector_store.update_metadata(
            id=chunk_ref["chunk_id"],
            metadata=boundaries
        )

    for module_ref in modules:
        module_data = fetch_json(resolve_relative(
            base_url, module_ref["module_filename"]
        ))
        boundaries = {
            "evidence_refs": module_data.get("evidence_refs", []),
            "boundary_clauses": module_data.get("boundary_clauses", []),
            "code_refs": module_data.get("code_refs", [])
        }
        vector_store.update_metadata(
            id=f"module:{module_ref['module_filename']}",
            metadata=boundaries
        )

边界声明的格式:

{
  "boundary_clauses": [
    {
      "type": "status_boundary",
      "description": "该功能目前处于规划阶段,尚未在主网上线",
      "status": "planned"
    }
  ],
  "evidence_refs": [
    {
      "type": "code",
      "url": "https://msgchain.org/whitepaper/proof/consensus.rs",
      "description": "共识模块源代码参考"
    }
  ],
  "code_refs": [
    {
      "repo": "msgchain/msgchain-core",
      "path": "x/consensus/types/round_robin.go",
      "tag": "v1.0.0"
    }
  ]
}

3.8 完整 6 步管线代码框架

import json
from typing import List, Dict
from urllib.parse import urljoin

class RAGIngestPipeline:
    def __init__(self, base_url: str, vector_store, embedding_model):
        self.base_url = base_url
        self.vector_store = vector_store
        self.embedding_model = embedding_model
        self.manifest = None
        self.modules = []
        self.chunks = []

    def resolve(self, path: str) -> str:
        return urljoin(self.base_url + "/", path)

    def fetch_json(self, url: str) -> dict:
        import requests
        resp = requests.get(url)
        resp.raise_for_status()
        return resp.json()

    def run(self):
        # Step 1
        manifest = self.fetch_json(self.resolve("../whitepaper_manifest.json"))
        self.manifest = manifest

        # Step 2
        module_index = self.fetch_json(self.resolve("module_exports/index.json"))
        self.modules = module_index["modules"]

        # Step 3
        chunk_index = self.fetch_json(self.resolve("module_chunks/index.json"))
        self.chunks = chunk_index["chunks"]

        # Step 4
        for c in self.chunks:
            chunk_data = self.fetch_json(self.resolve(c["chunk_id"]))
            self.vector_store.upsert(
                id=c["chunk_id"],
                vector=self.embedding_model.embed(chunk_data["text"]),
                metadata={
                    "chunk_id": c["chunk_id"],
                    "module_filename": c["module_filename"],
                    "chunk_public_url": c["chunk_public_url"],
                    "evidence_refs": chunk_data.get("evidence_refs", []),
                    "boundary_clauses": chunk_data.get("boundary_clauses", [])
                }
            )

        # Step 5
        for m in self.modules:
            module_data = self.fetch_json(self.resolve(m["module_filename"]))
            self.vector_store.upsert(
                id=f"module:{m['module_filename']}",
                vector=self.embedding_model.embed(module_data.get("abstract", m["module_title"])),
                metadata={
                    "module_filename": m["module_filename"],
                    "module_title": m["module_title"],
                    "status": m["status"],
                    "status_label": m["status_label"],
                    "group": m["group"],
                    "group_label": m["group_label"],
                    "tags": m["tags"],
                    "module_public_url": m["module_public_url"]
                }
            )

        # Step 6
        for c in self.chunks:
            chunk_data = self.fetch_json(self.resolve(c["chunk_id"]))
            self.vector_store.update_metadata(
                id=c["chunk_id"],
                metadata={
                    "evidence_refs": chunk_data.get("evidence_refs", []),
                    "boundary_clauses": chunk_data.get("boundary_clauses", [])
                }
            )

4. 元数据策略(推荐字段与索引设计)

4.1 推荐元数据字段

官方 RAG Ingest Flow 定义了以下推荐元数据字段。在向量数据库中,这些字段应作为过滤条件或排序依据。

字段 类型 用途 是否必须
module_filename string 模块文件名,用于模块关联 是
module_title string 模块标题,用于显示和路由 是
status string 模块状态(implemented/partial/planned) 是
status_label string 状态的可读标签 推荐
group string 分组标识 是
group_label string 分组名称 推荐
tags string[] 标签数组 推荐
chunk_id string 分块唯一标识 是(chunk 文档)
module_public_url string 模块公开 URL 是
chunk_public_url string 分块公开 URL 推荐(chunk 文档)
evidence_refs object[] 证据引用列表 推荐
boundary_clauses object[] 边界声明列表 推荐

4.2 索引设计

4.2.1 向量索引

4.2.2 标量索引

为支持高效的元数据过滤,建议为以下字段创建标量索引:

-- Chroma/Pinecone 中通过 metadata filtering 实现
-- Milvus 中通过标量字段索引实现

-- 必须创建的索引
CREATE INDEX idx_status ON msg_whitepaper(status);
CREATE INDEX idx_group ON msg_whitepaper(group);
CREATE INDEX idx_module_filename ON msg_whitepaper(module_filename);

-- 推荐的索引
CREATE INDEX idx_tags ON msg_whitepaper(tags);
CREATE INDEX idx_chunk_id ON msg_whitepaper(chunk_id);

4.2.3 复合过滤场景

场景 1:仅检索已实现模块

vector_store.query(
    vector=query_embedding,
    filter={"status": {"$eq": "implemented"}},
    top_k=10
)

场景 2:按分组检索

vector_store.query(
    vector=query_embedding,
    filter={"group": {"$eq": "consensus"}},
    top_k=10
)

场景 3:按标签筛选后检索

vector_store.query(
    vector=query_embedding,
    filter={
        "$and": [
            {"status": {"$ne": "planned"}},
            {"tags": {"$contains": "dilithium-5"}}
        ]
    },
    top_k=10
)

4.3 元数据演化策略

随着白皮书版本的更新,元数据可能发生变化:

变更类型 处理方式
新增模块 在下次完整摄取时自动发现
模块状态变更 更新对应的 metadata.status 字段
新增标签 更新 metadata.tags 数组
模块废弃 标记 metadata.status = "deprecated",保留文档但降低权重

版本管理:

建议在向量数据库中为每个文档记录 version 和 generated_at 字段,以便追踪版本变化:

document = {
    "metadata": {
        # ... 其他字段
        "version": manifest["schema_version"],
        "generated_at": manifest["generated_at_utc"]
    }
}

4.4 与 Bech32 地址的关联

当文档内容涉及链上地址时(如合约地址、验证者地址),建议统一使用 msg 前缀的 Bech32 格式:

msg1qypqxpq9qcrsszg2pvxq6rs0zqg3yyc5lzv7xu

在元数据中新增 related_addresses 字段以支持按地址检索:

{
  "related_addresses": [
    "msg1qypqxpq9qcrsszg2pvxq6rs0zqg3yyc5lzv7xu"
  ]
}

5. 检索策略(Chunk 优先、Rerank、降级)

5.1 官方检索策略概述

官方 RAG Ingest Flow 明确规定了三层检索策略:

primary:   chunk(分块)
secondary: module(模块)
rerank:    topic hints + status priority + shared tags + group affinity

5.2 第一层:Chunk 优先检索

Chunk 文档是检索的主要数据源。检索时优先返回 chunk 级别的结果:

def primary_retrieval(query: str, top_k: int = 20):
    query_vector = embedding_model.embed(query)
    results = vector_store.query(
        vector=query_vector,
        top_k=top_k,
        # 可选过滤:排除 planned 状态的 chunk
        filter={"status": {"$ne": "planned"}}
    )
    return results

为什么 chunk 优先?

  1. 粒度更细:chunk 是语义段落,比整个模块更精确
  2. 边界更清晰:每个 chunk 有自己的边界声明和证据引用
  3. 答案更准确:基于 chunk 的回答可以精确引用原文

5.3 第二层:Rerank 排序

官方定义了四个 rerank 维度:

维度 排序规则 权重
Topic Hints 与 retrieval_hints.json 中的主题提示匹配 高
Status Priority implemented > partial > planned 高
Shared Tags 查询中标签与文档标签的匹配数 中
Group Affinity 同组文档优先 低

Rerank 实现示例:

def rerank_results(query: str, results: List[Dict], hints: Dict = None) -> List[Dict]:
    def compute_score(doc):
        score = 0.0

        # 1. Topic Hints 匹配
        if hints:
            for topic in hints.get("topics", []):
                if topic["keyword"] in doc["text"]:
                    score += topic.get("boost", 0.5)

        # 2. Status Priority
        status_weights = {"implemented": 3.0, "partial": 2.0, "planned": 0.5}
        score += status_weights.get(doc["metadata"].get("status"), 1.0)

        # 3. Shared Tags
        query_tags = extract_tags_from_query(query)
        shared = set(query_tags) & set(doc["metadata"].get("tags", []))
        score += len(shared) * 1.5

        # 4. Group Affinity
        # 如果同一个 group 有多个结果,提升该 group 的分数
        score += 0.1  # 基础 group affinity 加成

        return score

    for doc in results:
        doc["_rerank_score"] = compute_score(doc)

    results.sort(key=lambda x: x["_rerank_score"], reverse=True)
    return results

5.4 第三层:降级回退

当 chunk 级别检索结果不理想时,降级到 module 级别的检索:

def fallback_retrieval(query: str, top_k: int = 5):
    # 降级到 module 级别
    query_vector = embedding_model.embed(query)
    results = vector_store.query(
        vector=query_vector,
        top_k=top_k,
        filter={"document_type": "module"}
    )
    return results

降级触发条件:

条件 说明
无 chunk 结果 向量检索返回空
所有结果相似度低于阈值 最高相似度 < 0.7
结果均不相关 经 rerank 后前 3 个结果仍然不匹配
查询明确针对模块级别 如 "consensus 模块包含什么"

5.5 完整检索流程

def search(query: str) -> Dict:
    # 第一阶段:chunk 优先检索
    chunk_results = primary_retrieval(query, top_k=20)

    if needs_fallback(query, chunk_results):
        # 降级到 module 检索
        module_results = fallback_retrieval(query, top_k=5)
        return format_response(module_results, source="module")
    else:
        # Rerank
        hints = load_retrieval_hints()
        reranked = rerank_results(query, chunk_results, hints)
        return format_response(reranked[:10], source="chunk")


def needs_fallback(query, results) -> bool:
    if not results:
        return True
    if results[0]["score"] < 0.7:
        return True
    return False

5.6 多查询融合

对于复杂问题,可以采用多查询融合策略:

def multi_query_search(queries: List[str], top_k: int = 10) -> List[Dict]:
    from collections import defaultdict

    scores = defaultdict(float)
    for q in queries:
        results = primary_retrieval(q, top_k=10)
        for r in results:
            scores[r["id"]] += r["score"]

    # 按总分排序
    ranked = sorted(scores.items(), key=lambda x: x[1], reverse=True)
    return [{"id": doc_id, "score": total} for doc_id, total in ranked[:top_k]]

5.7 检索结果格式化

def format_response(results, source="chunk"):
    formatted = []
    for r in results:
        entry = {
            "source": source,
            "confidence": "high" if r["score"] > 0.85 else "medium" if r["score"] > 0.7 else "low",
            "content": r["text"],
            "metadata": {
                "module_title": r["metadata"].get("module_title"),
                "status": r["metadata"].get("status"),
                "module_public_url": r["metadata"].get("module_public_url")
            }
        }
        if r["metadata"].get("boundary_clauses"):
            entry["boundaries"] = r["metadata"]["boundary_clauses"]
        if r["metadata"].get("evidence_refs"):
            entry["evidence"] = r["metadata"]["evidence_refs"]
        formatted.append(entry)
    return {
        "results": formatted,
        "total": len(formatted),
        "strategy": source
    }

6. 边界声明与服务策略

6.1 边界声明的必要性

在白皮书知识库中,不同模块处于不同的实现阶段。AI Agent 在检索和回答时,必须明确声明:

6.2 官方服务策略

官方 RAG Ingest Flow 定义了四条服务策略:

1. 返回 chunk 级别的结果以获得精确回答。
2. 如果多个 chunk 之间存在冲突,选择证据引用更强且边界声明更清晰的结果。
3. 引导用户返回 module_public_url 以获取完整上下文。
4. 避免将 planned 模块作为实现问题的首要证据。

6.3 策略 1:返回 Chunk 级别结果

def serve_answer(query: str) -> str:
    results = search(query)
    if not results["results"]:
        return "抱歉,我未能在白皮书中找到相关信息。"

    # 取 rerank 后的第一个结果
    top = results["results"][0]
    answer = top["content"]

    # 附加上下文
    if top.get("boundaries"):
        answer += f"\n\n> 边界声明:{top['boundaries']}"

    answer += f"\n\n[查看完整模块]({top['metadata']['module_public_url']})"
    return answer

6.4 策略 2:冲突解决

当多个 chunk 之间的内容冲突时:

def resolve_conflicts(conflicting_results: List[Dict]) -> Dict:
    def conflict_score(doc):
        score = 0
        # 更强的 evidence_refs
        score += len(doc["metadata"].get("evidence_refs", [])) * 2
        # 更清晰的 boundary_clauses
        clauses = doc["metadata"].get("boundary_clauses", [])
        for c in clauses:
            if c.get("type") == "status_boundary" and c.get("status") == "implemented":
                score += 5
        # 更高状态优先级
        status_weight = {"implemented": 10, "partial": 5, "planned": 1}
        score += status_weight.get(doc["metadata"].get("status"), 0)
        return score

    return max(conflicting_results, key=conflict_score)

6.5 策略 3:引导到完整上下文

def generate_llm_prompt(query: str, results: List[Dict]) -> str:
    context = ""
    for i, r in enumerate(results):
        context += f"[片段 {i + 1}]\n"
        context += f"内容:{r['content']}\n"
        context += f"状态:{r['metadata']['status_label']}\n"
        context += f"来源:{r['metadata']['module_public_url']}\n"
        if r.get("boundaries"):
            context += f"边界:{json.dumps(r['boundaries'], ensure_ascii=False)}\n"
        context += "\n"

    prompt = f"""基于以下白皮书知识片段回答用户的问题。

{context}

用户问题:{query}

请回答。在回答中请引用来源 URL,并在必要时说明当前功能的状态(已实现/部分实现/计划中)。
如果问题涉及计划中的功能,请明确说明该功能尚未上线。
"""
    return prompt

6.6 策略 4:避免 Planned 模块作为首要证据

def filter_primary_evidence(results: List[Dict]) -> List[Dict]:
    primary = [r for r in results if r["metadata"].get("status") != "planned"]
    secondary = [r for r in results if r["metadata"].get("status") == "planned"]

    if not primary:
        # 只有 planned 结果时仍返回,但标记为低置信度
        for r in secondary:
            r["warning"] = "该功能目前处于规划阶段,尚未上线"
        return secondary

    return primary

6.7 边界声明在向量库中的存储

BOUNDARY_DOCUMENT_TEMPLATE = {
    "id": "boundary:{module_filename}",
    "text": "",  # 不包含文本,仅用于元数据检索
    "metadata": {
        "document_type": "boundary",
        "module_filename": "...",  # 替换为实际值
        "module_title": "...",
        "boundary_clauses": [
            {
                "type": "status_boundary",
                "description": "该模块目前处于规划阶段",
                "status": "planned",
                "effective_from": "2025-01-01"
            }
        ],
        "evidence_refs": [
            {
                "type": "code",
                "url": "https://msgchain.org/whitepaper/proof/consensus.rs",
                "description": "共识代码参考实现"
            }
        ]
    }
}

7. 向量库选择与部署实践(Chroma/Milvus/Pinecone 集成)

7.1 向量库选型对比

特性 Chroma Milvus Pinecone
部署方式 嵌入式 / 客户端 独立服务 云托管
开源 是(Apache 2.0) 是(Apache 2.0) 否
标量过滤 支持 丰富 支持
元数据索引 有限 强大 良好
分布式 否 是 是
管理成本 低 中 高
适用场景 开发 / 小规模 生产 / 大规模 快速原型 / 中规模

7.2 Chroma 集成(推荐开发环境)

Chroma 是最轻量的向量数据库,适合开发阶段的 RAG 原型。

安装:

pip install chromadb

集成代码:

import chromadb
from chromadb.config import Settings

class ChromaVectorStore:
    def __init__(self, persist_directory="./chroma_db"):
        self.client = chromadb.Client(Settings(
            chroma_db_impl="duckdb+parquet",
            persist_directory=persist_directory
        ))
        self.collection = self.client.get_or_create_collection(
            name="msg_whitepaper",
            metadata={"hnsw:space": "cosine"}
        )

    def upsert(self, id, vector, metadata):
        self.collection.upsert(
            ids=[id],
            embeddings=[vector],
            metadatas=[metadata]
        )

    def update_metadata(self, id, metadata):
        self.collection.update(
            ids=[id],
            metadatas=[metadata]
        )

    def query(self, vector, top_k=10, filter=None):
        results = self.collection.query(
            query_embeddings=[vector],
            n_results=top_k,
            where=filter
        )
        return [
            {
                "id": results["ids"][0][i],
                "score": results["distances"][0][i],
                "metadata": results["metadatas"][0][i],
                "text": results["documents"][0][i] if results.get("documents") else ""
            }
            for i in range(len(results["ids"][0]))
        ]

7.3 Milvus 集成(推荐生产环境)

Milvus 适合大规模生产部署,支持分布式和丰富的标量索引。

安装:

pip install pymilvus

集成代码:

from pymilvus import (
    connections, Collection, FieldSchema, CollectionSchema,
    DataType, utility
)

class MilvusVectorStore:
    def __init__(self, host="localhost", port="19530"):
        connections.connect(host=host, port=port)
        self.collection_name = "msg_whitepaper"

        if utility.has_collection(self.collection_name):
            self.collection = Collection(self.collection_name)
        else:
            self._create_collection()

    def _create_collection(self):
        fields = [
            FieldSchema(name="id", dtype=DataType.VARCHAR, is_primary=True, max_length=255),
            FieldSchema(name="vector", dtype=DataType.FLOAT_VECTOR, dim=1536),
            FieldSchema(name="module_filename", dtype=DataType.VARCHAR, max_length=255),
            FieldSchema(name="module_title", dtype=DataType.VARCHAR, max_length=255),
            FieldSchema(name="status", dtype=DataType.VARCHAR, max_length=50),
            FieldSchema(name="status_label", dtype=DataType.VARCHAR, max_length=100),
            FieldSchema(name="group", dtype=DataType.VARCHAR, max_length=100),
            FieldSchema(name="group_label", dtype=DataType.VARCHAR, max_length=100),
            FieldSchema(name="tags", dtype=DataType.JSON),
            FieldSchema(name="module_public_url", dtype=DataType.VARCHAR, max_length=512),
            FieldSchema(name="chunk_public_url", dtype=DataType.VARCHAR, max_length=512),
            FieldSchema(name="evidence_refs", dtype=DataType.JSON),
            FieldSchema(name="boundary_clauses", dtype=DataType.JSON)
        ]
        schema = CollectionSchema(fields, description="MSG Chain Whitepaper RAG")
        self.collection = Collection(self.collection_name, schema)

        # 创建 IVF_FLAT 索引
        index_params = {
            "metric_type": "COSINE",
            "index_type": "IVF_FLAT",
            "params": {"nlist": 128}
        }
        self.collection.create_index("vector", index_params)

        # 创建标量索引
        self.collection.create_index("status", {"index_type": "INVERTED"})
        self.collection.create_index("group", {"index_type": "INVERTED"})

    def upsert(self, id, vector, metadata):
        entity = {
            "id": id,
            "vector": vector,
            **{k: v for k, v in metadata.items() if k in [
                "module_filename", "module_title", "status", "status_label",
                "group", "group_label", "tags", "module_public_url",
                "chunk_public_url", "evidence_refs", "boundary_clauses"
            ]}
        }
        self.collection.upsert([entity])

    def query(self, vector, top_k=10, filter=None):
        self.collection.load()
        expr = self._build_filter_expr(filter)
        results = self.collection.search(
            data=[vector],
            anns_field="vector",
            param={"metric_type": "COSINE", "params": {"nprobe": 10}},
            limit=top_k,
            expr=expr,
            output_fields=["module_title", "status", "module_public_url", "chunk_public_url"]
        )
        return [
            {
                "id": r.id,
                "score": r.score,
                "metadata": {k: r.entity.get(k) for k in r.entity.fields}
            }
            for r in results[0]
        ]

    def _build_filter_expr(self, filter):
        if not filter:
            return None
        # 将通用 filter 转换为 Milvus 表达式
        exprs = []
        for key, value in filter.items():
            if isinstance(value, dict):
                op = next(iter(value))
                val = value[op]
                if op == "$eq":
                    exprs.append(f"{key} == '{val}'")
                elif op == "$ne":
                    exprs.append(f"{key} != '{val}'")
            elif isinstance(value, str):
                exprs.append(f"{key} == '{value}'")
        return " && ".join(exprs)

7.4 Pinecone 集成(推荐快速原型)

Pinecone 是全托管的向量数据库,无需运维。

安装:

pip install pinecone-client

集成代码:

import pinecone

class PineconeVectorStore:
    def __init__(self, api_key, environment="us-west1-gcp"):
        pinecone.init(api_key=api_key, environment=environment)

        if "msg-whitepaper" not in pinecone.list_indexes():
            pinecone.create_index(
                name="msg-whitepaper",
                dimension=1536,
                metric="cosine",
                metadata_config={"indexed": [
                    "status", "group", "module_filename", "tags"
                ]}
            )

        self.index = pinecone.Index("msg-whitepaper")

    def upsert(self, id, vector, metadata):
        self.index.upsert(
            vectors=[(id, vector, metadata)]
        )

    def update_metadata(self, id, metadata):
        self.index.update(
            id=id,
            set_metadata=metadata
        )

    def query(self, vector, top_k=10, filter=None):
        results = self.index.query(
            vector=vector,
            top_k=top_k,
            filter=filter,
            include_metadata=True
        )
        return [
            {
                "id": r.id,
                "score": r.score,
                "metadata": r.metadata
            }
            for r in results.matches
        ]

7.5 Embedding 模型选择

模型 维度 适用场景
text-embedding-3-small 1536 通用场景,性价比高
text-embedding-3-large 3072 高精度场景
BAAI/bge-large-zh-v1.5 1024 中文优化,适合中文白皮书
shibing624/text2vec-base-chinese 768 轻量中文模型

中文白皮书的 embedding 建议:

由于 MSG Chain 白皮书同时支持中英文,建议:

7.6 部署架构

                    ┌─────────────┐
                    │   AI Agent  │
                    │    (LLM)    │
                    └──────┬──────┘
                           │
                    ┌──────▼──────┐
                    │   Retriever │
                    │  (Rerank)   │
                    └──────┬──────┘
                           │
              ┌────────────┼────────────┐
              │            │            │
        ┌─────▼────┐ ┌────▼────┐ ┌────▼────┐
        │ Vector   │ │Metadata │ │  Cache  │
        │  DB      │ │ Filter  │ │ (Redis) │
        └─────┬────┘ └─────────┘ └─────────┘
              │
     ┌────────┼─────────┐
     │        │         │
  Chroma  Milvus   Pinecone

8. 增量更新策略

8.1 更新触发条件

白皮书内容的更新可能来自以下原因:

触发条件 变更范围 urgency
白皮书版本发布 可能涉及多个模块 高
模块状态变更 单个模块的 implemented / partial / planned 状态变化 中
内容勘误 单个 chunk 的文本修正 低
新增模块 新模块加入,对应 module 和 chunk 中
废弃模块 旧模块标记废弃 低

8.2 全量更新 vs 增量更新

全量更新:

def full_refresh():
    pipeline = RAGIngestPipeline(
        base_url="https://msgchain.org/whitepaper/integration_examples/",
        vector_store=vector_store,
        embedding_model=embedding_model
    )

    # 清空已有数据
    vector_store.reset()

    # 完整运行 6 步流程
    pipeline.run()

适用场景:首次部署、重大版本变更、数据损坏恢复。

增量更新:

def incremental_update():
    # Step 1: 检查当前版本
    manifest = fetch_json("https://msgchain.org/whitepaper/integration_examples/whitepaper_manifest.json")
    current_version = get_stored_version()

    if manifest["schema_version"] != current_version:
        # 版本变更,需要全量更新
        full_refresh()
        return

    # Step 2: 比较模块变更
    old_modules = get_stored_module_index()
    new_modules = fetch_json("https://msgchain.org/whitepaper/integration_examples/module_exports/index.json")

    changes = diff_modules(old_modules, new_modules["modules"])
    for change in changes:
        if change["type"] == "modified":
            update_module(change["module"])
        elif change["type"] == "added":
            add_module(change["module"])
        elif change["type"] == "removed":
            remove_module(change["module_filename"])

    # Step 3: 比较 chunk 变更
    old_chunks = get_stored_chunk_index()
    new_chunks = fetch_json("https://msgchain.org/whitepaper/integration_examples/module_chunks/index.json")

    changes = diff_chunks(old_chunks, new_chunks["chunks"])
    for change in changes:
        if change["type"] == "modified":
            update_chunk(change["chunk"])
        elif change["type"] == "added":
            add_chunk(change["chunk"])
        elif change["type"] == "removed":
            remove_chunk(change["chunk_id"])

8.3 变更检测算法

def diff_modules(old_list, new_list):
    old_map = {m["module_filename"]: m for m in old_list}
    new_map = {m["module_filename"]: m for m in new_list}

    changes = []

    # 检测修改和添加
    for filename, new_mod in new_map.items():
        if filename in old_map:
            old_mod = old_map[filename]
            if old_mod != new_mod:
                changes.append({
                    "type": "modified",
                    "module": new_mod
                })
        else:
            changes.append({
                "type": "added",
                "module": new_mod
            })

    # 检测删除
    for filename in old_map:
        if filename not in new_map:
            changes.append({
                "type": "removed",
                "module_filename": filename
            })

    return changes

8.4 版本状态追踪

class VersionTracker:
    def __init__(self, state_file="./rag_state.json"):
        self.state_file = state_file
        self.state = self._load_state()

    def _load_state(self):
        try:
            with open(self.state_file) as f:
                return json.load(f)
        except FileNotFoundError:
            return {
                "version": None,
                "last_sync": None,
                "module_hashes": {},
                "chunk_hashes": {}
            }

    def _save_state(self):
        with open(self.state_file, "w") as f:
            json.dump(self.state, f, indent=2)

    def should_update(self, manifest):
        if self.state["version"] != manifest["schema_version"]:
            return True
        # 检查生成时间
        if self.state["last_sync"] < manifest["generated_at_utc"]:
            return True
        return False

    def record_sync(self, manifest, modules, chunks):
        self.state["version"] = manifest["schema_version"]
        self.state["last_sync"] = manifest["generated_at_utc"]

        for m in modules:
            self.state["module_hashes"][m["module_filename"]] = hash_dict(m)

        for c in chunks:
            self.state["chunk_hashes"][c["chunk_id"]] = hash_dict(c)

        self._save_state()

8.5 定时同步策略

import schedule
import time

def sync_job():
    print(f"[{datetime.utcnow()}] 开始 RAG 知识库同步...")
    manifest = fetch_manifest()

    tracker = VersionTracker()
    if tracker.should_update(manifest):
        incremental_update()
        tracker.record_sync(manifest, modules, chunks)
        print("同步完成")
    else:
        print("无需更新")

# 每小时检查一次
schedule.every().hour.do(sync_job)

while True:
    schedule.run_pending()
    time.sleep(60)

8.6 变更通知(Webhook)

对于需要实时响应的场景,MSG Chain 提供 webhook 通知机制:

from flask import Flask, request

app = Flask(__name__)

@app.route("/webhook/whitepaper-update", methods=["POST"])
def handle_update():
    payload = request.json
    event_type = payload["event"]  # "manifest_update" / "module_update" / "chunk_update"

    if event_type == "manifest_update":
        full_refresh()
    elif event_type == "module_update":
        update_module(payload["module_filename"])
    elif event_type == "chunk_update":
        update_chunk(payload["chunk_id"])

    return {"status": "ok"}, 200

if __name__ == "__main__":
    app.run(port=8080)

9. 示例:构建一个 MSG Chain RAG 知识库

9.1 项目结构

msg-rag-knowledge-base/
├── config/
│   ├── settings.yaml        # 配置文件
│   └── manifest.json        # 本地缓存的 manifest
├── ingest/
│   ├── __init__.py
│   ├── pipeline.py          # 6 步摄取管线
│   ├── fetcher.py           # 数据获取模块
│   ├── embedder.py          # Embedding 模块
│   └── store.py             # 向量库抽象
├── retrieval/
│   ├── __init__.py
│   ├── searcher.py          # 检索入口
│   ├── reranker.py          # Rerank 实现
│   └── formatter.py         # 结果格式化
├── api/
│   ├── __init__.py
│   └── server.py            # REST API 服务
├── scripts/
│   ├── full_refresh.py      # 全量更新脚本
│   └── watch_sync.py        # 定时同步脚本
├── tests/
│   ├── test_pipeline.py
│   ├── test_retrieval.py
│   └── test_reranker.py
├── requirements.txt
└── README.md

9.2 配置示例(settings.yaml)

# MSG Chain RAG 知识库配置
# 网络: msg-chain-1 | 地址前缀: msg | 域名: msgchain.org

whitepaper:
  base_url: "https://msgchain.org/whitepaper/integration_examples/"
  schema_version: "v1"
  entry_points:
    agent_entry: "../agent_entry.json"
    manifest: "../whitepaper_manifest.json"
    module_index: "module_exports/index.json"
    chunk_index: "module_chunks/index.json"

vector_store:
  provider: "chroma"  # chroma | milvus | pinecone
  persist_directory: "./data/chroma_db"
  collection_name: "msg_whitepaper"
  distance_metric: "cosine"

embedding:
  model: "text-embedding-3-small"  # 也支持 BAAI/bge-large-zh-v1.5
  dimension: 1536
  batch_size: 64

retrieval:
  default_top_k: 10
  rerank_enabled: true
  fallback_threshold: 0.7
  status_weights:
    implemented: 3.0
    partial: 2.0
    planned: 0.5

sync:
  schedule_interval_hours: 1
  state_file: "./data/rag_state.json"
  webhook_port: 8080

9.3 完整示例:从零构建

步骤 1:安装依赖

pip install chromadb openai requests schedule pyyaml

步骤 2:初始化向量库

# scripts/init_db.py
import chromadb
from chromadb.config import Settings

client = chromadb.Client(Settings(
    chroma_db_impl="duckdb+parquet",
    persist_directory="./data/chroma_db"
))

collection = client.get_or_create_collection(
    name="msg_whitepaper",
    metadata={"hnsw:space": "cosine"}
)

print(f"向量库已初始化:{collection.name}")

步骤 3:运行全量摄取

# scripts/full_refresh.py
from ingest.pipeline import RAGIngestPipeline
from ingest.store import ChromaVectorStore
from ingest.embedder import OpenAIEmbedder

store = ChromaVectorStore(persist_directory="./data/chroma_db")
embedder = OpenAIEmbedder(model="text-embedding-3-small")

pipeline = RAGIngestPipeline(
    base_url="https://msgchain.org/whitepaper/integration_examples/",
    vector_store=store,
    embedding_model=embedder
)

pipeline.run()
print("全量摄取完成")

步骤 4:启动检索 API

# api/server.py
from flask import Flask, request, jsonify
from retrieval.searcher import RAGSearcher

app = Flask(__name__)
searcher = RAGSearcher(
    store=ChromaVectorStore(persist_directory="./data/chroma_db"),
    embedder=OpenAIEmbedder(model="text-embedding-3-small")
)

@app.route("/api/v1/search", methods=["POST"])
def search():
    data = request.json
    query = data["query"]
    top_k = data.get("top_k", 10)

    results = searcher.search(query, top_k=top_k)
    return jsonify(results)

@app.route("/api/v1/health", methods=["GET"])
def health():
    return jsonify({
        "status": "ok",
        "chain": "msg-chain-1",
        "bech32_prefix": "msg",
        "whitepaper_source": "https://msgchain.org/whitepaper/integration_examples/"
    })

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=8080)

步骤 5:测试检索

curl -X POST http://localhost:8080/api/v1/search \
  -H "Content-Type: application/json" \
  -d '{
    "query": "Dilithium-5 签名如何工作",
    "top_k": 5
  }'

预期响应:

{
  "results": [
    {
      "source": "chunk",
      "confidence": "high",
      "content": "Dilithium-5 是 MSG Chain 采用的后量子密码签名方案...",
      "metadata": {
        "module_title": "Quantum-Safe Core",
        "status": "implemented",
        "module_public_url": "https://msgchain.org/whitepaper/quantum-safe-core"
      },
      "evidence": [
        {
          "type": "code",
          "url": "https://msgchain.org/whitepaper/proof/dilithium5.rs",
          "description": "Dilithium-5 实现参考"
        }
      ]
    }
  ],
  "total": 1,
  "strategy": "chunk"
}

9.4 测试用例

# tests/test_pipeline.py
import pytest
from ingest.pipeline import RAGIngestPipeline

class MockVectorStore:
    def __init__(self):
        self.documents = {}

    def upsert(self, id, vector, metadata):
        self.documents[id] = {"vector": vector, "metadata": metadata}

    def update_metadata(self, id, metadata):
        if id in self.documents:
            self.documents[id]["metadata"].update(metadata)

    def query(self, vector, top_k=10, filter=None):
        return []

class MockEmbedder:
    def embed(self, text):
        return [0.0] * 1536

def test_pipeline_complete():
    store = MockVectorStore()
    embedder = MockEmbedder()

    pipeline = RAGIngestPipeline(
        base_url="https://msgchain.org/whitepaper/integration_examples/",
        vector_store=store,
        embedding_model=embedder
    )

    pipeline.run()
    assert len(store.documents) > 0
    for doc_id, doc in store.documents.items():
        assert "metadata" in doc
        assert "vector" in doc
# tests/test_reranker.py
import pytest
from retrieval.reranker import rerank_results

def test_status_priority():
    query = "consensus 实现"
    results = [
        {"id": "1", "score": 0.8, "metadata": {"status": "planned"}, "text": "..."},
        {"id": "2", "score": 0.7, "metadata": {"status": "implemented"}, "text": "..."},
    ]
    hints = {"topics": [{"keyword": "consensus", "boost": 0.5}]}

    reranked = rerank_results(query, results, hints)
    assert reranked[0]["id"] == "2"  # implemented 应该排第一

9.5 检索 API 性能测试

# tests/test_performance.py
import time
import statistics

def benchmark_search(searcher, queries, iterations=10):
    latencies = []
    for query in queries:
        for _ in range(iterations):
            start = time.time()
            searcher.search(query, top_k=10)
            elapsed = time.time() - start
            latencies.append(elapsed)

    return {
        "avg_latency_ms": statistics.mean(latencies) * 1000,
        "p50_ms": statistics.median(latencies) * 1000,
        "p95_ms": sorted(latencies)[int(len(latencies) * 0.95)] * 1000,
        "total_queries": len(latencies)
    }

10. 总结

10.1 核心要点回顾

主题 要点
分层入口 agent_entry.json → manifest → module_index → chunk_index
6 步流程 load_manifest → load_module_index → load_chunk_index → ingest_chunks → ingest_modules → attach_boundaries
元数据策略 12 个推荐字段,status/tags/group 用于过滤和排序
检索策略 chunk 优先 → rerank(4 维度)→ module 降级
服务策略 chunk 结果 + 边界声明 + 引导 URL + 避免 planned 作为主要证据
向量库 Chroma(开发)/ Milvus(生产)/ Pinecone(原型)
增量更新 版本 + hash 比对,定时同步或 webhook 通知

10.2 实现检查清单

首次部署 RAG 知识库时,请按以下清单逐一确认:

10.3 常见问题

Q: 如果 agent_entry.json 访问失败怎么办?

A: 可以直接从 whitepaper_manifest.json 开始入口。入口文件有缓存策略,建议在 Agent 端实现重试和降级逻辑。

Q: 如何确认自己的索引是最新版本?

A: 检查 whitepaper_manifest.json 中的 generated_at_utc 字段,与本地最后一次同步时间比对。

Q: Chunk 和 Module 文档应该存储在同一个 collection 中吗?

A: 建议存储在同一个 collection 中,通过 document_type 字段区分(chunk / module)。这样可以实现统一的检索入口。

Q: 中文和白皮书英文内容如何统一处理?

A: 建议使用多语言 embedding 模型(如 intfloat/multilingual-e5-large)。中文 chunk 和英文 chunk 可以直接在同一向量空间中检索。

Q: 如何处理 planned 模块?

A: 检索时通过 metadata 过滤 status != "planned",但如果只有 planned 结果可用,仍然返回并标记低置信度,同时添加明确的计划状态声明。

10.4 参考资源


本文档基于 MSG Chain 白皮书机器层 RAG Ingest Flow 官方规范编写,schema 版本 v1,元数据配置 public_stable。所有技术细节以 https://msgchain.org/whitepaper/integration_examples/rag_ingest_flow.json 的官方定义为准。