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
2 changes: 1 addition & 1 deletion crates/utopia-cli/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ struct ManifestDataDir {
/// not a side effect of a code change.
// 是迁移文件的**个数**,不是最大的编号(守卫 `schema_version_policy_compares_against_current`
// 按个数比):编号有空缺时两者不同——0071 由一个开放 PR 占着,0072 先落,个数是 71
const CURRENT_SCHEMA_VERSION: u32 = 76;
const CURRENT_SCHEMA_VERSION: u32 = 77;

fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
Expand Down
18 changes: 15 additions & 3 deletions crates/utopia-server/src/api/sources_routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -633,8 +633,8 @@ struct StatementsBody {
}

/// 门口的校验:形状对不对、有没有契约之外的键、每条陈述的引文格是不是空的。
/// 通过就把 `{e, s, n}` 按契约重新序列化成文档正文——存的是我们自己写出来的那份,
/// 不是调用方发来的字节,于是文档里没有信封、没有多余空白,块就是契约本身
/// 通过就按契约重新序列化成文档正文——存的是我们自己写出来的那份,不是调用方发来的
/// 字节:身份、日期(有的话)和 `{e, s, n}`,键序固定、没有多余空白,块就是契约本身
fn validate_statements_payload(raw: &[u8]) -> Result<(StatementsBody, Option<String>), String> {
if raw.len() > STATEMENTS_MAX_BYTES {
return Err(format!(
Expand Down Expand Up @@ -706,8 +706,20 @@ fn validate_statements_payload(raw: &[u8]) -> Result<(StatementsBody, Option<Str
return Err(format!("n[{i}][2] (quote) must be null"));
}
}
// 存下的文档带着这次观测的身份和日期,再是契约的三个数组(#900)。身份在正文里,
// 两次看到同一件事就是两篇正文不同的文档:库里一份内容只能有一篇(`documents_kb_sha_idx`),
// 文件型来源那条「同内容出现在新路径 = 改名」的识别也就永远碰不到它。解析器只读
// `e` / `s` / `n`,多出的两个键它不看
// 只有契约能通过 parse_open_response:这一步在门口跑一遍,抽取时不会再有别的答案
let content = serde_json::to_string_pretty(&json!({ "e": body.e, "s": body.s, "n": body.n }))
let mut stored = serde_json::Map::new();
stored.insert("external_id".into(), json!(body.external_id.trim()));
if let Some(at) = body.doc_time {
stored.insert("doc_time".into(), json!(at));
}
stored.insert("e".into(), body.e.clone());
stored.insert("s".into(), body.s.clone());
stored.insert("n".into(), body.n.clone());
let content = serde_json::to_string_pretty(&serde_json::Value::Object(stored))
.map_err(|e| format!("cannot serialise the contract: {e}"))?;
let parsed = utopia_extract::open::parse_open_response(&content)
.map_err(|e| format!("the contract does not parse: {e}"))?;
Expand Down
52 changes: 50 additions & 2 deletions crates/utopia-server/src/api/sources_statements_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,19 @@ fn the_door_refuses_what_the_contract_has_no_slot_for() {
let (_, content) = super::validate_statements_payload(ok.to_string().as_bytes())
.expect("a well-formed contract passes");
let content = content.expect("a non-tombstone yields the document text");
// 存的是我们重新序列化的那份:没有信封,能被抽取用的同一个解析器读回
assert!(!content.contains("external_id"));
// 存的是我们重新序列化的那份:带着这次观测的身份,能被抽取用的同一个解析器读回
// (它只读三个数组,身份和日期两个键它不看)
let stored: Value = serde_json::from_str(&content).unwrap();
assert_eq!(stored["external_id"], "obs-1");
assert!(
stored.get("doc_time").is_none(),
"no date was given, none is stored"
);
assert_eq!(
stored.as_object().unwrap().keys().collect::<Vec<_>>(),
["e", "external_id", "n", "s"],
"identity plus the three arrays, nothing else"
);
let parsed = utopia_extract::open::parse_open_response(&content).unwrap();
assert_eq!(parsed.statements.len(), 1);
assert_eq!(parsed.statements[0].phrase, "is on");
Expand Down Expand Up @@ -529,3 +540,40 @@ async fn the_push_token_can_be_viewed_and_rotated_like_an_api_source() -> anyhow
assert_eq!(status, StatusCode::OK, "{body}");
f.cleanup().await
}

/// 同一份载荷在新身份下是另一次观测(#900):两次看到杯子在桌上就是两篇文档,各带自己的
/// 日期。身份写在正文里,所以两篇正文不同,库里「一份内容一篇文档」的唯一性和文件型
/// 来源那条「同内容出现在新路径 = 改名」的识别都碰不到它
#[tokio::test]
async fn the_same_payload_under_a_new_identity_is_a_second_observation() -> anyhow::Result<()> {
let Some(f) = Fixture::new().await? else {
return Ok(());
};
let mut first = observation("08:14:03", "kitchen table");
first["external_id"] = json!("obs-1");
first["doc_time"] = json!("2026-09-23T08:14:03Z");
let (status, body) = f.push(f.source, &f.token, &first).await?;
assert_eq!(status, StatusCode::OK, "{body}");
assert_eq!(body["action"], "created");
let mut second = first.clone();
second["external_id"] = json!("obs-2");
second["doc_time"] = json!("2026-09-23T08:20:00Z");
let (status, body) = f.push(f.source, &f.token, &second).await?;
assert_eq!(status, StatusCode::OK, "{body}");
assert_eq!(body["action"], "created", "not a move");
for (key, when) in [
("statements:obs-1", "2026-09-23T08:14:03Z"),
("statements:obs-2", "2026-09-23T08:20:00Z"),
] {
let doc = documents::find_by_external_key(&f.pool, f.source, key)
.await?
.unwrap_or_else(|| panic!("{key} is its own document"));
assert_eq!(
doc.doc_time
.map(|t| t.to_rfc3339_opts(chrono::SecondsFormat::Secs, true)),
Some(when.to_string()),
"each observation keeps its own date"
);
}
f.cleanup().await
}
24 changes: 12 additions & 12 deletions crates/utopia-store/src/documents.rs
Original file line number Diff line number Diff line change
Expand Up @@ -155,8 +155,8 @@ pub async fn create_with_version_and_processing(
_ => AppError::Db(e),
})?;
sqlx::query(
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes)
VALUES ($1, $2, 1, $3, $4)",
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes, doc_time)
VALUES ($1, $2, 1, $3, $4, (SELECT doc_time FROM documents WHERE id = $2))",
)
.bind(Uuid::now_v7())
.bind(document.id)
Expand Down Expand Up @@ -210,10 +210,10 @@ pub async fn replace_content_and_enqueue_processing(
.execute(&mut *tx)
.await?;
sqlx::query(
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes)
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes, doc_time)
VALUES ($1, $2,
(SELECT coalesce(max(version), 0) + 1 FROM document_versions WHERE document_id = $2),
$3, $4)",
$3, $4, (SELECT doc_time FROM documents WHERE id = $2))",
)
.bind(Uuid::now_v7())
.bind(id)
Expand Down Expand Up @@ -285,10 +285,10 @@ pub async fn upsert_source_document_tx(
.execute(&mut **tx)
.await?;
sqlx::query(
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes)
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes, doc_time)
VALUES ($1, $2,
(SELECT coalesce(max(version), 0) + 1 FROM document_versions WHERE document_id = $2),
$3, $4)",
$3, $4, (SELECT doc_time FROM documents WHERE id = $2))",
)
.bind(Uuid::now_v7())
.bind(document.id)
Expand Down Expand Up @@ -340,10 +340,10 @@ pub async fn upsert_source_document_tx(
.execute(&mut **tx)
.await?;
sqlx::query(
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes)
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes, doc_time)
VALUES ($1, $2,
(SELECT coalesce(max(version), 0) + 1 FROM document_versions WHERE document_id = $2),
$3, $4)",
$3, $4, (SELECT doc_time FROM documents WHERE id = $2))",
)
.bind(Uuid::now_v7())
.bind(document.id)
Expand Down Expand Up @@ -390,8 +390,8 @@ pub async fn upsert_source_document_tx(
_ => AppError::Db(e),
})?;
sqlx::query(
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes)
VALUES ($1, $2, 1, $3, $4)",
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes, doc_time)
VALUES ($1, $2, 1, $3, $4, (SELECT doc_time FROM documents WHERE id = $2))",
)
.bind(Uuid::now_v7())
.bind(document.id)
Expand Down Expand Up @@ -727,10 +727,10 @@ pub async fn record_version(
size_bytes: i64,
) -> AppResult<()> {
sqlx::query(
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes)
"INSERT INTO document_versions (id, document_id, version, sha256, size_bytes, doc_time)
VALUES ($1, $2,
(SELECT coalesce(max(version), 0) + 1 FROM document_versions WHERE document_id = $2),
$3, $4)",
$3, $4, (SELECT doc_time FROM documents WHERE id = $2))",
)
.bind(Uuid::now_v7())
.bind(document_id)
Expand Down
59 changes: 48 additions & 11 deletions crates/utopia-store/src/materialize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,11 @@
//! 重算是集合运算,跑多少遍结果一样:先删「不再成立」的来源(陈述作废了、签名不再绑着、
//! 绑到了别的属性或反了方向、行本身作废了),再作废来源全空的类型化行,最后给「该有而
//! 没有」的(陈述, 绑定)对补上——有同断言的行就并进去,没有才新建。没有模型调用。
//!
//! 写下的行随后**对账**(#899):写路径上抽取和点头写完一条 state 事实都会沿它的唯一性
//! 方向重算时间线(`temporal::reconcile_new_fact`),物化出来的行是同一种观察,不该
//! 少这一步——否则一个函数型属性的两个值各开着一段,后一次观察关不上前一次。重算
//! 提交之后再对账:对账按时间线各自开事务、拿自己的锁,不在物化的事务里做

