diff --git a/crates/utopia-server/src/api/export_routes.rs b/crates/utopia-server/src/api/export_routes.rs index ff853c7cd..6c995a0ba 100644 --- a/crates/utopia-server/src/api/export_routes.rs +++ b/crates/utopia-server/src/api/export_routes.rs @@ -4,6 +4,10 @@ //! 一个十万条事实的库正是最需要导出的那种库,也正是「先拼成一个 String」会 //! 把服务打死的那种库。 //! +//! **同一份快照**:体检、词汇表与每一页查询跑在同一条只读 REPEATABLE READ +//! 事务里。若每页各自拿一条连接,词汇表发完之后才提交的规则会在派生页里 +//! 留下一条指向它的 `wasGeneratedBy`——图里就出现没有本体的引用。 +//! //! 中途出错只能截断——HTTP 头早就发出去了。所以错误进日志,而客户端拿到的是 //! 一份短了一截的文件;这比先攒后发要好,那种做法在同样的库上根本发不出来。 @@ -44,7 +48,9 @@ pub async fn export( })?; let names = Names::new(kb_id, q.base.as_deref()).map_err(AppError::Validation)?; - // 导出是一次「整个库离开这台机器」的动作,台账要记下(0014 的同一条理由) + // 导出是一次「整个库离开这台机器」的动作,台账要记下(0014 的同一条理由)。 + // 必须在快照事务之前写完并放掉连接:下面那条事务一占就是整个流的时长, + // 占着它再回头要连接,会把最小池(两条)上的并发导出互相饿死 let _ = utopia_store::audit::record( &state.pool, Some(kb_id), @@ -56,13 +62,24 @@ pub async fn export( ) .await; - let pool = state.pool.clone(); + // 整份导出占一条连接上的只读 REPEATABLE READ 事务:下面的预检与流里 + // 每一页查询读同一个快照,事务活到流结束 + let mut tx = state.pool.begin().await?; + sqlx::query("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ READ ONLY") + .execute(&mut *tx) + .await?; + + // 出处链越库一律拒导(0070 挡新行,这里挡存量坏行):宁可不发一个字节, + // 也不能把别库的对象安上本库的 IRI——那份文件看着完整、实则悬空。 + // 在流开始之前拦下:客户端拿到的是确定性的错误,不是一截断掉的文件 + utopia_store::export::provenance_integrity(&mut tx, kb_id).await?; + let stream = async_stream::try_stream! { let buf = SharedBuf::default(); let mut sink = Sink::new(format, buf.clone()); - let classes = utopia_store::export::classes(&pool, kb_id).await.map_err(io)?; - let relations = utopia_store::export::relations(&pool, kb_id).await.map_err(io)?; + let classes = utopia_store::export::classes(&mut tx, kb_id).await.map_err(io)?; + let relations = utopia_store::export::relations(&mut tx, kb_id).await.map_err(io)?; let vocab = rdf::vocabulary(&names, &classes, &relations); for c in &classes { rdf::emit_class(&mut sink, &vocab, c)?; @@ -74,7 +91,7 @@ pub async fn export( let mut after = None; loop { - let page = utopia_store::export::documents_page(&pool, kb_id, after).await.map_err(io)?; + let page = utopia_store::export::documents_page(&mut tx, kb_id, after).await.map_err(io)?; let Some(last) = page.last() else { break }; after = Some(last.id); for d in &page { @@ -85,7 +102,7 @@ pub async fn export( let mut after = None; loop { - let page = utopia_store::export::entities_page(&pool, kb_id, after).await.map_err(io)?; + let page = utopia_store::export::entities_page(&mut tx, kb_id, after).await.map_err(io)?; let Some(last) = page.last() else { break }; after = Some(last.id); for e in &page { @@ -99,7 +116,7 @@ pub async fn export( let now = chrono::Utc::now(); let mut after = None; loop { - let page = utopia_store::export::facts_page(&pool, kb_id, after).await.map_err(io)?; + let page = utopia_store::export::facts_page(&mut tx, kb_id, after).await.map_err(io)?; let Some(last) = page.last() else { break }; after = Some(last.id); for f in &page { @@ -110,7 +127,7 @@ pub async fn export( let mut after = None; loop { - let page = utopia_store::export::derived_page(&pool, kb_id, after).await.map_err(io)?; + let page = utopia_store::export::derived_page(&mut tx, kb_id, after).await.map_err(io)?; let Some(last) = page.last() else { break }; after = Some(last.id); for d in &page { diff --git a/crates/utopia-server/src/api/mcp_tests.rs b/crates/utopia-server/src/api/mcp_tests.rs index 5e2f92adb..13e6316e3 100644 --- a/crates/utopia-server/src/api/mcp_tests.rs +++ b/crates/utopia-server/src/api/mcp_tests.rs @@ -30,6 +30,11 @@ impl Fixture { return Ok(None); }; let pool = sqlx::PgPool::connect(&url).await?; + Self::with_pool(pool).await.map(Some) + } + + /// 同一份种子,池子由调用方给——连池参数的测试(比如最小池)走这里 + async fn with_pool(pool: sqlx::PgPool) -> anyhow::Result { utopia_store::db::migrate(&pool).await?; let dir = std::env::temp_dir().join(format!("utopia-mcp-{}", Uuid::now_v7())); let search = Arc::new(utopia_search::SearchIndex::open(&dir.join("search"))?); @@ -148,7 +153,7 @@ impl Fixture { &f.document.to_string(), &[(f.chunk.to_string(), "orchard ".repeat(120))], )?; - Ok(Some(f)) + Ok(f) } async fn request( @@ -307,6 +312,82 @@ async fn record_axis_subseconds_survive_authenticated_rdf_export() -> anyhow::Re result.and(cleanup) } +/// 每个导出的快照事务活满整个流——它占住一条连接直到文件发完。台账如果排在 +/// 事务之后写,就是在「已经占了一条」的情况下再向池子要第二条:支持的最小池 +/// (2 条连接)上两个并发导出会互相把对方的审计饿死到超时。所以顺序必须是: +/// 先写完台账、放掉连接,再开始占着不放的长事务。两个导出都该落得下一行 +/// kb.exported,而不是在等一条永远不会来的连接 +#[tokio::test] +async fn concurrent_exports_on_a_minimum_pool_still_record_their_audits() -> anyhow::Result<()> { + use axum::body::{to_bytes, Body}; + use axum::http::{Request, StatusCode}; + use tower::ServiceExt; + + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + // 支持的最小池:两条连接。短的 acquire 超时只是为了不让失败的探测等太久—— + // 断言不依赖时钟,依赖台账行在不在 + let pool = sqlx::postgres::PgPoolOptions::new() + .max_connections(2) + .acquire_timeout(std::time::Duration::from_millis(400)) + .connect(&url) + .await?; + let f = Fixture::with_pool(pool).await?; + + let auth = utopia_store::tokens::authenticate(&f.state.pool, &f.token).await?; + let jwt = crate::auth::issue_token(&f.state, auth.user_id)?; + let app = crate::api::router(f.state.clone(), &Default::default()); + let uri = format!("/api/v1/kbs/{}/export?format=turtle", f.kb); + + let export = |app: axum::Router| { + let uri = uri.clone(); + let jwt = jwt.clone(); + async move { + let response = app + .oneshot( + Request::builder() + .uri(uri) + .header("authorization", format!("Bearer {jwt}")) + .body(Body::empty())?, + ) + .await?; + let status = response.status(); + let bytes = to_bytes(response.into_body(), 8 * 1024 * 1024).await?; + Ok::<(StatusCode, axum::body::Bytes), anyhow::Error>((status, bytes)) + } + }; + // 一次性失败上限:真饿死也只是多等几秒,不该挂着不走 + let (a, b) = tokio::time::timeout(std::time::Duration::from_secs(30), async { + let (a, b) = tokio::join!(export(app.clone()), export(app)); + (a.unwrap(), b.unwrap()) + }) + .await?; + anyhow::ensure!(a.0 == StatusCode::OK, "export A rejected: {}", a.0); + anyhow::ensure!(b.0 == StatusCode::OK, "export B rejected: {}", b.0); + anyhow::ensure!(!a.1.is_empty() && !b.1.is_empty(), "export body empty"); + + // 两份导出,两行台账——任何一份的审计被池子饿死这里都露馅 + let audits: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM audit_events WHERE kb_id = $1 AND action = 'kb.exported'", + ) + .bind(f.kb) + .fetch_one(&f.state.pool) + .await?; + anyhow::ensure!( + audits == 2, + "two exports must each record kb.exported, got {audits}" + ); + + // 两条流发完之后连接都得回家:接着借满整个池(两条)都该立刻拿到—— + // 快照事务没放下的话,这里就会撞 acquire 超时 + let c1 = f.state.pool.acquire().await?; + let c2 = f.state.pool.acquire().await?; + drop(c2); + drop(c1); + f.clean().await +} + #[tokio::test] async fn refused_and_executed_calls_are_each_audited_once() -> anyhow::Result<()> { let Some(f) = Fixture::new().await? else { @@ -622,7 +703,8 @@ async fn entity_facts_keeps_identity_values_filters_and_both_clocks() -> anyhow: assert!(derived["rule_id"].is_null()); uuid(&derived["attribute_rule_id"]); // The same UUID is the RDF statement's identity, not a newly minted response ID. - let exported = utopia_store::export::facts_page(&f.state.pool, f.kb, None).await?; + let exported = + utopia_store::export::facts_page(&mut f.state.pool.begin().await?, f.kb, None).await?; assert!(exported .iter() .any(|r| r.id == uuid(&corrected["id"]) && r.documents == vec![f.document])); diff --git a/crates/utopia-server/src/rdf.rs b/crates/utopia-server/src/rdf.rs index 0da8f3fc8..bcb2f9bc4 100644 --- a/crates/utopia-server/src/rdf.rs +++ b/crates/utopia-server/src/rdf.rs @@ -748,6 +748,14 @@ mod tests { documents: vec![], quotes: vec![], quote_origins: vec![], + subject_kb: Some(kb()), + object_kb: Some(kb()), + predicate_kb: Some(kb()), + supersedes_kb: None, + foreign_document: false, + foreign_chunk: false, + subject_merged: false, + object_merged: false, } } @@ -1164,6 +1172,15 @@ mod tests { rule_name: Some("Gas-bearing well".into()), premises: vec![id(5)], premises_derived: Vec::new(), + subject_kb: Some(kb()), + object_kb: None, + predicate_kb: Some(kb()), + rule_kb: None, + attribute_rule_kb: Some(kb()), + foreign_fact_premise: false, + foreign_derived_premise: false, + subject_merged: false, + object_merged: false, }; let quads = export(Format::Turtle, |sink, names, vocab| { emit_derived(sink, names, vocab, &derived).unwrap(); @@ -1289,6 +1306,15 @@ mod tests { rule_name: None, premises: vec![id(5)], premises_derived: vec![id(6)], + subject_kb: Some(kb()), + object_kb: Some(kb()), + predicate_kb: Some(kb()), + rule_kb: Some(kb()), + attribute_rule_kb: None, + foreign_fact_premise: false, + foreign_derived_premise: false, + subject_merged: false, + object_merged: false, }; for format in [Format::Turtle, Format::JsonLd] { let quads = export(format, |sink, names, vocab| { diff --git a/crates/utopia-store/src/export.rs b/crates/utopia-store/src/export.rs index 33fdc130e..360aff38a 100644 --- a/crates/utopia-store/src/export.rs +++ b/crates/utopia-store/src/export.rs @@ -6,15 +6,285 @@ //! //! 除本体外一律**按 id 分页**:一个库的事实可以有几十万条,全读进内存再序列化 //! 会在最需要它的那种部署上炸掉。id 是 uuid v7,按它排序即按写入顺序排序。 +//! 第一页没有下界(`$2 IS NULL OR id > $2`):uuid 没有更小的哨兵可垫—— +//! `id > NIL` 会把主键恰为 NIL 的合法行永远挡在导出外,而指向它的 +//! (document_id, version) 定位器照样解析,留下一条没有本体的边。 +//! +//! **同一份快照**:所有取数口吃调用方的事务(`&mut Transaction`),路由侧把它 +//! 定成只读 REPEATABLE READ——体检、词汇表与每一页读的是同一个时刻的库。 +//! 若每页各自向连接池要一条连接,词汇表发完之后才提交的规则会在派生页里 +//! 留下一条指向它的 `wasGeneratedBy`——图里就出现没有本体的引用。 +//! +//! **逐页校验用的是留下的那几行自己**:每个 page 查询把被引行的 +//! kb_id 与行本体**原子地一并选出**,校验在内存里跑。不能「先取一页、再去库里 +//! 问一次」——第二次问的是另一个时刻的状态,留下的行早已不是它。 +//! +//! 校验范围只覆盖**这份导出真正解析的引用**:会被铸成本库 IRI 的、会被按 id +//! 进本库词汇表查的。导出还没读的边(规则表自身、陈述属性、时间提及、 +//! 文档版本定位器等)不在这里管——它们随各自的导出面一起带上自己的校验。 use chrono::{DateTime, Utc}; -use sqlx::PgPool; -use utopia_core::AppResult; +use sqlx::{Postgres, Transaction}; +use utopia_core::{AppError, AppResult}; use uuid::Uuid; /// 一次取多少行。够大以免把往返次数拉满,够小以免一页就撑爆内存。 pub const PAGE: i64 = 500; +/// 出处链越界的一类引用。同库外键只认 id、不认库:A 库的 +/// 行可以引用 B 库的对象,schema 什么都不拦。0070 的触发器挡新行;这里拦的是 +/// **存量坏行**与绕过触发器写进来的行。 +/// +/// 处置一律是**整份拒导**:把别库对象的 id 铸进本库 IRI +/// (`urn:utopia:kb:A:document:{B 的文档}`)等于伪造身份——那份文件看着 +/// 完整,实则悬空。被词汇表解析的引用(谓词、属性类型、实体类型)越库则 +/// 静默消失——坏行一样不许放行。少导一行、换个 IRI 都不在选项里 +#[derive(Debug, sqlx::FromRow)] +pub struct CrossKbViolation { + /// 哪条边:evidence.chunk | derivation.premise_fact | … + pub edge: String, + pub rows: i64, +} + +fn cross_kb_error(violations: &[CrossKbViolation]) -> AppError { + let detail = violations + .iter() + .map(|v| format!("{}: {} row(s)", v.edge, v.rows)) + .collect::>() + .join("; "); + AppError::invalid_detail( + "cross_kb_provenance", + "export refused: KB-scoped provenance points outside the knowledge base", + detail, + ) +} + +/// 引用指着的东西**在本库,却不在导出集里**(合并掉的实体是唯一会缺席的 +/// 实体——`entities_page` 滤掉 `merged_into` 非空的行)。同库但缺席的引用 +/// 不能换 IRI、也不能静默省略:整份拒导,与越库同一处置。 +fn unexported_error(violations: &[CrossKbViolation]) -> AppError { + let detail = violations + .iter() + .map(|v| format!("{}: {} row(s)", v.edge, v.rows)) + .collect::>() + .join("; "); + AppError::invalid_detail( + "unexported_target", + "export refused: a reference points at a row that is not in this KB's exported set", + detail, + ) +} + +fn tally(violations: &mut Vec, edge: &str, rows: i64) { + if rows > 0 { + violations.push(CrossKbViolation { + edge: edge.into(), + rows, + }); + } +} + +/// 引用列的判定:**留下的是谁,就查谁**。`ref_kb` 由 page 查询与行本体原子地 +/// 一并选出——别库、悬空(NULL)都不算本库,一律判违规 +fn foreign(ref_kb: Option, kb_id: Uuid) -> bool { + ref_kb != Some(kb_id) +} + +/// 导出前的出处体检。逐类数一遍越界引用,有一行就整份拒导。 +/// 越界有两类,各自一条错:**别库/悬空**(`cross_kb`)与**同库但不在导出集** +/// (`unexported`——合并掉的实体)。只报哪条边坏了、坏了几行——具体哪些行 +/// 坏是库里的事,不进面向导出的报错 +/// +/// 扫的边与导出面一一对应:只查这份导出会解析的引用(IRI 会铸出去的、 +/// 词汇表会按 id 查的)。**必须在导出用的那条事务里跑**(REPEATABLE READ): +/// 体检与每一页查询看的是同一个快照,先体检后换连接会在两个时刻之间 +/// 漏掉刚提交的坏行 +pub async fn provenance_integrity( + tx: &mut Transaction<'_, Postgres>, + kb_id: Uuid, +) -> AppResult<()> { + #[derive(sqlx::FromRow)] + struct ScanViolation { + edge: String, + kind: String, + rows: i64, + } + let violations: Vec = sqlx::query_as( + "SELECT edge, kind, COUNT(*) AS rows FROM ( + -- 证据的段落:quote_origins 按它 JOIN chunks 取 origin——别库/悬空 + -- 的段落会让引文来源静默消失 + SELECT 'evidence.chunk'::text AS edge, 'cross_kb'::text AS kind, + c.kb_id IS DISTINCT FROM f.kb_id AS bad + FROM fact_evidence e + JOIN facts f ON f.id = e.fact_id + LEFT JOIN chunks c ON c.id = e.chunk_id + WHERE f.kb_id = $1 + UNION ALL + -- 证据的文档指针:铸成 prov:wasDerivedFrom 的文档 IRI + SELECT 'evidence.document', 'cross_kb', d.kb_id IS DISTINCT FROM f.kb_id + FROM fact_evidence e + JOIN facts f ON f.id = e.fact_id + LEFT JOIN documents d ON d.id = e.document_id + WHERE f.kb_id = $1 AND e.document_id IS NOT NULL + UNION ALL + -- 派生前提:铸成 prov:used 的事实/派生 IRI + SELECT 'derivation.premise_fact', 'cross_kb', p.kb_id IS DISTINCT FROM d.kb_id + FROM fact_derivations fd + JOIN derived_facts d ON d.id = fd.derived_fact_id + LEFT JOIN facts p ON p.id = fd.premise_fact_id + WHERE d.kb_id = $1 AND fd.premise_fact_id IS NOT NULL + UNION ALL + SELECT 'derivation.premise_derived', 'cross_kb', p.kb_id IS DISTINCT FROM d.kb_id + FROM fact_derivations fd + JOIN derived_facts d ON d.id = fd.derived_fact_id + LEFT JOIN derived_facts p ON p.id = fd.premise_derived_id + WHERE d.kb_id = $1 AND fd.premise_derived_id IS NOT NULL + UNION ALL + -- 边上的属性:类型进词汇表按 id 查(查不着静默丢),实体值铸 IRI + SELECT 'qualifier.type', 'cross_kb', r.kb_id IS DISTINCT FROM f.kb_id + FROM fact_qualifiers q + JOIN facts f ON f.id = q.fact_id + LEFT JOIN relation_types r ON r.id = q.qualifier_type_id + WHERE f.kb_id = $1 + UNION ALL + SELECT 'qualifier.entity', 'cross_kb', e.kb_id IS DISTINCT FROM f.kb_id + FROM fact_qualifiers q + JOIN facts f ON f.id = q.fact_id + LEFT JOIN entities e ON e.id = q.entity_id + WHERE f.kb_id = $1 AND q.entity_id IS NOT NULL + UNION ALL + SELECT 'qualifier.entity(merged)', 'unexported', TRUE + FROM fact_qualifiers q + JOIN facts f ON f.id = q.fact_id + JOIN entities e ON e.id = q.entity_id AND e.merged_into IS NOT NULL + WHERE f.kb_id = $1 + UNION ALL + -- 事实本体:主语铸 entity IRI,谓词进词汇表,supersedes 铸 fact IRI + SELECT 'fact.subject', 'cross_kb', s.kb_id IS DISTINCT FROM f.kb_id + FROM facts f LEFT JOIN entities s ON s.id = f.subject_id + WHERE f.kb_id = $1 + UNION ALL + SELECT 'fact.subject(merged)', 'unexported', TRUE + FROM facts f JOIN entities s ON s.id = f.subject_id AND s.merged_into IS NOT NULL + WHERE f.kb_id = $1 + UNION ALL + SELECT 'fact.object', 'cross_kb', o.kb_id IS DISTINCT FROM f.kb_id + FROM facts f LEFT JOIN entities o ON o.id = f.object_id + WHERE f.kb_id = $1 AND f.object_id IS NOT NULL + UNION ALL + SELECT 'fact.object(merged)', 'unexported', TRUE + FROM facts f JOIN entities o ON o.id = f.object_id AND o.merged_into IS NOT NULL + WHERE f.kb_id = $1 + UNION ALL + SELECT 'fact.predicate', 'cross_kb', r.kb_id IS DISTINCT FROM f.kb_id + FROM facts f LEFT JOIN relation_types r ON r.id = f.predicate_id + WHERE f.kb_id = $1 AND f.predicate_id IS NOT NULL + UNION ALL + SELECT 'fact.supersedes', 'cross_kb', s.kb_id IS DISTINCT FROM f.kb_id + FROM facts f LEFT JOIN facts s ON s.id = f.supersedes + WHERE f.kb_id = $1 AND f.supersedes IS NOT NULL + UNION ALL + -- 派生本体:规则 id 铸成 wasGeneratedBy 的 Activity IRI + SELECT 'derived.subject', 'cross_kb', s.kb_id IS DISTINCT FROM d.kb_id + FROM derived_facts d LEFT JOIN entities s ON s.id = d.subject_id + WHERE d.kb_id = $1 + UNION ALL + SELECT 'derived.subject(merged)', 'unexported', TRUE + FROM derived_facts d JOIN entities s ON s.id = d.subject_id AND s.merged_into IS NOT NULL + WHERE d.kb_id = $1 + UNION ALL + SELECT 'derived.object', 'cross_kb', o.kb_id IS DISTINCT FROM d.kb_id + FROM derived_facts d LEFT JOIN entities o ON o.id = d.object_id + WHERE d.kb_id = $1 AND d.object_id IS NOT NULL + UNION ALL + SELECT 'derived.object(merged)', 'unexported', TRUE + FROM derived_facts d JOIN entities o ON o.id = d.object_id AND o.merged_into IS NOT NULL + WHERE d.kb_id = $1 + UNION ALL + SELECT 'derived.predicate', 'cross_kb', r.kb_id IS DISTINCT FROM d.kb_id + FROM derived_facts d LEFT JOIN relation_types r ON r.id = d.predicate_id + WHERE d.kb_id = $1 + UNION ALL + SELECT 'derived.rule', 'cross_kb', r.kb_id IS DISTINCT FROM d.kb_id + FROM derived_facts d LEFT JOIN rules r ON r.id = d.rule_id + WHERE d.kb_id = $1 AND d.rule_id IS NOT NULL + UNION ALL + SELECT 'derived.attribute_rule', 'cross_kb', r.kb_id IS DISTINCT FROM d.kb_id + FROM derived_facts d LEFT JOIN attribute_rules r ON r.id = d.attribute_rule_id + WHERE d.kb_id = $1 AND d.attribute_rule_id IS NOT NULL + UNION ALL + -- 实体的类进词汇表按 id 查 + SELECT 'entity.type', 'cross_kb', t.kb_id IS DISTINCT FROM e.kb_id + FROM entities e LEFT JOIN entity_types t ON t.id = e.type_id + WHERE e.kb_id = $1 AND e.type_id IS NOT NULL + UNION ALL + -- 类层级与互斥都进词汇表按 id 查 + SELECT 'class.parent', 'cross_kb', p.kb_id IS DISTINCT FROM c.kb_id + FROM entity_type_parents x + JOIN entity_types c ON c.id = x.child_id + LEFT JOIN entity_types p ON p.id = x.parent_id + WHERE c.kb_id = $1 + UNION ALL + SELECT 'class.disjoint', 'cross_kb', a.kb_id IS DISTINCT FROM dd.kb_id + FROM entity_type_disjoint dd + LEFT JOIN entity_types a ON a.id = dd.a_id + WHERE dd.kb_id = $1 + UNION ALL + SELECT 'class.disjoint', 'cross_kb', b.kb_id IS DISTINCT FROM dd.kb_id + FROM entity_type_disjoint dd + LEFT JOIN entity_types b ON b.id = dd.b_id + WHERE dd.kb_id = $1 + UNION ALL + -- domain/range 进词汇表按 id 查;inverse/sub_property 铸关系 IRI + SELECT 'relation.domain', 'cross_kb', t.kb_id IS DISTINCT FROM r.kb_id + FROM relation_type_domains x + JOIN relation_types r ON r.id = x.relation_type_id + LEFT JOIN entity_types t ON t.id = x.entity_type_id + WHERE r.kb_id = $1 + UNION ALL + SELECT 'relation.range', 'cross_kb', t.kb_id IS DISTINCT FROM r.kb_id + FROM relation_type_ranges x + JOIN relation_types r ON r.id = x.relation_type_id + LEFT JOIN entity_types t ON t.id = x.entity_type_id + WHERE r.kb_id = $1 + UNION ALL + SELECT 'relation.inverse', 'cross_kb', t.kb_id IS DISTINCT FROM r.kb_id + FROM relation_types r LEFT JOIN relation_types t ON t.id = r.inverse_of + WHERE r.kb_id = $1 AND r.inverse_of IS NOT NULL + UNION ALL + SELECT 'relation.sub_property', 'cross_kb', t.kb_id IS DISTINCT FROM r.kb_id + FROM relation_types r LEFT JOIN relation_types t ON t.id = r.sub_property_of + WHERE r.kb_id = $1 AND r.sub_property_of IS NOT NULL + ) refs WHERE bad GROUP BY edge, kind", + ) + .bind(kb_id) + .fetch_all(&mut **tx) + .await?; + let cross_kb: Vec = violations + .iter() + .filter(|v| v.kind == "cross_kb") + .map(|v| CrossKbViolation { + edge: v.edge.clone(), + rows: v.rows, + }) + .collect(); + let unexported: Vec = violations + .iter() + .filter(|v| v.kind == "unexported") + .map(|v| CrossKbViolation { + edge: v.edge.clone(), + rows: v.rows, + }) + .collect(); + if !cross_kb.is_empty() { + return Err(cross_kb_error(&cross_kb)); + } + if !unexported.is_empty() { + return Err(unexported_error(&unexported)); + } + Ok(()) +} + #[derive(Debug, Clone, sqlx::FromRow)] pub struct ExportClass { pub id: Uuid, @@ -46,6 +316,8 @@ pub struct ExportRelation { pub is_symmetric: bool, pub is_asymmetric: bool, pub is_irreflexive: bool, + /// 同表自指:导出 owl:inverseOf / rdfs:subPropertyOf。存量的合法性在 + /// 迁移的递延约束里管,这里只负责把集合带上(越库/悬空 → 拒导) pub inverse_of: Option, pub sub_property_of: Option, pub domains: Vec, @@ -57,6 +329,9 @@ pub struct ExportEntity { pub id: Uuid, pub canonical_name: String, pub type_id: Option, + /// type_id 指着的类的 kb(LEFT JOIN 一并选出)。别库/悬空 → 拒导: + /// 序列化按 id 进本库词汇表查类,查不着就是静默丢类型 + pub type_kb: Option, } #[derive(Debug, Clone, sqlx::FromRow)] @@ -87,6 +362,20 @@ pub struct ExportFact { pub quotes: Vec, /// 证据文字的来源(0040),去重:一条陈述的引文里有没有扫描、转写、看图描述来的 pub quote_origins: Vec, + /// 以下各列是被引行的 kb,与行本体原子地一并选出。 + /// 别库或悬空(NULL)的引用不许被铸成本库 IRI,也不许静默跳过 + pub subject_kb: Option, + pub object_kb: Option, + pub predicate_kb: Option, + pub supersedes_kb: Option, + /// documents[] 里是否有别库或悬空的文档指针 + pub foreign_document: bool, + /// 引文 JOIN 的段落里是否有别库或悬空的(quote_origins 的来源会静默消失) + pub foreign_chunk: bool, + /// 主语/宾语指着**已合并**的实体:同库但不在导出集(merged_into IS NOT NULL + /// 的行 entities_page 不导)。与行本体原子地一并选出 + pub subject_merged: bool, + pub object_merged: bool, } #[derive(Debug, Clone, sqlx::FromRow)] @@ -116,6 +405,20 @@ pub struct ExportDerived { /// 前提里是**另一条派生**的那些(0030)。与上面那列分开,是因为读回来的人 /// 要知道该去哪张表接着往下走;合成一列的话,一条链在导出里就断了 pub premises_derived: Vec, + /// 被引行的 kb,与行本体原子地一并选出 + pub subject_kb: Option, + pub object_kb: Option, + pub predicate_kb: Option, + /// rule_id / attribute_rule_id 指着的规则行的 kb——它铸成 wasGeneratedBy + /// 的 Activity IRI + pub rule_kb: Option, + pub attribute_rule_kb: Option, + /// premises[] / premises_derived[] 里是否有别库或悬空的前提(两种前提分开报边) + pub foreign_fact_premise: bool, + pub foreign_derived_premise: bool, + /// 主语/宾语指着已合并的实体:同库但不在导出集 + pub subject_merged: bool, + pub object_merged: bool, } #[derive(Debug, Clone, sqlx::FromRow)] @@ -129,8 +432,11 @@ pub struct ExportDocument { pub deleted_at: Option>, } -pub async fn classes(pool: &PgPool, kb_id: Uuid) -> AppResult> { - Ok(sqlx::query_as( +pub async fn classes( + tx: &mut Transaction<'_, Postgres>, + kb_id: Uuid, +) -> AppResult> { + let classes: Vec = sqlx::query_as( "SELECT t.id, t.key, t.label, t.description, t.iri, COALESCE(ARRAY(SELECT p.parent_id FROM entity_type_parents p WHERE p.child_id = t.id ORDER BY p.parent_id), '{}') AS parents, @@ -141,12 +447,29 @@ pub async fn classes(pool: &PgPool, kb_id: Uuid) -> AppResult> FROM entity_types t WHERE t.kb_id = $1 ORDER BY t.key", ) .bind(kb_id) - .fetch_all(pool) - .await?) + .fetch_all(&mut **tx) + .await?; + // 父类与互斥类稍后要进本库词汇表按 id 查——查不着就是被静默丢掉。 + // 能解析的就地解析:词汇表全集就在手里,不在集合里的引用就是越库/悬空 + let own: std::collections::HashSet = classes.iter().map(|c| c.id).collect(); + let mut violations = Vec::new(); + for c in &classes { + let bad_parents = c.parents.iter().filter(|p| !own.contains(p)).count() as i64; + let bad_disjoint = c.disjoint.iter().filter(|p| !own.contains(p)).count() as i64; + tally(&mut violations, "class.parent", bad_parents); + tally(&mut violations, "class.disjoint", bad_disjoint); + } + if !violations.is_empty() { + return Err(cross_kb_error(&violations)); + } + Ok(classes) } -pub async fn relations(pool: &PgPool, kb_id: Uuid) -> AppResult> { - Ok(sqlx::query_as( +pub async fn relations( + tx: &mut Transaction<'_, Postgres>, + kb_id: Uuid, +) -> AppResult> { + let relations: Vec = sqlx::query_as( "SELECT r.id, r.key, r.label, r.description, r.iri, r.kind, r.datatype, r.unit, r.temporal, r.functional, r.inverse_functional, r.is_transitive, r.is_symmetric, r.is_asymmetric, r.is_irreflexive, @@ -158,32 +481,80 @@ pub async fn relations(pool: &PgPool, kb_id: Uuid) -> AppResult = sqlx::query_scalar( + "SELECT COALESCE(ARRAY(SELECT id FROM entity_types WHERE kb_id = $1), '{}')", + ) + .bind(kb_id) + .fetch_one(&mut **tx) + .await?; + let own: std::collections::HashSet = own_types.into_iter().collect(); + let own_rel: std::collections::HashSet = relations.iter().map(|r| r.id).collect(); + let mut violations = Vec::new(); + for r in &relations { + let bad_domains = r.domains.iter().filter(|t| !own.contains(t)).count() as i64; + let bad_ranges = r.ranges.iter().filter(|t| !own.contains(t)).count() as i64; + tally(&mut violations, "relation.domain", bad_domains); + tally(&mut violations, "relation.range", bad_ranges); + if let Some(t) = r.inverse_of { + tally( + &mut violations, + "relation.inverse", + (!own_rel.contains(&t)) as i64, + ); + } + if let Some(t) = r.sub_property_of { + tally( + &mut violations, + "relation.sub_property", + (!own_rel.contains(&t)) as i64, + ); + } + } + if !violations.is_empty() { + return Err(cross_kb_error(&violations)); + } + Ok(relations) } /// 合并掉的实体不导出:它已经不是一个东西了,它的事实早已搬到留下的那个身上。 pub async fn entities_page( - pool: &PgPool, + tx: &mut Transaction<'_, Postgres>, kb_id: Uuid, after: Option, ) -> AppResult> { - Ok(sqlx::query_as( - "SELECT id, canonical_name, type_id FROM entities - WHERE kb_id = $1 AND merged_into IS NULL AND id > COALESCE($2, '00000000-0000-0000-0000-000000000000'::uuid) - ORDER BY id LIMIT $3", + let page: Vec = sqlx::query_as( + "SELECT e.id, e.canonical_name, e.type_id, t.kb_id AS type_kb + FROM entities e LEFT JOIN entity_types t ON t.id = e.type_id + WHERE e.kb_id = $1 AND e.merged_into IS NULL + AND ($2 IS NULL OR e.id > $2) + ORDER BY e.id LIMIT $3", ) .bind(kb_id) .bind(after) .bind(PAGE) - .fetch_all(pool) - .await?) + .fetch_all(&mut **tx) + .await?; + let mut violations = Vec::new(); + for e in &page { + if e.type_id.is_some() && foreign(e.type_kb, kb_id) { + tally(&mut violations, "entity.type", 1); + } + } + if !violations.is_empty() { + return Err(cross_kb_error(&violations)); + } + Ok(page) } /// **不过滤 `invalidated_at`。** 撤回的、被修正顶掉的、区间早已闭合的,全在里面 /// ——它们各自带着两根轴上的时刻,读的人自己判断当时成立不成立(0019、0020)。 pub async fn facts_page( - pool: &PgPool, + tx: &mut Transaction<'_, Postgres>, kb_id: Uuid, after: Option, ) -> AppResult> { @@ -202,10 +573,26 @@ pub async fn facts_page( ORDER BY e.chunk_id), '{{}}') AS quotes, COALESCE(ARRAY(SELECT DISTINCT c.origin FROM fact_evidence e JOIN chunks c ON c.id = e.chunk_id - WHERE e.fact_id = f.id - ORDER BY c.origin), '{{}}') AS quote_origins + WHERE e.fact_id = f.id ORDER BY c.origin), '{{}}') + AS quote_origins, + s.kb_id AS subject_kb, o.kb_id AS object_kb, + p.kb_id AS predicate_kb, sp.kb_id AS supersedes_kb, + EXISTS(SELECT 1 FROM fact_evidence e + LEFT JOIN documents ed ON ed.id = e.document_id + WHERE e.fact_id = f.id AND e.document_id IS NOT NULL + AND ed.kb_id IS DISTINCT FROM f.kb_id) AS foreign_document, + EXISTS(SELECT 1 FROM fact_evidence e + LEFT JOIN chunks ec ON ec.id = e.chunk_id + WHERE e.fact_id = f.id + AND ec.kb_id IS DISTINCT FROM f.kb_id) AS foreign_chunk, + (s.merged_into IS NOT NULL) AS subject_merged, + (o.merged_into IS NOT NULL) AS object_merged FROM facts f - WHERE f.kb_id = $1 AND f.id > COALESCE($2, '00000000-0000-0000-0000-000000000000'::uuid) + LEFT JOIN entities s ON s.id = f.subject_id + LEFT JOIN entities o ON o.id = f.object_id + LEFT JOIN relation_types p ON p.id = f.predicate_id + LEFT JOIN facts sp ON sp.id = f.supersedes + WHERE f.kb_id = $1 AND ($2 IS NULL OR f.id > $2) ORDER BY f.id LIMIT $3", holds_from = crate::world_axis::facts_holds_from("f"), holds_to = crate::world_axis::facts_holds_to("f"), @@ -213,12 +600,68 @@ pub async fn facts_page( .bind(kb_id) .bind(after) .bind(PAGE) - .fetch_all(pool) + .fetch_all(&mut **tx) .await?; - // 边上的属性另一张表(0037),按事实 id 一次取回补上 + // 留下的行逐个过:被铸成本库 IRI 的引用不许越库,进词汇表的引用不许查空。 + // 检查用的是随行原子选出的 ref_kb——不是再去库里问一次的另一个时刻 + let mut violations = Vec::new(); + let mut unexported = Vec::new(); + for f in &facts { + tally( + &mut violations, + "fact.subject", + foreign(f.subject_kb, kb_id) as i64, + ); + tally( + &mut unexported, + "fact.subject(merged)", + f.subject_merged as i64, + ); + if f.object_id.is_some() { + tally( + &mut violations, + "fact.object", + foreign(f.object_kb, kb_id) as i64, + ); + tally( + &mut unexported, + "fact.object(merged)", + f.object_merged as i64, + ); + } + if f.predicate_id.is_some() { + tally( + &mut violations, + "fact.predicate", + foreign(f.predicate_kb, kb_id) as i64, + ); + } + if f.supersedes.is_some() { + tally( + &mut violations, + "fact.supersedes", + foreign(f.supersedes_kb, kb_id) as i64, + ); + } + tally( + &mut violations, + "evidence.document", + f.foreign_document as i64, + ); + tally(&mut violations, "evidence.chunk", f.foreign_chunk as i64); + } + if !violations.is_empty() { + return Err(cross_kb_error(&violations)); + } + if !unexported.is_empty() { + return Err(unexported_error(&unexported)); + } + // 边上的属性另一张表(0037),按事实 id 一次取回补上——把属性类型的 kb 与 + // 实体值的 kb 一并选出:别库类型会被词汇表静默跳过,别库实体会被铸进本库 + // IRI,两种行都得在序列化之前拦下来 { let ids: Vec = facts.iter().map(|f| f.id).collect(); - let mut by_fact = crate::graph::fact_qualifiers_for(pool, &ids).await?; + let mut by_fact = qualifiers_for_export(tx, &ids, kb_id).await?; for f in facts.iter_mut() { if let Some(q) = by_fact.remove(&f.id) { f.qualifiers = q; @@ -228,12 +671,97 @@ pub async fn facts_page( Ok(facts) } +#[derive(sqlx::FromRow)] +struct ExportQualifierRow { + fact_id: Uuid, + qualifier_type_id: Uuid, + key: Option, + label: Option, + value: Option, + entity_id: Option, + entity_name: Option, + type_kb: Option, + entity_kb: Option, + entity_merged: bool, +} + +/// 导出专用的属性取数:与 graph::fact_qualifiers_for 同一形状,多选两列 kb。 +/// 那边用 INNER JOIN——被引类型不在时整行属性**静默消失**;导出不能这么吞。 +/// 别库的属性类型会被词汇表查空而静默跳过,别库的实体值会被铸进本库 +/// entity IRI——两种都按坏行拒 +async fn qualifiers_for_export( + tx: &mut Transaction<'_, Postgres>, + fact_ids: &[Uuid], + kb_id: Uuid, +) -> AppResult>> { + let mut out: std::collections::HashMap> = + std::collections::HashMap::new(); + if fact_ids.is_empty() { + return Ok(out); + } + let rows: Vec = sqlx::query_as( + "SELECT q.fact_id, q.qualifier_type_id, r.key, r.label, q.value, q.entity_id, + e.canonical_name AS entity_name, + r.kb_id AS type_kb, e.kb_id AS entity_kb, + (e.merged_into IS NOT NULL) AS entity_merged + FROM fact_qualifiers q + LEFT JOIN relation_types r ON r.id = q.qualifier_type_id + LEFT JOIN entities e ON e.id = q.entity_id + WHERE q.fact_id = ANY($1) + ORDER BY q.fact_id, r.key", + ) + .bind(fact_ids) + .fetch_all(&mut **tx) + .await?; + let mut violations = Vec::new(); + let mut unexported = Vec::new(); + for r in &rows { + tally( + &mut violations, + "qualifier.type", + foreign(r.type_kb, kb_id) as i64, + ); + if r.entity_id.is_some() { + tally( + &mut violations, + "qualifier.entity", + foreign(r.entity_kb, kb_id) as i64, + ); + tally( + &mut unexported, + "qualifier.entity(merged)", + r.entity_merged as i64, + ); + } + } + if !violations.is_empty() { + return Err(cross_kb_error(&violations)); + } + if !unexported.is_empty() { + return Err(unexported_error(&unexported)); + } + for r in rows { + out.entry(r.fact_id) + .or_default() + .push(utopia_core::models::FactQualifier { + qualifier_type_id: r.qualifier_type_id, + // 过了校验 r 必然在:key/label 不会取不到 + key: r.key.unwrap_or_default(), + label: r.label.unwrap_or_default(), + value: r.value, + entity_id: r.entity_id, + entity_name: r.entity_name, + }); + } + Ok(out) +} + pub async fn derived_page( - pool: &PgPool, + tx: &mut Transaction<'_, Postgres>, kb_id: Uuid, after: Option, ) -> AppResult> { - Ok(sqlx::query_as( + let page: Vec = sqlx::query_as( // **两个 LEFT JOIN。** 表拓宽之后(0021)派生可能没有实体宾语、 // 也可能来自业务规则而不是公理——内连接会把这类结论整条挡在导出之外, // 而 0020 承诺的正是「审计员不靠我们也能读全」 @@ -249,34 +777,112 @@ pub async fn derived_page( COALESCE(ARRAY(SELECT fd.premise_derived_id FROM fact_derivations fd WHERE fd.derived_fact_id = d.id AND fd.premise_derived_id IS NOT NULL - ORDER BY fd.seq), '{}') AS premises_derived + ORDER BY fd.seq), '{}') AS premises_derived, + s.kb_id AS subject_kb, o.kb_id AS object_kb, p.kb_id AS predicate_kb, + ru.kb_id AS rule_kb, ar.kb_id AS attribute_rule_kb, + EXISTS(SELECT 1 FROM fact_derivations fd + LEFT JOIN facts pf ON pf.id = fd.premise_fact_id + WHERE fd.derived_fact_id = d.id AND fd.premise_fact_id IS NOT NULL + AND pf.kb_id IS DISTINCT FROM d.kb_id) AS foreign_fact_premise, + EXISTS(SELECT 1 FROM fact_derivations fd + LEFT JOIN derived_facts pd ON pd.id = fd.premise_derived_id + WHERE fd.derived_fact_id = d.id AND fd.premise_derived_id IS NOT NULL + AND pd.kb_id IS DISTINCT FROM d.kb_id) AS foreign_derived_premise, + (s.merged_into IS NOT NULL) AS subject_merged, + (o.merged_into IS NOT NULL) AS object_merged FROM derived_facts d LEFT JOIN rules ru ON ru.id = d.rule_id LEFT JOIN attribute_rules ar ON ar.id = d.attribute_rule_id - WHERE d.kb_id = $1 AND d.id > COALESCE($2, '00000000-0000-0000-0000-000000000000'::uuid) + LEFT JOIN entities s ON s.id = d.subject_id + LEFT JOIN entities o ON o.id = d.object_id + LEFT JOIN relation_types p ON p.id = d.predicate_id + WHERE d.kb_id = $1 AND ($2 IS NULL OR d.id > $2) ORDER BY d.id LIMIT $3", ) .bind(kb_id) .bind(after) .bind(PAGE) - .fetch_all(pool) - .await?) + .fetch_all(&mut **tx) + .await?; + let mut violations = Vec::new(); + let mut unexported = Vec::new(); + for d in &page { + tally( + &mut violations, + "derived.subject", + foreign(d.subject_kb, kb_id) as i64, + ); + tally( + &mut unexported, + "derived.subject(merged)", + d.subject_merged as i64, + ); + if d.object_id.is_some() { + tally( + &mut violations, + "derived.object", + foreign(d.object_kb, kb_id) as i64, + ); + tally( + &mut unexported, + "derived.object(merged)", + d.object_merged as i64, + ); + } + // 派生的谓词非空(CHECK 保证),NULL 的 ref_kb 一样是越界/悬空 + tally( + &mut violations, + "derived.predicate", + foreign(d.predicate_kb, kb_id) as i64, + ); + if d.rule_id.is_some() { + tally( + &mut violations, + "derived.rule", + foreign(d.rule_kb, kb_id) as i64, + ); + } + if d.attribute_rule_id.is_some() { + tally( + &mut violations, + "derived.attribute_rule", + foreign(d.attribute_rule_kb, kb_id) as i64, + ); + } + tally( + &mut violations, + "derivation.premise_fact", + d.foreign_fact_premise as i64, + ); + tally( + &mut violations, + "derivation.premise_derived", + d.foreign_derived_premise as i64, + ); + } + if !violations.is_empty() { + return Err(cross_kb_error(&violations)); + } + if !unexported.is_empty() { + return Err(unexported_error(&unexported)); + } + Ok(page) } pub async fn documents_page( - pool: &PgPool, + tx: &mut Transaction<'_, Postgres>, kb_id: Uuid, after: Option, ) -> AppResult> { Ok(sqlx::query_as( "SELECT id, filename, external_key, doc_time, created_at, deleted_at FROM documents - WHERE kb_id = $1 AND id > COALESCE($2, '00000000-0000-0000-0000-000000000000'::uuid) + WHERE kb_id = $1 AND ($2 IS NULL OR id > $2) ORDER BY id LIMIT $3", ) .bind(kb_id) .bind(after) .bind(PAGE) - .fetch_all(pool) + .fetch_all(&mut **tx) .await?) } diff --git a/crates/utopia-store/tests/a_merged_target_stays_out_of_the_export.rs b/crates/utopia-store/tests/a_merged_target_stays_out_of_the_export.rs new file mode 100644 index 000000000..18a1f53fe --- /dev/null +++ b/crates/utopia-store/tests/a_merged_target_stays_out_of_the_export.rs @@ -0,0 +1,288 @@ +//! 同库不等于在导出集里:`entities_page` 滤掉 `merged_into` +//! 非空的行,但 merge 只改写 fact 的主语/宾语——属性里的实体值、派生的 +//! 主语/宾语仍可能指着已合并的行。序列化照铸它的 IRI 就是一条悬空边。 +//! +//! 判法:**同库但不在导出集**也是越界——与越库同一处置,整份拒导。不 +//! 重写指向留下的实体(那是另一个语义动作),也不静默省略。 + +use sqlx::PgPool; +use uuid::Uuid; + +struct Fixture { + org: Uuid, + kb: Uuid, + survivor: Uuid, + merged: Uuid, + rel: Uuid, + attr: Uuid, + rule: Uuid, + fact: Uuid, +} + +/// 一个库、留着的实体与已合并的实体、一条事实、一条指向已合并实体的 +/// 派生——merged 是合法状态(没有触发器拦它),要拦的是导出侧 +async fn seed(pool: &PgPool) -> anyhow::Result { + let (org, ws, kb) = (Uuid::now_v7(), Uuid::now_v7(), Uuid::now_v7()); + sqlx::query("INSERT INTO organizations (id, name) VALUES ($1, 'merged-test')") + .bind(org) + .execute(pool) + .await?; + sqlx::query("INSERT INTO workspaces (id, org_id, name) VALUES ($1, $2, 'merged-test')") + .bind(ws) + .bind(org) + .execute(pool) + .await?; + sqlx::query( + "INSERT INTO knowledge_bases (id, workspace_id, name) VALUES ($1, $2, 'merged-test')", + ) + .bind(kb) + .bind(ws) + .execute(pool) + .await?; + + let (survivor, merged) = (Uuid::now_v7(), Uuid::now_v7()); + sqlx::query("INSERT INTO entities (id, kb_id, canonical_name) VALUES ($1, $2, 'survivor')") + .bind(survivor) + .bind(kb) + .execute(pool) + .await?; + sqlx::query( + "INSERT INTO entities (id, kb_id, canonical_name, merged_into) VALUES ($1, $2, 'gone', $3)", + ) + .bind(merged) + .bind(kb) + .bind(survivor) + .execute(pool) + .await?; + + let (rel, attr) = (Uuid::now_v7(), Uuid::now_v7()); + sqlx::query( + "INSERT INTO relation_types (id, kb_id, key, label) VALUES ($1, $2, 'knows', 'knows')", + ) + .bind(rel) + .bind(kb) + .execute(pool) + .await?; + sqlx::query( + "INSERT INTO relation_types (id, kb_id, key, label, kind) VALUES ($1, $2, 'since', 'since', 'attribute')", + ) + .bind(attr) + .bind(kb) + .execute(pool) + .await?; + + let fact = Uuid::now_v7(); + sqlx::query( + "INSERT INTO facts (id, kb_id, subject_id, predicate_id, object_id, confidence) + VALUES ($1, $2, $3, $4, $3, 0.9)", + ) + .bind(fact) + .bind(kb) + .bind(survivor) + .bind(rel) + .execute(pool) + .await?; + + let rule = Uuid::now_v7(); + sqlx::query( + "INSERT INTO rules (id, kb_id, predicate_id, kind) VALUES ($1, $2, $3, 'transitive')", + ) + .bind(rule) + .bind(kb) + .bind(rel) + .execute(pool) + .await?; + + Ok(Fixture { + org, + kb, + survivor, + merged, + rel, + attr, + rule, + fact, + }) +} + +async fn cleanup(pool: &PgPool, f: &Fixture) -> anyhow::Result<()> { + sqlx::query("DELETE FROM knowledge_bases WHERE id = $1") + .bind(f.kb) + .execute(pool) + .await?; + sqlx::query("DELETE FROM organizations WHERE id = $1") + .bind(f.org) + .execute(pool) + .await?; + Ok(()) +} + +async fn insert_derived( + pool: &PgPool, + f: &Fixture, + subject: Uuid, + object: Uuid, +) -> anyhow::Result { + let id = Uuid::now_v7(); + sqlx::query( + "INSERT INTO derived_facts (id, kb_id, subject_id, predicate_id, object_id, rule_id) + VALUES ($1, $2, $3, $4, $5, $6)", + ) + .bind(id) + .bind(f.kb) + .bind(subject) + .bind(f.rel) + .bind(object) + .bind(f.rule) + .execute(pool) + .await?; + Ok(id) +} + +/// 已合并的实体不进导出集——它不在 entities 页里出现 +#[tokio::test] +async fn a_merged_entity_is_not_emitted() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + let f = seed(&pool).await?; + + let page = utopia_store::export::entities_page(&mut pool.begin().await?, f.kb, None).await?; + assert!(page.iter().any(|e| e.id == f.survivor)); + assert!( + !page.iter().any(|e| e.id == f.merged), + "merged 实体必须不在导出集里" + ); + + cleanup(&pool, &f).await +} + +/// 指着已合并实体的属性值:同库但缺席——体检与事实页都要拦 +#[tokio::test] +async fn a_qualifier_on_a_merged_entity_fails_closed() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + let f = seed(&pool).await?; + + sqlx::query( + "INSERT INTO fact_qualifiers (fact_id, qualifier_type_id, entity_id) + VALUES ($1, $2, $3)", + ) + .bind(f.fact) + .bind(f.attr) + .bind(f.merged) + .execute(&pool) + .await?; + + let err = utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.kb).await; + let msg = format!("{err:?}"); + assert!(err.is_err(), "qualifier→merged 的体检必须拒导"); + assert!( + msg.contains("qualifier.entity(merged)"), + "要报 qualifier.entity(merged): {msg}" + ); + assert!( + utopia_store::export::facts_page(&mut pool.begin().await?, f.kb, None) + .await + .is_err(), + "事实页也要拦下同一条边" + ); + + sqlx::query("DELETE FROM fact_qualifiers WHERE fact_id = $1") + .bind(f.fact) + .execute(&pool) + .await?; + cleanup(&pool, &f).await +} + +/// 派生的主语/宾语指着已合并的实体:同库但缺席 +#[tokio::test] +async fn a_derived_on_a_merged_entity_fails_closed() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + let f = seed(&pool).await?; + + let d_subj = insert_derived(&pool, &f, f.merged, f.survivor).await?; + let err = utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.kb).await; + let msg = format!("{err:?}"); + assert!(err.is_err(), "derived.subject→merged 的体检必须拒导"); + assert!( + msg.contains("derived.subject(merged)"), + "要报 derived.subject(merged): {msg}" + ); + sqlx::query("DELETE FROM derived_facts WHERE id = $1") + .bind(d_subj) + .execute(&pool) + .await?; + + let d_obj = insert_derived(&pool, &f, f.survivor, f.merged).await?; + let err = utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.kb).await; + let msg = format!("{err:?}"); + assert!(err.is_err(), "derived.object→merged 的体检必须拒导"); + assert!( + msg.contains("derived.object(merged)"), + "要报 derived.object(merged): {msg}" + ); + sqlx::query("DELETE FROM derived_facts WHERE id = $1") + .bind(d_obj) + .execute(&pool) + .await?; + + // 干净之后就放行——拒的是缺席的引用,不是库本身 + utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.kb).await?; + cleanup(&pool, &f).await +} + +/// 事实的主语/宾语指着已合并的实体(merge 没走到的旧写或绕过 store 的写): +/// 事实页与体检都要拦 +#[tokio::test] +async fn a_fact_on_a_merged_entity_fails_closed() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + let f = seed(&pool).await?; + + let bad_fact = Uuid::now_v7(); + sqlx::query( + "INSERT INTO facts (id, kb_id, subject_id, predicate_id, object_id, confidence) + VALUES ($1, $2, $3, $4, $5, 0.9)", + ) + .bind(bad_fact) + .bind(f.kb) + .bind(f.survivor) + .bind(f.rel) + .bind(f.merged) + .execute(&pool) + .await?; + + let err = utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.kb).await; + let msg = format!("{err:?}"); + assert!(err.is_err(), "fact.object→merged 的体检必须拒导"); + assert!( + msg.contains("fact.object(merged)"), + "要报 fact.object(merged): {msg}" + ); + assert!( + utopia_store::export::facts_page(&mut pool.begin().await?, f.kb, None) + .await + .is_err(), + "事实页也要拦下同一条边" + ); + + sqlx::query("DELETE FROM facts WHERE id = $1") + .bind(bad_fact) + .execute(&pool) + .await?; + utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.kb).await?; + cleanup(&pool, &f).await +} diff --git a/crates/utopia-store/tests/an_export_carries_the_whole_ledger.rs b/crates/utopia-store/tests/an_export_carries_the_whole_ledger.rs index 908b3e8e1..821f2de88 100644 --- a/crates/utopia-store/tests/an_export_carries_the_whole_ledger.rs +++ b/crates/utopia-store/tests/an_export_carries_the_whole_ledger.rs @@ -222,7 +222,7 @@ async fn an_export_reads_the_whole_ledger_not_the_current_view() -> anyhow::Resu let f = seed(&pool).await?; // 1. 事实:撤回的那条**在**。界面把它藏起来是对的,导出把它藏起来就是骗人 - let facts = utopia_store::export::facts_page(&pool, f.kb, None).await?; + let facts = utopia_store::export::facts_page(&mut pool.begin().await?, f.kb, None).await?; let ids: Vec = facts.iter().map(|x| x.id).collect(); assert!(ids.contains(&f.live)); assert!( @@ -244,7 +244,8 @@ async fn an_export_reads_the_whole_ledger_not_the_current_view() -> anyhow::Resu assert_eq!(bare.surface_predicate.as_deref(), Some("advises")); // 4. 实体:合并掉的那个是唯一该消失的东西 - let entities = utopia_store::export::entities_page(&pool, f.kb, None).await?; + let entities = + utopia_store::export::entities_page(&mut pool.begin().await?, f.kb, None).await?; let ids: Vec = entities.iter().map(|e| e.id).collect(); assert!(ids.contains(&f.kept)); assert!( @@ -253,21 +254,21 @@ async fn an_export_reads_the_whole_ledger_not_the_current_view() -> anyhow::Resu ); // 5. 文档:删掉的留着墓碑(#268)。抹掉出处等于抹掉证据链 - let docs = utopia_store::export::documents_page(&pool, f.kb, None).await?; + let docs = utopia_store::export::documents_page(&mut pool.begin().await?, f.kb, None).await?; let deleted = docs.iter().find(|d| d.id == f.deleted_doc).unwrap(); assert!(deleted.deleted_at.is_some()); // 6. 派生:带着规则和前提,审计顺着它走得到断言 - let derived = utopia_store::export::derived_page(&pool, f.kb, None).await?; + let derived = utopia_store::export::derived_page(&mut pool.begin().await?, f.kb, None).await?; let d = derived.iter().find(|d| d.id == f.derived).unwrap(); assert_eq!(d.rule, "transitive"); assert_eq!(d.premises, vec![f.live]); // 7. 词汇表:导入来的类留着原 IRI,公理位照抄 - let classes = utopia_store::export::classes(&pool, f.kb).await?; + let classes = utopia_store::export::classes(&mut pool.begin().await?, f.kb).await?; let person = classes.iter().find(|c| c.key == "person").unwrap(); assert_eq!(person.iri.as_deref(), Some("https://schema.org/Person")); - let relations = utopia_store::export::relations(&pool, f.kb).await?; + let relations = utopia_store::export::relations(&mut pool.begin().await?, f.kb).await?; assert!(relations .iter() .any(|r| r.key == "works_for" && r.functional)); diff --git a/crates/utopia-store/tests/an_open_statement_keeps_the_documents_words.rs b/crates/utopia-store/tests/an_open_statement_keeps_the_documents_words.rs index a7c59e09f..c24c985e2 100644 --- a/crates/utopia-store/tests/an_open_statement_keeps_the_documents_words.rs +++ b/crates/utopia-store/tests/an_open_statement_keeps_the_documents_words.rs @@ -226,7 +226,8 @@ async fn an_open_statement_shows_under_its_phrase_and_reuses_its_row() -> anyhow .expect("a path walks the open statement"); assert_eq!(direct.edges[0].predicate.as_deref(), Some("acquired")); - let exported = utopia_store::export::facts_page(&pool, f.kb, None).await?; + let exported = + utopia_store::export::facts_page(&mut pool.begin().await?, f.kb, None).await?; let x = exported .iter() .find(|x| x.id == fact) diff --git a/crates/utopia-store/tests/cross_kb_provenance_fails_closed.rs b/crates/utopia-store/tests/cross_kb_provenance_fails_closed.rs new file mode 100644 index 000000000..1e595b7cf --- /dev/null +++ b/crates/utopia-store/tests/cross_kb_provenance_fails_closed.rs @@ -0,0 +1,295 @@ +//! 出处链不许跨库(0070):新行被触发器挡下,存量坏行让导出整份拒绝。 +//! +//! `fact_evidence`/`chunks` 的外键只认 id 不认库——原生写入路径碰巧全是同库 +//! 构造,但 schema 什么也不拦。两层防: +//! 1. 触发器(0070)挡在一切写入路径下游,包括绕过 store 层的 SQL; +//! 2. 导出侧体检 + 逐页校验——存量坏行与绕过触发器进来的行,宁可整份拒导, +//! 也不能把别库对象的 id 铸进本库 IRI。 +//! +//! 坏行在测试里靠 `SET LOCAL session_replication_role='replica'` 制造:只关 +//! 本事务的触发器,不碰 catalog——`DISABLE TRIGGER` 是全局的,并行测试会把 +//! 对方断言的拒绝窗口撞没。行指着的东西都真实存在,只是不在同一个库—— +//! 正是线上会遇到的形态(比如从 0070 之前的备份恢复进来的旧行)。 + +use sqlx::{Acquire, PgPool}; +use uuid::Uuid; + +struct TwoKbs { + org: Uuid, + a: Uuid, + b: Uuid, + doc_a: Uuid, + doc_b: Uuid, + chunk_a: Uuid, + chunk_b: Uuid, + fact_a: Uuid, + fact_b: Uuid, +} + +/// 两个库、每库一份文档一段一实体一事实——跨库引用需要的合法零件 +async fn seed(pool: &PgPool) -> anyhow::Result { + let (org, ws, a, b) = ( + Uuid::now_v7(), + Uuid::now_v7(), + Uuid::now_v7(), + Uuid::now_v7(), + ); + let (doc_a, doc_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (chunk_a, chunk_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (ent_a, ent_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (fact_a, fact_b) = (Uuid::now_v7(), Uuid::now_v7()); + + sqlx::query("INSERT INTO organizations (id, name) VALUES ($1, 'crosskb-test')") + .bind(org) + .execute(pool) + .await?; + sqlx::query("INSERT INTO workspaces (id, org_id, name) VALUES ($1, $2, 'crosskb-test')") + .bind(ws) + .bind(org) + .execute(pool) + .await?; + for kb in [a, b] { + sqlx::query( + "INSERT INTO knowledge_bases (id, workspace_id, name) VALUES ($1, $2, 'crosskb-test')", + ) + .bind(kb) + .bind(ws) + .execute(pool) + .await?; + } + for (id, kb, name) in [(doc_a, a, "a.md"), (doc_b, b, "b.md")] { + sqlx::query( + "INSERT INTO documents (id, kb_id, filename, sha256, status, external_key) + VALUES ($1, $2, $3, $4, 'ready', $5)", + ) + .bind(id) + .bind(kb) + .bind(name) + .bind(format!("sha-{id}")) + .bind(format!("file:///{name}")) + .execute(pool) + .await?; + } + for (id, kb, doc) in [(chunk_a, a, doc_a), (chunk_b, b, doc_b)] { + sqlx::query( + "INSERT INTO chunks (id, kb_id, document_id, seq, text) VALUES ($1, $2, $3, 0, 'x')", + ) + .bind(id) + .bind(kb) + .bind(doc) + .execute(pool) + .await?; + } + for (id, kb) in [(ent_a, a), (ent_b, b)] { + sqlx::query("INSERT INTO entities (id, kb_id, canonical_name) VALUES ($1, $2, 'e')") + .bind(id) + .bind(kb) + .execute(pool) + .await?; + } + for (id, kb, subject, object) in [(fact_a, a, ent_a, ent_a), (fact_b, b, ent_b, ent_b)] { + sqlx::query("INSERT INTO facts (id, kb_id, subject_id, object_id, confidence) VALUES ($1, $2, $3, $4, 0.9)") + .bind(id) + .bind(kb) + .bind(subject) + .bind(object) + .execute(pool) + .await?; + } + Ok(TwoKbs { + org, + a, + b, + doc_a, + doc_b, + chunk_a, + chunk_b, + fact_a, + fact_b, + }) +} + +async fn cleanup(pool: &PgPool, f: &TwoKbs) -> anyhow::Result<()> { + for kb in [f.a, f.b] { + sqlx::query("DELETE FROM knowledge_bases WHERE id = $1") + .bind(kb) + .execute(pool) + .await?; + } + sqlx::query("DELETE FROM organizations WHERE id = $1") + .bind(f.org) + .execute(pool) + .await?; + Ok(()) +} + +#[tokio::test] +async fn new_cross_kb_writes_are_rejected() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + // 触发器是 0070 带来的:测试库可能还没迁移,这里先保证约束在场 + utopia_store::db::migrate(&pool).await?; + let f = seed(&pool).await?; + + // 事实引用别库段落:原生路径碰巧不会这么写,但 schema 从前不拦——现在拦 + let err = sqlx::query( + "INSERT INTO fact_evidence (fact_id, chunk_id, quote, document_id, doc_version) + VALUES ($1, $2, 'x', $3, 1)", + ) + .bind(f.fact_a) + .bind(f.chunk_b) + .bind(f.doc_b) + .execute(&pool) + .await; + assert!(err.is_err(), "fact→foreign chunk 必须被拒"); + + // 只坏冗余文档指针那一头:段落同库、文档别库 + let err = sqlx::query( + "INSERT INTO fact_evidence (fact_id, chunk_id, quote, document_id, doc_version) + VALUES ($1, $2, 'x', $3, 1)", + ) + .bind(f.fact_a) + .bind(f.chunk_a) + .bind(f.doc_b) + .execute(&pool) + .await; + assert!(err.is_err(), "evidence.document→foreign 必须被拒"); + + // 段落挂在别库文档下 + let err = sqlx::query( + "INSERT INTO chunks (id, kb_id, document_id, seq, text) + VALUES ($1, $2, $3, 0, 'x')", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.doc_b) + .execute(&pool) + .await; + assert!(err.is_err(), "chunk→foreign document 必须被拒"); + + // 权威写入路径同一条约束:add_evidence 走 store API 也一样被拒 + let err = utopia_store::graph::add_evidence(&pool, f.fact_a, f.chunk_b, Some("x"), None).await; + assert!(err.is_err(), "add_evidence 的跨库配对必须被拒"); + + // 过户:把文档挪到别的库,等于把指着它的行一次全变坏行 + let err = sqlx::query("UPDATE documents SET kb_id = $2 WHERE id = $1") + .bind(f.doc_a) + .bind(f.b) + .execute(&pool) + .await; + assert!(err.is_err(), "documents.kb_id 过户必须被拒"); + let err = sqlx::query("UPDATE facts SET kb_id = $2 WHERE id = $1") + .bind(f.fact_a) + .bind(f.b) + .execute(&pool) + .await; + assert!(err.is_err(), "facts.kb_id 过户必须被拒"); + + // 同库的正常写入不受影响(防误伤) + sqlx::query( + "INSERT INTO fact_evidence (fact_id, chunk_id, quote, document_id, doc_version) + VALUES ($1, $2, 'ok', $3, 1)", + ) + .bind(f.fact_a) + .bind(f.chunk_a) + .bind(f.doc_a) + .execute(&pool) + .await?; + utopia_store::graph::add_evidence(&pool, f.fact_b, f.chunk_b, Some("ok"), None).await?; + + cleanup(&pool, &f).await +} + +#[tokio::test] +async fn malformed_existing_rows_fail_the_export_closed() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + let f = seed(&pool).await?; + + // 存量坏行:触发器只拦落在它之后的写,这种行只能绕过它造——导出侧要接住。 + // SET LOCAL 只在本事务内关触发器,提交即恢复,不会撞掉并行测试的断言窗口 + let mut conn = pool.acquire().await?; + let mut tx = conn.begin().await?; + sqlx::query("SET LOCAL session_replication_role = 'replica'") + .execute(&mut *tx) + .await?; + + // 只坏段落那一头:fact_a → chunk_b(文档指针留 NULL,隔离变量) + sqlx::query("INSERT INTO fact_evidence (fact_id, chunk_id) VALUES ($1, $2)") + .bind(f.fact_a) + .bind(f.chunk_b) + .execute(&mut *tx) + .await?; + // 只坏冗余文档指针:fact_b → chunk_b 本身同库,指针却指 a 的文档 + sqlx::query("INSERT INTO fact_evidence (fact_id, chunk_id, document_id) VALUES ($1, $2, $3)") + .bind(f.fact_b) + .bind(f.chunk_b) + .bind(f.doc_a) + .execute(&mut *tx) + .await?; + // 段落挂在别库文档下 + let stray_chunk = Uuid::now_v7(); + sqlx::query( + "INSERT INTO chunks (id, kb_id, document_id, seq, text) VALUES ($1, $2, $3, 9, 'stray')", + ) + .bind(stray_chunk) + .bind(f.b) + .bind(f.doc_a) + .execute(&mut *tx) + .await?; + tx.commit().await?; + drop(conn); + + // 体检:库 A 坏在 evidence.chunk,库 B 坏在 evidence.document 与 chunk.document + let err_a = utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.a).await; + let msg_a = format!("{err_a:?}"); + assert!(err_a.is_err(), "库 A 的体检必须拒导"); + assert!( + msg_a.contains("evidence.chunk"), + "库 A 该报 evidence.chunk: {msg_a}" + ); + + let err_b = utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.b).await; + let msg_b = format!("{err_b:?}"); + assert!(err_b.is_err(), "库 B 的体检必须拒导"); + assert!( + msg_b.contains("evidence.document"), + "库 B 该报 evidence.document: {msg_b}" + ); + + // 逐页校验同样fail-closed:体检之后的 TOCTOU 坏行也不能漏出伪造 IRI。 + // 坏段落指针会顺着 facts_page 的 quote_origins 出去,坏文档指针顺着 + // documents[] 出去——两路都在事实页上拦 + assert!( + utopia_store::export::facts_page(&mut pool.begin().await?, f.a, None) + .await + .is_err(), + "fact_a 的 evidence.chunk 别库:事实页要拦" + ); + assert!( + utopia_store::export::facts_page(&mut pool.begin().await?, f.b, None) + .await + .is_err(), + "fact_b 的 evidence.document 别库:事实页要拦" + ); + + // 清掉坏行,导出立刻恢复——拒的是坏行,不是库本身 + sqlx::query("DELETE FROM fact_evidence WHERE fact_id IN ($1, $2)") + .bind(f.fact_a) + .bind(f.fact_b) + .execute(&pool) + .await?; + sqlx::query("DELETE FROM chunks WHERE id = $1") + .bind(stray_chunk) + .execute(&pool) + .await?; + utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.a).await?; + utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.b).await?; + + cleanup(&pool, &f).await +} diff --git a/crates/utopia-store/tests/export_surfaces.rs b/crates/utopia-store/tests/export_surfaces.rs new file mode 100644 index 000000000..49b66dee1 --- /dev/null +++ b/crates/utopia-store/tests/export_surfaces.rs @@ -0,0 +1,401 @@ +//! 导出取数面:单事务快照与游标边界。 +//! +//! 判定在导出取数层:`provenance_integrity` 与各 page 函数吃调用方的事务—— +//! 同一条连接里种行、读页、断言。除「中途提交」那条探针外(它要两条连接), +//! 所有种子与坏行都在一个**回滚的事务**里造:快照库上跑这组测试不会留下 +//! 任何一行。 +//! +//! 坏行靠 `SET LOCAL session_replication_role='replica'` 造:只关本事务的 +//! 触发器(0070 装上的那些也一并关),回滚即恢复——要模拟的正是绕过触发器 +//! 进来的存量坏行。 + +use sqlx::{Acquire, PgPool, Postgres, Transaction}; +use utopia_store::export; +use uuid::Uuid; + +struct Fixture { + a: Uuid, + b: Uuid, + doc_a: Uuid, + attr_a: Uuid, +} + +/// 两个库;A 库一份文档一段一实体一事实一条证据一属性,B 库空着备查。 +/// 全部在调用方的事务里落——回滚即清场,一个 DELETE 都不用 +async fn seed_tx(tx: &mut Transaction<'_, Postgres>) -> anyhow::Result { + let (org, ws, a, b) = ( + Uuid::now_v7(), + Uuid::now_v7(), + Uuid::now_v7(), + Uuid::now_v7(), + ); + let (doc_a, doc_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (chunk_a, ent_a, ent_b, fact_a, attr_a) = ( + Uuid::now_v7(), + Uuid::now_v7(), + Uuid::now_v7(), + Uuid::now_v7(), + Uuid::now_v7(), + ); + + sqlx::query("INSERT INTO organizations (id, name) VALUES ($1, 'export-test')") + .bind(org) + .execute(&mut **tx) + .await?; + sqlx::query("INSERT INTO workspaces (id, org_id, name) VALUES ($1, $2, 'export-test')") + .bind(ws) + .bind(org) + .execute(&mut **tx) + .await?; + for kb in [a, b] { + sqlx::query( + "INSERT INTO knowledge_bases (id, workspace_id, name) VALUES ($1, $2, 'export-test')", + ) + .bind(kb) + .bind(ws) + .execute(&mut **tx) + .await?; + } + for (id, kb, name) in [(doc_a, a, "a.md"), (doc_b, b, "b.md")] { + sqlx::query( + "INSERT INTO documents (id, kb_id, filename, sha256, status, external_key) + VALUES ($1, $2, $3, $4, 'ready', $5)", + ) + .bind(id) + .bind(kb) + .bind(name) + .bind(format!("sha-{id}")) + .bind(format!("file:///{name}")) + .execute(&mut **tx) + .await?; + } + sqlx::query( + "INSERT INTO chunks (id, kb_id, document_id, seq, text, doc_version) + VALUES ($1, $2, $3, 0, 'x', 1)", + ) + .bind(chunk_a) + .bind(a) + .bind(doc_a) + .execute(&mut **tx) + .await?; + for (id, kb) in [(ent_a, a), (ent_b, b)] { + sqlx::query("INSERT INTO entities (id, kb_id, canonical_name) VALUES ($1, $2, 'e')") + .bind(id) + .bind(kb) + .execute(&mut **tx) + .await?; + } + sqlx::query( + "INSERT INTO facts (id, kb_id, subject_id, object_id, confidence) + VALUES ($1, $2, $3, $4, 0.9)", + ) + .bind(fact_a) + .bind(a) + .bind(ent_a) + .bind(ent_a) + .execute(&mut **tx) + .await?; + sqlx::query( + "INSERT INTO fact_evidence (fact_id, chunk_id, document_id, doc_version) + VALUES ($1, $2, $3, 1)", + ) + .bind(fact_a) + .bind(chunk_a) + .bind(doc_a) + .execute(&mut **tx) + .await?; + sqlx::query( + "INSERT INTO relation_types (id, kb_id, key, label, kind, datatype) + VALUES ($1, $2, 'headcount', 'headcount', 'attribute', 'number')", + ) + .bind(attr_a) + .bind(a) + .execute(&mut **tx) + .await?; + Ok(Fixture { + a, + b, + doc_a, + attr_a, + }) +} + +/// 文档过户的极端形态(replica):文档挪到 B 之后,A 的证据行还指着它—— +/// A 的导出宁拒也不能把别库文档铸进本库 IRI;B 没指着它,照常放行 +#[tokio::test] +async fn a_reassigned_document_breaks_the_edges_that_point_at_it() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + + let mut tx = pool.begin().await?; + let f = seed_tx(&mut tx).await?; + sqlx::query("SET LOCAL session_replication_role = 'replica'") + .execute(&mut *tx) + .await?; + sqlx::query("UPDATE documents SET kb_id = $2 WHERE id = $1") + .bind(f.doc_a) + .bind(f.b) + .execute(&mut *tx) + .await?; + + // A 的导出整体拒:evidence.document 还指着那份已经不属于本库的文档 + let err = export::provenance_integrity(&mut tx, f.a).await; + let msg = format!("{err:?}"); + assert!(err.is_err(), "文档过户后 A 的出处链必须拒导"); + assert!( + msg.contains("evidence.document"), + "该报 evidence.document: {msg}" + ); + assert!(export::facts_page(&mut tx, f.a, None).await.is_err()); + + // B 收下了文档,但它没有任何指着 A 的行:它的体检照样过 + export::provenance_integrity(&mut tx, f.b).await?; + tx.rollback().await?; + Ok(()) +} + +/// 导出中途落下的写进不了这一份。tx 起 REPEATABLE READ 快照后, +/// 另一条连接提交一条新谓词+引用它的派生——本事务的词汇表页与派生页 +/// 都看不见它:没有半个进来的引用,也没有悬空的 wasGeneratedBy。 +/// 这条要两条连接,种子必须提交——清场照常走 +#[tokio::test] +async fn a_mid_stream_commit_stays_outside_the_snapshot() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + + let (org, ws, a) = (Uuid::now_v7(), Uuid::now_v7(), Uuid::now_v7()); + let ent_a = Uuid::now_v7(); + sqlx::query("INSERT INTO organizations (id, name) VALUES ($1, 'export-race')") + .bind(org) + .execute(&pool) + .await?; + sqlx::query("INSERT INTO workspaces (id, org_id, name) VALUES ($1, $2, 'export-race')") + .bind(ws) + .bind(org) + .execute(&pool) + .await?; + sqlx::query( + "INSERT INTO knowledge_bases (id, workspace_id, name) VALUES ($1, $2, 'export-race')", + ) + .bind(a) + .bind(ws) + .execute(&pool) + .await?; + sqlx::query("INSERT INTO entities (id, kb_id, canonical_name) VALUES ($1, $2, 'e')") + .bind(ent_a) + .bind(a) + .execute(&pool) + .await?; + + // 先埋一条谓词和一条引用它的派生,作为「快照内」基线 + let (rule0, pred0) = (Uuid::now_v7(), Uuid::now_v7()); + sqlx::query( + "INSERT INTO relation_types (id, kb_id, key, label, kind) + VALUES ($1, $2, 'p0', 'p0', 'relation')", + ) + .bind(pred0) + .bind(a) + .execute(&pool) + .await?; + sqlx::query( + "INSERT INTO rules (id, kb_id, predicate_id, kind) VALUES ($1, $2, $3, 'transitive')", + ) + .bind(rule0) + .bind(a) + .bind(pred0) + .execute(&pool) + .await?; + sqlx::query( + "INSERT INTO derived_facts (id, kb_id, subject_id, predicate_id, object_id, rule_id) + VALUES ($1, $2, $3, $4, $5, $6)", + ) + .bind(Uuid::now_v7()) + .bind(a) + .bind(ent_a) + .bind(pred0) + .bind(ent_a) + .bind(rule0) + .execute(&pool) + .await?; + + // 导出事务:只读 REPEATABLE READ,快照从第一条语句起钉死 + let mut conn = pool.acquire().await?; + let mut tx = conn.begin().await?; + sqlx::query("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ READ ONLY") + .execute(&mut *tx) + .await?; + export::provenance_integrity(&mut tx, a).await?; + let _ = export::entities_page(&mut tx, a, None).await?; + + // 中途:另一条连接提交一条新谓词+引用它的派生事实 + let (rule_late, pred_late, derived_late) = (Uuid::now_v7(), Uuid::now_v7(), Uuid::now_v7()); + sqlx::query( + "INSERT INTO relation_types (id, kb_id, key, label, kind) + VALUES ($1, $2, 'p_late', 'p_late', 'relation')", + ) + .bind(pred_late) + .bind(a) + .execute(&pool) + .await?; + sqlx::query( + "INSERT INTO rules (id, kb_id, predicate_id, kind) VALUES ($1, $2, $3, 'transitive')", + ) + .bind(rule_late) + .bind(a) + .bind(pred_late) + .execute(&pool) + .await?; + sqlx::query( + "INSERT INTO derived_facts (id, kb_id, subject_id, predicate_id, object_id, rule_id) + VALUES ($1, $2, $3, $4, $5, $6)", + ) + .bind(derived_late) + .bind(a) + .bind(ent_a) + .bind(pred_late) + .bind(ent_a) + .bind(rule_late) + .execute(&pool) + .await?; + + // 快照里:词汇表页、派生页都不该有中途进来的行——一致性是整份的 + let relations = export::relations(&mut tx, a).await?; + assert!( + !relations.iter().any(|r| r.id == pred_late), + "中途提交的谓词不许进这份导出" + ); + let derived = export::derived_page(&mut tx, a, None).await?; + assert!( + !derived.iter().any(|d| d.id == derived_late), + "中途提交的派生不许进这份导出——也就不会有悬空的 wasGeneratedBy" + ); + assert_eq!(derived.len(), 1, "快照内的那条还在"); + tx.rollback().await?; + drop(conn); + + // 新事务是新的快照:两条都该在 + let mut tx2 = pool.begin().await?; + let relations = export::relations(&mut tx2, a).await?; + assert!(relations.iter().any(|r| r.id == pred_late)); + let derived = export::derived_page(&mut tx2, a, None).await?; + assert_eq!(derived.len(), 2); + tx2.rollback().await?; + + sqlx::query("DELETE FROM knowledge_bases WHERE id = $1") + .bind(a) + .execute(&pool) + .await?; + sqlx::query("DELETE FROM organizations WHERE id = $1") + .bind(org) + .execute(&pool) + .await?; + Ok(()) +} + +/// NIL 是合法 uuid——schema 不拦它当主键。按序它排在最前;首页谓词若是 +/// `id > 哨兵`,这一行就永远进不了任何一页。而指向它的引用照样解析过去: +/// 节点缺席、边在场,导出里就悬一条没有本体的引用。所有按 id 翻页的 +/// 取数口——实体、事实、派生、文档——第一页都得把它翻出来 +#[tokio::test] +async fn a_nil_id_row_still_reaches_the_first_page() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + + let mut tx = pool.begin().await?; + let f = seed_tx(&mut tx).await?; + let nil = Uuid::nil(); + + // 每张走 id 游标的表各埋一行 NIL 主键;引用一律指回这些 NIL 行自己, + // 外键不因 NIL 失效——它们跟其他行一样合法 + sqlx::query("INSERT INTO entities (id, kb_id, canonical_name) VALUES ($1, $2, 'nil-e')") + .bind(nil) + .bind(f.a) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO documents (id, kb_id, filename, sha256, status, external_key) + VALUES ($1, $2, 'nil.md', 'nil-sha', 'ready', 'file:///nil.md')", + ) + .bind(nil) + .bind(f.a) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO facts (id, kb_id, subject_id, predicate_id, object_id, confidence) + VALUES ($1, $2, $3, $4, $3, 0.9)", + ) + .bind(nil) + .bind(f.a) + .bind(nil) + .bind(f.attr_a) + .execute(&mut *tx) + .await?; + let (pred_r, rule_r) = (Uuid::now_v7(), Uuid::now_v7()); + sqlx::query( + "INSERT INTO relation_types (id, kb_id, key, label, kind) + VALUES ($1, $2, 'rel_p', 'rel_p', 'relation')", + ) + .bind(pred_r) + .bind(f.a) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO rules (id, kb_id, predicate_id, kind) VALUES ($1, $2, $3, 'transitive')", + ) + .bind(rule_r) + .bind(f.a) + .bind(pred_r) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO derived_facts (id, kb_id, subject_id, predicate_id, object_id, rule_id) + VALUES ($1, $2, $3, $4, $3, $5)", + ) + .bind(nil) + .bind(f.a) + .bind(nil) + .bind(pred_r) + .bind(rule_r) + .execute(&mut *tx) + .await?; + + assert!( + export::entities_page(&mut tx, f.a, None) + .await? + .iter() + .any(|e| e.id == nil), + "entities 首页漏掉 NIL 行" + ); + assert!( + export::documents_page(&mut tx, f.a, None) + .await? + .iter() + .any(|d| d.id == nil), + "documents 首页漏掉 NIL 行" + ); + assert!( + export::facts_page(&mut tx, f.a, None) + .await? + .iter() + .any(|x| x.id == nil), + "facts 首页漏掉 NIL 行" + ); + assert!( + export::derived_page(&mut tx, f.a, None) + .await? + .iter() + .any(|d| d.id == nil), + "derived 首页漏掉 NIL 行" + ); + tx.rollback().await?; + Ok(()) +} diff --git a/crates/utopia-store/tests/exported_references_never_cross_a_kb.rs b/crates/utopia-store/tests/exported_references_never_cross_a_kb.rs new file mode 100644 index 000000000..15f36bc39 --- /dev/null +++ b/crates/utopia-store/tests/exported_references_never_cross_a_kb.rs @@ -0,0 +1,843 @@ +//! 写入侧:导出会触碰的每条边都被 0070 的外键/触发器挡住;导出侧只复查 +//! 这份导出真正解析的那些引用——别库/悬空的行宁可整份拒导,也不许 +//! 伪造 IRI 或静默丢语义。落到别库的下场分两种: +//! - 铸成本库 IRI 的引用(实体、事实、派生、文档、段落、规则)→ 伪造身份; +//! - 进本库词汇表按 id 查的引用(谓词、属性类型、实体类型、父类)→ 静默消失。 +//! +//! 两种都是坏行,触发器与导出侧一律拒。 +//! +//! 坏行靠 `SET LOCAL session_replication_role='replica'` 造——只关本事务的 +//! 触发器(含 FK 强制),提交即恢复;这同时允许造出「指着的行已不在」的 +//! 悬空引用,正是从 0070 之前的备份恢复进来的形态。 + +use sqlx::{Acquire, PgPool}; +use uuid::Uuid; + +struct Fixture { + org: Uuid, + a: Uuid, + b: Uuid, + doc_a: Uuid, + chunk_a: Uuid, + chunk_b: Uuid, + ent_a: Uuid, + ent_b: Uuid, + fact_a: Uuid, + fact_b: Uuid, + rel_a: Uuid, + rel_b: Uuid, + cls_a: Uuid, + cls_b: Uuid, + rule_a: Uuid, + rule_b: Uuid, + arule_a: Uuid, + arule_b: Uuid, + der_a: Uuid, + der_b: Uuid, +} + +/// 两库,每库一套能被引用的零件:文档/段落/实体/事实/谓词/类/公理规则/业务规则/派生 +async fn seed(pool: &PgPool) -> anyhow::Result { + let (org, ws, a, b) = ( + Uuid::now_v7(), + Uuid::now_v7(), + Uuid::now_v7(), + Uuid::now_v7(), + ); + let (doc_a, doc_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (chunk_a, chunk_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (ent_a, ent_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (fact_a, fact_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (rel_a, rel_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (cls_a, cls_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (rule_a, rule_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (arule_a, arule_b) = (Uuid::now_v7(), Uuid::now_v7()); + let (der_a, der_b) = (Uuid::now_v7(), Uuid::now_v7()); + + sqlx::query("INSERT INTO organizations (id, name) VALUES ($1, 'xkb-ref-test')") + .bind(org) + .execute(pool) + .await?; + sqlx::query("INSERT INTO workspaces (id, org_id, name) VALUES ($1, $2, 'xkb-ref-test')") + .bind(ws) + .bind(org) + .execute(pool) + .await?; + for kb in [a, b] { + sqlx::query( + "INSERT INTO knowledge_bases (id, workspace_id, name) VALUES ($1, $2, 'xkb-ref-test')", + ) + .bind(kb) + .bind(ws) + .execute(pool) + .await?; + } + sqlx::query( + "INSERT INTO documents (id, kb_id, filename, sha256, status, external_key) + VALUES ($1, $2, 'a.md', 'sha-a', 'ready', 'file:///a.md')", + ) + .bind(doc_a) + .bind(a) + .execute(pool) + .await?; + sqlx::query( + "INSERT INTO documents (id, kb_id, filename, sha256, status, external_key) + VALUES ($1, $2, 'b.md', 'sha-b', 'ready', 'file:///b.md')", + ) + .bind(doc_b) + .bind(b) + .execute(pool) + .await?; + sqlx::query( + "INSERT INTO chunks (id, kb_id, document_id, seq, text) VALUES ($1, $2, $3, 0, 'x')", + ) + .bind(chunk_a) + .bind(a) + .bind(doc_a) + .execute(pool) + .await?; + sqlx::query( + "INSERT INTO chunks (id, kb_id, document_id, seq, text) VALUES ($1, $2, $3, 0, 'x')", + ) + .bind(chunk_b) + .bind(b) + .bind(doc_b) + .execute(pool) + .await?; + for (id, kb) in [(ent_a, a), (ent_b, b)] { + sqlx::query("INSERT INTO entities (id, kb_id, canonical_name) VALUES ($1, $2, 'e')") + .bind(id) + .bind(kb) + .execute(pool) + .await?; + } + for (id, kb, s, o) in [(fact_a, a, ent_a, ent_a), (fact_b, b, ent_b, ent_b)] { + sqlx::query( + "INSERT INTO facts (id, kb_id, subject_id, object_id, confidence) + VALUES ($1, $2, $3, $4, 0.9)", + ) + .bind(id) + .bind(kb) + .bind(s) + .bind(o) + .execute(pool) + .await?; + } + for (id, kb) in [(rel_a, a), (rel_b, b)] { + sqlx::query("INSERT INTO relation_types (id, kb_id, key, label) VALUES ($1, $2, $3, 'r')") + .bind(id) + .bind(kb) + .bind(format!("r-{id}")) + .execute(pool) + .await?; + } + for (id, kb) in [(cls_a, a), (cls_b, b)] { + sqlx::query("INSERT INTO entity_types (id, kb_id, key, label) VALUES ($1, $2, $3, 'c')") + .bind(id) + .bind(kb) + .bind(format!("c-{id}")) + .execute(pool) + .await?; + } + for (id, kb, pred) in [(rule_a, a, rel_a), (rule_b, b, rel_b)] { + sqlx::query( + "INSERT INTO rules (id, kb_id, predicate_id, kind) VALUES ($1, $2, $3, 'transitive')", + ) + .bind(id) + .bind(kb) + .bind(pred) + .execute(pool) + .await?; + } + for (id, kb, ty, pred) in [(arule_a, a, cls_a, rel_a), (arule_b, b, cls_b, rel_b)] { + sqlx::query( + "INSERT INTO attribute_rules (id, kb_id, name, subject_type_id, conclusion, + conclude_predicate_id, conclude_value) + VALUES ($1, $2, $3, $4, 'attribute', $5, '42'::jsonb)", + ) + .bind(id) + .bind(kb) + .bind(format!("ar-{id}")) + .bind(ty) + .bind(pred) + .execute(pool) + .await?; + } + for (id, kb, s, p, o, r) in [ + (der_a, a, ent_a, rel_a, ent_a, rule_a), + (der_b, b, ent_b, rel_b, ent_b, rule_b), + ] { + sqlx::query( + "INSERT INTO derived_facts (id, kb_id, subject_id, predicate_id, object_id, rule_id) + VALUES ($1, $2, $3, $4, $5, $6)", + ) + .bind(id) + .bind(kb) + .bind(s) + .bind(p) + .bind(o) + .bind(r) + .execute(pool) + .await?; + } + Ok(Fixture { + org, + a, + b, + doc_a, + chunk_a, + chunk_b, + ent_a, + ent_b, + fact_a, + fact_b, + rel_a, + rel_b, + cls_a, + cls_b, + rule_a, + rule_b, + arule_a, + arule_b, + der_a, + der_b, + }) +} + +async fn cleanup(pool: &PgPool, f: &Fixture) -> anyhow::Result<()> { + for kb in [f.a, f.b] { + sqlx::query("DELETE FROM knowledge_bases WHERE id = $1") + .bind(kb) + .execute(pool) + .await?; + } + sqlx::query("DELETE FROM organizations WHERE id = $1") + .bind(f.org) + .execute(pool) + .await?; + Ok(()) +} + +/// 新的跨库写:每条边都被它自己的触发器当场挡下 +#[tokio::test] +async fn new_cross_kb_writes_are_rejected_on_every_exported_edge() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + let f = seed(&pool).await?; + + // —— 派生的前提:两种前提,两种越库形态 + for (premise_fact, premise_derived, what) in [ + (Some(f.fact_b), None, "foreign fact premise"), + (None, Some(f.der_b), "foreign derived premise"), + ] { + let err = sqlx::query( + "INSERT INTO fact_derivations (derived_fact_id, premise_fact_id, premise_derived_id, seq) + VALUES ($1, $2, $3, 0)", + ) + .bind(f.der_a) + .bind(premise_fact) + .bind(premise_derived) + .execute(&pool) + .await; + assert!(err.is_err(), "derivation premise: {what} 必须被拒"); + } + + // —— 边上的属性:别库类型会被词汇表静默跳过,别库实体会被铸进本库 IRI + let err = sqlx::query( + "INSERT INTO fact_qualifiers (fact_id, qualifier_type_id, value) + VALUES ($1, $2, '\"lit\"'::jsonb)", + ) + .bind(f.fact_a) + .bind(f.rel_b) + .execute(&pool) + .await; + assert!(err.is_err(), "qualifier.type→foreign 必须被拒"); + let err = sqlx::query( + "INSERT INTO fact_qualifiers (fact_id, qualifier_type_id, entity_id) + VALUES ($1, $2, $3)", + ) + .bind(f.fact_a) + .bind(f.rel_a) + .bind(f.ent_b) + .execute(&pool) + .await; + assert!(err.is_err(), "qualifier.entity→foreign 必须被拒"); + + // —— 事实本体的五个引用列(from_statement_id 是同表自指:immediate 端在 + // BEFORE 触发器上,目标早已在场时当场拒) + for (col, foreign_id, what) in [ + ("subject_id", f.ent_b, "fact.subject"), + ("object_id", f.ent_b, "fact.object"), + ("predicate_id", f.rel_b, "fact.predicate"), + ("supersedes", f.fact_b, "fact.supersedes"), + ("from_statement_id", f.fact_b, "fact.from_statement"), + ] { + let err = sqlx::query(&format!( + "INSERT INTO facts (id, kb_id, subject_id, {col}) VALUES ($1, $2, $3, $4)" + )) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(if col == "subject_id" { + f.ent_b + } else { + f.ent_a + }) + .bind(foreign_id) + .execute(&pool) + .await; + assert!(err.is_err(), "{what}→foreign 必须被拒"); + } + + // —— 陈述来源(typed_fact_sources,0068):行自己没有 kb 列,归属按所属 + // fact 的库判——A 的事实吃 B 的陈述,序列化会铸出指着别库陈述的边 + let err = sqlx::query("INSERT INTO typed_fact_sources (fact_id, statement_id) VALUES ($1, $2)") + .bind(f.fact_a) + .bind(f.fact_b) + .execute(&pool) + .await; + assert!(err.is_err(), "factsource.statement→foreign 必须被拒"); + + // —— 开放陈述的属性(statement_qualifiers,0061):行没有 kb 列,归属按 + // 所属 fact 的库判;别库的实体值会被铸进本库 IRI + let err = sqlx::query( + "INSERT INTO statement_qualifiers (fact_id, role, entity_id) + VALUES ($1, 'as', $2)", + ) + .bind(f.fact_a) + .bind(f.ent_b) + .execute(&pool) + .await; + assert!(err.is_err(), "squalifier.entity→foreign 必须被拒"); + + // —— 时间提及(0061/0064):提及自己的 kb 必须与它指的事实、段落同属一库 + let err = sqlx::query( + "INSERT INTO time_mentions (id, kb_id, fact_id, chunk_id, text, char_start) + VALUES ($1, $2, $3, $4, '去年', 0)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.fact_b) // 事实在别库 + .bind(f.chunk_a) + .execute(&pool) + .await; + assert!(err.is_err(), "timemention.fact→foreign 必须被拒"); + let err = sqlx::query( + "INSERT INTO time_mentions (id, kb_id, fact_id, chunk_id, text, char_start) + VALUES ($1, $2, $3, $4, '去年', 0)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.fact_a) + .bind(f.chunk_b) // 段落在别库 + .execute(&pool) + .await; + assert!(err.is_err(), "timemention.chunk→foreign 必须被拒"); + + // —— 类型与短语绑定(0065/0066):绑定的类/属性以绑定行自己的库为准 + let err = sqlx::query( + "INSERT INTO type_bindings (id, kb_id, kind_word, status, type_id) + VALUES ($1, $2, 'corp', 'bound', $3)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.cls_b) + .execute(&pool) + .await; + assert!(err.is_err(), "binding.type→foreign 必须被拒"); + let err = sqlx::query( + "INSERT INTO phrase_bindings (id, kb_id, phrase, subject_type_id, status) + VALUES ($1, $2, 'runs', $3, 'none')", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.cls_b) + .execute(&pool) + .await; + assert!(err.is_err(), "pbinding.subject_type→foreign 必须被拒"); + let err = sqlx::query( + "INSERT INTO phrase_bindings (id, kb_id, phrase, relation_type_id, + direction, status) + VALUES ($1, $2, 'runs', $3, 'forward', 'bound')", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.rel_b) + .execute(&pool) + .await; + assert!(err.is_err(), "pbinding.relation→foreign 必须被拒"); + + // —— 派生事实本体的五个引用列 + for (col, foreign_id, what) in [ + ("subject_id", f.ent_b, "derived.subject"), + ("object_id", f.ent_b, "derived.object"), + ("predicate_id", f.rel_b, "derived.predicate"), + ("rule_id", f.rule_b, "derived.rule"), + ("attribute_rule_id", f.arule_b, "derived.attribute_rule"), + ] { + // attribute_rule 与 rule 互斥:测 attribute_rule 那行不带 rule_id + let err = if col == "attribute_rule_id" { + sqlx::query( + "INSERT INTO derived_facts (id, kb_id, subject_id, predicate_id, attribute_rule_id) + VALUES ($1, $2, $3, $4, $5)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.ent_a) + .bind(f.rel_a) + .bind(foreign_id) + .execute(&pool) + .await + } else { + sqlx::query(&format!( + "INSERT INTO derived_facts (id, kb_id, subject_id, predicate_id, rule_id, {col}) + VALUES ($1, $2, $3, $4, $5, $6)" + )) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(if col == "subject_id" { + f.ent_b + } else { + f.ent_a + }) + .bind(if col == "predicate_id" { + f.rel_b + } else { + f.rel_a + }) + .bind(if col == "rule_id" { f.rule_b } else { f.rule_a }) + .bind(foreign_id) + .execute(&pool) + .await + }; + assert!(err.is_err(), "{what}→foreign 必须被拒"); + } + + // —— 实体的类、类层级、互斥、domain/range + let err = sqlx::query( + "INSERT INTO entities (id, kb_id, canonical_name, type_id) VALUES ($1, $2, 'e', $3)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.cls_b) + .execute(&pool) + .await; + assert!(err.is_err(), "entity.type→foreign 必须被拒"); + + let err = sqlx::query("INSERT INTO entity_type_parents (child_id, parent_id) VALUES ($1, $2)") + .bind(f.cls_a) + .bind(f.cls_b) + .execute(&pool) + .await; + assert!(err.is_err(), "class.parent→foreign 必须被拒"); + + let err = + sqlx::query("INSERT INTO entity_type_disjoint (kb_id, a_id, b_id) VALUES ($1, $2, $3)") + .bind(f.a) + .bind(f.cls_a) + .bind(f.cls_b) + .execute(&pool) + .await; + assert!(err.is_err(), "class.disjoint→foreign 必须被拒"); + + for (table, what) in [ + ("relation_type_domains", "relation.domain"), + ("relation_type_ranges", "relation.range"), + ] { + let err = sqlx::query(&format!( + "INSERT INTO {table} (relation_type_id, entity_type_id) VALUES ($1, $2)" + )) + .bind(f.rel_a) + .bind(f.cls_b) + .execute(&pool) + .await; + assert!(err.is_err(), "{what}→foreign 必须被拒"); + } + + // —— 关系的同表自指与边属性声明:owl:inverseOf / + // rdfs:subPropertyOf 与 qualifier 列表都是语义,别库的不许进来 + let err = sqlx::query( + "INSERT INTO relation_types (id, kb_id, key, label, inverse_of) VALUES ($1, $2, 'inv', 'inv', $3)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.rel_b) + .execute(&pool) + .await; + assert!(err.is_err(), "relation.inverse→foreign 必须被拒"); + let err = sqlx::query( + "INSERT INTO relation_types (id, kb_id, key, label, sub_property_of) VALUES ($1, $2, 'sub', 'sub', $3)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.rel_b) + .execute(&pool) + .await; + assert!(err.is_err(), "relation.sub_property→foreign 必须被拒"); + let err = sqlx::query( + "INSERT INTO relation_type_qualifiers (relation_type_id, qualifier_type_id) VALUES ($1, $2)", + ) + .bind(f.rel_a) + .bind(f.rel_b) + .execute(&pool) + .await; + assert!(err.is_err(), "relation.qualifier→foreign 必须被拒"); + + // —— 规则本体:公理编在哪个谓词上、业务规则看什么类得什么结论, + // 现在都在导出里——它们的引用同样不许跨库 + let err = sqlx::query( + "INSERT INTO rules (id, kb_id, predicate_id, kind) VALUES ($1, $2, $3, 'transitive')", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.rel_b) + .execute(&pool) + .await; + assert!(err.is_err(), "rule.predicate→foreign 必须被拒"); + // 每条结论形状自己的必填列(CHECK 钉死了 typing/attribute/computed 的列组), + // 一个引用列一种合法形状 + let err = sqlx::query( + "INSERT INTO attribute_rules (id, kb_id, name, subject_type_id, conclusion, + conclude_predicate_id, conclude_value) + VALUES ($1, $2, 'ar', $3, 'attribute', $4, '42'::jsonb)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.cls_b) // subject_type 别库 + .bind(f.rel_a) + .execute(&pool) + .await; + assert!(err.is_err(), "arule.subject_type→foreign 必须被拒"); + let err = sqlx::query( + "INSERT INTO attribute_rules (id, kb_id, name, subject_type_id, conclusion, + conclude_type_id) + VALUES ($1, $2, 'ar', $3, 'typing', $4)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.cls_a) + .bind(f.cls_b) // conclude_type 别库 + .execute(&pool) + .await; + assert!(err.is_err(), "arule.conclude_type→foreign 必须被拒"); + let err = sqlx::query( + "INSERT INTO attribute_rules (id, kb_id, name, subject_type_id, conclusion, + conclude_predicate_id, conclude_value) + VALUES ($1, $2, 'ar', $3, 'attribute', $4, '42'::jsonb)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.cls_a) + .bind(f.rel_b) // conclude_predicate 别库 + .execute(&pool) + .await; + assert!(err.is_err(), "arule.conclude_predicate→foreign 必须被拒"); + + // —— kb 过户:把已被引用的行挪到别的库,等于把指着它的行一次全变坏行 + for (table, id, what) in [ + ("entities", f.ent_a, "entities.kb_id"), + ("entity_types", f.cls_a, "entity_types.kb_id"), + ("relation_types", f.rel_a, "relation_types.kb_id"), + ("derived_facts", f.der_a, "derived_facts.kb_id"), + ("rules", f.rule_a, "rules.kb_id"), + ("attribute_rules", f.arule_a, "attribute_rules.kb_id"), + ] { + let err = sqlx::query(&format!("UPDATE {table} SET kb_id = $2 WHERE id = $1")) + .bind(id) + .bind(f.b) + .execute(&pool) + .await; + assert!(err.is_err(), "{what} 过户必须被拒"); + } + + // —— 防误伤:同库的合法写入照常 + sqlx::query( + "INSERT INTO fact_derivations (derived_fact_id, premise_fact_id, seq) + VALUES ($1, $2, 0)", + ) + .bind(f.der_a) + .bind(f.fact_a) + .execute(&pool) + .await?; + sqlx::query( + "INSERT INTO fact_qualifiers (fact_id, qualifier_type_id, entity_id) + VALUES ($1, $2, $3)", + ) + .bind(f.fact_a) + .bind(f.rel_a) + .bind(f.ent_a) + .execute(&pool) + .await?; + + cleanup(&pool, &f).await +} + +/// 存量坏行:体检报出正确的边,对应的页读取同样拒 +#[tokio::test] +async fn malformed_rows_fail_every_exported_edge_closed() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + let f = seed(&pool).await?; + + let mut conn = pool.acquire().await?; + let mut tx = conn.begin().await?; + sqlx::query("SET LOCAL session_replication_role = 'replica'") + .execute(&mut *tx) + .await?; + + // 每条边造一行坏行——全部落在库 A 身上 + sqlx::query( + "INSERT INTO fact_derivations (derived_fact_id, premise_fact_id, seq) + VALUES ($1, $2, 0)", + ) + .bind(f.der_a) + .bind(f.fact_b) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO fact_derivations (derived_fact_id, premise_derived_id, seq) + VALUES ($1, $2, 1)", + ) + .bind(f.der_a) + .bind(f.der_b) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO fact_qualifiers (fact_id, qualifier_type_id, value) + VALUES ($1, $2, '\"lit\"'::jsonb)", + ) + .bind(f.fact_a) + .bind(f.rel_b) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO fact_qualifiers (fact_id, qualifier_type_id, entity_id) + VALUES ($1, $2, $3)", + ) + .bind(f.fact_a) + .bind(f.rel_a) + .bind(f.ent_b) + .execute(&mut *tx) + .await?; + // 事实的 supersedes 指向别库事实 + sqlx::query("UPDATE facts SET supersedes = $2 WHERE id = $1") + .bind(f.fact_a) + .bind(f.fact_b) + .execute(&mut *tx) + .await?; + // 陈述来源指着别库陈述(typed_fact_sources 与 from_statement_id 两条路) + sqlx::query("INSERT INTO typed_fact_sources (fact_id, statement_id) VALUES ($1, $2)") + .bind(f.fact_a) + .bind(f.fact_b) + .execute(&mut *tx) + .await?; + sqlx::query("UPDATE facts SET from_statement_id = $2 WHERE id = $1") + .bind(f.fact_a) + .bind(f.fact_b) + .execute(&mut *tx) + .await?; + // 开放陈述的属性值指着别库实体 + sqlx::query( + "INSERT INTO statement_qualifiers (fact_id, role, entity_id) + VALUES ($1, 'as', $2)", + ) + .bind(f.fact_a) + .bind(f.ent_b) + .execute(&mut *tx) + .await?; + // 时间提及两种坏法:归属在 A 却指着 B 的段落(体检扫得见——行按自己的 + // kb 计数);归属在 B 却挂在 A 的事实上(按事实取回时,页校验拦它) + sqlx::query( + "INSERT INTO time_mentions (id, kb_id, fact_id, chunk_id, text, char_start) + VALUES ($1, $2, $3, $4, '去年', 0)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.fact_a) + .bind(f.chunk_b) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO time_mentions (id, kb_id, fact_id, chunk_id, text, char_start) + VALUES ($1, $2, $3, $4, '前年', 2)", + ) + .bind(Uuid::now_v7()) + .bind(f.b) + .bind(f.fact_a) + .bind(f.chunk_a) + .execute(&mut *tx) + .await?; + // 类型绑定指着别库的类 + sqlx::query( + "INSERT INTO type_bindings (id, kb_id, kind_word, status, type_id) + VALUES ($1, $2, 'corp', 'bound', $3)", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.cls_b) + .execute(&mut *tx) + .await?; + // 短语绑定指着别库的属性 + sqlx::query( + "INSERT INTO phrase_bindings (id, kb_id, phrase, relation_type_id, + direction, status) + VALUES ($1, $2, 'runs', $3, 'forward', 'bound')", + ) + .bind(Uuid::now_v7()) + .bind(f.a) + .bind(f.rel_b) + .execute(&mut *tx) + .await?; + // 派生的公理规则换库 + sqlx::query("UPDATE derived_facts SET rule_id = $2 WHERE id = $1") + .bind(f.der_a) + .bind(f.rule_b) + .execute(&mut *tx) + .await?; + // 实体的类换库 + sqlx::query("UPDATE entities SET type_id = $2 WHERE id = $1") + .bind(f.ent_a) + .bind(f.cls_b) + .execute(&mut *tx) + .await?; + // 类层级与 domain/range + sqlx::query("INSERT INTO entity_type_parents (child_id, parent_id) VALUES ($1, $2)") + .bind(f.cls_a) + .bind(f.cls_b) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO relation_type_domains (relation_type_id, entity_type_id) VALUES ($1, $2)", + ) + .bind(f.rel_a) + .bind(f.cls_b) + .execute(&mut *tx) + .await?; + tx.commit().await?; + drop(conn); + + let err = utopia_store::export::provenance_integrity(&mut pool.begin().await?, f.a).await; + let msg = format!("{err:?}"); + assert!(err.is_err(), "库 A 的体检必须拒导"); + // 体检只报这份导出真正解析的边:陈述来源、开放陈述属性、时间提及、 + // 类型/短语绑定这些行上面同样埋了坏行,但导出现在不读它们——它们随 + // 各自的导出面一起带上自己的校验 + for edge in [ + "derivation.premise_fact", + "derivation.premise_derived", + "qualifier.type", + "qualifier.entity", + "fact.supersedes", + "derived.rule", + "entity.type", + "class.parent", + "relation.domain", + ] { + assert!(msg.contains(edge), "体检该报 {edge}: {msg}"); + } + + // 逐页校验:坏行落在哪页,哪页就拒——不能漏出伪造 IRI + assert!( + utopia_store::export::facts_page(&mut pool.begin().await?, f.a, None) + .await + .is_err() + ); + assert!( + utopia_store::export::derived_page(&mut pool.begin().await?, f.a, None) + .await + .is_err() + ); + assert!( + utopia_store::export::entities_page(&mut pool.begin().await?, f.a, None) + .await + .is_err() + ); + assert!(utopia_store::export::classes(&mut pool.begin().await?, f.a) + .await + .is_err()); + assert!( + utopia_store::export::relations(&mut pool.begin().await?, f.a) + .await + .is_err() + ); + + cleanup(&pool, &f).await +} + +/// 留下的行才参与校验:被引的行在两次读之间消失,判违规的是 +/// **随页原子选出的归属**,不是事后另一个时刻的 JOIN——悬空引用照样拒 +#[tokio::test] +async fn retained_row_validation_survives_dangling_and_late_state() -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + utopia_store::db::migrate(&pool).await?; + let f = seed(&pool).await?; + + // 合法证据先行 + sqlx::query("INSERT INTO fact_evidence (fact_id, chunk_id, document_id) VALUES ($1, $2, $3)") + .bind(f.fact_a) + .bind(f.chunk_a) + .bind(f.doc_a) + .execute(&pool) + .await?; + utopia_store::export::facts_page(&mut pool.begin().await?, f.a, None).await?; + + // 绕过触发器把文档删掉:evidence 行还指着它——悬空不是「不存在所以跳过」 + let mut conn = pool.acquire().await?; + let mut tx = conn.begin().await?; + sqlx::query("SET LOCAL session_replication_role = 'replica'") + .execute(&mut *tx) + .await?; + sqlx::query("DELETE FROM documents WHERE id = $1") + .bind(f.doc_a) + .execute(&mut *tx) + .await?; + tx.commit().await?; + drop(conn); + + let err = utopia_store::export::facts_page(&mut pool.begin().await?, f.a, None).await; + let msg = format!("{err:?}"); + assert!(err.is_err(), "悬空 document 指针必须拒"); + assert!( + msg.contains("evidence.document"), + "该报 evidence.document: {msg}" + ); + + // 同理:前提行被删,derived 的前提数组仍留着死指针——逐页校验要拦 + sqlx::query( + "INSERT INTO fact_derivations (derived_fact_id, premise_fact_id, seq) + VALUES ($1, $2, 0)", + ) + .bind(f.der_a) + .bind(f.fact_a) + .execute(&pool) + .await?; + let mut conn = pool.acquire().await?; + let mut tx = conn.begin().await?; + sqlx::query("SET LOCAL session_replication_role = 'replica'") + .execute(&mut *tx) + .await?; + sqlx::query("DELETE FROM facts WHERE id = $1") + .bind(f.fact_a) + .execute(&mut *tx) + .await?; + tx.commit().await?; + drop(conn); + + let err = utopia_store::export::derived_page(&mut pool.begin().await?, f.a, None).await; + let msg = format!("{err:?}"); + assert!(err.is_err(), "悬空前提必须拒"); + assert!( + msg.contains("derivation.premise_fact"), + "该报 derivation.premise_fact: {msg}" + ); + + cleanup(&pool, &f).await +}