Skip to content
Merged
5 changes: 5 additions & 0 deletions bench/feature-off-declarations/6655.json
Original file line number Diff line number Diff line change
@@ -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."
}
203 changes: 160 additions & 43 deletions crates/sparq-engine/src/exec.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1929,18 +1929,19 @@ pub(crate) mod functions {
static ACTIVE: RefCell<Option<Arc<FunctionRegistry>>> = 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<Arc<FunctionRegistry>>);
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
Expand Down Expand Up @@ -2072,16 +2073,16 @@ pub(crate) mod aggregates {
static ACTIVE: RefCell<Option<Arc<CustomAggregateRegistry>>> = const { RefCell::new(None) };
}

pub(crate) struct Guard;
pub(crate) struct Guard(Option<Arc<CustomAggregateRegistry>>);
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))]
Expand Down Expand Up @@ -2276,16 +2277,16 @@ pub(crate) mod spatial {
static ACTIVE: RefCell<Option<Arc<dyn SpatialProvider>>> = const { RefCell::new(None) };
}

pub(crate) struct Guard;
pub(crate) struct Guard(Option<Arc<dyn SpatialProvider>>);
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<dyn SpatialProvider>) -> 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.
Expand Down Expand Up @@ -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
Expand All @@ -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(());
}
}
Expand All @@ -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(())
}

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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];

Expand Down Expand Up @@ -10231,8 +10255,11 @@ fn try_theta_antijoin(
let all_cands: Vec<usize> = (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,
Expand Down Expand Up @@ -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<Term> {
functions::lookup(F).map(|f| f(&[]).unwrap())
}

fn lit(tag: &str) -> Option<Term> {
Some(Term::Literal(oxrdf::Literal::new_simple_literal(tag)))
}

#[test]
fn nested_function_registry_restores_the_outer_one() {
crate::with_functions(&registry("outer"), || {
assert_eq!(call(), lit("outer"));
crate::with_functions(&registry("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("<http://ex/s> <http://ex/p> <http://ex/o> .", "turtle").unwrap();
let r = crate::query(&g, "SELECT ?v WHERE { BIND(<http://ex/f>() 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(&registry("outer"), || {
let unwound = std::panic::catch_unwind(|| {
crate::with_functions(&registry("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<Term>]| Ok(lit(tag)));
reg
};
let call = || aggregates::lookup(F).map(|f| f(&[]).unwrap());
crate::aggregate::with_aggregates(&reg("outer"), || {
crate::aggregate::with_aggregates(&reg("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<Vec<Term>> {
None
}
fn is_indexed(&self, _: &Term) -> bool {
false
}
}

#[test]
fn nested_spatial_index_restores_the_outer_one() {
let outer: Arc<dyn crate::SpatialProvider> = Arc::new(NoIndex);
let inner: Arc<dyn crate::SpatialProvider> = Arc::new(NoIndex);
let is = |want: &Arc<dyn crate::SpatialProvider>| {
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());
}
}
50 changes: 50 additions & 0 deletions crates/sparq-engine/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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, "<http://ex/s{i}> <http://ex/p> \"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 <http://ex/p> ?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]
Expand Down
Loading
Loading