use sqlx::PgPool;
use utopia_core::AppResult;
Expand All @@ -38,6 +43,10 @@ pub struct Outcome {
pub merged: u64,
/// 规则算出来的隐含行(0044 决定 3 第五片),新建的
pub implied: u64,
/// 对账自动闭合而改写出来的修正行数(#899)
pub corrected: u64,
/// 对账裁不了、交给人的冲突数
pub conflicts: u32,
}

/// 一条该物化的(陈述, 绑定)对,连陈述上要抄的东西。
Expand Down Expand Up @@ -71,8 +80,24 @@ pub async fn materialize(pool: &PgPool, kb_id: Uuid) -> AppResult<Outcome> {
.bind(kb_id.to_string())
.execute(&mut *tx)
.await?;
let outcome = materialize_in_tx(&mut tx, kb_id).await?;
let (outcome, written) = materialize_in_tx(&mut tx, kb_id).await?;
tx.commit().await?;
reconcile_written(pool, kb_id, outcome, &written).await
}

/// 这一轮写下(新建或并入)的类型化行沿各自的唯一性时间线对账(#899)。只有 state 且
/// 声明了唯一性的谓词有时间线,`timelines_of` 自己筛;其余的行这里是空转
async fn reconcile_written(
pool: &PgPool,
kb_id: Uuid,
mut outcome: Outcome,
written: &[Uuid],
) -> AppResult<Outcome> {
if !written.is_empty() {
let report = crate::temporal::reconcile_facts(pool, kb_id, written).await?;
outcome.corrected = report.corrected.len() as u64;
outcome.conflicts = report.conflicts;
}
Ok(outcome)
}

