diff --git a/README.md b/README.md index 86cd9a2..8c10af0 100644 --- a/README.md +++ b/README.md @@ -43,6 +43,8 @@ fn main() -> Result<(), String> { `Backend` is the lower-level compilation/execution interface. `Backend::install_composition(®istry, &manifest)` validates exact registry pins and typed bindings, then installs the composition and returns its `CompositionResolution` with generated input/output relation names. Use `ProgramInstance` for public-interface enforcement or registered operations. Existing MCP queries and deltas contain DDlog row text; library callers can use `Backend::query_typed` and `export_inputs` for typed rows. `why` returns direct rule-variable witnesses, not recursive proof trees or confidence/provenance. +For large maintained outputs, use `Backend::query_typed_bounded` with exact positional filters, row/JSON-byte limits and an explicit continuation. It streams only the selected relation and drains its native dump with bounded host accumulation. `apply_without_deltas` acknowledges mutations without transporting unrelated derived changes. See [bounded output reads](docs/bounded-reads.md) for limits and scan costs. + Library callers can explicitly save and restore pure-program local checkpoints. See [checkpoint contracts](docs/checkpoints.md). Compatible output-schema changes recompile and replay retained inputs; input schemas stay fixed when data is retained. A runnable library example is [`examples/program.rs`](examples/program.rs). Installation invokes native compilation; pure parsing, lowering and registry validation do not. The workspace does not require the MCP feature for library use. diff --git a/docs/bounded-reads.md b/docs/bounded-reads.md new file mode 100644 index 0000000..1f1ba14 --- /dev/null +++ b/docs/bounded-reads.md @@ -0,0 +1,105 @@ +# Bounded output reads + +The lower-level Rust library supports pages from one maintained native output: + +```rust +use ddlog_runtime::{Backend, BoundedQuery}; +use serde_json::json; +use std::collections::BTreeMap; + +fn read(backend: &mut Backend) -> Result<(), String> { + let mut query = BoundedQuery { + filters: BTreeMap::from([(0, json!("selected-key"))]), + max_rows: 20, + max_bytes: 64 * 1024, + continuation: None, + }; + loop { + let page = backend.query_typed_bounded("visible", &query)?; + // Consume this page before requesting another; retaining every page + // would move an unbounded snapshot into the application. + println!("{} rows, {} JSON bytes", page.rows.len(), page.bytes); + match page.continuation { + Some(cursor) => query.continuation = Some(cursor), + None => break, + } + } + Ok(()) +} +``` + +`filters` maps zero-based field positions to exact typed values. Every filter +must match. String comparisons are case-sensitive and integers must fit signed +64 bits; wrong types or positions fail before issuing a native command. +Only declared output relations are accepted. Composition callers translate +exported ports through the installed `CompositionResolution`. The backend +does not enforce `ProgramInstance` public interfaces or registered-operation +admission; these remain the caller's responsibility at that lower boundary. +This addition does not alter MCP tools, shared-owner messages or native drivers. + +`max_rows` must be between 1 and `MAX_QUERY_ROWS` (10,000). `max_bytes` must be +between 2 and `MAX_QUERY_BYTES` (4 MiB). `QueryPage.bytes` counts the exact +serialized JSON of `rows`, including the outer brackets, row brackets, commas +and escaped string bytes. It does not count page metadata or the cursor. An +empty complete result is two bytes. Defaults are 100 rows and 64 KiB. + +`truncated` is true only when an additional matching row was observed; it always +comes with a continuation that advances past the returned rows. If the first +matching row cannot fit the requested byte limit, the request returns an error +with the required JSON page size. It never returns an empty, non-advancing +continuation in this case. Increase the limit if the row fits the hard cap. +An individual native record, including its newline, must fit +`MAX_NATIVE_RECORD_BYTES` (4 MiB). An over-limit or malformed record encountered +while selecting the page returns an error; it is never silently skipped to +present the page as complete. The transport drains the rest of the response +before returning a record/size error, preserving the healthy live owner. +Loss of framing or native I/O still disables the owner. + +Cursors are opaque, serializable and bound to one live native owner, revision, +predicate and exact filter set. Page limits may change between requests. +Acknowledged mutations, program replacement, a new owner and checkpoint +restore invalidate prior cursors. Reads do not advance the revision. Cursors +are transport positions, not durable checkpoints or authorization tokens. +Native relation order is stable within an unchanged live revision; callers +must not treat it as a product ranking or rely on the same order after restore. + +The pinned DDlog CLI supports `dump R_selected;` but has no cursor/limit command. +The runtime therefore scans and drains the entire selected relation for each +page, retaining only a capped native record and the bounded result page. It +decodes and compares rows until it finds the page and evidence of another +match, then drains the remaining bytes without accumulating them. It does not +dump unrelated relations, reconstruct a derived cache, evaluate rules in the +host, or call `query_typed` and truncate its complete vector. Limits bound +returned payload and host accumulation, **not** native output work, total bytes +read, or scan latency. Pagination rescans the selected relation. The native +CLI can block, and no operation timeout is introduced here. + +For writers that need selected output reads rather than a complete delta dump, +`apply_without_deltas(&changes)` returns the acknowledged `version` and +`revision` using a plain native commit. Validation, set semantics, failure +handling, acknowledged input ownership and checkpoint behavior match `apply`. +Lost acknowledgement leaves inputs/revision unadvanced and disables the owner; +it does not authorize mutation replay. Candidate installation and checkpoint +restore already replay inputs with plain commits. The original `apply`, +`query`, `query_typed`, and MCP operations preserve their existing 4 MiB +aggregate transport and delta contracts. + +The controlled fixture regressions in `tests/memory_runtime.rs` exercise +bounded pages and filtering, byte accounting, cursor binding, malformed and +over-limit native records, lost commit acknowledgement, and large selected and +unrelated outputs across replacement/checkpoint restore. The fixture is a +transport simulator and is not evidence of native rule evaluation. + +Native acceptance is separately opt-in with an operator-configured driver: + +```sh +DDLOG_RUNTIME_NATIVE_BUILD=/absolute/native-build-driver \ + cargo test --locked --test bounded_native -- --ignored --nocapture +``` + +It installs a real graph with 100 selected 64 KiB rows and 100 unrelated 64 KiB +rows, verifies every selected key across pages, rejects a native record above +4 MiB without losing the owner, then checks mutation, replay and checkpoint +reopen. Build logs and the checkpoint remain under the printed temporary +directory for inspection. The driver controls whether compilation is fresh or +an exact native artifact is reused; preserve that evidence separately. diff --git a/docs/checkpoints.md b/docs/checkpoints.md index b888c5a..9c425c9 100644 --- a/docs/checkpoints.md +++ b/docs/checkpoints.md @@ -2,9 +2,16 @@ The runtime library owns acknowledged input state. `Backend::export_inputs()` returns every input relation, including empty ones, as typed rows. -`query_typed(predicate)` validates and decodes native output records. Existing -string query APIs remain available. Native responses remain bounded to 4 MiB; -typed reads do not introduce pagination or push down caller-side filters. +`query_typed(predicate)` validates and decodes complete native output records. +Its existing string transport has a 4 MiB aggregate response limit; exceeding +that limit disables the owner because the command is no longer synchronized. +Use `query_typed_bounded(predicate, &BoundedQuery)` for bounded retrieval from +larger maintained outputs. See [bounded output reads](bounded-reads.md). +`apply_without_deltas(changes)` commits the same input transaction without +transporting unrelated output deltas. Candidate installation and checkpoint +replay also use plain commits, so their acknowledgements do not grow with the +derived output snapshot. Existing string queries and delta-returning mutation +APIs remain available with their original contracts. `Backend::install_composition(&mut self, registry: &ProcessorRegistry, manifest: &CompositionManifest) -> Result` resolves diff --git a/src/bounded.rs b/src/bounded.rs new file mode 100644 index 0000000..0081126 --- /dev/null +++ b/src/bounded.rs @@ -0,0 +1,214 @@ +//! Bounded transport reads from a single maintained native output relation. +//! The native CLI has no limit/cursor command: every page drains that relation's +//! dump. Selection here only compares decoded columns, never evaluates rules. +use crate::{rows, Backend, Result}; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use std::collections::BTreeMap; +use std::io::Write; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::{SystemTime, UNIX_EPOCH}; + +pub const MAX_QUERY_ROWS: usize = 10_000; +pub const MAX_QUERY_BYTES: usize = 4 * 1024 * 1024; +/// Maximum native record size, including its trailing newline. +pub const MAX_NATIVE_RECORD_BYTES: usize = 4 * 1024 * 1024; +static OWNER_SEQUENCE: AtomicU64 = AtomicU64::new(0); + +pub(crate) fn owner_identity() -> String { + format!( + "{}-{}-{}", + std::process::id(), + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_nanos(), + OWNER_SEQUENCE.fetch_add(1, Ordering::Relaxed) + ) +} + +/// A page request with exact, zero-based positional equality filters. +/// All filters must match. Limits apply after filtering, and `max_bytes` counts +/// the exact JSON encoding of the returned `rows`, including outer brackets. +/// Neither limit bounds native scan/drain time; the CLI dumps the full selected +/// relation on each request. Unrelated relations are never requested. +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct BoundedQuery { + pub filters: BTreeMap, + pub max_rows: usize, + pub max_bytes: usize, + pub continuation: Option, +} + +impl Default for BoundedQuery { + fn default() -> Self { + Self { + filters: BTreeMap::new(), + max_rows: 100, + max_bytes: 64 * 1024, + continuation: None, + } + } +} + +/// Continuation in native relation order, bound to one live owner, revision, +/// predicate and filter set. It is invalid after a mutation, replacement or +/// checkpoint restore. Its encoding is opaque to callers, not an authorization +/// token. Page limits may be changed while continuing the same query. +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct QueryCursor { + owner: String, + revision: u64, + predicate: String, + filters: BTreeMap, + offset: u64, +} + +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct QueryPage { + pub rows: Vec>, + pub bytes: usize, + /// True only after observing an additional matching row. + pub truncated: bool, + pub continuation: Option, + pub revision: u64, +} + +struct CountBytes(usize); +impl Write for CountBytes { + fn write(&mut self, bytes: &[u8]) -> std::io::Result { + self.0 += bytes.len(); + Ok(bytes.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } +} + +impl Backend { + /// Read a bounded page from one maintained output without accumulating its + /// full snapshot. A row too large for an empty page, malformed native record, + /// or over-limit native record returns an error after draining the response; + /// the healthy owner and acknowledged inputs remain usable. An I/O/framing + /// failure still disables the owner, as with other native commands. + pub fn query_typed_bounded( + &mut self, + predicate: &str, + query: &BoundedQuery, + ) -> Result { + if query.max_rows == 0 || query.max_rows > MAX_QUERY_ROWS { + return Err(format!("max_rows must be in 1..={MAX_QUERY_ROWS}")); + } + if query.max_bytes < 2 || query.max_bytes > MAX_QUERY_BYTES { + return Err(format!("max_bytes must be in 2..={MAX_QUERY_BYTES}")); + } + let schema = self.schema.get(predicate).ok_or("Unknown relation")?; + if schema.input { + return Err("Bounded queries accept output relations only".into()); + } + for (position, value) in &query.filters { + let field = schema + .fields + .get(*position) + .ok_or_else(|| format!("Filter position {position} exceeds relation arity"))?; + match field.as_str() { + "int" if value.as_i64().is_some() => (), + "string" if value.as_str().is_some() => (), + _ => return Err(format!("Filter position {position} requires {field}")), + } + } + let runtime = self.runtime.as_mut().ok_or("Install a program first")?; + let offset = if let Some(cursor) = &query.continuation { + if cursor.owner != runtime.identity + || cursor.revision != self.revision + || cursor.predicate != predicate + || cursor.filters != query.filters + { + return Err("Continuation does not match the live owner, revision or query".into()); + } + cursor.offset + } else { + 0 + }; + let mut page = QueryPage { + rows: Vec::new(), + bytes: 2, + truncated: false, + continuation: None, + revision: self.revision, + }; + let mut matched = 0_u64; + let mut error = None; + let exchange = runtime.exchange_stream(&format!("dump R_{predicate};"), |line| { + // The native stream must still be drained once the page is full or + // its consumer rejects a record. No derived rows are cached. + if error.is_some() || page.truncated { + return; + } + let mut accept = |line: Result<&str>| -> Result<()> { + let mut decoded = rows::decode_rows(line?, predicate, &schema.fields)?; + let Some(row) = decoded.pop() else { + return Ok(()); + }; + if !query + .filters + .iter() + .all(|(position, value)| row[*position] == *value) + { + return Ok(()); + } + matched = matched.checked_add(1).ok_or("Row offset exhausted")?; + if matched <= offset { + return Ok(()); + } + let mut count = CountBytes(0); + serde_json::to_writer(&mut count, &row).map_err(|e| e.to_string())?; + if page.rows.is_empty() && count.0 + 2 > query.max_bytes { + return Err(format!( + "Matching row requires {} JSON page bytes, exceeding max_bytes={}; increase the limit", + count.0 + 2, query.max_bytes + )); + } + let bytes = page.bytes + usize::from(!page.rows.is_empty()) + count.0; + if page.rows.len() == query.max_rows || bytes > query.max_bytes { + page.truncated = true; + return Ok(()); + } + page.bytes = bytes; + page.rows.push(row); + Ok(()) + }; + if let Err(reason) = accept(line) { + error = Some(reason); + } + }); + if let Err(error) = exchange { + self.runtime = None; + self.failed = true; + return Err(format!( + "Runtime unavailable; reconcile outstanding work: {error}" + )); + } + if let Some(error) = error { + return Err(error); + } + if matched < offset { + return Err("Continuation offset exceeds the selected relation".into()); + } + if page.truncated { + page.continuation = Some(QueryCursor { + owner: runtime.identity.clone(), + revision: self.revision, + predicate: predicate.to_string(), + filters: query.filters.clone(), + offset: offset + .checked_add(page.rows.len() as u64) + .ok_or("Row offset exhausted")?, + }); + } + Ok(page) + } +} diff --git a/src/lib.rs b/src/lib.rs index 5a8aad4..8cf1612 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,7 +6,11 @@ pub use lemmalog_syntax as syntax; pub mod instance; pub use instance::ProgramInstance; +mod bounded; mod checkpoint; +pub use bounded::{ + BoundedQuery, QueryCursor, QueryPage, MAX_NATIVE_RECORD_BYTES, MAX_QUERY_BYTES, MAX_QUERY_ROWS, +}; pub mod composition; #[cfg(all(feature = "mcp", unix))] pub mod host; @@ -34,6 +38,7 @@ struct Runtime { input: ChildStdin, output: BufReader, sequence: u64, + identity: String, group: processes::Group, } impl Runtime { @@ -51,6 +56,7 @@ impl Runtime { output: BufReader::new(child.stdout.take().unwrap()), child, sequence: 0, + identity: bounded::owner_identity(), group, }) } @@ -82,6 +88,59 @@ impl Runtime { result.push_str(&line); } } + + /// Read one record at a time and always drain the response to its marker. + /// Consumer/record errors do not desynchronize the live owner. I/O errors + /// still fail the exchange because the command boundary is then uncertain. + fn exchange_stream( + &mut self, + commands: &str, + mut consume: impl FnMut(Result<&str>), + ) -> Result<()> { + self.sequence += 1; + let marker = format!("LEMMALOG_END_{}", self.sequence); + writeln!(self.input, "{commands}\necho {marker};") + .and_then(|_| self.input.flush()) + .map_err(|e| e.to_string())?; + let mut line = Vec::new(); + loop { + line.clear(); + let read = self + .output + .by_ref() + .take((MAX_NATIVE_RECORD_BYTES + 1) as u64) + .read_until(b'\n', &mut line) + .map_err(|e| e.to_string())?; + if read == 0 { + return Err("DDlog exited before completing the request".into()); + } + if line.len() > MAX_NATIVE_RECORD_BYTES { + if !line.ends_with(b"\n") { + loop { + let bytes = self.output.fill_buf().map_err(|e| e.to_string())?; + if bytes.is_empty() { + return Err("DDlog exited before completing the request".into()); + } + let end = bytes.iter().position(|byte| *byte == b'\n'); + let length = end.map_or(bytes.len(), |end| end + 1); + self.output.consume(length); + if end.is_some() { + break; + } + } + } + consume(Err(format!( + "DDlog record exceeded {MAX_NATIVE_RECORD_BYTES} bytes; bounded query rejected" + ))); + continue; + } + match std::str::from_utf8(&line) { + Ok(text) if text.trim() == marker => return Ok(()), + Ok(text) => consume(Ok(text)), + Err(error) => consume(Err(format!("Invalid UTF-8 in DDlog record: {error}"))), + } + } + } } impl Drop for Runtime { fn drop(&mut self) { @@ -262,6 +321,15 @@ impl Backend { Ok(format!("R_{predicate}({})", rendered.join(", "))) } pub fn apply(&mut self, changes: &Value) -> Result { + self.apply_inner(changes, true) + } + /// Commit inputs without transporting output deltas. Input acknowledgement, + /// revision advancement and checkpoint semantics are identical to `apply`. + /// Use this when the caller will read selected outputs separately. + pub fn apply_without_deltas(&mut self, changes: &Value) -> Result { + self.apply_inner(changes, false) + } + fn apply_inner(&mut self, changes: &Value, dump_deltas: bool) -> Result { if self.runtime.is_none() { return Err("Install a program first".into()); } @@ -288,12 +356,32 @@ impl Backend { for (_, fact) in staged.keys().filter(|key| !self.facts.contains_key(*key)) { commands.push_str(&format!("insert {fact};\n")); } - commands.push_str("commit dump_changes;"); - match self.runtime.as_mut().unwrap().exchange(&commands) { + commands.push_str(if dump_deltas { + "commit dump_changes;" + } else { + "commit;" + }); + let response = self + .runtime + .as_mut() + .unwrap() + .exchange(&commands) + .and_then(|response| { + if !dump_deltas && !response.trim().is_empty() { + Err(format!("DDlog rejected transaction: {response}")) + } else { + Ok(response) + } + }); + match response { Ok(delta) => { self.facts = staged; self.revision = revision; - Ok(json!({"version":self.version,"deltas":delta})) + Ok(if dump_deltas { + json!({"version":self.version,"deltas":delta}) + } else { + json!({"version":self.version,"revision":self.revision}) + }) } Err(error) => { self.runtime = None; diff --git a/tests/bounded_native.rs b/tests/bounded_native.rs new file mode 100644 index 0000000..6ff00ff --- /dev/null +++ b/tests/bounded_native.rs @@ -0,0 +1,129 @@ +#![cfg(unix)] +//! Explicit native acceptance, never run with the transport simulator. +use ddlog_runtime::{Backend, BoundedQuery, MAX_NATIVE_RECORD_BYTES, MAX_QUERY_BYTES}; +use serde_json::json; +use std::collections::BTreeMap; +use std::path::PathBuf; + +#[test] +#[ignore = "requires DDLOG_RUNTIME_NATIVE_BUILD pointing to an operator-configured native driver"] +fn native_large_outputs_support_bounded_reads_plain_commits_and_reopen() { + let driver = PathBuf::from(std::env::var("DDLOG_RUNTIME_NATIVE_BUILD").unwrap()); + let root = std::env::temp_dir().join(format!("ddlog-bounded-native-{}", std::process::id())); + std::fs::create_dir(&root).unwrap(); + let rules = "echo(N,S) :- source(N,S). unrelated(S) :- unused(S)."; + let schemas = json!({ + "source":{"input":true,"fields":["int","string"]}, + "unused":{"input":true,"fields":["string"]}, + "echo":{"input":false,"fields":["int","string"]}, + "unrelated":{"input":false,"fields":["string"]} + }); + let mut live = Backend::new(root.join("live"), driver.clone()); + live.install(rules, schemas.clone()).unwrap(); + let payload = "x".repeat(65536); + let changes: Vec<_> = (0..100) + .flat_map(|index| { + [ + json!({"op":"insert","predicate":"source","values":[index,payload]}), + json!({"op":"insert","predicate":"unused","values":[format!("{index}:{payload}")]}), + ] + }) + .collect(); + live.apply_without_deltas(&json!(changes)).unwrap(); + let selected = BoundedQuery { + filters: BTreeMap::from([(0, json!(50)), (1, json!(payload))]), + max_rows: 1, + max_bytes: 128 * 1024, + continuation: None, + }; + let exact = live.query_typed_bounded("echo", &selected).unwrap(); + assert_eq!(exact.rows, vec![vec![json!(50), json!(payload)]]); + assert!(!exact.truncated); + let mut query = BoundedQuery { + max_rows: 31, + max_bytes: 2 * 1024 * 1024, + ..Default::default() + }; + let mut keys = Vec::new(); + loop { + let page = live.query_typed_bounded("echo", &query).unwrap(); + assert!(page.rows.len() <= query.max_rows); + assert!(page.bytes <= query.max_bytes); + assert_eq!(page.bytes, serde_json::to_vec(&page.rows).unwrap().len()); + keys.extend(page.rows.iter().map(|row| row[0].as_i64().unwrap())); + assert_eq!(page.truncated, page.continuation.is_some()); + query.continuation = page.continuation; + if query.continuation.is_none() { + break; + } + } + assert_eq!(keys.len(), 100); + let unique = keys + .iter() + .copied() + .collect::>(); + assert_eq!( + unique.len(), + keys.len(), + "continuations must not repeat rows" + ); + keys.sort_unstable(); + assert_eq!(keys, (0..100).collect::>()); + // Even one record beyond the transport cap is drained without losing the + // owner; the following mutation must still reach the same native graph. + let enormous = "z".repeat(MAX_NATIVE_RECORD_BYTES + 1); + live.apply_without_deltas(&json!([ + {"op":"insert","predicate":"source","values":[-1,enormous]} + ])) + .unwrap(); + let error = live + .query_typed_bounded( + "echo", + &BoundedQuery { + max_bytes: MAX_QUERY_BYTES, + ..Default::default() + }, + ) + .unwrap_err(); + assert!(error.contains("record exceeded"), "{error}"); + assert_eq!(live.health(), "ready"); + live.apply_without_deltas(&json!([ + {"op":"delete","predicate":"source","values":[-1,enormous]}, + {"op":"insert","predicate":"source","values":[-2,"a\n\"b\"\\é🚀"]} + ])) + .unwrap(); + let escaped = BoundedQuery { + filters: BTreeMap::from([(1, json!("a\n\"b\"\\é🚀"))]), + ..Default::default() + }; + assert_eq!( + live.query_typed_bounded("echo", &escaped).unwrap().rows, + vec![vec![json!(-2), json!("a\n\"b\"\\é🚀")]] + ); + live.install(rules, schemas).unwrap(); + let checkpoint = root.join("checkpoint.json"); + live.save_checkpoint(&checkpoint, json!({"native":"bounded"})) + .unwrap(); + let revision = live.revision(); + drop(live); + let mut reopened = Backend::new(root.join("reopened"), driver); + assert_eq!( + reopened.restore_checkpoint(&checkpoint).unwrap(), + json!({"native":"bounded"}) + ); + assert_eq!(reopened.revision(), revision); + assert_eq!( + reopened + .query_typed_bounded("echo", &selected) + .unwrap() + .rows, + exact.rows + ); + assert_eq!( + reopened.query_typed_bounded("echo", &escaped).unwrap().rows, + vec![vec![json!(-2), json!("a\n\"b\"\\é🚀")]] + ); + assert_eq!(reopened.health(), "ready"); + println!("Native bounded acceptance artifacts: {}", root.display()); + println!("Verified 100 x 64KiB selected rows, 100 x 64KiB unrelated rows, exact filters, paginated completeness, oversized record recovery, candidate replay, checkpoint reopen"); +} diff --git a/tests/fixtures/memory_fake_runtime.py b/tests/fixtures/memory_fake_runtime.py index e2696ce..b38567f 100644 --- a/tests/fixtures/memory_fake_runtime.py +++ b/tests/fixtures/memory_fake_runtime.py @@ -13,6 +13,9 @@ facts, staged = {}, {} for raw in sys.stdin: command = raw.strip() + if command.startswith(("commit", "dump")): + with (control / "commands").open("a") as log: + log.write(command + "\n") if command == "start;": staged = {name: set(rows) for name, rows in facts.items()} elif command.startswith(("insert R_", "delete R_")): @@ -27,12 +30,17 @@ os._exit(42) if (control / "fail_replay").exists() and command == "commit;": print("error: simulated replay rejection", flush=True) + if (control / "large_deltas").exists() and command == "commit dump_changes;": + for i in range(100): + print('R_unrelated{.f0 = "' + "x" * 65536 + '"}: +1', flush=True) elif command.startswith("dump R_"): name = command[len("dump R_"):-1] if (control / "malformed_query").exists(): print("unexpected native output", flush=True) + elif (control / "oversized_query").exists(): + print('R_' + name + '{.f0 = 1, .f1 = "' + "x" * (4 * 1024 * 1024) + '"}', flush=True) else: - for row in sorted(facts.get("source", set())): + for row in sorted(facts.get("unused" if name == "unrelated" else "source", set())): values = json.loads(row)[:arity[name]] fields = ", ".join(f".f{i} = {json.dumps(v, ensure_ascii=False)}" for i, v in enumerate(values)) print("R_" + name + "{" + fields + "}", flush=True) diff --git a/tests/memory_runtime.rs b/tests/memory_runtime.rs index 8254c20..e511a36 100644 --- a/tests/memory_runtime.rs +++ b/tests/memory_runtime.rs @@ -2,9 +2,10 @@ //! Controlled transport/compiler fixtures, never claimed as native DDlog proof. use ddlog_runtime::composition::CompositionManifest; use ddlog_runtime::registry::{ProcessorDefinition, ProcessorRegistry, ProcessorVersion}; -use ddlog_runtime::Backend; +use ddlog_runtime::{Backend, BoundedQuery, MAX_QUERY_BYTES, MAX_QUERY_ROWS}; use serde_json::{json, Value}; use sha2::{Digest, Sha256}; +use std::collections::BTreeMap; use std::fs; use std::os::unix::fs::PermissionsExt; use std::path::{Path, PathBuf}; @@ -471,3 +472,269 @@ fn failed_publication_does_not_replace_target_or_leave_partial_file() { .to_string_lossy() .ends_with(".tmp"))); } + +#[test] +fn bounded_pages_filter_exact_positions_account_bytes_and_continue_without_gaps() { + let fixture = Fixture::new(); + let mut live = initial(&fixture); + live.apply_without_deltas(&json!([ + {"op":"insert","predicate":"source","values":[1,"a\n\"b\"\\é"]}, + {"op":"insert","predicate":"source","values":[2,"other"]}, + {"op":"insert","predicate":"source","values":[3,"a\n\"b\"\\é"]} + ])) + .unwrap(); + let revision = live.revision(); + let expected = live + .query_typed("echo") + .unwrap() + .into_iter() + .filter(|row| row[1] == json!("a\n\"b\"\\é")) + .collect::>(); + let mut query = BoundedQuery { + filters: BTreeMap::from([(1, json!("a\n\"b\"\\é"))]), + max_rows: 1, + ..Default::default() + }; + let mut actual = Vec::new(); + for index in 0..expected.len() { + let page = live.query_typed_bounded("echo", &query).unwrap(); + assert_eq!(page.rows, vec![expected[index].clone()]); + assert_eq!(page.bytes, serde_json::to_vec(&page.rows).unwrap().len()); + assert_eq!(page.revision, revision); + assert_eq!(page.truncated, index + 1 < expected.len()); + assert_eq!(page.truncated, page.continuation.is_some()); + actual.extend(page.rows); + // Cursors survive serialization, but remain opaque to their caller. + query.continuation = serde_json::from_value(json!(page.continuation)).unwrap(); + } + assert_eq!(actual, expected); + assert_eq!(live.revision(), revision); + query.continuation = None; + query.max_rows = 10; + query.max_bytes = serde_json::to_vec(&vec![expected[0].clone()]) + .unwrap() + .len(); + let first = live.query_typed_bounded("echo", &query).unwrap(); + assert_eq!(first.rows.len(), 1); + assert_eq!(first.bytes, query.max_bytes); + assert!(first.truncated); + query.continuation = first.continuation; + query.max_bytes = MAX_QUERY_BYTES; + let rest = live.query_typed_bounded("echo", &query).unwrap(); + assert_eq!(rest.rows, expected[1..]); + assert!(!rest.truncated); + query.continuation = None; + query.filters = BTreeMap::from([(0, json!(2)), (1, json!("other"))]); + assert_eq!( + live.query_typed_bounded("echo", &query).unwrap().rows, + vec![vec![json!(2), json!("other")]] + ); + query.filters.insert(1, json!("Other")); + let empty = live.query_typed_bounded("echo", &query).unwrap(); + assert!(empty.rows.is_empty()); + assert_eq!(empty.bytes, 2); + assert!(!empty.truncated); +} + +#[test] +fn bounded_read_errors_drain_without_poisoning_owner_and_cursors_bind_snapshot() { + let fixture = Fixture::new(); + let mut live = initial(&fixture); + live.apply_without_deltas(&json!([ + {"op":"insert","predicate":"source","values":[9,"other"]} + ])) + .unwrap(); + let query = BoundedQuery { + max_rows: 1, + ..Default::default() + }; + let cursor = live + .query_typed_bounded("echo", &query) + .unwrap() + .continuation; + assert!(cursor.is_some()); + let mut continued = query.clone(); + continued.continuation = cursor; + let mut changed = continued.clone(); + changed.filters.insert(0, json!(9)); + assert!(live + .query_typed_bounded("echo", &changed) + .unwrap_err() + .contains("Continuation")); + let mut small = query.clone(); + small.max_bytes = 2; + assert!(live + .query_typed_bounded("echo", &small) + .unwrap_err() + .contains("Matching row requires")); + for flag in ["malformed_query", "oversized_query"] { + fixture.flag(flag); + let error = live.query_typed_bounded("echo", &query).unwrap_err(); + if flag == "oversized_query" { + assert!(error.contains("record exceeded"), "{error}"); + } + fixture.unflag(flag); + assert_eq!(live.health(), "ready"); + assert_eq!( + live.query_typed_bounded("echo", &continued).unwrap().rows, + vec![vec![json!(9), json!("other")]] + ); + } + let path = fixture.root.join("bounded-checkpoint.json"); + live.save_checkpoint(&path, Value::Null).unwrap(); + let mut restored = fixture.backend("other-owner"); + restored.restore_checkpoint(&path).unwrap(); + assert_eq!(restored.revision(), live.revision()); + assert!(restored + .query_typed_bounded("echo", &continued) + .unwrap_err() + .contains("Continuation")); + live.apply_without_deltas(&json!([])).unwrap(); + assert!(live + .query_typed_bounded("echo", &continued) + .unwrap_err() + .contains("Continuation")); + assert_eq!(live.health(), "ready"); +} + +#[test] +fn invalid_bounded_requests_fail_before_native_commands() { + let fixture = Fixture::new(); + let mut live = initial(&fixture); + let commands = fs::read_to_string(fixture.root.join("commands")).unwrap(); + for invalid in [ + BoundedQuery { + max_rows: 0, + ..Default::default() + }, + BoundedQuery { + max_rows: MAX_QUERY_ROWS + 1, + ..Default::default() + }, + BoundedQuery { + max_bytes: 1, + ..Default::default() + }, + BoundedQuery { + max_bytes: MAX_QUERY_BYTES + 1, + ..Default::default() + }, + BoundedQuery { + filters: BTreeMap::from([(2, json!(1))]), + ..Default::default() + }, + BoundedQuery { + filters: BTreeMap::from([(0, json!("1"))]), + ..Default::default() + }, + BoundedQuery { + filters: BTreeMap::from([(0, json!(1.0))]), + ..Default::default() + }, + BoundedQuery { + filters: BTreeMap::from([(1, Value::Null)]), + ..Default::default() + }, + ] { + assert!(live.query_typed_bounded("echo", &invalid).is_err()); + } + for predicate in ["source", "unknown", "echo; dump"] { + assert!(live + .query_typed_bounded(predicate, &BoundedQuery::default()) + .is_err()); + } + assert_eq!( + fs::read_to_string(fixture.root.join("commands")).unwrap(), + commands + ); + assert_eq!(live.health(), "ready"); +} + +#[test] +fn large_current_and_unrelated_outputs_allow_bounded_reads_replacement_and_reopen() { + let fixture = Fixture::new(); + let mut live = fixture.backend("large"); + let mut schema = schemas(); + schema["unrelated"] = json!({"input":false,"fields":["string"]}); + let rules = "echo(N,S) :- source(N,S). unrelated(S) :- unused(S)."; + live.install(rules, schema.clone()).unwrap(); + fixture.flag("large_deltas"); + fs::write(fixture.root.join("commands"), "").unwrap(); + let payload = "x".repeat(65536); + let changes: Vec<_> = (0..100) + .flat_map(|index| { + [ + json!({"op":"insert","predicate":"source","values":[index,payload]}), + json!({"op":"insert","predicate":"unused","values":[format!("{index}:{payload}")]}), + ] + }) + .collect(); + // Each relation's native snapshot exceeds the old aggregate response cap. + assert!(100 * payload.len() > MAX_QUERY_BYTES); + let receipt = live.apply_without_deltas(&json!(changes)).unwrap(); + assert_eq!(receipt["revision"], live.revision()); + let selected = BoundedQuery { + filters: BTreeMap::from([(0, json!(50))]), + max_rows: 1, + max_bytes: 128 * 1024, + continuation: None, + }; + assert_eq!( + live.query_typed_bounded("echo", &selected).unwrap().rows, + vec![vec![json!(50), json!(payload)]] + ); + let first = live + .query_typed_bounded( + "echo", + &BoundedQuery { + max_rows: 1, + max_bytes: 128 * 1024, + ..Default::default() + }, + ) + .unwrap(); + assert!(first.truncated); + assert_eq!(live.health(), "ready"); + // Candidate replay and checkpoint restore also use plain commit, so neither + // transports all derived deltas while acknowledging large retained inputs. + live.install(rules, schema).unwrap(); + let path = fixture.root.join("large-checkpoint.json"); + live.save_checkpoint(&path, json!({"large":true})).unwrap(); + let mut reopened = fixture.backend("large-reopened"); + assert_eq!( + reopened.restore_checkpoint(&path).unwrap(), + json!({"large":true}) + ); + assert_eq!(reopened.revision(), live.revision()); + let page = reopened.query_typed_bounded("echo", &selected).unwrap(); + assert_eq!(page.rows, vec![vec![json!(50), json!(payload)]]); + assert!(!page.truncated); + assert_eq!(reopened.health(), "ready"); + let commands = fs::read_to_string(fixture.root.join("commands")).unwrap(); + assert!(!commands.contains("dump_changes")); + assert!( + commands + .lines() + .all(|line| line == "commit;" || line == "dump R_echo;"), + "{commands}" + ); +} + +#[test] +fn lost_plain_commit_ack_does_not_advance_inputs_or_allow_checkpoint() { + let fixture = Fixture::new(); + let mut live = initial(&fixture); + let revision = live.revision(); + fixture.flag("die_on_commit"); + assert!(live + .apply_without_deltas(&json!([ + {"op":"insert","predicate":"source","values":[1,"uncertain"]} + ])) + .is_err()); + assert_eq!(live.revision(), revision); + assert_eq!(live.health(), "failed"); + assert!(live.export_inputs().is_err()); + assert!(live + .save_checkpoint(&fixture.root.join("uncertain.json"), Value::Null) + .is_err()); +}