diff --git a/bench/feature-off-declarations/6655.json b/bench/feature-off-declarations/6655.json new file mode 100644 index 0000000000..dac4b67079 --- /dev/null +++ b/bench/feature-off-declarations/6655.json @@ -0,0 +1,5 @@ +{ + "pr": 6655, + "date": "2026-10-05", + "reason": "Engine correctness fixes (#4467, #4438, #4239, #4158, #4160) in always-compiled crates/sparq-engine/src/exec.rs and update.rs: the functions and spatial registry guards now restore the previous registry on drop, the Join arm evaluates a left SERVICE once, the single-pattern SELECT-JSON stream skips its terminator after a budget trip, the SIP theta anti-join breaks its outer loop on exhaustion, and LOAD freshens document blank nodes. These sites are outside every `#[cfg(feature = ...)]` block and compile into the feature-OFF sparq-wasm bundle, so its bytes move by design. No feature gate is added or widened, no dependency is added, and there is no `unsafe`. Bundle SIZE remains governed separately by the metrics.wasm_bundle_bytes floor ratchet in bench.yml." +} diff --git a/crates/sparq-engine/src/exec.rs b/crates/sparq-engine/src/exec.rs index 3758b6f03e..b2024e5483 100644 --- a/crates/sparq-engine/src/exec.rs +++ b/crates/sparq-engine/src/exec.rs @@ -1929,18 +1929,19 @@ pub(crate) mod functions { static ACTIVE: RefCell>> = const { RefCell::new(None) }; } - /// Uninstalls the registry when the installing entry point returns (also on - /// error/unwind, so a poisoned thread never leaks a stale registry). - pub(crate) struct Guard; + /// Restores the PREVIOUS registry when the installing entry point returns (also + /// on error/unwind). Restoring rather than clearing keeps a nested install from + /// unregistering the outer scope's functions for the rest of that scope (#4467). + pub(crate) struct Guard(Option>); impl Drop for Guard { fn drop(&mut self) { - ACTIVE.with(|a| a.borrow_mut().take()); + let prev = self.0.take(); + ACTIVE.with(|a| *a.borrow_mut() = prev); } } pub(crate) fn install(fns: &FunctionRegistry) -> Guard { - ACTIVE.with(|a| *a.borrow_mut() = Some(Arc::new(fns.clone()))); - Guard + Guard(ACTIVE.with(|a| a.borrow_mut().replace(Arc::new(fns.clone())))) } /// Snapshot of the installed registry for the rayon-parallel branches @@ -2072,16 +2073,16 @@ pub(crate) mod aggregates { static ACTIVE: RefCell>> = const { RefCell::new(None) }; } - pub(crate) struct Guard; + pub(crate) struct Guard(Option>); impl Drop for Guard { fn drop(&mut self) { - ACTIVE.with(|a| a.borrow_mut().take()); + let prev = self.0.take(); + ACTIVE.with(|a| *a.borrow_mut() = prev); } } pub(crate) fn install(reg: &CustomAggregateRegistry) -> Guard { - ACTIVE.with(|a| *a.borrow_mut() = Some(Arc::new(reg.clone()))); - Guard + Guard(ACTIVE.with(|a| a.borrow_mut().replace(Arc::new(reg.clone())))) } #[cfg_attr(not(feature = "parallel"), allow(dead_code))] @@ -2276,16 +2277,16 @@ pub(crate) mod spatial { static ACTIVE: RefCell>> = const { RefCell::new(None) }; } - pub(crate) struct Guard; + pub(crate) struct Guard(Option>); impl Drop for Guard { fn drop(&mut self) { - ACTIVE.with(|a| a.borrow_mut().take()); + let prev = self.0.take(); + ACTIVE.with(|a| *a.borrow_mut() = prev); } } pub(crate) fn install(idx: Arc) -> Guard { - ACTIVE.with(|a| *a.borrow_mut() = Some(idx)); - Guard + Guard(ACTIVE.with(|a| a.borrow_mut().replace(idx))) } /// The installed spatial index, if any. @@ -2899,7 +2900,8 @@ fn single_pattern_scan_json_emit( // deadline-only budget, the now-past wall clock: sets the sticky flag the // caller's `budget::check(0)` converts into the budget error (a chunk skipped // above means the deadline is globally past, so this fires deterministically). - let _ = budget::exhausted(frags.iter().map(|(n, _)| n).sum()); + let total: usize = frags.iter().map(|(n, _)| n).sum(); + budget::exhausted(total); // Accumulate into `pending` and hand a chunk to `emit` at each flush boundary // (byte-identical concatenation to the old `emit_chunk` Vec layout — only the // chunk *boundaries* differ, and the concat is what the byte-identity contract @@ -2915,14 +2917,26 @@ fn single_pattern_scan_json_emit( } wrote = true; pending.push_str(&f); - if flush.is_some_and(|n| pending.len() >= n) - && emit(std::mem::take(&mut pending)).is_break() - { - return Some(()); + if flush.is_some_and(|n| pending.len() >= n) { + if emit(std::mem::take(&mut pending)).is_break() { + return Some(()); + } + // A cancellation (possibly set by the sink itself) or a deadline that + // passed while emitting stops the stream here, rechecked per chunk. + if budget::exhausted(total) { + return Some(()); + } } } - pending.push_str("]}}"); - let _ = emit(pending); + // Never close the document over a result the budget cut short (#4239): the + // caller reports the abort, and a sink that saw `]}}` would hold a complete- + // looking but truncated body. Rechecked here, not cached from before emission. + if !budget::exhausted(total) { + pending.push_str("]}}"); + } + if !pending.is_empty() { + let _ = emit(pending); + } return Some(()); } } @@ -2941,13 +2955,20 @@ fn single_pattern_scan_json_emit( } written += 1; write_row(row, &mut s); - if flush.is_some_and(|n| s.len() >= n) && emit(std::mem::take(&mut s)).is_break() { + if flush.is_some_and(|n| s.len() >= n) + && (emit(std::mem::take(&mut s)).is_break() || budget::exhausted(written)) + { return Some(()); } } - let _ = budget::exhausted(written); // final row-count gate (sticky) - s.push_str("]}}"); - let _ = emit(s); + // Final row-count gate (sticky). An exhausted budget leaves the document unclosed + // (#4239), as above. + if !budget::exhausted(written) { + s.push_str("]}}"); + } + if !s.is_empty() { + let _ = emit(s); + } Some(()) } @@ -5393,26 +5414,29 @@ fn eval_graph_pattern_inner(graph: &Graph, local: &mut LocalVocab, p: &GraphPatt Ok(b) } GraphPattern::Join { left, right } => { - let l = eval_graph_pattern(graph, local, left)?; - // Bind-join pushdown: if the RIGHT side is a SERVICE and the left has - // already bound its join variables, push those bindings to the remote as a - // VALUES block instead of materialising the whole remote relation. Join is + // Bind-join pushdown: if one side is a SERVICE and the other has already + // bound its join variables, push those bindings to the remote as a VALUES + // block instead of materialising the whole remote relation. Join is // symmetric, so try either side as the SERVICE. [OPUS-4.8] (sq-sjkj) + // + // SERVICE on the left (and not on the right): evaluate the right FIRST and + // the SERVICE at most once — evaluating it eagerly and again on a declined + // pushdown fetched the endpoint (or ran a local handler) twice (#4438). #[cfg(feature = "service")] + if matches!(left.as_ref(), GraphPattern::Service { .. }) + && !matches!(right.as_ref(), GraphPattern::Service { .. }) { - if let Some(r) = try_bound_join_service(graph, local, &l, right)? { - return Ok(join_bindings(l, r)); - } - // Symmetric: SERVICE on the left, bindings produced by the right. - if matches!(left.as_ref(), GraphPattern::Service { .. }) { - let r = eval_graph_pattern(graph, local, right)?; - if let Some(sl) = try_bound_join_service(graph, local, &r, left)? { - return Ok(join_bindings(r, sl)); - } - // Fall through with the already-evaluated right; recompute left verbatim. - let l2 = eval_graph_pattern(graph, local, left)?; - return Ok(join_bindings(l2, r)); + let r = eval_graph_pattern(graph, local, right)?; + if let Some(sl) = try_bound_join_service(graph, local, &r, left)? { + return Ok(join_bindings(r, sl)); } + let l = eval_graph_pattern(graph, local, left)?; + return Ok(join_bindings(l, r)); + } + let l = eval_graph_pattern(graph, local, left)?; + #[cfg(feature = "service")] + if let Some(r) = try_bound_join_service(graph, local, &l, right)? { + return Ok(join_bindings(l, r)); } // Sideways information passing (SIP): when the already-evaluated `l` is // SMALL, evaluate the big `right` child CORRELATED on it — seeding scans @@ -10163,7 +10187,7 @@ fn try_theta_antijoin( // ---- SIP-seed anti-join (small correlation cardinality / literal keys) ---- let mut fired = false; - for key in &order { + 'groups: for key in &order { let members = &groups[key]; let ri0 = members[0]; @@ -10231,8 +10255,11 @@ fn try_theta_antijoin( let all_cands: Vec = (0..b_prime.rows.len()).collect(); for &ri in members { let lrow = &left_b.rows[ri]; + // Stop the whole SIP strategy, not just this correlation group, so it ends + // exactly like the hash strategy: no further seeded `B'` is evaluated and the + // caller's operator-exit check raises the budget error (#4158). if budget::exhausted(result_rows.len()) { - break; + break 'groups; } let matched = antijoin_row_matches( graph, local, lrow, &b_prime, &all_cands, &shared, &out_src, &tmp_vars, &checks, @@ -22241,3 +22268,93 @@ mod capped_rhs_tests { ); } } + +/// #4467 — the scoped registry guards restore the registry the install replaced +/// (rather than clearing it), so a nested install hands the outer scope its own +/// registry back, on normal return and on unwind alike. +#[cfg(test)] +mod scoped_registry_tests { + use super::*; + use std::sync::Arc; + + const F: &str = "http://ex/f"; + + fn registry(tag: &'static str) -> crate::FunctionRegistry { + let mut reg = crate::FunctionRegistry::new(); + reg.register(F, move |_: &[Term]| Ok(Term::Literal(oxrdf::Literal::new_simple_literal(tag)))); + reg + } + + fn call() -> Option { + functions::lookup(F).map(|f| f(&[]).unwrap()) + } + + fn lit(tag: &str) -> Option { + Some(Term::Literal(oxrdf::Literal::new_simple_literal(tag))) + } + + #[test] + fn nested_function_registry_restores_the_outer_one() { + crate::with_functions(®istry("outer"), || { + assert_eq!(call(), lit("outer")); + crate::with_functions(®istry("inner"), || assert_eq!(call(), lit("inner"))); + assert_eq!(call(), lit("outer"), "outer registry lost after the inner scope"); + let g = sparq_core::Graph::load_str(" .", "turtle").unwrap(); + let r = crate::query(&g, "SELECT ?v WHERE { BIND(() AS ?v) }").unwrap(); + assert_eq!(r.rows[0][0], lit("outer")); + }); + assert_eq!(call(), None, "registry leaked past the outermost scope"); + } + + #[test] + fn nested_function_registry_restores_the_outer_one_on_unwind() { + crate::with_functions(®istry("outer"), || { + let unwound = std::panic::catch_unwind(|| { + crate::with_functions(®istry("inner"), || panic!("inner scope panics")) + }); + assert!(unwound.is_err()); + assert_eq!(call(), lit("outer")); + }); + assert_eq!(call(), None); + } + + #[cfg(feature = "window-functions")] + #[test] + fn nested_aggregate_registry_restores_the_outer_one() { + let reg = |tag: &'static str| { + let mut reg = crate::CustomAggregateRegistry::new(); + reg.register(F, move |_: &[Option]| Ok(lit(tag))); + reg + }; + let call = || aggregates::lookup(F).map(|f| f(&[]).unwrap()); + crate::aggregate::with_aggregates(®("outer"), || { + crate::aggregate::with_aggregates(®("inner"), || assert_eq!(call(), Some(lit("inner")))); + assert_eq!(call(), Some(lit("outer"))); + }); + assert_eq!(call(), None); + } + + struct NoIndex; + impl crate::SpatialProvider for NoIndex { + fn candidates(&self, _: &crate::SpatialQuery) -> Option> { + None + } + fn is_indexed(&self, _: &Term) -> bool { + false + } + } + + #[test] + fn nested_spatial_index_restores_the_outer_one() { + let outer: Arc = Arc::new(NoIndex); + let inner: Arc = Arc::new(NoIndex); + let is = |want: &Arc| { + spatial::active().is_some_and(|a| std::ptr::addr_eq(Arc::as_ptr(&a), Arc::as_ptr(want))) + }; + crate::with_spatial_index(outer.clone(), || { + crate::with_spatial_index(inner.clone(), || assert!(is(&inner))); + assert!(is(&outer), "outer spatial index lost after the inner scope"); + }); + assert!(spatial::active().is_none()); + } +} diff --git a/crates/sparq-engine/src/lib.rs b/crates/sparq-engine/src/lib.rs index c0b2316f35..dd54f5e9f3 100644 --- a/crates/sparq-engine/src/lib.rs +++ b/crates/sparq-engine/src/lib.rs @@ -2581,6 +2581,56 @@ mod tests { assert!(e.contains("query budget exceeded (max-rows)"), "got: {e}"); } + /// #4239 — a row cap that trips on the streaming scan fast path must not close the + /// JSON document: the sink never receives the `]}}` terminator of a truncated result. + #[test] + fn budget_tripped_select_json_stream_leaves_the_document_unclosed() { + use std::ops::ControlFlow; + let b = QueryBudget { max_rows: Some(1), ..QueryBudget::unlimited() }; + let mut body = String::new(); + let e = query_json_stream_with_budget(&g(), "SELECT * WHERE { ?s ?p ?o }", &b, |c| { + body.push_str(&c); + ControlFlow::Continue(()) + }) + .unwrap_err(); + assert!(e.contains("query budget exceeded (max-rows)"), "got: {e}"); + assert!(!body.ends_with("]}}"), "truncated stream was closed: {body}"); + // Under the cap, the same query streams a complete document. + let mut whole = String::new(); + query_json_stream_with_budget(&g(), "SELECT * WHERE { ?s ?p ?o }", &QueryBudget::unlimited(), |c| { + whole.push_str(&c); + ControlFlow::Continue(()) + }) + .unwrap(); + assert!(whole.ends_with("]}}")); + } + + /// A cancellation raised by the sink after its first chunk stops a large parallel + /// stream: no later chunk and no `]}}`, so the aborted body never looks complete. + #[test] + fn cancel_from_the_first_chunk_stops_a_parallel_select_json_stream() { + use std::fmt::Write as _; + use std::ops::ControlFlow; + use std::sync::{atomic::{AtomicBool, Ordering}, Arc}; + let mut nt = String::new(); + for i in 0..60_000 { + writeln!(nt, " \"v{i}\" .").unwrap(); + } + let graph = Graph::load_str(&nt, "ntriples").unwrap(); + let flag = Arc::new(AtomicBool::new(false)); + let budget = QueryBudget::cancelled_by(flag.clone()); + let (mut body, mut chunks) = (String::new(), 0); + let result = query_json_stream_with_budget(&graph, "SELECT ?s ?o WHERE { ?s ?o }", &budget, |c| { + chunks += 1; + body.push_str(&c); + flag.store(true, Ordering::Relaxed); + ControlFlow::Continue(()) + }); + assert_eq!(result.unwrap_err(), "query budget exceeded (cancelled)"); + assert_eq!(chunks, 1, "chunks kept flowing after cancellation"); + assert!(!body.ends_with("]}}"), "a cancelled stream was closed"); + } + // [GPT-6] Refusal precedes any output, including empty-result headers, on both // the scan fast path and the general evaluator. Cancellation works on wasm too. #[test] diff --git a/crates/sparq-engine/src/update.rs b/crates/sparq-engine/src/update.rs index b270c2f938..e4a64e9b2e 100644 --- a/crates/sparq-engine/src/update.rs +++ b/crates/sparq-engine/src/update.rs @@ -203,20 +203,21 @@ type Subst<'a> = dyn Fn(&Variable) -> Option + 'a; /// Fresh-blank-node state for ONE solution row: per SPARQL, every blank node in an INSERT /// template is instantiated FRESH per solution (same label, same row → same fresh node; /// different rows — and different operations in one request — get DIFFERENT nodes, hence -/// the process-wide counter). +/// a random id per node). struct FreshBnodes { map: FxHashMap, } -static FRESH_BNODE_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); - impl FreshBnodes { fn get(&mut self, label: &str) -> Term { if let Some(t) = self.map.get(label) { return t.clone(); } - let n = FRESH_BNODE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed); - let t = Term::BlankNode(BlankNode::new_unchecked(format!("fb{n}"))); + // A random 128-bit id (oxrdf's own fresh node), never a per-process counter: a counter + // restarts at 0, so its `_:fb0` could already be in a store (persisted by an earlier + // process, or loaded from a document using that label) and would merge with it. The + // RNG works on wasm too (the browser `getrandom` backend is already configured). + let t = Term::BlankNode(BlankNode::default()); self.map.insert(label.to_string(), t.clone()); t } @@ -410,7 +411,11 @@ fn load_document(source: &str) -> Result { _ => "turtle", }; let g = Graph::load_str(&text, format).map_err(|e| format!("LOAD {source}: {e}"))?; - Ok(decode_triples(&g)) + // LOAD merges the document into the destination (SPARQL 1.1 Update §3.1.5), and an RDF + // merge standardises the incoming blank nodes apart from those already in the store: one + // fresh node per document label, never a store's (or an earlier LOAD's) `_:b` (#4160). + let mut fresh = FreshBnodes { map: FxHashMap::default() }; + Ok(decode_triples(&g).iter().map(|t| t.each_ref().map(|x| freshen_term(x, &mut fresh))).collect()) } // --- the rebuild path ---------------------------------------------------------------------------- @@ -1660,6 +1665,54 @@ mod tests { let one_subj = crate::count(&g3, "SELECT DISTINCT ?s WHERE { ?s ?p ?o . FILTER(isBlank(?s)) }").unwrap(); assert_eq!(one_subj, 1, "same label in one INSERT DATA op is one node"); } + + /// #4160 — LOAD is an RDF merge: the document's blank nodes are standardised apart from + /// the store's and from every other LOAD's, on the rebuild and delta-overlay paths alike, + /// while one label within one document stays one node. + #[test] + fn load_blank_nodes_are_fresh() { + let dir = std::env::temp_dir().join(format!("sparq_load_4160_{}", std::process::id())); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + std::fs::write(dir.join("a.nt"), "_:b0 .\n_:b0 .\n").unwrap(); + std::fs::write(dir.join("b.nt"), "_:b0 .\n").unwrap(); + let (a, b) = (dir.join("a.nt"), dir.join("b.nt")); + let req = format!("LOAD ; LOAD ; LOAD ", a.display(), b.display(), a.display()); + let blank_subjects = "SELECT DISTINCT ?s WHERE { ?s ?p ?o . FILTER(isBlank(?s)) }"; + // The store already holds a `_:b0` of its own. + let src = "_:b0 ."; + + let g = Graph::load_str(src, "ntriples").unwrap(); + let g = with_load_base(dir.clone(), || update(&g, &req)).unwrap(); + // existing + a.nt + b.nt + a.nt again: four distinct nodes, a.nt's two triples on one. + assert_eq!(count(&g), 6); + assert_eq!(crate::count(&g, blank_subjects).unwrap(), 4, "LOAD conflated blank nodes"); + + let mut g2 = Graph::load_str(src, "ntriples").unwrap(); + with_load_base(dir.clone(), || update_in_place(&mut g2, &req)).unwrap(); + assert_eq!(count(&g2), 6); + assert_eq!(crate::count(&g2, blank_subjects).unwrap(), 4, "in-place LOAD conflated blank nodes"); + + let _ = std::fs::remove_dir_all(&dir); + } + + /// Fresh labels must not collide with labels a store already holds, such as the + /// `_:fbN` an earlier process minted before its counter restarted at 0. + #[test] + fn fresh_blank_nodes_avoid_labels_already_stored() { + use std::fmt::Write; + let mut src = String::new(); + for n in 0..4096 { + writeln!(src, "_:fb{n} .").unwrap(); + } + let blank_subjects = "SELECT DISTINCT ?s WHERE { ?s ?p ?o . FILTER(isBlank(?s)) }"; + let req = "INSERT DATA { _:b0 . }"; + let g = update(&Graph::load_str(&src, "ntriples").unwrap(), req).unwrap(); + assert_eq!(crate::count(&g, blank_subjects).unwrap(), 4097, "INSERT DATA reused a stored blank node"); + let mut g2 = Graph::load_str(&src, "ntriples").unwrap(); + update_in_place(&mut g2, req).unwrap(); + assert_eq!(crate::count(&g2, blank_subjects).unwrap(), 4097, "in-place INSERT DATA reused a stored blank node"); + } } /// [OPUS-4.8] (sq-qu8o) Pins the SPARQL 1.1 Update SILENT / non-SILENT error contract per diff --git a/crates/sparq-engine/tests/federation_features/service_local.rs b/crates/sparq-engine/tests/federation_features/service_local.rs index d5aaf7d9c1..2bafd6292d 100644 --- a/crates/sparq-engine/tests/federation_features/service_local.rs +++ b/crates/sparq-engine/tests/federation_features/service_local.rs @@ -410,3 +410,32 @@ fn without_the_service_feature_an_unregistered_iri_is_unsupported() { .expect("SILENT degrades to the join identity"); assert_eq!(res.rows.len(), 2); } + +/// #4438 — a SERVICE on the LEFT of a join whose bind-join pushdown declines (a local +/// handler always declines it) runs the handler ONCE, not once eagerly and once more +/// on the fall-through. +#[test] +fn service_on_the_left_of_a_join_runs_the_handler_once() { + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::Arc; + + let calls = Arc::new(AtomicUsize::new(0)); + let mut reg = LocalServiceRegistry::new(); + let counter = calls.clone(); + reg.register(HANDLER, move |_req| { + counter.fetch_add(1, Ordering::SeqCst); + Ok(LocalServiceRows::new( + vec![var("s"), var("name")], + vec![vec![Some(iri("http://ex/alice")), Some(lit("Alice"))]], + )) + }); + let g = local_graph(); + let q = format!( + "PREFIX ex: \n\ + SELECT ?s ?name WHERE {{ SERVICE <{}> {{ ?s ex:name ?name }} . ?s a ex:Person }}", + HANDLER + ); + let res = with_local_services(®, || query(&g, &q)).expect("query"); + assert_eq!(column(&res, "name"), vec!["\"Alice\"".to_string()]); + assert_eq!(calls.load(Ordering::SeqCst), 1, "handler ran more than once"); +}