Skip to content

Storyline 三表 Lance 存储

StorylineLanceStore 是 pChronicle 的 Storyline-native 规范化物理表示。它与 events.lance 原始事件日志并列存在,不替代后者。

Projection contract and closed loop

Storyline retains the Hub interchange contract (path A), while the three-table store is a rebuildable silver projection of canonical events.lance (path B). The two uses share a schema, but not write identity: interchange imports and direct replace_storyline calls carry no canonical lineage. Only the events projector may publish a CURRENT with projection lineage.

events.lance (source of truth)
  ├─ project build/rebuild ─► runs + steps + tool_calls + objects
  ├─ project sync ──────────► replace only sessions touched by the append suffix
  └─ Catalog fallback ──────► project a pinned events snapshot when missing or stale

CURRENT pins exact Lance versions for all four tables and records the source URI and source identity, fact_version, fact_rows, the source layout revision at build time, projector and recipe identity, recipe hash, and completeness. fact_version and fact_rows are the freshness watermark. Compaction changes only the layout revision and does not stale a projection. Direct document writes clear lineage; maintenance preserves it.

Operational commands:

# The destination must be empty; build from one pinned fact snapshot.
pchronicle project build --from ./run/events.lance --output ./run/storyline

# status reads CURRENT only; verify compares it with the current canonical fact watermark.
pchronicle project status --from ./run/storyline
pchronicle project verify --from ./run/storyline --source ./run/events.lance

# Read only the new append suffix, then reload and replace the affected sessions.
pchronicle project sync --from ./run/storyline --source ./run/events.lance

# Continuously sync and periodically verify outside the Gateway hot path; emit one JSONL event per loop.
pchronicle project watch --from ./run/storyline --source ./run/events.lance

# Write a new table generation before switching CURRENT.
pchronicle project rebuild --from ./run/events.lance --output ./run/storyline

Catalog merges a lineage-linked sidecar with the events source into one logical source. When sources.projection_status is fresh, normalized queries use the three tables. When it is stale, Catalog hides the sidecar and falls back to a deterministic projection of the pinned events snapshot. projection_generation exposes the generation actually selected. A Storyline document store without lineage is never inferred to be a projection of canonical events.

project watch is the recommended operational runner under systemd, Kubernetes, or another supervisor. It applies capped exponential backoff, shuts down on process signals, and supports --exit-on-error when the supervisor owns retries. It deliberately stays outside Gateway capture, so projection failures cannot block canonical event writes.

本文只负责三表物理 schema、内容层、Snapshot 发布、查询接入和维护语义。事实源与 projection ownership 见轨迹存储,用户查询流程见 Dataset 查询指南

这是 pChronicle 唯一的规范化三表模型。旧的 ATIF sessions / steps / tool_callsNormalizedStore 和内存联表视图已经删除。ATIF 仍作为输入输出格式存在,但查询时先转换为 Storyline,再投影到本页定义的 runs / steps / tool_calls schema,不再维护第二套表结构。

表模型

粒度 逻辑主键 外键
runs.lance 每个 Storyline 一行 session_id
steps.lance 每个 turn 一行 (session_id, step_id) session_id → runs
tool_calls.lance 每个 tool call 一行 (session_id, tool_call_id) (session_id, step_id) → steps

run_id 是 Run 分组键;一个 Run 可以包含主 Story 和多个 subagent Story,因此 runs.lance 中可能有多行共享同一个 run_idsession_id 才是 Storyline 文档的 唯一键。

常用 JSON 值(message、arguments、metrics、extra 等)以 UTF-8 JSON 列保存;身份、 顺序、类型、时间和性能字段使用独立的 Arrow 标量列,便于过滤和分析。

steps.timestamp 是规范化到 UTC 的 Timestamp(Millisecond, "UTC")。写入端拒绝无法 解析为 RFC3339 的非空时间;原始时区偏移和亚毫秒精度不会写入物理表,读取为 Storyline / ATIF 时统一编码成带 Z 后缀、精度不超过毫秒的 UTC 字符串(整秒省略小数部分)。SQL 排序、范围过滤和时间聚合直接使用 timestamp。这是一次物理 schema 变更;旧三表投影 需要从 canonical events 或原始交换文件重建,不对既有 Lance generation 做原地类型猜测。

tool result 不再留在 step 的 observation JSON 中。写入时根据 observation.results[].source_call_id 关联到对应 tool call,并保存到该行的 results_json。缺失或错误的关联会拒绝整次写入。

大块内容层

设计目标与边界

