Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,8 @@ fn main() -> Result<(), String> {

`Backend` is the lower-level compilation/execution interface. `Backend::install_composition(&registry, &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.
Expand Down
105 changes: 105 additions & 0 deletions docs/bounded-reads.md
Original file line number Diff line number Diff line change
@@ -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.
13 changes: 10 additions & 3 deletions docs/checkpoints.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<CompositionResolution, String>` resolves
Expand Down
214 changes: 214 additions & 0 deletions src/bounded.rs
Original file line number Diff line number Diff line change
@@ -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<usize, Value>,
pub max_rows: usize,
pub max_bytes: usize,
pub continuation: Option<QueryCursor>,
}

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<usize, Value>,
offset: u64,
}

#[derive(Clone, Debug, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct QueryPage {
pub rows: Vec<Vec<Value>>,
pub bytes: usize,
/// True only after observing an additional matching row.
pub truncated: bool,
pub continuation: Option<QueryCursor>,
pub revision: u64,
}

struct CountBytes(usize);
impl Write for CountBytes {
fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
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<QueryPage> {
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)
}
}
Loading
Loading