本节摘要:解析器只管「把万物变文本」,从文本到
viking://空间里的分层知识库,中间隔着一条完整的编排流水线。本节走两趟:第一趟看会话日志入库——ingest/orchestrator.py的回填编排与poller.py的增量轮询,靠「持久化读游标」做到崩溃自愈;第二趟看资源入库——ResourceService的异步任务化(解析→分块→嵌入→写入→生成分层 sidecar),WatchScheduler定时轮询源目录自动重灌,scrapy+trafilatura 递归爬取网页。第 3 章讲的 L0/L1 sidecar 就在这条流水线的语义 DAG 里落地:自底向上,先文件摘要、后目录摘要,一层层向上长。
内容来源:原项目源码 openviking/ingest/(orchestrator.py、poller.py、replay.py、sources/)、openviking/parse/accessors/(web_crawler/、web_importer.py)、openviking/resource/(watch_scheduler.py、watch_manager.py)、openviking/service/resource_service.py、openviking/storage/queuefs/(semantic_dag.py、semantic_sidecar.py)
⚠️ 注意:目录里有两个「编排」别混淆——
openviking/ingest/编排的是会话日志回放(把 Claude Code 等工具的本地会话记录灌进 OpenViking),openviking/resource/与service/resource_service.py编排的才是资源入库(文档/仓库/网页)。两者最终都汇入同一套写入队列与语义 sidecar 生成管线。
ingest/orchestrator.py 的开篇 docstring 定了位:「回填编排:把发现的每个会话从 cursor 回放到 end,然后提交;增量 watch 模式由 poller.py 负责」。主干 backfill_source()(57-102 行)是一个朴素而结实的循环:
refs = list(source.discover_sessions()) # 1. 发现会话 for ref in refs: if since and ref.started_at and ref.started_at < since: stats.skipped += 1 continue if not dry_run: if reset: await self.replayer.reset_session(name, ref) else: await self.replayer.reconcile(name, ref) added = await self._backfill_one(name, source, ref, dry_run=dry_run) stats.sessions += 1 stats.messages += added if not dry_run: if await self.replayer.commit_if_needed( name, ref, harness_cfg.commit.keep_recent_count ): stats.committed += 1
三步曲:发现(discover_sessions 扫描源工具的会话记录)→ 回放(_backfill_one 从存储的游标处循环 read_messages 批量搬运)→ 提交(commit_if_needed 触发第 5 章的记忆提取)。有一行注释值得单独摘出来:「added == 0 也要 commit——上次崩溃可能留下已追加未提交的消息」。回填的敌人从来不是慢,是重复:同一条消息灌两遍,记忆提取就会产出双份事实。dry_run 模式只数不写,reset 模式推倒重来,BackfillStats(sessions/messages/committed/skipped/errors) 把五本账记全,失败按会话隔离——一个坏会话不拖垮整批。
数据源在 sources/ 下:claude_code.py 读 ~/.claude/projects/<slug>/<uuid>.jsonl(逐行 JSON,user/assistant 轮次、cwd、gitBranch 都在记录里),同类还有 codex/cursor/hermes/openclaw/opencode。每个源是一个 LogSource 适配器,把各家格式归一成 NormalizedMessage——又是「先归一」的套路。claude_code.py 的文件头注释就是一份微型格式文档:每条记录有顶层 type,会话轮次是 type in {user, assistant},content 既可能是字符串也可能是块列表(text/tool_use/tool_result),_extract_text 只取 text 块——工具调用的中间产物不进记忆,对话文本才是沉淀的原料。
回放的核心 replay.py 有一段教科书式的 docstring:
"""SessionReplayer drives the canonical recipe in idempotent batches: reconcile() -> for each <=100-message batch: set_pending -> append -> confirm -> commit_if_needed Each batch's intent (target cursor + size + the server's message count *before* the append) is persisted BEFORE the append. If the process crashes mid-append, the next run ``reconcile()`` compares the server's current message count against that baseline to decide whether the batch landed — confirming it (no re-append) or dropping the intent (re-read from the confirmed cursor). """
配方拆开看:批次上限 _BATCH = 100(服务端批量接口的硬顶);每批先落盘意图(目标游标 + 条数 + 追加前服务端的消息数基线),再真正 append。若进程在 append 中途暴毙,下次 reconcile() 拿服务端当前消息数与基线比对:对得上就 confirm(不重灌),对不上就丢弃意图、从已确认游标重读。游标只在追加被确认后才前进——at-least-once 与 at-most-once 之间的窄门,靠「意图先行 + 基线比对」走过去。
增量侧的 poller.py 是一个 WatchScheduler 风格的 asyncio 轮询环:
class IngestPoller: def __init__(self, sources, replayer): ... self._dirty: Dict[Tuple[str, str], _Dirty] = {} self.poll_interval = min( (cfg.poll_interval_seconds for _, cfg, _ in sources), default=5.0) async def run(self) -> None: while self._running: await self._tick() await self._commit_idle() ...
不依赖文件系统事件,每 tick 重扫会话、从游标增量读取、追加、空闲或 token 阈值到点就提交。docstring 点破了选型理由:「持久化读游标让轮询天然正确且自愈——错过的 tick、睡眠、重启,下次从 cursor 读到 EOF 即可」。事件驱动快但易丢,轮询慢但稳,会话日志这种「追加为主的低频写入」场景,稳字当先。
资源侧的编排主角是 service/resource_service.py(两千余行)。对外入口 add_resource(),内部先构造 ContentTargetSpec(决定落进 viking://resources 还是用户空间),再走 _submit_resource_ingestion() → _execute_resource_ingestion(),后者调 process_resource(..., defer_post_processing=True)——解析同步做、后处理异步化,这是全链路的基调:
result = await self._resource_processor.process_resource( path=path, ctx=ctx, scope="resources", to=target.to, parent=target.parent, build_index=build_index, summarize=summarize, defer_post_processing=True, ...) ... prepared = result.pop("_post_process", None) deferred_lock = result.pop("_resource_lock", None)
defer_post_processing=True 意味着:文件树先落盘(persist_temp_tree,第 2 章的临时目录转正),随后的嵌入与语义处理排队异步执行,_monitor_queue_processing 盯队列状态、wait_processed(timeout) 供调用方按需等待,ov add 命令的「处理中→就绪」状态机就建在这上面。队列本体在 storage/queuefs/(第 7 章细讲):嵌入队列按 chunk 批量算向量,语义队列跑 sidecar 生成,各有并发上限。入库时自动生成 L0/L1 sidecar 的第 3 章分层,正是在这里落地——storage/semantic_sidecar.py 的 write_semantic_sidecars() 与 queuefs/semantic_dag.py 接力:
@dataclass class DirNode: """Directory node state for DAG execution.""" uri: str children_dirs: List[str] file_paths: List[str] file_summaries: List[Optional[Dict[str, str]]] children_abstracts: List[Optional[Dict[str, str]]] pending: int dispatched: bool = False overview_scheduled: bool = False
生成顺序是自底向上的事件驱动 DAG:叶子文件先出摘要,目录节点的 pending 计数清零后调度自己的 .abstract.md(汇拢子项);子目录摘要齐了再排 .overview.md。一个节点完成即唤醒父节点——「懒分派」让成千上万个目录的摘要任务在有限并发下自然排产,而不是全量铺开。messages.jsonl 这类会话归档文件显式跳过(「总结它只浪费 token」),sidecar 的 freshness 元数据(第 3 章 OKF 头)也在这里写。
「入库一次」之外,资源还有「持续同步」的形态。resource/watch_scheduler.py 的 WatchScheduler 每 60 秒(可配)检查到期任务,经 asyncio.Semaphore(max_concurrency=4) 限并发后调 ResourceService 重灌:
class WatchScheduler: DEFAULT_CHECK_INTERVAL = 60.0 def __init__(self, resource_service, viking_fs=None, check_interval=60.0, max_concurrency=4, ...):
任务本体由 WatchManager 持久化在 viking_fs 里:target URI、interval、认证状态一应俱全;飞书源带 OAuth token 自动刷新(feishu_watch_auth),Git 源带 HTTP 认证;URI 前缀改写有专门的 UriMutationCoordinator 协调存量任务。配合上一节的 FeishuAccessor,「每天早上自动同步这个飞书知识库」就是一个 watch 任务。
网页侧,accessors/web_crawler/ 是一个完整的 scrapy 工程:OpenVikingWebSpider(scrapy_spider.py)从根 URL 出发递归爬取,跳过 .css/.js/.png 等资产扩展名;遇到 PDF/MD/TXT/DOC 等下载类链接单独收集成 CrawledDownload 直接入库;正文抽取则交给 HTML 解析器里的 trafilatura(html.py 159-181 行:trafilatura.extract(...) 抽正文,trafilatura.extract_metadata(...) 抽元数据,连 <noscript> 占位符这种边角都有处理)。scrapy 管调度去重与 politeness,trafilatura 管正文降噪——爬取引擎与抽取引擎各用所长。轻量场景另有 web_importer.py 单页导入(明确注释「刻意不用 trafilatura,保持轻量」),按需取用。
把两趟漫游叠起来看编排层的全貌。批量:回填模式按会话成批搬运,资源入库按 chunk 成批嵌入,ovpack 导入按目录成批重建索引——批是吞吐的单位。异步:defer_post_processing 让 HTTP 请求早返回、队列慢慢消化;poller 与 WatchScheduler 都是「短 tick + 持久游标」的节奏;语义 DAG 事件驱动按需分派。而所有异步路径殊途同归:写入 viking_fs 文件层 → 嵌入进向量层 → 自底向上生成 L0/L1 sidecar → 向量化 sidecar。解析定树(上一节),编排定序(本节),语义层随后生长——第 4 章的目录递归检索,搜的就是这趟流水线种出来的两层索引。
💡 漫游要点:编排层的三个关键词是游标、意图、异步。会话回填用持久化游标 + 意图先行的批次配方(set_pending→append→confirm→commit)在崩溃下做到不重不漏;轮询优于事件,因为「cursor→EOF」天然自愈。资源入库把重活全部队列化(嵌入队列/语义队列/自底向上的 sidecar DAG),HTTP 早返回、状态可查询;watch 调度器把「一次入库」变成「持续同步」,scrapy+trafilatura 把整个网站变成一棵可检索的树。第 3 章的分层不是独立功能,而是这条流水线的最后一站。
_backfill_one 循环 read_messages)→ commit_if_needed;added==0 也提交以防「已追加未提交」;BackfillStats 五本账,失败按会话隔离。下一节:
第 7 章 · 01 storage 双层与虚拟 FS——流水线的尽头是存储:VikingFS 虚拟文件系统与向量库如何分层协作,rm/mv如何级联清理向量索引,queuefs 队列与路径锁又长什么样。