Agent 轨迹中的长 reasoning、工具输出、源代码、日志和多模态载荷会让列式表出现少量超大 cell。若把它们与身份、顺序、类型和指标一起内联,常规过滤与聚合也要承受更大的 fragment、 page cache 和解码开销。pChronicle 因此在三表之外增加共享内容层,但保持三个约束:

  1. runs / steps / tool_calls 的 Arrow schema 和 SQL 结果不变;内容层是内部物理优化。
  2. 小值继续内联,只有达到阈值的 UTF-8/JSON cell 才外置,避免所有读取都退化成 KV lookup。
  3. 内容按原始字节寻址并跨 Storyline 复用;不把轨迹主键、生命周期或业务去重混入内容层。

当前实现没有定制 Lance 文件格式或私有索引类型,而是组合 Lance Blob v2、普通 BTree scalar index 和 DataFusion execution node。这样可以得到需要的延迟物化能力,同时把维护 面限制在 pChronicle 自己的协议与执行计划中。

内部描述符协议

超过默认 64 KiB 的内容列在三表中暂时编码为:

<RS>PCHRONICLE-CONTENT:<type>:<codec>:<blake3-256>:<raw_length>:<preview-base64url>
字段 当前编码 作用
magic PCHRONICLE-CONTENT 严格识别内部引用
logical type u / j / b UTF-8、JSON;binary 标签已保留给后续二进制列
codec i / z identity 或 Zstd
content id 64 位十六进制 BLAKE3-256 对未压缩原始字节寻址、校验和跨轨迹复用
raw length u64 解压后长度校验,也允许无 payload 的代价判断
preview URL-safe Base64 默认最多 256 个 UTF-8 字节的安全前缀

描述符只允许存在于内部物理列。用户原文若恰好以 magic 开头,也会被强制外置,读取时再 恢复为原文,从而消除“用户字符串被误认为引用”的歧义。公开的读取、SQL、转换和导出 API 必须返回完整值或显式 preview,不能泄露描述符。

objects.lance 使用以下物理列:

作用
content_id BLAKE3 内容地址;建立 BTree index
logical_type, media_type 逻辑类型和 MIME 提示
raw_length, stored_length, codec 完整性检查和存储代价
preview 无 Blob I/O 的安全预览
payload Lance Blob v2,保存 identity/Zstd 字节
created_at_ms 对象创建时间

写入、复用与发布

写入端对候选 cell 依次执行:

原始 UTF-8/JSON
  ├─ 小于阈值 ───────────────────────────────► 原值内联
  └─ 达到阈值 / 命中 magic
       ├─ BLAKE3(raw bytes) + UTF-8 preview
       ├─ Zstd;没有净收益则保留 identity
       ├─ batch 内按 content_id 合并并检查碰撞
       ├─ BTree 批量查询 objects.lance,跳过已存在对象
       └─ 先提交对象 version,再写三表 descriptor,最后发布 CURRENT

对象必须先于引用持久化;CURRENT 同时固定三张业务表和对象表的精确 Lance version。 任一步骤失败都不会发布新快照:允许留下不可达对象,但不会发布悬空引用或跨表半提交。 跨轨迹复用只依赖内容地址,不依赖 session 生命周期,因此同一长文本在不同 Run 中只保存 一次。同一写入批次内若 content id 相同但 codec、原始长度或存储字节不一致,会拒绝写入。

对象层在普通写入期间保持 append-only,GC 不进入写入热路径。显式 maintain 会只扫描三表 的内容引用列,计算当前快照的可达 content id,并清理不可达 payload。生产环境仍需要把对象 增长率、不可达字节和维护耗时纳入指标。

查询期延迟物化

StorylineDataSource 先让 Lance 完成业务表的 projection、可安全谓词、scalar index、limit 和并行扫描,再在计划中插入 ContentHydrationExec

  • 查询不引用内容列时,不打开 objects.lance payload;
  • 只收集投影中实际出现的 content id,以最多 512 个为一组走 BTree lookup;
  • 根据 row address 批量读取 Blob,解压后验证长度与 BLAKE3,再恢复原 Utf8 列;
  • 内容列谓词不能作用于描述符,必须保留在 hydration 之后由 DataFusion 计算;
  • Preview 模式只返回描述符中的 UTF-8 前缀,零 payload I/O,并拒绝内容列谓词,避免把 preview 错当成完整值。

因此大内容的成本只由真正读取这些内容的查询承担;身份过滤、计数、分组和指标分析仍沿用 紧凑的三表列式路径。

提交布局

root/
├── CURRENT
├── objects.lance/
└── generations/
    └── <table-generation>/
        ├── runs.lance/
        ├── steps.lance/
        └── tool_calls.lance/

