Skip to content
Closed
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
72 changes: 65 additions & 7 deletions crates/utopia-server/src/api/export_routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,10 @@
//! 一个十万条事实的库正是最需要导出的那种库,也正是「先拼成一个 String」会
//! 把服务打死的那种库。
//!
//! **同一份快照**:体检、词汇表与每一页查询跑在同一条只读 REPEATABLE READ
//! 事务里。若每页各自拿一条连接,词汇表发完之后才提交的规则会在派生页里
//! 留下一条指向它的 `wasGeneratedBy`——图里就出现没有本体的引用。
//!
//! 中途出错只能截断——HTTP 头早就发出去了。所以错误进日志,而客户端拿到的是
//! 一份短了一截的文件;这比先攒后发要好,那种做法在同样的库上根本发不出来。

Expand Down Expand Up @@ -44,6 +48,18 @@ pub async fn export(
})?;
let names = Names::new(kb_id, q.base.as_deref()).map_err(AppError::Validation)?;

// 整份导出占一条连接上的只读 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?;

// 导出是一次「整个库离开这台机器」的动作,台账要记下(0014 的同一条理由)
let _ = utopia_store::audit::record(
&state.pool,
Expand All @@ -56,25 +72,31 @@ pub async fn export(
)
.await;

let pool = state.pool.clone();
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)?;
}
for r in &relations {
rdf::emit_relation(&mut sink, &vocab, r)?;
}
// 公理与业务规则整份进词汇表区:引它们的派生才有 prov:wasGeneratedBy 可指
for r in &utopia_store::export::rules(&mut tx, kb_id).await.map_err(io)? {
rdf::emit_rule(&mut sink, &names, &vocab, r)?;
}
for r in &utopia_store::export::attribute_rules(&mut tx, kb_id).await.map_err(io)? {
rdf::emit_attribute_rule(&mut sink, &names, &vocab, r)?;
}
yield axum::body::Bytes::from(buf.take());

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 {
Expand All @@ -85,7 +107,31 @@ 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::document_versions_page(&mut tx, kb_id, after)
.await
.map_err(io)?;
let Some(last) = page.last() else { break };
after = Some(last.id);
for v in &page {
rdf::emit_docversion(&mut sink, &names, v)?;
}
yield axum::body::Bytes::from(buf.take());
}

let mut after = None;
loop {
let page = utopia_store::export::chunks_page(&mut tx, kb_id, after).await.map_err(io)?;
let Some(last) = page.last() else { break };
after = Some(last.id);
for c in &page {
rdf::emit_chunk(&mut sink, &names, c)?;
}
yield axum::body::Bytes::from(buf.take());
}

let mut after = None;
loop {
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 {
Expand All @@ -99,7 +145,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 {
Expand All @@ -108,9 +154,21 @@ pub async fn export(
yield axum::body::Bytes::from(buf.take());
}

// 证据游标是复合的 (fact_id, chunk_id)——与取数页同一键
let mut after: Option<(Uuid, Uuid)> = None;
loop {
let page = utopia_store::export::evidence_page(&mut tx, kb_id, after).await.map_err(io)?;
let Some(last) = page.last() else { break };
after = Some((last.fact_id, last.chunk_id));
for e in &page {
rdf::emit_evidence(&mut sink, &names, e)?;
}
yield axum::body::Bytes::from(buf.take());
}

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 {
Expand Down
3 changes: 2 additions & 1 deletion crates/utopia-server/src/api/mcp_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -622,7 +622,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]));
Expand Down
Loading