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
33 changes: 25 additions & 8 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,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),
Expand All @@ -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)?;
Expand All @@ -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 {
Expand All @@ -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 {
Expand All @@ -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 {
Expand All @@ -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 {
Expand Down
86 changes: 84 additions & 2 deletions crates/utopia-server/src/api/mcp_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Self> {
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"))?);
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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]));
Expand Down
26 changes: 26 additions & 0 deletions crates/utopia-server/src/rdf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
}

Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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| {
Expand Down
Loading