首次导入创建三张规范化 Lance dataset、共享的 objects.lance 和标量索引。后续 replace_storyline 不再读取或重写全库,而是按各表主键执行 merge-upsert,并只删除指定 session_id 中已经不再存在的旧键。每次替换 会产生一个新的逻辑 snapshot; CURRENT 是一段 JSON,记录逻辑 snapshot id、物理 table_generation、三张表以及对象表 各自精确的 Lance version id。对象先持久化,三张业务表随后写入,最后才更新 CURRENT; 因此失败最多留下不可达对象,不会发布悬空引用或跨表半提交。

阈值、preview 长度和 Zstd level 可通过 StorylineContentOptions 配置;三表 schema 不变。

Lance MVCC 的旧版本默认保留,便于已打开的 reader 固定快照及故障恢复。频繁增量更新 会积累 fragment、delete file 和未合并的索引增量。普通 replace 不执行 index refresh 或 compaction,避免某次写请求出现维护型长尾;生产环境通过 maintain 显式执行三表并行 compaction、补齐/刷新索引、内容 GC 和按保留期 vacuum。维护产生的三个新 version 仍先原子更新 CURRENT,之后才回收旧版本。CURRENT 必须是包含全部精确版本的 JSON 指针,不读取旧的 纯文本 generation 指针。

本地写入通过进程内锁和文件锁串行化;对象存储通过 CURRENT 的 ETag/version 条件更新 执行 optimistic CAS。stale commit 不能移动 CURRENTStorylineLanceStore 在 CAS 冲突后 直接返回错误,不会重新读取、merge 或自动重试。调用方若选择重试,必须从最新 snapshot 重新 开始完整 replace。上层 lease 可减少冲突,但不改变这一失败语义。

Rust API

let store = StorylineLanceStore::open(path).await?;
store.replace_storyline(&storyline).await?;
let restored = store.get_storyline_full("session-id").await?;
let report = store.maintain(&LanceMaintenanceOptions::default()).await?;

replace_storylinesession_id 为边界替换三张表中的相关行,同时保留同一 store 内的其他 Storyline。

get_storyline_full 明确表示会读取三表并恢复该 Storyline 的全部内容。未被 CLI 或 Web 使用的 store-local 分页 API 已删除;产品层的列表、分页和投影统一由 Catalog、Warehouse API 和 DataFusion query 承担,避免维护第二套不可达的读取协议。

首次导入和替换都并行写三张表。Arrow 行按最多 8192 行一批懒编码并流入 Lance,避免 导入大型语料时同时保留整表的 Arrow 副本。CURRENT 只解析一次;DataSource 随后把每张 表直接打开到指针指定的 version,不再先验证、再重复打开同一 dataset。

生产环境通过 StorylineLanceStore::maintain Rust API 执行维护;公共 CLI 不提供维护命令。

DataFusion datasource

StorylineDataSource 在打开时固定 CURRENT 中的三个业务表 version 和对象表 version,并把三张 dataset 注册为 runsstepstool_calls。即使写入端随后切换 CURRENT,已经打开 的查询仍使用同一份三表快照。

let source = StorylineDataSource::open(path).await?;
let ctx = source.session_context()?;
let rows = ctx
    .sql("SELECT step_id, source FROM steps WHERE session_id = 's-1' ORDER BY step_id")
    .await?
    .collect()
    .await?;

Datasource 使用 Lance 原生 DataFusion execution plan,支持列裁剪、谓词和 limit 下推, 并采用 unordered physical scan 允许并行读取;有顺序要求的查询必须显式使用 ORDER BY step_id, call_index。未引用大内容列的查询不会打开 Blob;引用内容列时在 Lance 投影/安全谓词/limit 之后批量恢复。针对内容列的谓词不下推到内部引用,而是在恢复 后由 DataFusion 求值,确保 SQL 语义不变。内部引用不会由 pChronicle 的读取、查询、导出 API 返回;直接绕过 pChronicle 扫描底层 Lance 文件属于诊断接口,不在该保证内。

预览 UI 可把 StorylineDataSourceOptions::content_read_mode 设为 StorylineContentReadMode::Preview。该模式直接从描述符返回 UTF-8 安全的短 preview,零 Blob payload I/O;为避免把 preview 当成完整值产生错误结果,内容列谓词在 preview 模式 下会被明确拒绝。

首次创建 table generation 时建立以下标量索引:

BTree Bitmap
runs session_id, run_id
steps session_id, timestamp effective_kind, source
tool_calls session_id, tool_call_id function_name

