diff --git a/.github/workflows/checks.yml b/.github/workflows/checks.yml index c109319..5e3f11c 100644 --- a/.github/workflows/checks.yml +++ b/.github/workflows/checks.yml @@ -16,6 +16,7 @@ jobs: - name: Formatting and transport-free contracts run: | cargo fmt --all --check + cargo clippy --locked --workspace --all-targets --all-features -- -D warnings cargo test --locked --workspace - name: MCP and controlled worker contracts run: | diff --git a/README.md b/README.md index 8c10af0..d8a7486 100644 --- a/README.md +++ b/README.md @@ -81,6 +81,9 @@ One registered operation consumes/returns strings, with explicit submission, cla ## Verification +The [runtime contract](docs/specification.md) maps supported behavior to executable +tests and distinguishes local checkpoints from unimplemented storage features. + ```sh cargo fmt --all --check cargo test --locked --workspace diff --git a/crates/syntax/src/lib.rs b/crates/syntax/src/lib.rs index b56cbe6..332c9ef 100644 --- a/crates/syntax/src/lib.rs +++ b/crates/syntax/src/lib.rs @@ -2,9 +2,9 @@ pub mod ast; pub use ast::{parse_program, Atom, Clause, CmpOp, Expr, Lit, ParseError}; -/// Aggregate functions usable in rule HEAD arguments only -/// (`kit_count(P, count(K))`). Lowered internally to a temp relation plus -/// a group-by fold; the head predicate completes before any reader. +/// Aggregate syntax recognized in rule heads, such as `kit_count(P, count(K))`. +/// This crate only parses the syntax; each consumer decides which constructs it +/// supports. DDlog Runtime currently rejects aggregates during lowering. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum AggFn { Count, diff --git a/docs/hardening.md b/docs/hardening.md new file mode 100644 index 0000000..27eaf0c --- /dev/null +++ b/docs/hardening.md @@ -0,0 +1,75 @@ +# Hardening baseline — September 6, 2026 + +Initially assessed `f763612129a0275e53b34869a428c12c1a5bf7a9`, then rebased and +revalidated against merged `dc50b2f43c430ff406c6145e180afb0cf0584449`, including +bounded reads. All work is in an isolated checkout. This is a measured baseline and targeted cleanup, not a +production certification or a claim that all hotspots have been eliminated. + +## Changes and validation + +- Separated registry operations and pinned activation from semantic dispatch. + Existing health/pin admission guards remain in the shared dispatch path. +- Added registered-request contracts for claim requirements, stale completion, + identical/conflicting completion, and lost native acknowledgment. They test + the real Rust admission implementation with explicitly simulated transport. +- Added signed-integer extrema and whole-batch rejection checks relevant to M0. +- Corrected syntax documentation that implied aggregate execution support. +- Added the runtime specification with explicit tests and unsupported guarantees. +- Simplified three Clippy findings without changing public APIs or generated code. +- Final instrumented Rust suite: 75 passed, one native-compiler test ignored. +- Final Python suite: 58 passed. Unix socket access required an escalation; + the initial sandbox denial occurred before any runtime test executed. +- Strict all-target/all-feature Clippy passed. After checkpoint ownership was + released, a local type alias removed its test-only type-complexity warning. + CI now runs strict Clippy alongside the existing contract suites. + +## Measurement + +Tools: rustc 1.94.1, cargo-llvm-cov 0.9.1, Lizard 1.24.0. Runtime source line +coverage is 2983/3275 = **91.08%** on the rebased product. On the earlier +like-for-like baseline it improved from 2755/3065 = **89.89%** to 2795/3067 = +**91.13%** before the bounded-read merge changed the measured population. These totals +exclude integration tests, the separate syntax crate and generated/native engine +code. Host subprocess profiles from the Python suite are included. Forced process +termination can omit profiling data, so an uncovered function is a review target, +not proof that no test executes it. + +The diagnostic uses Lizard's Rust cyclomatic complexity and executable-line +coverage in each function's source range from LLVM LCOV. Formula: +`CC² × (1 − coverage)³ + CC`. This is a line-based CRAP estimate, not branch +coverage or formal correctness evidence. No hard score gate has been chosen. +The main dispatch estimate changed from 163.2 (CC106) to 49.4 (CC45); responsibility +extraction reduces local complexity but does not remove total system complexity. + +| Remaining hotspot | CC | Line coverage | CRAP estimate | +| --- | ---: | ---: | ---: | +| `src/instance.rs::execute` | 45 | 87.1% | 49.4 | +| `src/host.rs::standalone` | 6 | 0.0% | 42.0 | +| `src/lower.rs::lower_clauses_with_operators` | 41 | 92.3% | 41.8 | +| `src/instance.rs::execute_registry` | 41 | 93.3% | 41.5 | +| `src/instance.rs::install_processor` | 22 | 68.3% | 37.4 | +| `src/composition.rs::manifest` | 35 | 92.4% | 35.5 | +| `src/host.rs::host` | 32 | 88.8% | 33.4 | +| `src/host.rs::connection` | 26 | 87.9% | 27.2 | + +## Reproduction + +Install cargo-llvm-cov 0.9.1 and Rust llvm-tools-preview. Run Rust coverage with +`cargo llvm-cov --locked --workspace --all-features --no-report`. Run the Python +suite against the instrumented `target/llvm-cov-target/debug/lemmalog-ddlog-mcp` +by setting `LEMMALOG_DDLOG_MCP`, with `LLVM_PROFILE_FILE` pointing into that coverage +target using distinct process/module placeholders `%p-%m.profraw`. Then export +`cargo llvm-cov report --lcov --output-path coverage.lcov`. In an environment with +Lizard 1.24.0, run `python scripts/complexity_coverage.py coverage.lcov`. + +Do not report Rust-only coverage as transport coverage. Do not report this fixture +suite as native compilation or formal validation of a Lean-to-dataflow compiler. + +## Remaining work before a broad clean bill + +The bounded-read branch is now included. Review compiler, composition and host +hotspots in context rather than splitting them merely to reduce a score. Make +native compiler validation available as explicit evidence, and decide which +coverage/complexity checks belong in CI after the baseline stabilizes. The current +library has local pure checkpoints, not integrated generic Iceberg persistence. +No world, ECS, mission, deployment or external-effect recovery features were added. diff --git a/docs/specification.md b/docs/specification.md new file mode 100644 index 0000000..c125e6b --- /dev/null +++ b/docs/specification.md @@ -0,0 +1,62 @@ +# Runtime contract + +This document describes the supported runtime, not planned world or storage +abstractions. The library and MCP expose the same program admission rules; +transport does not define program semantics. Focused guides below describe +wire shapes and operational prerequisites. + +| Boundary | Required behavior | Executable evidence | +| --- | --- | --- | +| Language | Explicit `int`/`string` schemas; positive recursion and safe stratified negation; reject unsupported constructs and negative cycles before compilation | `tests/ddlog_lowering.rs`, `tests/recursive_lowering.rs` | +| Definitions | Immutable versions preserve identity and lineage; conditional pointer and lifecycle updates reject stale callers | `tests/processor_registry.rs` | +| Composition | Resolve exact versions, isolate private names, validate bindings, and preserve public interfaces through nesting | `tests/processor_composition_registry.rs` | +| Activation | Compile a candidate before replacement; compatible retained inputs replay; failed replacement preserves the prior usable program | `tests/memory_runtime.rs` | +| Input changes | Validate the complete input transaction before execution; acknowledged state changes only after completion; uncertain native failure disables continued use | `tests/memory_runtime.rs` | +| Bounded reads | Bound returned rows/bytes, bind continuation to the backend revision and request, drain selected native output after local limit errors; no indexed-query latency guarantee | `tests/memory_runtime.rs`, [bounded reads](bounded-reads.md) | +| Registered requests | Claim before completion; preserve late results with explicit freshness; identical completion is idempotent; conflicting completion fails; uncertain settlement never implies permission to repeat provider work | `tests/registered_requests.rs` | +| Shared access | Connections share one owner; disconnect does not stop it; pinned versions and exported ports remain enforced | `tests/test_shared_host.py`, `tests/upstream_compatibility.rs` | +| Checkpoints | Explicit pure-backend checkpoint; integrity-checked restore into a fresh backend; reconstruct outputs from acknowledged inputs | `tests/memory_runtime.rs`, [checkpoint contract](checkpoints.md) | +| Compatibility | Preserve previously published definition hashes, native source fixtures, and MCP tool schemas unless a deliberate compatibility change is declared | `tests/upstream_compatibility.rs` | + +## What an acknowledgment establishes + +A completed input operation means the native program acknowledged its transaction. +It does not mean that inputs are durable, an external provider completed, or an +Iceberg catalog published anything. Program installation version and backend +revision are distinct; callers must not use a version number as a transaction ID. + +An explicit checkpoint has its own durability outcome. Directory-sync failure +after replacement is uncertain publication and requires inspection before retry. +Checkpoint integrity is corruption detection, not authentication. Caller metadata +does not automatically reconstruct registry pins, public interfaces, or provider +admission state. See [checkpoints](checkpoints.md) for supported contents and limits. + +## Deliberate limits + +- The parser recognizes more syntax than the runtime supports. Parsing is not + compilation, and compilation is not a mathematical correctness proof. +- Registered-operation programs do not compose with ordinary program nodes. + Imported native operators and registered operations cannot use checkpoint + format 1. Ordinary inference state is session-local. +- One host owns one graph instance. No fleet scheduling, world semantics, + automatic crash recovery, distributed transaction, or provider exactly-once + guarantee is supplied. +- Generic Arrow/Iceberg persistence is not integrated. Separate experiments are + evidence for their scoped scenarios, not shipped runtime guarantees. +- `why` provides direct rule witnesses, not a Lean proof, recursive provenance, + or a certificate for a source-to-dataflow compiler. + +## Validation scope + +Default Rust and Python suites validate runtime contracts using controlled native +transport fixtures. They do not execute the actual DDlog compiler or certify the +upstream engine. [Native acceptance](building.md#native-acceptance) is a separate +configured run. Record whether native evidence is a fresh compilation or reuse of +a hash-verified executable. + +Coverage must include the instrumented MCP binary exercised by Python, otherwise +host coverage is understated. Generated native sources, Rust host code, syntax +code, and Python worker coverage are separate populations. Report the population +and skipped tests with any percentage. Complexity/coverage scores prioritize +review; they do not prove correctness or justify splitting code solely to improve +a metric. diff --git a/scripts/complexity_coverage.py b/scripts/complexity_coverage.py new file mode 100644 index 0000000..be2502c --- /dev/null +++ b/scripts/complexity_coverage.py @@ -0,0 +1,45 @@ +#!/usr/bin/env python3 +"""Diagnostic CRAP estimates from Lizard Rust complexity and LLVM LCOV lines. + +Requires lizard==1.24.0. Run from the repository root, passing an LCOV file. +This uses executable-line coverage, not branch coverage. No score is invented +for functions without coverage mapping. Generated native sources are excluded. +""" +import json +from pathlib import Path +import sys + +import lizard + + +def report(path): + coverage = {} + current = None + for line in Path(path).read_text().splitlines(): + if line.startswith("SF:"): + current = str(Path(line[3:]).resolve()) + coverage.setdefault(current, {}) + elif line.startswith("DA:") and current is not None: + number, hits, *_ = line[3:].split(",") + number, hits = int(number), int(hits) + coverage[current][number] = max(hits, coverage[current].get(number, 0)) + rows = [] + for source in sorted(Path("src").glob("*.rs")): + lines = coverage.get(str(source.resolve()), {}) + for function in lizard.analyze_file(str(source)).function_list: + hits = [count for number, count in lines.items() + if function.start_line <= number <= function.end_line] + fraction = sum(count > 0 for count in hits) / len(hits) if hits else None + complexity = function.cyclomatic_complexity + rows.append({ + "file": str(source), "function": function.name, + "line": function.start_line, "complexity": complexity, + "line_coverage": fraction, + "crap_estimate": (complexity ** 2 * (1 - fraction) ** 3 + complexity + if fraction is not None else None), + }) + return sorted(rows, key=lambda row: row["crap_estimate"] or 0, reverse=True) + + +if __name__ == "__main__": + print(json.dumps(report(sys.argv[1]), indent=2)) diff --git a/src/composition.rs b/src/composition.rs index 88a8d1d..933125f 100644 --- a/src/composition.rs +++ b/src/composition.rs @@ -100,7 +100,7 @@ pub fn validate_interface(program: &ProgramDefinition) -> Result<()> { if !outputs.insert(name.clone()) { return Err(format!("Duplicate interface output {name}")); } - if !schemas.get(name).is_some_and(|schema| !schema.input) { + if !matches!(schemas.get(name), Some(schema) if !schema.input) { return Err(format!( "Interface output {name} must name a declared derived relation" )); diff --git a/src/instance.rs b/src/instance.rs index ef786d8..b75746e 100644 --- a/src/instance.rs +++ b/src/instance.rs @@ -50,88 +50,7 @@ impl ProgramInstance { ); } if name.starts_with("processor_") && name != "processor_install" { - let registry = self - .registry - .as_ref() - .ok_or("Processor registry is not configured")?; - if name == "processor_list" || name == "processor_search" { - let limit = a - .get("limit") - .map(|value| value.as_u64().ok_or("limit must be an integer")) - .transpose()? - .unwrap_or(20); - let limit = usize::try_from(limit).map_err(|e| e.to_string())?; - let after = a.get("after").map(|_| string(a, "after")).transpose()?; - let include_archived = a - .get("include_archived") - .map(|value| value.as_bool().ok_or("include_archived must be a boolean")) - .transpose()? - .unwrap_or(false); - let page = if name == "processor_search" { - registry.search(string(a, "query")?, limit, after, include_archived)? - } else { - registry.list(limit, after, include_archived)? - }; - return serde_json::to_value(page).map_err(|e| e.to_string()); - } - if name == "processor_archive" || name == "processor_restore" { - let expected_revision = a["expected_revision"] - .as_u64() - .ok_or("Missing nonnegative expected_revision")?; - let lifecycle = if name == "processor_archive" { - registry.archive( - string(a, "processor_id")?, - string(a, "expected_version")?, - expected_revision, - )? - } else { - registry.restore( - string(a, "processor_id")?, - string(a, "expected_version")?, - expected_revision, - )? - }; - return serde_json::to_value(lifecycle).map_err(|e| e.to_string()); - } - let provenance = || -> Result, String> { - a.get("git_provenance") - .filter(|v| !v.is_null()) - .map(|v| serde_json::from_value(v.clone()).map_err(|e| e.to_string())) - .transpose() - }; - let definition = || -> Result { - // Select the shape before deserialization so missing/unknown - // fields remain visible instead of an opaque untagged-enum error. - if a["definition"].get("composition").is_some() { - serde_json::from_value(a["definition"].clone()) - .map(ProcessorDefinition::Composition) - .map_err(|e| e.to_string()) - } else { - serde_json::from_value(a["definition"].clone()) - .map(ProcessorDefinition::Program) - .map_err(|e| e.to_string()) - } - }; - let record = match name { - "processor_create" => registry.create(definition()?, provenance()?)?, - "processor_publish" => registry.publish( - string(a, "processor_id")?, - definition()?, - string(a, "expected_version")?, - provenance()?, - )?, - "processor_fork" => registry.fork( - string(a, "processor_id")?, - string(a, "version")?, - provenance()?, - )?, - "processor_get" => registry.get( - string(a, "processor_id")?, - a.get("version").map(|_| string(a, "version")).transpose()?, - )?, - _ => return Err("Unknown tool".into()), - }; - return serde_json::to_value(record).map_err(|e| e.to_string()); + return self.execute_registry(name, a); } if self.instance_id.is_some() && self.backend.health() == "failed" @@ -149,97 +68,7 @@ impl ProgramInstance { return Err("Instance is pinned to an immutable processor version; create a new instance to select another version".into()); } match name { - "processor_install" => { - if self.backend.health() != "uninitialized" || self.agent.is_some() { - return Err("Select a processor only in a fresh instance".into()); - } - self.registry - .as_ref() - .ok_or("Processor registry is not configured")? - .ensure_active(string(a, "processor_id")?)?; - let record = self - .registry - .as_ref() - .ok_or("Processor registry is not configured")? - .get( - string(a, "processor_id")?, - a.get("version").map(|_| string(a, "version")).transpose()?, - )?; - let mut result = match record.definition { - ProcessorDefinition::Composition(definition) => { - let compiled = self - .registry - .as_ref() - .unwrap() - .compile_composition(&definition.composition)?; - let result = self - .backend - .install_source(compiled.source, compiled.schemas)?; - self.interface = Some(PublicInterface { - inputs: compiled.resolution.inputs.clone(), - outputs: compiled.resolution.outputs.clone(), - }); - self.composition = Some(compiled.resolution); - result - } - ProcessorDefinition::Program(definition) => { - let result = if let Some(binding) = definition.operation { - if !definition.operators.is_empty() { - return Err("Typed operators cannot be combined with a registered operation; put the operator in a separate pure program".into()); - } - let operation = self - .operations - .get(&binding.name) - .ok_or("Pinned operation is not registered on this host")? - .clone(); - if operation.version != binding.version - || operation.description != binding.description - { - return Err( - "Pinned operation definition does not match this host registry" - .into(), - ); - } - let (agent, result) = AgentProgram::install( - &mut self.backend, - &binding.name, - operation, - &definition.rules, - definition.schemas, - )?; - self.agent = Some(agent); - result - } else { - self.backend.install_with_operators( - &definition.rules, - definition.schemas, - &definition.operators, - )? - }; - self.interface = definition.interface.map(|interface| PublicInterface { - inputs: interface - .inputs - .into_iter() - .map(|name| (name.clone(), name)) - .collect(), - outputs: interface - .outputs - .into_iter() - .map(|name| (name.clone(), name)) - .collect(), - }); - result - } - }; - self.processor = - Some(json!({"processor_id":record.processor_id,"version":record.version})); - result["processor"] = self.processor.clone().unwrap(); - if let Some(composition) = &self.composition { - result["composition"] = - serde_json::to_value(composition).map_err(|e| e.to_string())?; - } - Ok(result) - } + "processor_install" => self.install_processor(a), "agent_operations" => Ok( json!({"operations":self.operations.iter().map(|(name,op)|json!({"name":name,"version":op.version,"description":op.description,"input":"string","output":"string"})).collect::>()}), ), @@ -356,6 +185,180 @@ impl ProgramInstance { _ => Err("Unknown tool".into()), } } + + fn execute_registry(&self, name: &str, a: &Value) -> Result { + let registry = self + .registry + .as_ref() + .ok_or("Processor registry is not configured")?; + if name == "processor_list" || name == "processor_search" { + let limit = a + .get("limit") + .map(|value| value.as_u64().ok_or("limit must be an integer")) + .transpose()? + .unwrap_or(20); + let limit = usize::try_from(limit).map_err(|e| e.to_string())?; + let after = a.get("after").map(|_| string(a, "after")).transpose()?; + let include_archived = a + .get("include_archived") + .map(|value| value.as_bool().ok_or("include_archived must be a boolean")) + .transpose()? + .unwrap_or(false); + let page = if name == "processor_search" { + registry.search(string(a, "query")?, limit, after, include_archived)? + } else { + registry.list(limit, after, include_archived)? + }; + return serde_json::to_value(page).map_err(|e| e.to_string()); + } + if name == "processor_archive" || name == "processor_restore" { + let expected_revision = a["expected_revision"] + .as_u64() + .ok_or("Missing nonnegative expected_revision")?; + let lifecycle = if name == "processor_archive" { + registry.archive( + string(a, "processor_id")?, + string(a, "expected_version")?, + expected_revision, + )? + } else { + registry.restore( + string(a, "processor_id")?, + string(a, "expected_version")?, + expected_revision, + )? + }; + return serde_json::to_value(lifecycle).map_err(|e| e.to_string()); + } + let provenance = || -> Result, String> { + a.get("git_provenance") + .filter(|v| !v.is_null()) + .map(|v| serde_json::from_value(v.clone()).map_err(|e| e.to_string())) + .transpose() + }; + let definition = || -> Result { + // Select the shape before deserialization so missing/unknown + // fields remain visible instead of an opaque untagged-enum error. + if a["definition"].get("composition").is_some() { + serde_json::from_value(a["definition"].clone()) + .map(ProcessorDefinition::Composition) + .map_err(|e| e.to_string()) + } else { + serde_json::from_value(a["definition"].clone()) + .map(ProcessorDefinition::Program) + .map_err(|e| e.to_string()) + } + }; + let record = match name { + "processor_create" => registry.create(definition()?, provenance()?)?, + "processor_publish" => registry.publish( + string(a, "processor_id")?, + definition()?, + string(a, "expected_version")?, + provenance()?, + )?, + "processor_fork" => registry.fork( + string(a, "processor_id")?, + string(a, "version")?, + provenance()?, + )?, + "processor_get" => registry.get( + string(a, "processor_id")?, + a.get("version").map(|_| string(a, "version")).transpose()?, + )?, + _ => return Err("Unknown tool".into()), + }; + serde_json::to_value(record).map_err(|e| e.to_string()) + } + + fn install_processor(&mut self, a: &Value) -> Result { + if self.backend.health() != "uninitialized" || self.agent.is_some() { + return Err("Select a processor only in a fresh instance".into()); + } + self.registry + .as_ref() + .ok_or("Processor registry is not configured")? + .ensure_active(string(a, "processor_id")?)?; + let record = self + .registry + .as_ref() + .ok_or("Processor registry is not configured")? + .get( + string(a, "processor_id")?, + a.get("version").map(|_| string(a, "version")).transpose()?, + )?; + let mut result = match record.definition { + ProcessorDefinition::Composition(definition) => { + let compiled = self + .registry + .as_ref() + .unwrap() + .compile_composition(&definition.composition)?; + let result = self + .backend + .install_source(compiled.source, compiled.schemas)?; + self.interface = Some(PublicInterface { + inputs: compiled.resolution.inputs.clone(), + outputs: compiled.resolution.outputs.clone(), + }); + self.composition = Some(compiled.resolution); + result + } + ProcessorDefinition::Program(definition) => { + let result = if let Some(binding) = definition.operation { + if !definition.operators.is_empty() { + return Err("Typed operators cannot be combined with a registered operation; put the operator in a separate pure program".into()); + } + let operation = self + .operations + .get(&binding.name) + .ok_or("Pinned operation is not registered on this host")? + .clone(); + if operation.version != binding.version + || operation.description != binding.description + { + return Err( + "Pinned operation definition does not match this host registry".into(), + ); + } + let (agent, result) = AgentProgram::install( + &mut self.backend, + &binding.name, + operation, + &definition.rules, + definition.schemas, + )?; + self.agent = Some(agent); + result + } else { + self.backend.install_with_operators( + &definition.rules, + definition.schemas, + &definition.operators, + )? + }; + self.interface = definition.interface.map(|interface| PublicInterface { + inputs: interface + .inputs + .into_iter() + .map(|name| (name.clone(), name)) + .collect(), + outputs: interface + .outputs + .into_iter() + .map(|name| (name.clone(), name)) + .collect(), + }); + result + } + }; + self.processor = Some(json!({"processor_id":record.processor_id,"version":record.version})); + result["processor"] = self.processor.clone().unwrap(); + if let Some(composition) = &self.composition { + result["composition"] = serde_json::to_value(composition).map_err(|e| e.to_string())?; + } + Ok(result) + } } /// Only these declared ports can be addressed through the ordinary fact tools. diff --git a/src/processes.rs b/src/processes.rs index b68036c..72a3da8 100644 --- a/src/processes.rs +++ b/src/processes.rs @@ -40,10 +40,8 @@ impl ProcessControl { } pub fn track(&self, pid: u32) -> Group { self.inner.groups.lock().unwrap().insert(pid); - if self.stopped() { - if self.inner.detached { - kill_group(pid); - } + if self.stopped() && self.inner.detached { + kill_group(pid); } Group { pid, diff --git a/src/registry.rs b/src/registry.rs index 12dbb25..649e34d 100644 --- a/src/registry.rs +++ b/src/registry.rs @@ -420,7 +420,7 @@ impl ProcessorRegistry { continue; }; if validate_processor_id(&identity).is_ok() - && after.map_or(true, |cursor| identity.as_str() > cursor) + && (after.is_none() || after.is_some_and(|cursor| identity.as_str() > cursor)) && entry.file_type().map_err(io_error)?.is_dir() { identities.push(identity); diff --git a/tests/memory_runtime.rs b/tests/memory_runtime.rs index e511a36..e6fb077 100644 --- a/tests/memory_runtime.rs +++ b/tests/memory_runtime.rs @@ -372,7 +372,8 @@ fn corrupt_and_unsupported_checkpoints_fail_before_compilation_or_activation() { .restore_checkpoint(&path) .unwrap_err() .contains("integrity")); - let edits: Vec> = vec![ + type CheckpointEdit = Box; + let edits: Vec = vec![ Box::new(|state| state["format_version"] = json!(999)), Box::new(|state| state["program_version"] = json!(999)), Box::new(|state| { diff --git a/tests/registered_requests.rs b/tests/registered_requests.rs new file mode 100644 index 0000000..fb3b647 --- /dev/null +++ b/tests/registered_requests.rs @@ -0,0 +1,151 @@ +#![cfg(unix)] +//! Admission/settlement contracts using simulated transport, not native evaluation. +use ddlog_runtime::{AgentProgram, Backend, Operation}; +use serde_json::json; +use std::fs; +use std::os::unix::fs::PermissionsExt; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicU64, Ordering}; +static NEXT: AtomicU64 = AtomicU64::new(0); +struct Fixture { + root: PathBuf, +} +impl Fixture { + fn new() -> Self { + let root = std::env::temp_dir().join(format!( + "ddlog-requests-test-{}-{}", + std::process::id(), + NEXT.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir(&root).unwrap(); + let template = + Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/memory_fake_runtime.py"); + let script = format!("#!/usr/bin/env python3\nimport sys\nfrom pathlib import Path\nroot=Path({})\nif (root/'reject_build').exists(): sys.exit(91)\nsource=Path({}).read_text().replace('__CONTROL__', {})\nPath(sys.argv[2]).write_text(source)\nPath(sys.argv[2]).chmod(0o700)\n", json!(root), json!(template), json!(serde_json::to_string(&root).unwrap())); + fs::write(root.join("build.py"), script).unwrap(); + fs::set_permissions(root.join("build.py"), fs::Permissions::from_mode(0o700)).unwrap(); + Self { root } + } + fn backend(&self, name: &str) -> Backend { + Backend::new(self.root.join(name), self.root.join("build.py")) + } + fn flag(&self, name: &str) { + fs::write(self.root.join(name), "").unwrap(); + } + fn unflag(&self, name: &str) { + fs::remove_file(self.root.join(name)).unwrap(); + } +} +impl Drop for Fixture { + fn drop(&mut self) { + let _ = fs::remove_dir_all(&self.root); + } +} + +fn install(f: &Fixture) -> (AgentProgram, Backend) { + let mut backend = f.backend("live"); + let (agent, _) = AgentProgram::install( + &mut backend, + "infer", + Operation { + version: "v1".into(), + description: "fixture".into(), + }, + "answer(E,R,O) :- agent_result(E,R,O).", + json!({"answer":{"input":false,"fields":["string","int","string"]}}), + ) + .unwrap(); + (agent, backend) +} +#[test] +fn completion_requires_claim_and_preserves_current_and_historical_results() { + let f = Fixture::new(); + let (mut agent, mut backend) = install(&f); + let old = agent.submit(&mut backend, "entity", 1, "first").unwrap()["request_id"] + .as_str() + .unwrap() + .to_owned(); + assert!(agent + .complete(&mut backend, &old, "answer") + .unwrap_err() + .contains("claimed")); + agent.claim(&mut backend, &old).unwrap(); + assert!(agent.claim(&mut backend, &old).is_err()); + let new = agent.submit(&mut backend, "entity", 2, "second").unwrap()["request_id"] + .as_str() + .unwrap() + .to_owned(); + assert_eq!( + agent.complete(&mut backend, &old, "old result").unwrap()["fresh"], + false + ); + let revision = backend.revision(); + assert_eq!( + agent.complete(&mut backend, &old, "old result").unwrap()["duplicate"], + true + ); + assert_eq!(backend.revision(), revision); + assert!(agent + .complete(&mut backend, &old, "different") + .unwrap_err() + .contains("Conflicting")); + agent.claim(&mut backend, &new).unwrap(); + let done = agent.complete(&mut backend, &new, "new result").unwrap(); + assert_eq!(done["fresh"], true); + assert_eq!(done["duplicate"], false); + assert!(agent.submit(&mut backend, "entity", 1, "first").is_err()); + assert!(agent.submit(&mut backend, "entity", 2, "changed").is_err()); + assert!(agent.complete(&mut backend, "missing", "answer").is_err()); +} +#[test] +fn lost_completion_ack_keeps_request_unsettled_and_disables_runtime_reuse() { + let f = Fixture::new(); + let (mut agent, mut backend) = install(&f); + let id = agent.submit(&mut backend, "entity", 1, "input").unwrap()["request_id"] + .as_str() + .unwrap() + .to_owned(); + agent.claim(&mut backend, &id).unwrap(); + f.flag("die_on_commit"); + assert!(agent.complete(&mut backend, &id, "result").is_err()); + assert_eq!(backend.health(), "failed"); + assert_eq!(agent.status()["requests"][0]["status"], "claimed"); + f.unflag("die_on_commit"); + assert!(agent.complete(&mut backend, &id, "result").is_err()); + assert_eq!(agent.status()["requests"][0]["status"], "claimed"); +} + +#[test] +fn signed_integer_boundaries_and_invalid_tail_preserve_complete_input_state() { + let f = Fixture::new(); + let mut backend = f.backend("typed"); + backend + .install( + "echo(N,S) :- source(N,S).", + json!({ + "source":{"input":true,"fields":["int","string"]}, + "echo":{"input":false,"fields":["int","string"]} + }), + ) + .unwrap(); + backend + .apply(&json!([ + {"op":"insert","predicate":"source","values":[i64::MIN,"low"]}, + {"op":"insert","predicate":"source","values":[i64::MAX,"high"]} + ])) + .unwrap(); + let before = serde_json::to_value(backend.export_inputs().unwrap()).unwrap(); + let revision = backend.revision(); + for invalid in [json!(u64::MAX), json!(1.5), json!(true), json!("1")] { + assert!(backend + .apply(&json!([ + {"op":"delete","predicate":"source","values":[i64::MIN,"low"]}, + {"op":"insert","predicate":"source","values":[invalid,"invalid"]} + ])) + .is_err()); + assert_eq!(backend.revision(), revision); + assert_eq!( + serde_json::to_value(backend.export_inputs().unwrap()).unwrap(), + before + ); + } +}