diff --git a/crates/sparq-canon/README.md b/crates/sparq-canon/README.md index fd2cf18b24..4e9f6a3e18 100644 --- a/crates/sparq-canon/README.md +++ b/crates/sparq-canon/README.md @@ -38,10 +38,10 @@ let map = sparq_canon::issued_identifiers(&[q]).unwrap(); // issuer map return a `CanonicalGraph` (sorted canonical N-Quads lines + re-parsed canonical triples) — what the ZK per-graph commitment pipeline consumes (`leaf_index = line index`). -- **Fail-closed on poison graphs** — RDFC-1.0's pathological blow-ups hit the - HNDQ call-limit guard and surface as `CanonError::Canonicalization`; RDF 1.2 - triple terms are outside the standard data model, so the standard paths fail - closed with `CanonError::TripleTerm` unless `rdf12-triple-terms` is enabled. +- **Fail-closed on poison graphs** — RDFC-1.0's pathological blow-ups hit the HNDQ call-limit guard + (`CanonError::Canonicalization`). RDF 1.2 triple terms and directional literals (`"…"@en--ltr`) are outside + RDFC-1.0: the standard paths fail closed with `CanonError::TripleTerm` / `CanonError::DirectionalLiteral` + unless the `rdf12-triple-terms` profile is enabled. - **W3C-conformant** — validated against the official [rdf-canon test suite] (eval + issued-map + negative cases, SHA-256 and SHA-384) through this crate's own public API (`tests/rdf_canon_suite.rs`). diff --git a/crates/sparq-canon/src/lib.rs b/crates/sparq-canon/src/lib.rs index aaa0531ac6..05a05471a7 100644 --- a/crates/sparq-canon/src/lib.rs +++ b/crates/sparq-canon/src/lib.rs @@ -54,7 +54,10 @@ //! //! Triple terms (`Term::Triple`, the RDF-1.2 `<<( s p o )>>` object) are //! **outside** the W3C RDFC-1.0 data model, so the standard paths above fail -//! closed with [`CanonError::TripleTerm`]. Enabling the **opt-in, off-by-default** +//! closed with [`CanonError::TripleTerm`]. RDF-1.2 directional-language +//! literals (`"…"@en--ltr`) are likewise outside that RDF-1.1 data model and +//! fail closed with [`CanonError::DirectionalLiteral`] on the standard paths. +//! Enabling the **opt-in, off-by-default** //! `rdf12-triple-terms` cargo feature adds a *separate, clearly non-standard* v2 //! profile (`canonicalize_rdf12`, `canonicalize_triples_rdf12`, …) that //! natively re-implements the RDFC-1.0 algorithm over oxrdf-0.3 and **descends @@ -209,6 +212,14 @@ pub enum CanonError { /// present (not feature-gated), per the `NestedBlankNode` precedent. /// [FABLE-5] sq-x3oj2. TripleTermDepthExceeded, + /// An RDF 1.2 directional-language literal (`"…"@lang--ltr` / `--rtl`, + /// [`oxrdf::BaseDirection`]) reached a **standard** (`rdf-canon`-backed) + /// entry point. RDFC-1.0 is defined over RDF 1.1, and the standard path's + /// oxrdf-0.2 bridge cannot represent a base direction, so these paths fail + /// closed with this typed error instead of a generic [`CanonError::Bridge`]. + /// The opt-in, non-standard `rdf12-triple-terms` profile (`canonicalize_rdf12` + /// and siblings) canonicalizes directional literals natively. GitHub #5359. + DirectionalLiteral, /// Bridge serialization/parse failure (should not happen for RDFC-1.0-model /// content; surfaced rather than swallowed). Bridge(String), @@ -242,6 +253,13 @@ impl std::fmt::Display for CanonError { triple terms are a stack-overflow vector for recursive descent)" ) } + CanonError::DirectionalLiteral => { + write!( + f, + "RDF 1.2 directional-language literals are outside the W3C RDFC-1.0 \ + data model; enable the `rdf12-triple-terms` profile to canonicalize them" + ) + } CanonError::Bridge(e) => write!(f, "oxrdf bridge error: {e}"), CanonError::Canonicalization(e) => write!(f, "RDFC-1.0 canonicalization failed: {e}"), } @@ -422,6 +440,7 @@ pub fn graph_triples(g: &Graph) -> Result, CanonError> { /// literal word `DEFAULT`, so the default graph is emitted explicitly as a /// 3-term line. fn bridge_to_02(dataset: &[Quad]) -> Result, CanonError> { + reject_unbridgeable(dataset.iter().map(|q| &q.object))?; let doc = serialize_quads(dataset)?; parse_02(&doc) } @@ -479,6 +498,7 @@ fn serialize_quads(dataset: &[Quad]) -> Result { } fn bridge_triples_to_02(triples: &[Triple]) -> Result, CanonError> { + reject_unbridgeable(triples.iter().map(|t| &t.object))?; let mut doc = String::new(); #[cfg(feature = "bridge-lowcopy")] use std::fmt::Write as _; @@ -538,6 +558,29 @@ fn parse_canonical(canonical: &str) -> Result { Ok(CanonicalGraph { lines, triples }) } +/// Fails closed on object terms the RDF-1.1 oxrdf-0.2 bridge cannot carry, +/// before serialization: triple terms ([`CanonError::TripleTerm`], which takes +/// precedence) and directional-language literals +/// ([`CanonError::DirectionalLiteral`]). Literals only occur in object +/// position, and triple terms are rejected outright, so the top-level objects +/// are the only places a base direction can appear. GitHub #5359. +fn reject_unbridgeable<'a>( + objects: impl Iterator, +) -> Result<(), CanonError> { + let mut directional = false; + for o in objects { + match o { + oxrdf::Term::Triple(_) => return Err(CanonError::TripleTerm), + oxrdf::Term::Literal(l) if l.direction().is_some() => directional = true, + _ => {} + } + } + if directional { + return Err(CanonError::DirectionalLiteral); + } + Ok(()) +} + fn contains_triple_term(t: &Triple) -> bool { matches!(t.object, oxrdf::Term::Triple(_)) } diff --git a/crates/sparq-canon/tests/directional_literal.rs b/crates/sparq-canon/tests/directional_literal.rs new file mode 100644 index 0000000000..a53800d1cb --- /dev/null +++ b/crates/sparq-canon/tests/directional_literal.rs @@ -0,0 +1,71 @@ +//! GitHub #5359: RDF 1.2 directional-language literals (`"…"@en--ltr`) are +//! outside RDFC-1.0's RDF-1.1 data model and cannot cross the standard path's +//! oxrdf-0.2 bridge. Every standard entry point must fail closed with the typed +//! `CanonError::DirectionalLiteral`, not a generic `CanonError::Bridge`. + +use oxrdf::{BaseDirection, BlankNode, GraphName, Literal, NamedNode, Quad, Triple}; +use sparq_canon::CanonError; + +fn directional_triple() -> Triple { + Triple::new( + BlankNode::new("b0").unwrap(), + NamedNode::new("http://ex/p").unwrap(), + Literal::new_directional_language_tagged_literal("hello", "en", BaseDirection::Ltr) + .unwrap(), + ) +} + +fn directional_quad() -> Quad { + let t = directional_triple(); + Quad::new(t.subject, t.predicate, t.object, GraphName::DefaultGraph) +} + +fn assert_directional(r: Result) { + match r { + Err(CanonError::DirectionalLiteral) => {} + other => panic!("expected CanonError::DirectionalLiteral, got {other:?}"), + } +} + +#[test] +fn standard_dataset_paths_reject_directional_literal_with_typed_error() { + let ds = [directional_quad()]; + assert_directional(sparq_canon::canonicalize(&ds)); + assert_directional(sparq_canon::canonicalize_quads(&ds)); + assert_directional(sparq_canon::issue_quads(&ds)); + assert_directional(sparq_canon::issued_identifiers(&ds)); +} + +#[test] +fn standard_single_graph_paths_reject_directional_literal_with_typed_error() { + let ts = [directional_triple()]; + assert_directional(sparq_canon::canonicalize_triples(&ts)); + assert_directional(sparq_canon::issue_triples(&ts)); +} + +#[test] +fn nquads_text_entry_point_rejects_directional_literal_with_typed_error() { + let doc = "_:b0 \"hello\"@en--rtl .\n"; + assert_directional(sparq_canon::canonicalize_nquads(doc)); +} + +#[test] +fn directional_literal_error_names_the_rdf12_profile() { + let msg = CanonError::DirectionalLiteral.to_string(); + assert!(msg.contains("rdf12-triple-terms"), "{msg}"); +} + +#[test] +fn plain_language_literal_still_canonicalizes() { + let doc = "_:b0 \"hello\"@en .\n"; + let canon = sparq_canon::canonicalize_nquads(doc).unwrap(); + assert_eq!(canon, "_:c14n0 \"hello\"@en .\n"); +} + +/// The opt-in non-standard profile canonicalizes directional literals natively. +#[cfg(feature = "rdf12-triple-terms")] +#[test] +fn rdf12_profile_canonicalizes_directional_literal() { + let canon = sparq_canon::rdf12::canonicalize_rdf12(&[directional_quad()]).unwrap(); + assert_eq!(canon, "_:c14n0 \"hello\"@en--ltr .\n"); +} diff --git a/crates/sparq-gpu/examples/gpu_bench.rs b/crates/sparq-gpu/examples/gpu_bench.rs index 47c6ffb0c8..4aa9c0f065 100644 --- a/crates/sparq-gpu/examples/gpu_bench.rs +++ b/crates/sparq-gpu/examples/gpu_bench.rs @@ -36,6 +36,9 @@ impl Rng { } const PAR_CHUNK: usize = 64 * 1024; +/// A kernel call only fails if the device stalls past `sparq_gpu::POLL_TIMEOUT`; +/// the panic then unwinds safely (a stalled `Gpu` leaks rather than waits on drop). +const GPU_OK: &str = "GPU kernel stalled past POLL_TIMEOUT"; /// A named benchmark leg returning a checksum (asserted equal across legs). type Variant<'a> = (&'a str, Box u64 + 'a>); @@ -122,7 +125,7 @@ fn main() { let col: Vec = (0..n).map(|_| rng.u32()).collect(); let (lo, hi) = (0u32, u32::MAX / 8); let resident = gpu.upload_u32(&col); - let _ = gpu.filter_count_u32(&resident, lo, hi); // warm pipeline/caches + let _ = gpu.filter_count_u32(&resident, lo, hi).expect(GPU_OK); // warm pipeline/caches let mut variants: Vec = vec![ ("cpu1", Box::new(|| cpu::filter_count_u32(&col, lo, hi))), @@ -136,13 +139,13 @@ fn main() { ), ( "gpu resident", - Box::new(|| gpu.filter_count_u32(&resident, lo, hi)), + Box::new(|| gpu.filter_count_u32(&resident, lo, hi).expect(GPU_OK)), ), ( "gpu e2e", Box::new(|| { gpu.write_u32(&resident, &col); - gpu.filter_count_u32(&resident, lo, hi) + gpu.filter_count_u32(&resident, lo, hi).expect(GPU_OK) }), ), ]; @@ -163,7 +166,7 @@ fn main() { let col: Vec = (0..n).map(|_| rng.u32() as f64 / 1e3).collect(); let t = u32::MAX as f64 / 1e3 * 0.875; let resident = gpu.upload_f64(&col); - let _ = gpu.filter_count_f64_gt(&resident, t); + let _ = gpu.filter_count_f64_gt(&resident, t).expect(GPU_OK); let mut variants: Vec = vec![ ("cpu1", Box::new(|| cpu::filter_count_f64_gt(&col, t))), @@ -177,13 +180,13 @@ fn main() { ), ( "gpu resident", - Box::new(|| gpu.filter_count_f64_gt(&resident, t)), + Box::new(|| gpu.filter_count_f64_gt(&resident, t).expect(GPU_OK)), ), ( "gpu e2e", Box::new(|| { gpu.write_f64(&resident, &col); - gpu.filter_count_f64_gt(&resident, t) + gpu.filter_count_f64_gt(&resident, t).expect(GPU_OK) }), ), ]; @@ -212,7 +215,7 @@ fn main() { let slots = cpu::build_hash_table(&build_keys, &payloads); let table = gpu.upload_table(&slots); let probe_col = gpu.upload_u32(&probe); - let _ = gpu.hash_probe(&table, &probe_col); + let _ = gpu.hash_probe(&table, &probe_col).expect(GPU_OK); let mut variants: Vec = vec![ ( @@ -235,7 +238,7 @@ fn main() { ( "gpu resident", Box::new(|| { - let (m, s) = gpu.hash_probe(&table, &probe_col); + let (m, s) = gpu.hash_probe(&table, &probe_col).expect(GPU_OK); m.wrapping_add(s) }), ), @@ -244,7 +247,7 @@ fn main() { Box::new(|| { gpu.write_table(&table, &slots); gpu.write_u32(&probe_col, &probe); - let (m, s) = gpu.hash_probe(&table, &probe_col); + let (m, s) = gpu.hash_probe(&table, &probe_col).expect(GPU_OK); m.wrapping_add(s) }), ), @@ -263,12 +266,15 @@ fn main() { continue; } const G: u32 = 256; + const KEYS_IN_RANGE: &str = "keys are generated `% G`"; let mut rng = Rng(0x0F0F_F0F0_1337_4242); let keys: Vec = (0..n).map(|_| rng.u32() % G).collect(); let vals: Vec = (0..n).map(|_| rng.u32()).collect(); let keys_col = gpu.upload_u32(&keys); let vals_col = gpu.upload_u32(&vals); - let _ = gpu.group_aggregate(&keys_col, &vals_col, G); + let _ = gpu + .group_aggregate(&keys_col, &vals_col, G) + .expect(KEYS_IN_RANGE); let fold = |rows: Vec<(u64, u64)>| -> u64 { rows.iter().fold(0u64, |acc, (c, s)| { @@ -286,7 +292,7 @@ fn main() { let mut variants: Vec = vec![ ( "cpu1", - Box::new(|| fold(cpu::group_aggregate(&keys, &vals, G))), + Box::new(|| fold(cpu::group_aggregate(&keys, &vals, G).expect(KEYS_IN_RANGE))), ), ( "cpuN", @@ -294,21 +300,29 @@ fn main() { let rows = keys .par_chunks(PAR_CHUNK) .zip(vals.par_chunks(PAR_CHUNK)) - .map(|(k, v)| cpu::group_aggregate(k, v, G)) + .map(|(k, v)| cpu::group_aggregate(k, v, G).expect(KEYS_IN_RANGE)) .reduce(|| vec![(0, 0); G as usize], merge); fold(rows) }), ), ( "gpu resident", - Box::new(|| fold(gpu.group_aggregate(&keys_col, &vals_col, G))), + Box::new(|| { + fold( + gpu.group_aggregate(&keys_col, &vals_col, G) + .expect(KEYS_IN_RANGE), + ) + }), ), ( "gpu e2e", Box::new(|| { gpu.write_u32(&keys_col, &keys); gpu.write_u32(&vals_col, &vals); - fold(gpu.group_aggregate(&keys_col, &vals_col, G)) + fold( + gpu.group_aggregate(&keys_col, &vals_col, G) + .expect(KEYS_IN_RANGE), + ) }), ), ]; diff --git a/crates/sparq-gpu/src/cpu.rs b/crates/sparq-gpu/src/cpu.rs index aec8604b19..0f32103e3d 100644 --- a/crates/sparq-gpu/src/cpu.rs +++ b/crates/sparq-gpu/src/cpu.rs @@ -5,7 +5,7 @@ //! the example — the library carries no thread-pool dependency). They are //! deliberately the same tight scalar loops a sparq scan thread runs today. -use crate::{hash32, EMPTY_KEY}; +use crate::{hash32, GroupKeyOutOfRange, EMPTY_KEY}; /// Counts elements with `lo <= v < hi`. pub fn filter_count_u32(col: &[u32], lo: u32, hi: u32) -> u64 { @@ -60,7 +60,8 @@ pub fn hash_probe(slots: &[[u32; 2]], probe: &[u32]) -> (u64, u64) { let mut sum = 0u64; for &k in probe { let mut slot = (hash32(k) & mask) as usize; - loop { + // Bounded at mask + 1 steps, mirroring the kernel (GitHub #4603). + for _ in 0..=mask { let [key, payload] = slots[slot]; if key == EMPTY_KEY { break; @@ -75,16 +76,22 @@ pub fn hash_probe(slots: &[[u32; 2]], probe: &[u32]) -> (u64, u64) { (matches, sum) } -/// COUNT + SUM GROUP BY; `keys[i] < groups`. Returns `(count, sum)` per group. -pub fn group_aggregate(keys: &[u32], vals: &[u32], groups: u32) -> Vec<(u64, u64)> { +/// COUNT + SUM GROUP BY. Returns `(count, sum)` per group, or +/// [`GroupKeyOutOfRange`] if any `keys[i] >= groups` — the same contract as +/// [`crate::Gpu::group_aggregate`] (GitHub #4603). +pub fn group_aggregate( + keys: &[u32], + vals: &[u32], + groups: u32, +) -> Result, GroupKeyOutOfRange> { assert_eq!(keys.len(), vals.len()); let mut out = vec![(0u64, 0u64); groups as usize]; for (&k, &v) in keys.iter().zip(vals) { - let e = &mut out[k as usize]; + let e = out.get_mut(k as usize).ok_or(GroupKeyOutOfRange)?; e.0 += 1; e.1 += u64::from(v); } - out + Ok(out) } // [OPUS-4.8] sq-goay: unit tests for the CPU reference itself. @@ -284,15 +291,33 @@ mod tests { assert_eq!(s, 10 + 100 + 1); } + /// GitHub #4603: a table with no EMPTY_KEY slot must not spin forever; the + /// walk is bounded at one pass over the table. + #[test] + fn hash_probe_terminates_on_a_full_table() { + let slots = [[1u32, 10], [2, 20], [3, 30], [4, 40]]; + assert_eq!(hash_probe(&slots, &[3, 9]), (1, 30)); + } + // ---- group_aggregate ------------------------------------------------- + /// GitHub #4603: an out-of-range key is a typed error, not a panic. + #[test] + fn group_aggregate_rejects_out_of_range_key() { + assert_eq!( + group_aggregate(&[0, 3], &[1, 1], 3), + Err(GroupKeyOutOfRange) + ); + assert_eq!(group_aggregate(&[600], &[1], 3), Err(GroupKeyOutOfRange)); + } + #[test] fn group_aggregate_brute_force_random() { let mut rng = Rng(0x9999_7777_5555_3333); let groups = 37u32; let keys: Vec = (0..6_000).map(|_| rng.u32() % groups).collect(); let vals: Vec = (0..6_000).map(|_| rng.u32()).collect(); - let got = group_aggregate(&keys, &vals, groups); + let got = group_aggregate(&keys, &vals, groups).unwrap(); // Independent oracle, per group. for g in 0..groups { @@ -317,7 +342,7 @@ mod tests { let n = 100usize; let keys = vec![0u32; n]; let vals = vec![u32::MAX; n]; - let got = group_aggregate(&keys, &vals, 1); + let got = group_aggregate(&keys, &vals, 1).unwrap(); assert_eq!(got[0], (n as u64, n as u64 * u64::from(u32::MAX))); assert!(got[0].1 > u64::from(u32::MAX), "sum overflows u32"); } diff --git a/crates/sparq-gpu/src/lib.rs b/crates/sparq-gpu/src/lib.rs index a63ca8af79..ec3d405912 100644 --- a/crates/sparq-gpu/src/lib.rs +++ b/crates/sparq-gpu/src/lib.rs @@ -44,6 +44,90 @@ pub const EMPTY_KEY: u32 = u32::MAX; /// under WebGPU's 16 KiB minimum guarantee). pub const MAX_GROUPS: u32 = 512; +/// Upper bound on how long [`Gpu`] waits for one kernel submission before +/// giving up with [`GpuError::Stalled`] instead of blocking the host thread +/// forever. Every kernel is O(n) with bounded per-element work, so a healthy +/// device finishes far inside this; hitting it means a wedged device. +/// GitHub #4603. +pub const POLL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60); + +/// A GROUP BY key was `>= groups` (the documented `keys[i] < groups` +/// precondition was violated). Returned by both [`Gpu::group_aggregate`] and +/// the CPU oracle [`cpu::group_aggregate`], so the two sides agree on malformed +/// input instead of the GPU silently dropping the row while the CPU panics. +/// GitHub #4603. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct GroupKeyOutOfRange; + +impl std::fmt::Display for GroupKeyOutOfRange { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str("GROUP BY key out of range: every key must be < groups") + } +} + +impl std::error::Error for GroupKeyOutOfRange {} + +/// Why a [`Gpu`] kernel call failed. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum GpuError { + /// The device did not finish a submission within [`POLL_TIMEOUT`] (or the + /// poll itself failed). The [`Gpu`] is now permanently stalled: every later + /// kernel call returns this error without touching the device, and dropping + /// the [`Gpu`] leaks its device and queue instead of waiting on them (see + /// the `Drop` impl). Recreate a [`Gpu`] to retry. + Stalled, + /// See [`GroupKeyOutOfRange`] (only [`Gpu::group_aggregate`]). + GroupKeyOutOfRange, +} + +impl std::fmt::Display for GpuError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + GpuError::Stalled => write!( + f, + "GPU device stalled: a submission did not complete within {POLL_TIMEOUT:?}" + ), + GpuError::GroupKeyOutOfRange => GroupKeyOutOfRange.fmt(f), + } + } +} + +impl std::error::Error for GpuError {} + +impl From for GpuError { + fn from(_: GroupKeyOutOfRange) -> Self { + GpuError::GroupKeyOutOfRange + } +} + +/// Keeps `handle`'s shared resource alive for the rest of the process by +/// forgetting one extra strong reference to it, so its destructor never runs. +/// Used on wgpu's `Device`/`Queue` (both `Arc`-backed `Clone` handles) once a +/// submission has timed out — see `impl Drop for Gpu`. +fn leak_handle(handle: &T) { + std::mem::forget(handle.clone()); +} + +/// Validates a host-built open-addressing table before it reaches the device: +/// the slot count must be a non-zero power of two (the kernel masks with +/// `len - 1`) and at least one slot must be `EMPTY_KEY`, so every probe walk has +/// a terminating sentinel. Panics otherwise. GitHub #4603. +fn check_table(slots: &[[u32; 2]]) { + assert!( + slots.len().is_power_of_two(), + "hash table slot count must be a non-zero power of two, got {}", + slots.len() + ); + assert!( + u32::try_from(slots.len()).is_ok(), + "hash table slot count must fit in u32" + ); + assert!( + slots.iter().any(|s| s[0] == EMPTY_KEY), + "hash table must contain at least one EMPTY_KEY slot (load factor < 1)" + ); +} + /// The 32-bit mixer used by both the CPU reference and the WGSL kernel (they /// must agree bit-for-bit so both sides walk identical probe sequences). #[inline] @@ -58,7 +142,54 @@ pub fn hash32(mut x: u32) -> u32 { #[cfg(test)] mod tests { - use super::{hash32, EMPTY_KEY}; + use super::{check_table, hash32, leak_handle, GpuError, GroupKeyOutOfRange, EMPTY_KEY}; + + #[test] + fn leak_handle_keeps_the_shared_resource_alive() { + let handle = std::rc::Rc::new(()); + let weak = std::rc::Rc::downgrade(&handle); + leak_handle(&handle); + drop(handle); + // The destructor of the shared value never runs: a stalled wgpu queue is + // never torn down (and so never waited on). + assert!(weak.upgrade().is_some()); + } + + #[test] + fn gpu_error_display_and_from() { + assert_eq!( + GpuError::from(GroupKeyOutOfRange), + GpuError::GroupKeyOutOfRange + ); + assert_eq!( + GpuError::GroupKeyOutOfRange.to_string(), + GroupKeyOutOfRange.to_string() + ); + assert!(GpuError::Stalled.to_string().contains("stalled")); + } + + #[test] + fn check_table_accepts_a_built_table() { + check_table(&crate::cpu::build_hash_table(&[1, 2, 3], &[4, 5, 6])); + } + + #[test] + #[should_panic(expected = "EMPTY_KEY slot")] + fn check_table_rejects_a_full_table() { + check_table(&[[1, 0], [2, 0], [3, 0], [4, 0]]); + } + + #[test] + #[should_panic(expected = "power of two")] + fn check_table_rejects_a_non_power_of_two_table() { + check_table(&[[EMPTY_KEY, 0], [EMPTY_KEY, 0], [EMPTY_KEY, 0]]); + } + + #[test] + #[should_panic(expected = "power of two")] + fn check_table_rejects_an_empty_table() { + check_table(&[]); + } // [GPT-5.6] sq-pz5rf: Exercise the mixer directly, including values around // the hash-table sentinel and a bounded sample of the full u32 domain. @@ -199,7 +330,12 @@ fn main(@builtin(local_invocation_id) lid: vec3, if (i < params.n) { let k = probe[i]; var slot = hash32(k) & params.mask; + // Bounded walk: at most mask + 1 slots (every slot once), so a table + // with no EMPTY_KEY slot cannot spin forever (GitHub #4603). + var steps = 0u; loop { + if (steps > params.mask) { break; } + steps += 1u; let e = table[slot]; if (e.x == 0xFFFFFFFFu) { break; } if (e.x == k) { @@ -231,7 +367,9 @@ fn main(@builtin(local_invocation_id) lid: vec3, } "#; -/// COUNT + SUM GROUP BY: keys must already be in [0, g), g <= MAX_GROUPS. +/// COUNT + SUM GROUP BY: keys must already be in [0, g), g <= MAX_GROUPS. A key +/// `>= g` is not accumulated; it sets the out-of-range flag word `out[3g]`, which +/// the host turns into [`GroupKeyOutOfRange`] (GitHub #4603). /// Two-level atomics: every thread accumulates into workgroup-shared per-group /// counters (the compute-dense part), then each workgroup flushes its g partial /// rows to global atomics once. u64 sums via the same lo/carry/hi emulation. @@ -240,7 +378,7 @@ const MAX_G: u32 = 512u; struct Params { n: u32, g: u32, _p0: u32, _p1: u32 }; @group(0) @binding(0) var keys: array; @group(0) @binding(1) var vals: array; -@group(0) @binding(2) var out: array>; // [count;g][sum_lo;g][sum_hi;g] +@group(0) @binding(2) var out: array>; // [count;g][sum_lo;g][sum_hi;g][oob_flag] @group(0) @binding(3) var params: Params; var s_cnt: array, MAX_G>; @@ -263,9 +401,13 @@ fn main(@builtin(local_invocation_id) lid: vec3, if (i < params.n) { let k = keys[i]; let v = vals[i]; - atomicAdd(&s_cnt[k], 1u); - let old = atomicAdd(&s_lo[k], v); - if (old + v < old) { atomicAdd(&s_hi[k], 1u); } + if (k < params.g) { + atomicAdd(&s_cnt[k], 1u); + let old = atomicAdd(&s_lo[k], v); + if (old + v < old) { atomicAdd(&s_hi[k], 1u); } + } else { + atomicOr(&out[3u * params.g], 1u); + } } workgroupBarrier(); g = lid.x; @@ -342,6 +484,24 @@ pub struct Gpu { filter_f64: wgpu::ComputePipeline, hash_probe: wgpu::ComputePipeline, group_agg: wgpu::ComputePipeline, + /// Set once a submission fails to complete within [`POLL_TIMEOUT`]. + stalled: std::sync::atomic::AtomicBool, +} + +/// GitHub #4603: a timed-out submission may still be pending on the device, +/// and wgpu-core 30's `Drop for Queue` calls the HAL's `wait_for_idle()` with +/// no timeout — so a plain drop of a stalled [`Gpu`] (including during panic +/// unwinding) would block forever on the very submission we gave up on. Once +/// stalled, we deliberately leak the device and queue instead: their +/// destructors never run, and the driver reclaims them at process exit. A +/// healthy [`Gpu`] drops normally. +impl Drop for Gpu { + fn drop(&mut self) { + if self.stalled.load(std::sync::atomic::Ordering::Acquire) { + leak_handle(&self.queue); + leak_handle(&self.device); + } + } } impl Gpu { @@ -407,6 +567,7 @@ impl Gpu { queue, info, max_storage_bytes, + stalled: std::sync::atomic::AtomicBool::new(false), }) } @@ -427,8 +588,12 @@ impl Gpu { } /// Uploads a host-built open-addressing table (see [`cpu::build_hash_table`]). + /// + /// # Panics + /// If the slot count is not a non-zero power of two, or no slot is + /// `EMPTY_KEY` (a full table would give the probe walk no sentinel). pub fn upload_table(&self, slots: &[[u32; 2]]) -> HashTable { - debug_assert!(slots.len().is_power_of_two()); + check_table(slots); HashTable { buf: self.storage_buffer(bytemuck::cast_slice(slots)), mask: (slots.len() - 1) as u32, @@ -451,6 +616,9 @@ impl Gpu { pub fn write_table(&self, table: &HashTable, slots: &[[u32; 2]]) { assert_eq!(table.mask as usize + 1, slots.len()); + // No `check_table` scan here: this is the timed host->device refill in the + // e2e benchmark legs, and an extra O(n) host pass would skew them. A full + // table written here still cannot hang — the probe walk is bounded. self.queue .write_buffer(&table.buf, 0, bytemuck::cast_slice(slots)); } @@ -468,7 +636,7 @@ impl Gpu { // -- kernels -- /// Counts elements with `lo <= v < hi` in a resident u32 column. - pub fn filter_count_u32(&self, col: &ColU32, lo: u32, hi: u32) -> u64 { + pub fn filter_count_u32(&self, col: &ColU32, lo: u32, hi: u32) -> Result { let r = self.run( &self.filter_u32, &[col.buf.as_entire_binding()], @@ -480,15 +648,15 @@ impl Gpu { }, col.len, 1, - ); - u64::from(r[0]) + )?; + Ok(u64::from(r[0])) } /// Counts elements with `v > t` (NaN-excluding, IEEE semantics, exact) in a /// resident f64 column. - pub fn filter_count_f64_gt(&self, col: &ColF64, t: f64) -> u64 { + pub fn filter_count_f64_gt(&self, col: &ColF64, t: f64) -> Result { if t.is_nan() { - return 0; // v > NaN is false for every v; no dispatch needed + return Ok(0); // v > NaN is false for every v; no dispatch needed } let key = f64_order_key(if t == 0.0 { 0.0 } else { t }); // canonicalise -0.0 let r = self.run( @@ -502,15 +670,15 @@ impl Gpu { }, col.len, 1, - ); - u64::from(r[0]) + )?; + Ok(u64::from(r[0])) } /// Probes every id in `probe` against the resident `table`; returns /// (number of matches, exact u64 sum of matched payloads). Both are exact /// u64s on the device (lo/carry/hi atomic emulation) — duplicate build keys /// can push a high-fanout join past u32::MAX matches. - pub fn hash_probe(&self, table: &HashTable, probe: &ColU32) -> (u64, u64) { + pub fn hash_probe(&self, table: &HashTable, probe: &ColU32) -> Result<(u64, u64), GpuError> { let r = self.run( &self.hash_probe, &[table.buf.as_entire_binding(), probe.buf.as_entire_binding()], @@ -522,16 +690,22 @@ impl Gpu { }, probe.len, 4, - ); - ( + )?; + Ok(( (u64::from(r[1]) << 32) | u64::from(r[0]), (u64::from(r[3]) << 32) | u64::from(r[2]), - ) + )) } - /// COUNT + SUM GROUP BY over resident columns. `keys[i]` must be `< groups` - /// and `groups <= MAX_GROUPS`. Returns `(count, sum)` per group. - pub fn group_aggregate(&self, keys: &ColU32, vals: &ColU32, groups: u32) -> Vec<(u64, u64)> { + /// COUNT + SUM GROUP BY over resident columns. `groups <= MAX_GROUPS`. + /// Returns `(count, sum)` per group, or [`GpuError::GroupKeyOutOfRange`] if + /// any `keys[i] >= groups` (detected on the device, no host scan). + pub fn group_aggregate( + &self, + keys: &ColU32, + vals: &ColU32, + groups: u32, + ) -> Result, GpuError> { assert!( (1..=MAX_GROUPS).contains(&groups), "groups must be in 1..={MAX_GROUPS}" @@ -547,21 +721,26 @@ impl Gpu { d: 0, }, keys.len, - 3 * groups as usize, - ); + 3 * groups as usize + 1, + )?; let g = groups as usize; - (0..g) + if r[3 * g] != 0 { + return Err(GpuError::GroupKeyOutOfRange); + } + Ok((0..g) .map(|i| { let count = u64::from(r[i]); let sum = (u64::from(r[2 * g + i]) << 32) | u64::from(r[g + i]); (count, sum) }) - .collect() + .collect()) } /// Shared dispatch path: binds `inputs` at 0.., a zero-initialised result /// buffer of `result_words` u32s next, the uniform params last; dispatches a /// 2D-tiled grid covering `n` elements; reads the result back synchronously. + /// Fails with [`GpuError::Stalled`] (and marks `self` stalled) if the + /// submission does not complete within [`POLL_TIMEOUT`]. fn run( &self, pipeline: &wgpu::ComputePipeline, @@ -569,8 +748,12 @@ impl Gpu { params: Params, n: u32, result_words: usize, - ) -> Vec { + ) -> Result, GpuError> { + use std::sync::atomic::Ordering; use wgpu::util::DeviceExt; + if self.stalled.load(Ordering::Acquire) { + return Err(GpuError::Stalled); + } let result_bytes = (result_words * 4) as u64; // Freshly created buffers are zero-initialised per the WebGPU spec. let result = self.device.create_buffer(&wgpu::BufferDescriptor { @@ -633,7 +816,7 @@ impl Gpu { cp.dispatch_workgroups(gx, gy, 1); } enc.copy_buffer_to_buffer(&result, 0, &readback, 0, result_bytes); - self.queue.submit(Some(enc.finish())); + let submission = self.queue.submit(Some(enc.finish())); let slice = readback.slice(..); let (tx, rx) = std::sync::mpsc::channel(); @@ -646,9 +829,22 @@ impl Gpu { // wgpu 30: `get_mapped_range()` now returns `Result` // instead of `BufferView` directly — unwrap because a mapping error here is // a GPU protocol violation (we waited for the map to succeed above). - self.device - .poll(wgpu::PollType::wait_indefinitely()) - .expect("device poll failed"); + // GitHub #4603: wait for this submission with a bound ([`POLL_TIMEOUT`]) + // rather than `wait_indefinitely()`, so a wedged device fails instead of + // hanging the host thread. The failure is a returned error, not a panic: + // we mark `self` stalled first so neither the caller's eventual drop nor + // a panic unwind runs wgpu's unbounded queue-idle wait (see `Drop`). + if self + .device + .poll(wgpu::PollType::Wait { + submission_index: Some(submission), + timeout: Some(POLL_TIMEOUT), + }) + .is_err() + { + self.stalled.store(true, Ordering::Release); + return Err(GpuError::Stalled); + } rx.recv() .expect("map_async callback dropped") .expect("readback map failed"); @@ -656,7 +852,7 @@ impl Gpu { bytemuck::cast_slice(&slice.get_mapped_range().expect("get_mapped_range failed")) .to_vec(); readback.unmap(); - out + Ok(out) } } diff --git a/crates/sparq-gpu/tests/correctness.rs b/crates/sparq-gpu/tests/correctness.rs index 2e1b3e624e..3ff58ba165 100644 --- a/crates/sparq-gpu/tests/correctness.rs +++ b/crates/sparq-gpu/tests/correctness.rs @@ -6,7 +6,7 @@ //! exact by construction (bit-pattern total-order comparison, not float math on //! the device), so it is asserted exactly too — including NaN/±inf/-0.0 edges. -use sparq_gpu::{cpu, Gpu, EMPTY_KEY, MAX_GROUPS}; +use sparq_gpu::{cpu, Gpu, GpuError, GroupKeyOutOfRange, EMPTY_KEY, MAX_GROUPS}; /// Deterministic xorshift64* stream (same generator the wgpu spike used). struct Rng(u64); @@ -63,7 +63,7 @@ fn filter_u32_matches_cpu() { ] { let expect = cpu::filter_count_u32(&col, lo, hi); assert_eq!( - gpu.filter_count_u32(&resident, lo, hi), + gpu.filter_count_u32(&resident, lo, hi).unwrap(), expect, "lo={lo} hi={hi}" ); @@ -111,7 +111,11 @@ fn filter_f64_matches_cpu_including_ieee_edges() { 5e-324, ] { let expect = cpu::filter_count_f64_gt(&col, t); - assert_eq!(gpu.filter_count_f64_gt(&resident, t), expect, "t={t}"); + assert_eq!( + gpu.filter_count_f64_gt(&resident, t).unwrap(), + expect, + "t={t}" + ); } } @@ -144,7 +148,7 @@ fn hash_probe_matches_cpu() { let probe_col = gpu.upload_u32(&probe); let expect = cpu::hash_probe(&slots, &probe); - assert_eq!(gpu.hash_probe(&table, &probe_col), expect); + assert_eq!(gpu.hash_probe(&table, &probe_col).unwrap(), expect); } #[test] @@ -161,14 +165,57 @@ fn group_aggregate_matches_cpu() { let keys_col = gpu.upload_u32(&keys); let vals_col = gpu.upload_u32(&vals); - let expect = cpu::group_aggregate(&keys, &vals, groups); - let got = gpu.group_aggregate(&keys_col, &vals_col, groups); + let expect = cpu::group_aggregate(&keys, &vals, groups).unwrap(); + let got = gpu.group_aggregate(&keys_col, &vals_col, groups).unwrap(); assert_eq!(got, expect, "groups={groups}"); let total: u64 = got.iter().map(|(c, _)| c).sum(); assert_eq!(total, N as u64, "every row lands in exactly one group"); } } +/// GitHub #4603: an out-of-range GROUP BY key must be rejected identically by +/// the GPU and the CPU oracle — never silently dropped on the GPU. Covers a key +/// in `groups..MAX_GROUPS` (lands in a workgroup slot that is never flushed) and +/// one `>= MAX_GROUPS` (past the shared-memory array), at a row deep inside a +/// multi-workgroup dispatch. +#[test] +fn group_aggregate_rejects_out_of_range_keys_like_cpu() { + let Some(gpu) = gpu_or_skip("group_aggregate_out_of_range") else { + return; + }; + let groups = 7u32; + for bad in [groups, MAX_GROUPS - 1, MAX_GROUPS, u32::MAX] { + let mut keys: Vec = (0..10_000u32).map(|i| i % groups).collect(); + keys[5_000] = bad; + let vals = vec![1u32; keys.len()]; + let keys_col = gpu.upload_u32(&keys); + let vals_col = gpu.upload_u32(&vals); + assert_eq!( + gpu.group_aggregate(&keys_col, &vals_col, groups), + Err(GpuError::GroupKeyOutOfRange), + "gpu, bad key {bad}" + ); + assert_eq!( + cpu::group_aggregate(&keys, &vals, groups), + Err(GroupKeyOutOfRange), + "cpu, bad key {bad}" + ); + } +} + +/// GitHub #4603: `upload_table` rejects a table with no EMPTY_KEY slot (in +/// release builds too) instead of handing the probe walk a table it could +/// never leave. +#[test] +fn upload_table_rejects_a_full_table() { + let Some(gpu) = gpu_or_skip("upload_table_full") else { + return; + }; + let full = [[1u32, 0], [2, 0], [3, 0], [4, 0]]; + let r = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| gpu.upload_table(&full))); + assert!(r.is_err(), "a full table must be rejected"); +} + /// [OPUS-4.8] sq-goay: randomized differential oracle. /// /// The tests above pin each kernel against the CPU reference on ONE fixed @@ -206,7 +253,7 @@ fn differential_sweep_all_kernels() { let b = rng.u32(); let (lo, hi) = if a <= b { (a, b) } else { (b, a) }; assert_eq!( - gpu.filter_count_u32(&resident, lo, hi), + gpu.filter_count_u32(&resident, lo, hi).unwrap(), cpu::filter_count_u32(&col, lo, hi), "filter_u32 n={n} lo={lo} hi={hi}" ); @@ -222,7 +269,7 @@ fn differential_sweep_all_kernels() { for _ in 0..4 { let t = (rng.u32() as f64 - (u32::MAX / 2) as f64) / 1e3; assert_eq!( - gpu.filter_count_f64_gt(&resident, t), + gpu.filter_count_f64_gt(&resident, t).unwrap(), cpu::filter_count_f64_gt(&col, t), "filter_f64 n={n} t={t}" ); @@ -250,7 +297,7 @@ fn differential_sweep_all_kernels() { let table = gpu.upload_table(&slots); let probe_col = gpu.upload_u32(&probe); assert_eq!( - gpu.hash_probe(&table, &probe_col), + gpu.hash_probe(&table, &probe_col).unwrap(), cpu::hash_probe(&slots, &probe), "hash_probe n={n}" ); @@ -265,7 +312,7 @@ fn differential_sweep_all_kernels() { let vals_col = gpu.upload_u32(&vals); assert_eq!( gpu.group_aggregate(&keys_col, &vals_col, groups), - cpu::group_aggregate(&keys, &vals, groups), + cpu::group_aggregate(&keys, &vals, groups).map_err(GpuError::from), "group_aggregate n={n} groups={groups}" ); } diff --git a/crates/sparq-mcp/src/nlq.rs b/crates/sparq-mcp/src/nlq.rs index cd858a93c4..9645e6b212 100644 --- a/crates/sparq-mcp/src/nlq.rs +++ b/crates/sparq-mcp/src/nlq.rs @@ -15,7 +15,10 @@ //! default config both tools build, because it false-positives on a legitimately-empty //! answer. So an ungrounded predicate/class IRI is accepted by BOTH — `ask` executes it to //! zero rows, `nl_query` hands it back — and that parity is pinned by -//! `both_tools_accept_an_ungrounded_predicate`. [OPUS-5] +//! `both_tools_accept_an_ungrounded_predicate`. Both tools derive their config from one +//! shared base, and [`run_nl_query_with`] honours `check_dictionary` exactly as +//! [`sparq_nlq::Nlq::ask`] does, so turning it on cannot make the tools diverge +//! (GitHub #4833). //! //! This is the second of the two complementary grounding tools chosen on the 2026-06-23 //! design call (`shapes` is the no-LLM structured one): instead of handing the client a @@ -146,7 +149,21 @@ pub fn nl_query(graph: &Graph, question: &str) -> Result { /// fails at runtime (or matches nothing) is returned as-is. That is the honest trade the /// tool advertises: cheaper and reviewable, but only *syntactically* validated. pub fn run_nl_query(graph: &Graph, question: &str, llm: Box) -> Result { - let config = NlqConfig::default(); + run_nl_query_with(graph, question, base_config(), llm) +} + +/// [`run_nl_query`] under an explicit [`NlqConfig`] — the config seam that keeps the two +/// tools from drifting (GitHub #4833). It applies the same pre-execution checks +/// [`sparq_nlq::Nlq::ask`] does under the same config, in the same order: the question +/// guard, the `spargebra` parse, the forbidden-construct refusal, and — when +/// `config.check_dictionary` is on — the dictionary-grounding repair +/// ([`sparq_nlq::constrain::unknown_terms`]). Only execution is skipped. +pub fn run_nl_query_with( + graph: &Graph, + question: &str, + config: NlqConfig, + llm: Box, +) -> Result { // `Nlq` owns its `Llm` and does not lend it out, so build the grounding and repair // prompts through a prompt-only instance (neither `prompt_for` nor `repair_prompt_for` // calls the model) and drive the caller's model directly. Constructing it runs the one @@ -170,11 +187,21 @@ pub fn run_nl_query(graph: &Graph, question: &str, llm: Box) -> Result< Err(e) => (Some(q), e.to_string()), Ok(parsed) => { let forbidden = sparq_nlq::guard::forbidden_constructs(&parsed, &config.guard); - if forbidden.is_empty() { - return Ok(render_translation(question, &q, round)); + if !forbidden.is_empty() { + let msg = sparq_nlq::guard::forbidden_repair_message(&forbidden); + (Some(q), msg) + } else { + let unknowns = if config.check_dictionary { + sparq_nlq::constrain::unknown_terms(graph, &parsed) + } else { + Vec::new() + }; + if unknowns.is_empty() { + return Ok(render_translation(question, &q, round)); + } + let msg = sparq_nlq::constrain::dictionary_repair_message(&unknowns); + (Some(q), msg) } - let msg = sparq_nlq::guard::forbidden_repair_message(&forbidden); - (Some(q), msg) } }, }; @@ -247,7 +274,7 @@ pub fn run_ask( fn config_from_budget(budget: &QueryBudget) -> NlqConfig { let mut config = NlqConfig { max_rows: budget.max_rows, - ..NlqConfig::default() + ..base_config() }; // Translate the budget's absolute deadline into a per-query duration (the loop builds // a fresh deadline at execution time). Use the remaining time; fall back to the @@ -263,6 +290,13 @@ fn config_from_budget(budget: &QueryBudget) -> NlqConfig { config } +/// The single [`NlqConfig`] both tools start from: `ask` layers the server budget on top +/// ([`config_from_budget`]) and `nl_query` uses it as-is, so a change to the pre-execution +/// checks (e.g. turning `check_dictionary` on) reaches both tools at once. GitHub #4833. +fn base_config() -> NlqConfig { + NlqConfig::default() +} + /// Render the answer as a structured JSON object: the executed SPARQL, the result rows /// (faithfully, never paraphrased into prose the model could distort), the repair count, /// and — with `nlq/citations` (the `sparq-nlq/citations` feature) — in-graph citations. @@ -551,6 +585,47 @@ ex:bob rdf:type ex:Person ; ex:name "Bob" . assert_eq!(v["repairs"], 0); } + /// GitHub #4833: with `check_dictionary` on, `nl_query` must apply the same + /// dictionary-grounding repair `Nlq::ask` does under the same config — an ungrounded + /// predicate is fed back as a repair signal and, when the model repairs it, the + /// grounded query is returned. Both tools must reject the never-repaired case. + #[test] + fn nl_query_honours_check_dictionary_like_ask() { + const UNGROUNDED: &str = "PREFIX ex: SELECT ?s WHERE { ?s ex:nosuch ?o }"; + let config = NlqConfig { + check_dictionary: true, + max_repair_rounds: 1, + ..base_config() + }; + let g = graph(); + + // A model that only ever emits the ungrounded query: both tools reject it. + let stuck = || Box::new(FnLlm(|_| Ok(format!("```sparql\n{}\n```", UNGROUNDED)))); + let ask_err = sparq_nlq::Nlq::with_config(&g, stuck(), config.clone()) + .ask("anything") + .expect_err("ask repairs then rejects an ungrounded predicate"); + assert!(ask_err.to_string().contains("nosuch"), "{ask_err}"); + let err = run_nl_query_with(&g, "anything", config.clone(), stuck()) + .expect_err("nl_query must reject what ask rejects under the same config"); + assert!(err.contains("nosuch"), "{err}"); + + // A model that repairs on the dictionary signal: the grounded query comes back. + let repairing = Box::new(FnLlm(|prompt| { + // Only the repair prompt names the ungrounded IRI. + let q = if prompt.contains("http://ex/nosuch") { + COUNT_QUERY + } else { + UNGROUNDED + }; + Ok(format!("```sparql\n{}\n```", q)) + })); + let out = run_nl_query_with(&g, "anything", config, repairing) + .expect("the repaired, grounded query is returned"); + let v: Value = serde_json::from_str(&out).unwrap(); + assert_eq!(v["sparql"], COUNT_QUERY); + assert_eq!(v["repairs"], 1); + } + #[test] fn run_nl_query_surfaces_an_untranslatable_question_as_a_tool_error() { let llm = Box::new(FnLlm(|_| Ok("sorry, I cannot help".to_string()))); diff --git a/crates/sparq-terse/src/transpile.rs b/crates/sparq-terse/src/transpile.rs index 528bbbc73e..35832d015d 100644 --- a/crates/sparq-terse/src/transpile.rs +++ b/crates/sparq-terse/src/transpile.rs @@ -305,7 +305,7 @@ fn k_spans(src: &str) -> Vec { b'K' => { // Standalone sigil: the preceding char (if any) must not be an identifier // char (so `?xK:type` or `fooK:type` are NOT the keyword token). - let prev_ok = i == 0 || !is_ident_char(bytes[i - 1]); + let prev_ok = !preceded_by_name(bytes, i); if prev_ok && i + 1 < n && bytes[i + 1] == b':' { if let Some((name, end)) = parse_k_name(bytes, i + 2) { out.push(KSpan { name, start: i, end }); @@ -371,7 +371,7 @@ fn declares_prefix_k(src: &str) -> bool { } _ => { // At a token boundary, try to match `PREFIX`/`@prefix` then `K` `:`. - let boundary = i == 0 || !is_ident_char(bytes[i - 1]); + let boundary = !preceded_by_name(bytes, i); if boundary { if let Some(after) = match_prefix_kw(bytes, i) { let k = skip_ws(bytes, after); @@ -417,7 +417,7 @@ fn match_word_ci(bytes: &[u8], start: usize, word: &[u8]) -> Option { } } // Must be followed by a non-identifier byte (or end of input). - if end < bytes.len() && is_ident_char(bytes[end]) { + if end < bytes.len() && is_name_byte(bytes[end]) { return None; } Some(end) @@ -498,7 +498,7 @@ fn scan_v_constructs( b'V' => { // Must be a standalone `V` token: the preceding char (if any) must not be // part of an identifier (so `?fooV(` or `abcV(` are NOT the construct). - let prev_ok = i == 0 || !is_ident_char(bytes[i - 1]); + let prev_ok = !preceded_by_name(bytes, i); if prev_ok { if let Some((phrase, end)) = parse_v_call(bytes, i) { if let Some(t) = f(&phrase, i, end) { @@ -522,6 +522,30 @@ fn is_ident_char(b: u8) -> bool { b.is_ascii_alphanumeric() || b == b'_' || b == b':' || b == b'?' || b == b'$' } +/// `true` if `b` can be part of a SPARQL name: an ASCII identifier byte, or any +/// non-ASCII byte. Outside strings, IRIs and comments (which the scanners skip), +/// a non-ASCII character can only be a `PN_CHARS` name character — SPARQL +/// whitespace and punctuation are all ASCII — so every UTF-8 lead/continuation +/// byte counts as a name byte. GitHub #4662. +fn is_name_byte(b: u8) -> bool { + is_ident_char(b) || !b.is_ascii() +} + +/// `true` if the token starting at `bytes[i]` is glued to a preceding name, so +/// it is NOT a standalone `K` / `V` / `PREFIX` token: the previous byte is a +/// name byte ([`is_name_byte`]), or the previous two bytes are a `PN_LOCAL_ESC` +/// (`\` + one of `_~.-!$&'()*+,;=/?#@%`), which continues a local name. +/// GitHub #4662. +fn preceded_by_name(bytes: &[u8], i: usize) -> bool { + if i == 0 { + return false; + } + if is_name_byte(bytes[i - 1]) { + return true; + } + i >= 2 && bytes[i - 2] == b'\\' && b"_~.-!$&'()*+,;=/?#@%".contains(&bytes[i - 1]) +} + /// At `bytes[start] == b'V'`, tries to parse `V("phrase")` / `V('phrase')` allowing /// whitespace around the parens and the string. On success returns `(phrase, end)` where /// `end` is the byte index just past the closing `)`. The phrase is the *unescaped* string @@ -653,6 +677,36 @@ mod tests { assert!(e.warnings.is_empty()); } + /// GitHub #4662: a `K` that ends a Unicode variable name (`?aéK`) is not the + /// keyword sigil — the byte before it is a UTF-8 continuation byte, which is + /// part of the name, not a delimiter. Valid SPARQL must pass through unchanged. + #[test] + fn k_after_unicode_name_char_is_not_a_sigil() { + let q = "PREFIX : SELECT ?s WHERE { ?aéK:x ?s }"; + let e = terse_to_sparql(q).expect("valid SPARQL passes through"); + assert_eq!(e.canonical_sparql, q); + assert!(e.keywords.is_empty()); + } + + /// GitHub #4662: a `K` after a `PN_LOCAL_ESC` (`\-`) continues the local + /// name (`ex:a\-K:x` is the local name `a-K:x`), so it is not the sigil. + #[test] + fn k_after_pn_local_esc_is_not_a_sigil() { + let q = "PREFIX ex: SELECT ?s WHERE { ?s ex:p ex:a\\-K:x }"; + let e = terse_to_sparql(q).expect("valid SPARQL passes through"); + assert_eq!(e.canonical_sparql, q); + assert!(e.keywords.is_empty()); + } + + /// GitHub #4662: the `V(` scanner shares the boundary predicate, so a `V` + /// glued to a Unicode name is not the construct either. + #[test] + fn v_after_unicode_name_char_is_not_a_construct() { + let q = "SELECT ?s WHERE { ?s ?aéV(\"x\") }"; + let err = terse_to_sparql(q).expect_err("not valid SPARQL"); + assert!(matches!(err, TerseError::CanaryFailed { .. }), "{err:?}"); + } + #[test] fn canary_rejects_invalid_sparql() { let bad = "SELECT ?s WHERE { ?s ?p"; // unbalanced — does not parse diff --git a/skills/agent-tools/SKILL.md b/skills/agent-tools/SKILL.md index d45be460df..c44a289248 100644 --- a/skills/agent-tools/SKILL.md +++ b/skills/agent-tools/SKILL.md @@ -251,7 +251,10 @@ same `nlq` feature: the query may still fail at runtime or match nothing. Note that sparq-nlq's dictionary-grounding constraint (`NlqConfig::check_dictionary`) is opt-in and **off** in the default config, for `ask` as well as `nl_query` — so an ungrounded predicate/class - IRI is accepted by **both** (`ask` just executes it to zero rows). Same backend, same + IRI is accepted by **both** (`ask` just executes it to zero rows). Both tools start + from one shared config, and the library entry point `run_nl_query_with(graph, question, + config, llm)` takes an explicit `NlqConfig`: with `check_dictionary` on it applies the + same dictionary repair `ask` does, so the two cannot drift. Same backend, same fail-closed "not configured" error. These are **ergonomics / grounding aids pending measurement** — *not* a token-saving diff --git a/skills/gpu-kernels/SKILL.md b/skills/gpu-kernels/SKILL.md index 4cd1897ee4..05812a473a 100644 --- a/skills/gpu-kernels/SKILL.md +++ b/skills/gpu-kernels/SKILL.md @@ -55,17 +55,20 @@ let Some(gpu) = Gpu::new() else { return; }; // no adapter → skip // FILTER + count: how many elements fall in [lo, hi]? let col = gpu.upload_u32(&[1u32, 5, 9, 3, 7]); -let n = gpu.filter_count_u32(&col, 4, 8); // -> 2 (5 and 7) +let n = gpu.filter_count_u32(&col, 4, 8)?; // -> 2 (5 and 7) // Hash-join probe against a resident open-addressing table. -let table = gpu.upload_table(&[[42, 100], [7, 200]]); // [key, payload] +// Build it on the host with `cpu::build_hash_table` (load factor <= 0.5): `upload_table` +// panics on a table that is not a power of two or has no EMPTY_KEY slot. +let slots = sparq_gpu::cpu::build_hash_table(&[42, 7], &[100, 200]); // keys, payloads +let table = gpu.upload_table(&slots); let probe = gpu.upload_u32(&[42u32, 7, 7, 1]); -let (matches, payload_sum) = gpu.hash_probe(&table, &probe); +let (matches, payload_sum) = gpu.hash_probe(&table, &probe)?; // GROUP BY COUNT+SUM (keys pre-densified to 0..groups, groups ≤ MAX_GROUPS = 512). let keys = gpu.upload_u32(&[0u32, 1, 0, 1, 1]); let vals = gpu.upload_u32(&[10u32, 20, 30, 40, 50]); -let per_group = gpu.group_aggregate(&keys, &vals, 2); // Vec<(count, sum)> +let per_group = gpu.group_aggregate(&keys, &vals, 2)?; // Vec<(count, sum)> let _ = (n, matches, payload_sum, per_group); ``` @@ -76,15 +79,24 @@ let _ = (n, matches, payload_sum, per_group); `upload_table` (build), and `write_u32` / `write_f64` / `write_table` (re-fill an existing buffer in place). Resident types: `ColU32`, `ColF64`, `HashTable`. - Kernels (all reduce **on-device** and read back O(1)/O(groups) bytes): - - `filter_count_u32(&col, lo, hi) -> u64` — count elements in `[lo, hi]`. - - `filter_count_f64_gt(&col, t) -> u64` — count `> t`. WGSL has **no f64**, so the + - `filter_count_u32(&col, lo, hi) -> Result` — count elements in `[lo, hi]`. + - `filter_count_f64_gt(&col, t) -> Result` — count `> t`. WGSL has **no f64**, so the f64 kernel compares IEEE-754 bit patterns mapped to a monotonic u64 key (sign-flip trick): **exact, NaN-correct, no float math on the device** for comparisons. - - `hash_probe(&table, &probe) -> (u64, u64)` — `(matches, payload_sum)` against a + - `hash_probe(&table, &probe) -> Result<(u64, u64), GpuError>` — `(matches, payload_sum)` against a resident linear-probing table (load ≤ 0.5; u64 sum via 32-bit atomic carry). - - `group_aggregate(&keys, &vals, groups) -> Vec<(u64, u64)>` — `(count, sum)` per - group with two-level (workgroup-shared → global) atomics. -- Constants / introspection: `EMPTY_KEY` (`u32::MAX`), `MAX_GROUPS` (`512`), + - `group_aggregate(&keys, &vals, groups) -> Result, GpuError>` + — `(count, sum)` per group with two-level (workgroup-shared → global) atomics. A key + `>= groups` is flagged on the device and returned as `GpuError::GroupKeyOutOfRange`; + the CPU oracle `cpu::group_aggregate` returns `GroupKeyOutOfRange` (it no longer + panics; `GpuError: From`). +- The probe walk is bounded at one pass over the table, and each submission is waited on + for at most `POLL_TIMEOUT` (60 s) rather than indefinitely. A timeout returns + `GpuError::Stalled` (not a panic) and marks the `Gpu` stalled: later kernel calls fail + fast with the same error, and dropping it **leaks** its wgpu device + queue on purpose, + because wgpu's queue destructor waits for idle with no timeout and would hang on the + stalled submission (including during a panic unwind). Recreate a `Gpu` to retry. +- Constants / introspection: `EMPTY_KEY` (`u32::MAX`), `MAX_GROUPS` (`512`), `POLL_TIMEOUT`, `Gpu::max_storage_bytes` (the `max_storage_buffer_binding_size` cap on resident column size).