这些索引针对按 Story/Run 定位、tool-call 查找和类型过滤。step_id 在每个 Storyline 内 从小值重新计数,全局选择性低,因此不单独建立 BTree;组合条件先用 session_id 定位到 单个 Storyline,再过滤很短的 step 范围。DataFusion 的索引谓词会下推为 Lance ScalarIndexQuery

StorylineDataSourceOptions 可显式控制 use_scalar_indexesscan_in_order;默认配置 面向在线分析查询启用索引、关闭物理顺序。关闭索引主要用于 benchmark、诊断或极小表 的全扫描对照。

统一查询引擎

ChronicleQueryEngine 是对外的只读 SQL 门面。Lance 与 ATIF 后端注册完全相同的 runsstepstool_calls 表,因此查询语句不需要随物理格式改变:

let lance = ChronicleQueryEngine::open_lance("./storyline-store").await?;
let batches = lance.query(
    "SELECT session_id, step_id, source FROM steps WHERE step_id >= 10"
).await?;

let atif = ChronicleQueryEngine::open_atif("./trajectories.ndjson")?;
let jsonl = atif.query_jsonl(
    "SELECT source, COUNT(*) AS steps FROM steps GROUP BY source ORDER BY source"
).await?;

query 返回 Arrow RecordBatch,适合服务端继续处理;dataframe 返回 lazy DataFrame, 适合追加 DataFusion 变换或查看计划;query_jsonl 用于 CLI/API 边界。调用者也可通过 context() 取得 SessionContext 注册 UDF 或额外表。

AtifDataSource 接受单个 ATIF JSON 对象、JSON 数组、每行一个完整 trajectory 的 JSONL/NDJSON,以及包含这些 ATIF 文档的目录。文件路径默认注册为按文件 lazy 的 StreamingTable:manifest 在打开时冻结路径和文件身份,scan 才读取命中文件;目录按 稳定顺序发现,每个文件是独立 partition,并以固定大小 Arrow batch 提供背压。显式 from_json / from_trajectories 因为调用者已经持有完整输入,仍使用 MemTable

pChronicle + JSON 投影查询快路径

旧路径把每份 JSON 完整解析为格式对象,再执行 ATIF → Storyline → 三表行 → 全宽 Arrow 后交给 SQL。即使查询只需要 sourceCOUNT(*),也会构造 message、reasoning、metrics、 tool calls 等未使用的大字段。新路径把优化边界前移到 TableProvider::scan

SQL / DataFrame
  → DataFusion projection + filters
  → FileScanSpec
      ├─ _file_ = / IN / LIKE:manifest 文件裁剪
      ├─ session_id:trajectory 裁剪
      ├─ step_id / source:step 裁剪
      └─ projected column set
  → BufRead / serde streaming decoder
      └─ DeserializeSeed + Visitor + IgnoredAny
  → 只为命中行解码被引用字段
  → projected Arrow RecordBatch
  → DataFusion 保留 inexact filter 再次校验

当前 fast path 的适用范围是 ATIF 单对象 JSON、JSON 数组(包括 pretty JSON),以及每行一个对象的 JSONL/NDJSON, 目标表为 steps,并且物理计划存在严格列裁剪。它有意保持保守:

输入/查询 执行路径
ATIF object/pretty object + projected steps reader-backed seeded projected decoder
ATIF array/pretty array + projected steps fill_buf 结构扫描 + 有界 element buffer + seeded from_slice
ATIF JSONL/NDJSON + projected steps BufRead 逐记录、有界 record buffer
_file_session_idstep_idsource 的安全简单谓词 可提前裁剪,DataFusion 仍复核
SELECT * 完整规范化 fallback
runs / tool_calls 完整规范化 fallback
OpenAI-message / ACTF 完整规范化 fallback
无法证明安全的表达式、OR/函数/跨列条件 不预裁剪,由 DataFusion 求值

DeserializeSeed 把查询 projection 和安全谓词传入 Visitor;未引用字段交给 IgnoredAny 做语法扫描,不构造 Value/Storyline。JSONL/NDJSON 以 BufRead 逐记录读取; JSON array 的结构扫描器识别字符串和转义,在不构造 DOM 的情况下提取单个 trajectory, 再通过 slice decoder 执行投影解析。单条 JSONL 记录或 array element 由 max_record_bytes 限制;单对象直接从 reader 解码。三种路径都不先复制整文件。 Arrow encoder 也只创建投影列,COUNT(*) 使用合法的零列 batch。轻量路径 校验 JSON、必需字段、重复 session、命中文档内的重复 step 和当前表内约束;跨表引用 完整性仍由导入路径或完整 fallback 负责。这一边界使临时查询不承担导入语义,同时不降低 SQL 结果正确性。

