第 6 章 · 02 ingest 编排与 sidecar 生成


第 6 章 · 02 ingest 编排与 sidecar 生成

本节摘要:解析器只管「把万物变文本」,从文本到 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 生成管线。

学习目标

  1. 读懂 IngestOrchestrator 的回填三步:发现会话 → 游标回放 → 空闲提交;理解 BackfillStats 的统计口径。
  2. 说清 IngestPoller 的轮询循环与「持久化游标自愈」:错过的 tick 为什么不会丢消息。
  3. 理解 SessionReplayer 的崩溃安全配方:set_pending → append → confirm → commit,意图先落盘。
  4. 走通资源入库异步链:add_resource 入队 → process_resource → 嵌入/语义双队列 → 写入 viking:// → sidecar 生成。
  5. 掌握语义 DAG 的自底向上生成顺序,以及 WatchScheduler 的定时重灌与网页递归爬取。

一、会话回填:orchestrator 的游标三步曲

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 的意图先行与 poller 的自愈轮询

回放的核心 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.pywrite_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 头)也在这里写。

四、轮询源目录与递归网页:watch 与 web_crawler

「入库一次」之外,资源还有「持续同步」的形态。resource/watch_scheduler.pyWatchScheduler 每 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 章的分层不是独立功能,而是这条流水线的最后一站。

本节要点回顾

  • ingest/ 与 resource/ 两个编排:前者回放会话日志(claude_code/codex/cursor/hermes/openclaw/opencode 六类 LogSource,claude_code 只取 text 块),后者负责资源入库;共享写入队列与语义管线。
  • IngestOrchestrator 三步:discover_sessions → 游标回放(_backfill_one 循环 read_messages)→ commit_if_needed;added==0 也提交以防「已追加未提交」;BackfillStats 五本账,失败按会话隔离。
  • SessionReplayer 配方:每批 ≤100 条,意图(目标游标+条数+服务端消息数基线)先落盘再 append,崩溃后按基线比对决定 confirm 还是重读;游标只在确认后前进。
  • IngestPoller:5 秒级轮询、_dirty 跟踪、空闲/token 阈值提交;持久化游标使错过 tick、重启都不丢消息。
  • 资源入库异步链:add_resource → ContentTargetSpec → process_resource(defer_post_processing=True)→ persist_temp_tree → 嵌入/语义双队列 → wait_processed 可选等待;L0/L1 由 write_semantic_sidecars + semantic_dag 自底向上生成,文件摘要→目录 abstract→父层 overview,事件驱动懒分派。
  • WatchScheduler:60 秒检查、Semaphore(4) 限并发、WatchManager 持久化任务、飞书 OAuth/Git 认证自动续期;web_crawler 用 scrapy 递归爬 + trafilatura 抽正文,下载类链接单独收集,轻量单页走 web_importer。
  • 排产哲学三条:批量是吞吐的单位(会话批/嵌入批/导入批),异步是响应的单位(defer_post_processing/队列消化),游标是一致性的单位(所有断点续传都靠它)。

下一节:第 7 章 · 01 storage 双层与虚拟 FS——流水线的尽头是存储:VikingFS 虚拟文件系统与向量库如何分层协作,rm/mv 如何级联清理向量索引,queuefs 队列与路径锁又长什么样。


作者与出处
原作者: 灏天文库
整理: 灏天文库整理
本站整理收录,版权归原作者/开源协议所有;欢迎通过原文链接访问源仓库。
发布者: 作者: 灏天文库 转发
评论区 (0)
U