Expand All @@ -96,15 +121,19 @@ pub async fn try_materialize(pool: &PgPool, kb_id: Uuid) -> AppResult<Option<Out
tx.rollback().await?;
return Ok(None);
}
let outcome = materialize_in_tx(&mut tx, kb_id).await?;
let (outcome, written) = materialize_in_tx(&mut tx, kb_id).await?;
tx.commit().await?;
Ok(Some(outcome))
Ok(Some(
reconcile_written(pool, kb_id, outcome, &written).await?,
))
}

async fn materialize_in_tx(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
kb_id: Uuid,
) -> AppResult<Outcome> {
) -> AppResult<(Outcome, Vec<Uuid>)> {
// 这一轮写下的行(新建的和并入的),提交后对账
let mut written: Vec<Uuid> = Vec::new();
// 1. 删不再成立的来源:陈述死了、行死了、签名没绑着、属性或方向变了、陈述带了 mood
sqlx::query(&format!(
"DELETE FROM typed_fact_sources src
Expand Down Expand Up @@ -248,6 +277,7 @@ async fn materialize_in_tx(
}
_ => continue,
};
written.push(fact);
if new {
added += 1;
sqlx::query("UPDATE facts SET from_statement_id = $2 WHERE id = $1 AND from_statement_id IS NULL")
Expand Down Expand Up @@ -311,13 +341,18 @@ async fn materialize_in_tx(
}
// 3b. 已批准的规则算隐含行(0044 决定 3 第五片)。读数只查缓存:缓存里没有的这一轮
// 不算,`read_phrases` 填上之后再来。短语规则按陈述触发,类别词规则按实体触发
let implied = imply_in_tx(tx, kb_id).await?;
Ok(Outcome {
retired,
added,
merged,
implied,
})
let implied = imply_in_tx(tx, kb_id, &mut written).await?;
Ok((
Outcome {
retired,
added,
merged,
implied,
corrected: 0,
conflicts: 0,
},
written,
))
}

#[derive(sqlx::FromRow)]
Expand All @@ -341,6 +376,7 @@ struct Implied {
async fn imply_in_tx(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
kb_id: Uuid,
written: &mut Vec<Uuid>,
) -> AppResult<u64> {
// 短语规则:签名下活着的、没 mood 的陈述;宾语是读数的答案(缓存里的实体或值),
// 没有读数时就是陈述的宾语。已经有活着的隐含行以这条陈述为来源的不再算
Expand Down Expand Up @@ -425,6 +461,7 @@ async fn imply_in_tx(
d.confidence,
)
.await?;
written.push(fact);
if new {
implied += 1;
sqlx::query("UPDATE facts SET implied = TRUE WHERE id = $1")
Expand Down
2 changes: 1 addition & 1 deletion crates/utopia-store/src/materialize_delivery_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ async fn try_materialize(pool: &PgPool, kb_id: Uuid) -> AppResult<Option<Outcome
tx.rollback().await?;
return Ok(None);
}
let outcome = materialize_in_tx(&mut tx, kb_id).await?;
let (outcome, _written) = materialize_in_tx(&mut tx, kb_id).await?;
tx.commit().await?;
Ok(Some(outcome))
}
Expand Down
Loading
Loading