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
153 changes: 153 additions & 0 deletions crates/utopia-server/src/api/mcp_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1632,3 +1632,156 @@ async fn rule_matches_keep_materialized_intervals_and_count_rows() -> anyhow::Re
);
f.clean().await
}

// Reuse the authenticated ledger fixture so RDF exercises the same stored records
// as structured MCP reads, including evidence and retracted history.
#[tokio::test]
async fn rdf_export_preserves_unbound_literal_objects() -> anyhow::Result<()> {
use axum::body::{to_bytes, Body};
use axum::http::{Request, StatusCode};
use oxrdf::{vocab::rdf, Literal, Term};
use tower::ServiceExt;

let Some(f) = Fixture::new().await? else {
return Ok(());
};
async fn check(f: &Fixture) -> anyhow::Result<()> {
let value = json!({"value": "待复检"});
let (statement, _) = utopia_store::graph::insert_open_statement(
&f.state.pool,
f.kb,
f.subject,
"状态",
utopia_store::graph::FactObject::Value(&value),
Some("2026-01-01T00:00:00Z".parse()?),
0.9,
)
.await?;
sqlx::query("UPDATE facts SET recorded_at='2026-02-01' WHERE id=$1")
.bind(statement)
.execute(&f.state.pool)
.await?;
sqlx::query(
"INSERT INTO fact_evidence(fact_id,chunk_id,document_id,doc_version,quote)
VALUES ($1,$2,$3,1,'设备 A 待复检')",
)
.bind(statement)
.bind(f.chunk)
.bind(f.document)
.execute(&f.state.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 names = crate::rdf::Names::new(f.kb, None).map_err(anyhow::Error::msg)?;
let stmt = names.fact(statement);
let mut formats = Vec::new();
for retracted in [false, true] {
if retracted {
sqlx::query("UPDATE facts SET invalidated_at='2026-03-01' WHERE id=$1")
.bind(statement)
.execute(&f.state.pool)
.await?;
}
// Snapshot every KB-scoped business table, including queues and adoption
// records. Request audit is deliberately excluded from this read-only check.
let tables: Vec<String> = sqlx::query_scalar(
"SELECT table_name FROM information_schema.columns
WHERE table_schema='public' AND column_name='kb_id'
AND table_name <> 'audit_events' ORDER BY table_name",
)
.fetch_all(&f.state.pool)
.await?;
let snapshot = async {
let mut rows = Vec::new();
for table in &tables {
let sql = format!("SELECT COALESCE(jsonb_agg(to_jsonb(t) ORDER BY to_jsonb(t)::text), '[]'::jsonb) FROM \"{}\" t WHERE kb_id=$1", table.replace('"', "\"\""));
rows.push(
sqlx::query_scalar::<_, Value>(&sql)
.bind(f.kb)
.fetch_one(&f.state.pool)
.await?,
);
}
Ok::<_, anyhow::Error>(rows)
};
let before = snapshot.await?;
let extra_sql = "SELECT jsonb_build_array(
(SELECT COALESCE(jsonb_agg(to_jsonb(e) ORDER BY to_jsonb(e)::text), '[]')
FROM fact_evidence e JOIN facts f ON f.id=e.fact_id WHERE f.kb_id=$1),
(SELECT COALESCE(jsonb_agg(to_jsonb(j) ORDER BY j.id), '[]') FROM jobs j
WHERE payload->>'kb_id'=$1::text OR payload->>'document_id' IN
(SELECT id::text FROM documents WHERE kb_id=$1)))";
let extra_before: Value = sqlx::query_scalar(extra_sql)
.bind(f.kb)
.fetch_one(&f.state.pool)
.await?;
for format in ["turtle", "jsonld"] {
let response = app
.clone()
.oneshot(
Request::builder()
.uri(format!("/api/v1/kbs/{}/export?format={format}", f.kb))
.header("authorization", format!("Bearer {jwt}"))
.body(Body::empty())?,
)
.await?;
anyhow::ensure!(response.status() == StatusCode::OK, "export rejected");
let bytes = to_bytes(response.into_body(), 1024 * 1024).await?;
let format = if format == "turtle" {
oxrdfio::RdfFormat::Turtle
} else {
oxrdfio::RdfFormat::JsonLd {
profile: oxrdfio::JsonLdProfileSet::empty(),
}
};
let quads = oxrdfio::RdfParser::from_format(format)
.for_slice(&bytes)
.collect::<Result<std::collections::HashSet<_>, _>>()?;
anyhow::ensure!(
quads.iter().any(|q| q.subject == stmt.clone().into()
&& q.predicate == rdf::OBJECT
&& q.object == Term::Literal(Literal::new_simple_literal("待复检"))),
"unbound statement lost its rdf:object in authenticated export"
);
anyhow::ensure!(
!quads
.iter()
.any(|q| q.subject == stmt.clone().into() && q.predicate == rdf::PREDICATE),
"invented a bound predicate"
);
anyhow::ensure!(
quads.iter().any(|q| q.subject == stmt.clone().into()
&& q.predicate.as_str() == "http://www.w3.org/ns/prov#wasDerivedFrom"
&& q.object == names.document(f.document).into()),
"lost evidence source"
);
formats.push(quads);
}
anyhow::ensure!(
formats[formats.len() - 1] == formats[formats.len() - 2],
"formats disagree"
);
let extra_after: Value = sqlx::query_scalar(extra_sql)
.bind(f.kb)
.fetch_one(&f.state.pool)
.await?;
anyhow::ensure!(
extra_before == extra_after,
"export changed evidence or jobs"
);
for (table, expected) in tables.iter().zip(before) {
let sql = format!("SELECT COALESCE(jsonb_agg(to_jsonb(t) ORDER BY to_jsonb(t)::text), '[]'::jsonb) FROM \"{}\" t WHERE kb_id=$1", table.replace('"', "\"\""));
let actual: Value = sqlx::query_scalar(&sql)
.bind(f.kb)
.fetch_one(&f.state.pool)
.await?;
anyhow::ensure!(actual == expected, "export changed {table}");
}
}
Ok(())
}
let result = check(&f).await;
let cleanup = f.clean().await;
result.and(cleanup)
}
103 changes: 99 additions & 4 deletions crates/utopia-server/src/rdf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -453,10 +453,11 @@ pub fn emit_fact(
let predicate = f.predicate_id.and_then(|p| vocab.relation(p)).cloned();
let object: Option<Term> = match (f.object_id, &f.object_value) {
(Some(o), _) => Some(names.entity(o).into()),
(None, Some(v)) => f.predicate_id.map(|p| {
let (datatype, _) = vocab.literal_shape(p);
literal_value(v, datatype).into()
}),
(None, Some(v)) => {
// An unbound statement still has an object; only its datatype is unknown.
let datatype = f.predicate_id.and_then(|p| vocab.literal_shape(p).0);
Some(literal_value(v, datatype).into())
}
_ => None,
};

Expand Down Expand Up @@ -801,6 +802,100 @@ mod tests {
const OBJ: &str = "<urn:utopia:kb:01a06dc4-f40a-7013-b09f-1b499e2e7441:entity:0b0b0b0b-0b0b-0b0b-0b0b-0b0b0b0b0b0b>";
const WORKS_FOR: &str = "https://schema.org/worksFor";

#[test]
fn unbound_literal_objects_survive_both_formats() {
for value in [
serde_json::json!({"value": "待复检"}),
serde_json::json!({"value": "quote: \" and slash: \\"}),
serde_json::json!({"value": ""}),
serde_json::json!({"value": 0}),
serde_json::json!({"value": false}),
serde_json::json!({"value": null, "summary": "not specified"}),
serde_json::Value::Null,
] {
let mut f = fact(5);
f.predicate_id = None;
f.surface_predicate = Some("状态".into());
f.object_id = None;
f.object_value = Some(value.clone());
f.documents = vec![id(20)];
f.quotes = vec!["设备 A 待复检".into()];
f.supersedes = Some(id(6));
for retracted in [false, true] {
f.invalidated_at = retracted.then(|| at("2026-02-01T00:00:00Z"));
let mut sets = Vec::new();
for format in [Format::Turtle, Format::JsonLd] {
let quads = export(format, |sink, names, vocab| {
emit_fact(sink, names, vocab, &f, at("2026-06-01T00:00:00Z")).unwrap();
});
assert_eq!(
objects(&quads, STMT, rdf::OBJECT.as_str()),
vec![literal_value(&value, None).to_string()],
"unbound statement lost its literal object: {value}"
);
assert!(objects(&quads, STMT, rdf::PREDICATE.as_str()).is_empty());
assert!(!quads.iter().any(|q| q.subject.to_string() == SUBJ));
sets.push(quads.into_iter().collect::<std::collections::HashSet<_>>());
}
assert_eq!(sets[0], sets[1]);
}
}
}

#[test]
fn bound_literal_datatypes_survive_both_formats() {
for (datatype, value, expected) in [
(
"number",
serde_json::json!(0),
Literal::new_typed_literal("0", xsd::DECIMAL),
),
(
"text",
serde_json::json!("待复检"),
Literal::new_simple_literal("待复检"),
),
(
"bool",
serde_json::json!(false),
Literal::new_typed_literal("false", xsd::BOOLEAN),
),
] {
let mut f = fact(5);
f.predicate_id = Some(id(4));
f.object_id = None;
f.object_value = Some(serde_json::json!({"value": value}));
for format in [Format::Turtle, Format::JsonLd] {
let quads = export(format, |sink, names, _| {
let mut property = relation(4, "value", None, "attribute");
property.datatype = Some(datatype.into());
let vocab = vocabulary(names, &[], &[property]);
emit_fact(sink, names, &vocab, &f, at("2026-06-01T00:00:00Z")).unwrap();
});
assert_eq!(
objects(&quads, STMT, rdf::OBJECT.as_str()),
vec![expected.to_string()]
);
assert!(quads.iter().any(|q| q.subject.to_string() == SUBJ
&& q.object == Term::Literal(expected.clone())));
}
}
}

#[test]
fn an_absent_object_is_not_an_empty_literal() {
let mut f = fact(5);
f.predicate_id = None;
f.object_id = None;
f.object_value = None;
for format in [Format::Turtle, Format::JsonLd] {
let quads = export(format, |sink, names, vocab| {
emit_fact(sink, names, vocab, &f, at("2026-06-01T00:00:00Z")).unwrap();
});
assert!(objects(&quads, STMT, rdf::OBJECT.as_str()).is_empty());
}
}

#[test]
fn an_imported_class_keeps_its_own_iri() {
let quads = export(Format::Turtle, |_, _, _| {});
Expand Down
Loading