Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions crates/persisting-dlcapt/src/tlv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -218,6 +218,13 @@ fn encode_block(
kind: &str,
) -> Result<String> {
let timestamp = Utc::now().to_rfc3339();
let storyline_id =
i64::try_from(seq).context("tlv sequence exceeds Storyline turn id range")?;
let storyline_source = if speaker == "assistant" {
"agent"
} else {
speaker
};
let mut fields = BTreeMap::new();
fields.insert("agent_id".to_string(), json!(record.agent_id));
fields.insert("call_id".to_string(), json!(record.call_id));
Expand All @@ -232,6 +239,18 @@ fn encode_block(
fields.insert("trace_id".to_string(), json!(record.call_id));
fields.insert("turn".to_string(), json!(record.turn));
fields.insert("v".to_string(), json!(BLOCK_FORMAT_VERSION));
fields.insert("message_encoding".to_string(), json!("text"));
fields.insert("step_id".to_string(), json!(storyline_id));
fields.insert(
"storyline".to_string(),
json!({
"id": storyline_id,
"kind": kind,
"ts": timestamp,
"src": storyline_source,
"model": record.model,
}),
);

if speaker == "assistant" {
fields.insert("status".to_string(), json!(record.status_code));
Expand Down Expand Up @@ -290,6 +309,11 @@ fn format_document_preamble(session_id: &str, agent_id: &str, turns: u64) -> Res
"session: {session_id}\n",
"agent: {agent_id}\n",
"turns: {turns}\n",
"storyline:\n",
" session: {session_id}\n",
" agent:\n",
" id: {agent_id}\n",
" name: {agent_id}\n",
"client:\n",
" peer: ''\n",
" peer_port: 0\n",
Expand Down
12 changes: 6 additions & 6 deletions crates/persisting-gateway/tests/agenticmd_bridge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,12 +65,12 @@ fn capture_turns_roundtrip_through_public_storyline_api() {
}

#[test]
fn legacy_fixture_still_imports_as_storyline() {
let story = decode_agenticmd(fixture()).expect("legacy AgenticMD import parse");
assert_eq!(story.turns.len(), 2);
assert_eq!(story.turns[0].source, "user");
assert_eq!(story.turns[0].message, json!("你好"));
assert_eq!(story.turns[1].source, "agent");
fn legacy_fixture_without_authoritative_storyline_is_rejected() {
let error = decode_agenticmd(fixture()).unwrap_err();
assert_eq!(
error.to_string(),
"missing authoritative Storyline metadata"
);
}

#[test]
Expand Down
12 changes: 6 additions & 6 deletions crates/persisting-gateway/tests/agenticmd_golden.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,12 +53,12 @@ fn generated_agenticmd_preserves_golden_storyline_semantics() {
}

#[test]
fn checked_in_legacy_golden_remains_readable() {
fn checked_in_legacy_golden_without_authoritative_storyline_is_rejected() {
let fixture = std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/agenticmd/demo-run-001.md");
let story = decode_agenticmd(&std::fs::read_to_string(&fixture).unwrap()).unwrap();
assert_eq!(story.session_id, "demo-run-001");
assert_eq!(story.turns.len(), 2);
assert_eq!(story.turns[0].message, json!("你好"));
assert_eq!(story.turns[1].message, json!("你好!有什么可以帮你的?"));
let error = decode_agenticmd(&std::fs::read_to_string(&fixture).unwrap()).unwrap_err();
assert_eq!(
error.to_string(),
"missing authoritative Storyline metadata"
);
}
6 changes: 4 additions & 2 deletions crates/persisting-pchronicle-cli/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1414,8 +1414,10 @@ async fn export_filters_complete_trajectories_and_streams_finite_json() -> Resul
let mut stderr = Vec::new();
run(cli, false, &mut stdout, &mut stderr).await?;

let rows: Value = serde_json::from_slice(&stdout)?;
let rows = rows.as_array().context("OpenAI export must be an array")?;
let document: Value = serde_json::from_slice(&stdout)?;
let rows = document["session_steps"]
.as_array()
.context("OpenAI export must contain a session_steps array")?;
assert_eq!(rows.len(), 1);
assert_eq!(rows[0]["session_id"], "training-002");
assert!(String::from_utf8(stderr)?.contains("trajectories=1"));
Expand Down
35 changes: 32 additions & 3 deletions crates/persisting-pchronicle-cli/tests/import_export_roundtrip.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,7 @@ async fn import_export_roundtrip_is_byte_identical_and_reimportable() -> Result<
}

#[tokio::test]
async fn forced_storyline_roundtrip_is_canonical_json_byte_identical() -> Result<()> {
async fn forced_storyline_roundtrip_is_canonical_and_reimport_stable() -> Result<()> {
let temp = tempfile::tempdir()?;
for fixture in EXAMPLE_FIXTURES {
let format = fixture.name;
Expand Down Expand Up @@ -115,10 +115,39 @@ async fn forced_storyline_roundtrip_is_canonical_json_byte_identical() -> Result
.await?;
assert!(exported_output.stderr_text()?.contains("exact=false"));

let reimported = temp.path().join(format!("{format}-storyline-reimported"));
run_cli([
"import",
"--from",
exported.to_str().unwrap(),
"--output",
reimported.to_str().unwrap(),
"--format",
format,
])
.await?;

let reexported = temp
.path()
.join(format!("{format}-storyline-reexport.json"));
let reexported_output = run_cli([
"export",
"--from",
reimported.to_str().unwrap(),
"--output",
reexported.to_str().unwrap(),
"--format",
format,
"--where",
"TRUE",
])
.await?;
assert!(reexported_output.stderr_text()?.contains("exact=false"));

assert_eq!(
canonical_json_bytes(&reexported)?,
canonical_json_bytes(&exported)?,
canonical_json_bytes(&input)?,
"Storyline round-trip canonical JSON differs for {format}"
"Storyline canonical JSON is not reimport-stable for {format}"
);
}
Ok(())
Expand Down
4 changes: 4 additions & 0 deletions crates/persisting-pchronicle/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,10 @@ required-features = ["search"]
name = "storyline_lance_roundtrip"
required-features = ["lance-store"]

[[test]]
name = "conversion_semantics"
required-features = ["lance-store"]

[[test]]
name = "atif_lance_corpus"
required-features = ["lance-store"]
Expand Down
36 changes: 23 additions & 13 deletions crates/persisting-pchronicle/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,8 +31,8 @@ events.lance ──单向投影──► StorylineDocument ──► Storyline
| `Storyline` | Storyline 三表 Lance | `runs`、`steps`、`tool_calls` | 权威二进制表示 |
| `AgenticMd` | Markdown 文件 | `runs`、`steps`、`tool_calls` | 可读编码,可双向转换 |
| `Atif` | ATIF JSON/JSONL/NDJSON | `runs`、`steps`、`tool_calls` | ATIF v1.7 对齐,可双向转换 |
| `OpenaiMsg` | OpenAI message corpus JSON | `runs`、`steps`、`tool_calls` | 通过分层 residual 无损往返 |
| `Actf` | ACTF JSON | `runs`、`steps`、`tool_calls` | 通过分层 residual 无损往返 |
| `OpenaiMsg` | OpenAI message corpus JSON | `runs`、`steps`、`tool_calls` | 通过分层 unknown fields 无损往返 |
| `Actf` | ACTF JSON | `runs`、`steps`、`tool_calls` | 通过分层 unknown fields 无损往返 |

统一读取入口是 `document::open_document`;返回的 `DocumentSource` 隐藏具体 provider,
并提供有预算上限的物化、逐条 Storyline 回调和 DataFusion 注册。写入仍使用
Expand All @@ -48,17 +48,27 @@ ACTF → Storyline Lance → ACTF
OpenAI Msg → Storyline Lance → OpenAI Msg
```

保真内容包括 ATIF 的 missing/null/value 三态、嵌套 subagent 顺序、`trajectory_id` 与
run-scoped `session_id` 的独立身份、RFC3339 原始偏移与亚毫秒精度,以及 ACTF/OpenAI 的
未知字段、数组顺序、attempt 分组和多 session 关系。外围格式无法映射到正式 Storyline
字段的内容保存在对应语义层级的受控 residual 中;不会保存完整原始对象副本。跨格式转换
只保证目标格式能够表达的语义。Canonical Event → Storyline 是有意的有损规范化投影,
不属于上述无损承诺。
已建模的 Storyline 语义(包括嵌套 subagent 顺序、`trajectory_id` 与 run-scoped
`session_id` 的独立身份、RFC3339 原始偏移与亚毫秒精度,以及 ACTF/OpenAI 的数组顺序、
attempt 分组和多 session 关系)按其规范化表示保存。已知字段的 missing/null 区别,以及
输入的物理容器形态(例如 ATIF 顶层单对象与单元素数组),都会被规范化,因而不作为
往返保真承诺。

ATIF 的顶层单对象/数组形态与 root 顺序作为格式无关的 Storyline 集合语义随 Lance
持久化,不依赖进程内 sidecar;Lance 内部另用 `storage_ordinal` 维护全局稳定读取顺序,
不会把多次增量写入都退化到 document id 排序。无法用任何 Storyline 表达的空 ATIF 数组
或空 OpenAI 信封会 fail closed,而不是接受后在导出时静默丢失容器字段。
源格式中 Storyline 未建模的键保存在受控 unknown fields:键名是带命名空间的精确
[RFC 6901 JSON Pointer](https://www.rfc-editor.org/rfc/rfc6901);未知字段值不保存完整原始对象
副本。未知键即使值为 `null` 也会保留。写回同一格式时,目标格式的规范字段优先;若
unknown field 与它们冲突,编码会 fail closed,而不会覆盖目标字段或静默丢弃冲突。

跨格式、多跳转换使用保留的 version-1 `_storyline` envelope 携带这些 unknown fields,确保目标
格式不能直接表示的源语义仍可在后续转换中恢复。每条 trajectory 跨所有来源默认最多
4,096 个 unknown fields、最多 1 MiB;任一上限溢出都会拒绝整条 Storyline,而非截断或只
保留部分未知字段。Canonical Event → Storyline 是有意的有损规范化投影,不属于上述无损
承诺。

Storyline Lance 的 `objects.lance` 可用于 unknown field 值的内部去重/卸载优化;它从不出现在
公共 Storyline 模型或任何公共 wire 输出中。Lance 内部另用 `storage_ordinal` 维护全局稳定
读取顺序,不会把多次增量写入都退化到 document id 排序。无法用任何 Storyline 表达的空
ATIF 数组或空 OpenAI 信封会 fail closed,而不是接受后在导出时静默丢失容器字段。

## DataFusion 能力

Expand All @@ -70,7 +80,7 @@ ATIF 的顶层单对象/数组形态与 root 顺序作为格式无关的 Storyli
| Storyline Lance | 是 | expression-dependent | 是 | 是 | 否 | 是 |
| ATIF | 是 | inexact | 是 | 否 | 是 | 否 |
| OpenAI Msg | 是 | unsupported | 是 | 否 | 否 | 否 |
| ACTF | 是 | unsupported | 是 | 否 | | 否 |
| ACTF | 是 | inexact | 是 | 否 | | 否 |
| AgenticMD | 是 | unsupported | 否 | 否 | 否 | 否 |

Canonical Event 保留 Lance projection/filter/limit pushdown、scalar index、pinned manifest
Expand Down
Loading
Loading