查询指标额外报告 projected_filesstreamed_recordsstreaming_buffer_peak_bytes、 scanned/pruned documents、scanned/pruned/emitted rows 和 projected_arrow_bytes,用来区分 “源字节扫描”“输入缓冲”“JSON 字段物化”和“Arrow 输出”四个成本。仓库 benchmark 报告 median/P95、rows/s、独立进程峰值 RSS,以及计数 allocator 观测到的 allocation calls/bytes; 这些是指定 corpus、查询和机器的回归数据,不是跨环境 SLA。

该路径仍需顺序扫描命中文件的全部 JSON 字节,不是文件内索引。一次性或受控批次查询可 直接使用 JSON;超大、远端或反复查询的数据应先转换为 Lance,利用 snapshot、列裁剪、 并行 fragment scan 和 scalar index。

ATIF 导入同样默认走 AtifReader。空 store 使用一个 producer 单遍完成 校验、Storyline 规范化和三表拆分,再经三条有界 Arrow channel 并行创建三个 Lance dataset;已有 store 则以最多 256 个 Storyline 为一个增量替换批次。两种路径都在所有 输入和三表写入成功后才原子切换一次 CURRENT

CLI 使用相同引擎,输出稳定的 JSONL:

pchronicle query ./trajectories.ndjson \
  'SELECT source, COUNT(*) AS steps FROM dataset.steps GROUP BY source ORDER BY source'

# 含 CURRENT 的三表 store 根目录会被 auto 识别为 Lance
pchronicle query ./storyline-store \
  'SELECT step_id, source FROM dataset.steps WHERE session_id = '\''s-1'\'' ORDER BY step_id'

# OpenAI/ACTF 目录直接查询;_file_ 为查询期相对路径列,不写入 Lance
pchronicle query ./openai-data \
  "SELECT _file_, COUNT(*) FROM dataset.steps WHERE _file_ LIKE 'batch/%' GROUP BY _file_"

查询是只读的;SQL 可以使用 SELECT、CTE、JOIN、聚合和 DataFusion 内置函数,但不通过 这个门面执行 DDL/DML。Lance 引擎打开时固定 CURRENT 指向的三个版本,从而保证一次 查询会话内三张表来自同一快照。

仓库使用 Criterion.rs + hyperfine 的统一 benchmark runner。Criterion 负责 CPU-bound 转换、events→Storyline 和三表 split/reconstruct 微基准;canonical event append、投影 build/sync/verify、Lance/DataFusion 生命周期、JSON streaming 与 RSS 场景由 hyperfine 重复执行独立进程,最终生成统一 JSON、Markdown 和 HTML:

# PR/local smoke workload
just benchmark-pchronicle

# larger nightly workload
just benchmark-pchronicle nightly target/pchronicle-benchmark/nightly

# compare two raw reports produced on the same testbed
just benchmark-pchronicle-compare \
  target/pchronicle-benchmark/main/raw-report.json \
  target/pchronicle-benchmark/current/raw-report.json

raw-report.json$["measurements"]... JSONPath 地址保存原始指标和环境, bencher.json 是历史平台使用的扁平投影, report.md 写入 GitHub Actions Job Summary,report.html 与 Criterion 明细作为 artifact。

JSON 对照使用单个 NDJSON 文件,避免大量小文件打开开销。ATIF steps 直接查询会把 DataFusion projection 和可安全预裁剪的 session_idstep_idsource 谓词传给 projected decoder:未引用 JSON 字段只做语法扫描,不构造 Storyline/三表对象,Arrow batch 也只包含执行计划需要的列。object、array、pretty JSON 和 JSONL/NDJSON 共用流式 projection decoder;SELECT * 和其他格式仍走完整规范化 fallback。轻量路径执行 JSON、必需字段和表内约束校验,跨表引用完整性由导入或完整 fallback 校验。预解析内存 JSON 对照只计算查询逻辑,用来区分产品工作流与纯内存遍历成本。 benchmark 还单独输出 DataSource 冷打开并执行 SQL、get_storyline_full 点查和单 Storyline 替换的延迟,避免 warm SQL 吞吐掩盖在线读写路径的写放大。

性能结论不应写成“Lance 在所有规模和查询上必然更快”:显式构造的 MemTable 或预解析 内存 JSON 在小数据下仍可能更快。默认 ATIF streaming 解决的是内存上界,不提供物理 索引;Lance 的主要优势仍是更小的物理体积、近乎常数的 datasource 打开时间,以及列 裁剪、并行扫描和选择性索引收益。

相关文档