Skip to main content

vole_document/field/
dag.rs

1//! The procedural seed DAG: bounded closure traversal and materialization.
2//!
3//! Every node names a **bounded, versioned materializer**; there is no arbitrary
4//! execution (ADR-0025). A materializer is a pure function of the field's exact
5//! descriptor, the seed store, and the node's dependency outputs. Unknown kinds
6//! or materializer versions fail closed.
7//!
8//! Exact node kinds (`Q_ref`) resolve to byte-identical spans of the source via
9//! the *existing* partial-materialization machinery, so a narrow observation does
10//! not force a whole-document reconstruction. Derived node kinds (`Q_gen`, e.g.
11//! decoded streams, text runs, previews) are deterministic projections and are
12//! always labelled as such.
13
14use crate::container::{ObjectSource, ParsedDescriptor};
15use crate::error::{Error, Result};
16use crate::limits::Limits;
17use crate::materialize::observation::{select_ops, selection_references, serve_selection};
18use crate::store::{NodeId, SeedStore};
19
20use super::derive;
21use super::node::{NodeKind, SeedNode, read_object_params, read_span_params, read_u32_params};
22
23/// Supplies exact source byte ranges to a seed-DAG materialization.
24///
25/// The full implementation is a parsed descriptor
26/// ([`impl SourceServer for ParsedDescriptor`]); a partial loader that reads only
27/// the records a query needs is the other. Every materializer that resolves an
28/// exact `Q_ref` node goes through this one method, so the DAG logic cannot
29/// diverge between the complete and partial sources.
30pub trait SourceServer {
31    /// Bytes of the source range `[offset, offset + len)`.
32    fn serve_range(&self, offset: u64, len: u64, limits: Limits) -> Result<Vec<u8>>;
33    /// The whole reconstructed source (only a `DocumentExact` node needs this).
34    fn serve_document(&self, limits: Limits) -> Result<Vec<u8>>;
35}
36
37impl SourceServer for ParsedDescriptor {
38    fn serve_range(&self, offset: u64, len: u64, limits: Limits) -> Result<Vec<u8>> {
39        serve_source_range(self, offset, len, limits)
40    }
41
42    fn serve_document(&self, limits: Limits) -> Result<Vec<u8>> {
43        crate::materialize::materialize(self, limits)
44    }
45}
46
47/// Hard cap on the number of nodes one materialization may evaluate.
48pub const MAX_EVAL_NODES: u64 = 1 << 20;
49/// Hard cap on total intermediate+output bytes one materialization may produce.
50pub const MAX_EVAL_BYTES: u64 = 1 << 32;
51
52/// A bounded evaluation budget, shared across a recursive materialization.
53#[derive(Debug, Clone)]
54pub struct EvalBudget {
55    /// Maximum nodes evaluated.
56    pub max_nodes: u64,
57    /// Maximum total bytes produced (intermediate + final).
58    pub max_bytes: u64,
59    /// Nodes evaluated so far.
60    pub nodes: u64,
61    /// Bytes produced so far.
62    pub produced: u64,
63}
64
65impl Default for EvalBudget {
66    fn default() -> Self {
67        EvalBudget {
68            max_nodes: MAX_EVAL_NODES,
69            max_bytes: MAX_EVAL_BYTES,
70            nodes: 0,
71            produced: 0,
72        }
73    }
74}
75
76impl EvalBudget {
77    fn charge_node(&mut self) -> Result<()> {
78        self.nodes = self
79            .nodes
80            .checked_add(1)
81            .ok_or_else(|| Error::resource_limit("seed evaluation node count overflow"))?;
82        if self.nodes > self.max_nodes {
83            return Err(Error::resource_limit(format!(
84                "seed evaluation exceeded {} nodes",
85                self.max_nodes
86            )));
87        }
88        Ok(())
89    }
90
91    pub(crate) fn charge_bytes(&mut self, n: u64) -> Result<()> {
92        self.produced = self
93            .produced
94            .checked_add(n)
95            .ok_or_else(|| Error::resource_limit("seed evaluation byte count overflow"))?;
96        if self.produced > self.max_bytes {
97            return Err(Error::resource_limit(format!(
98                "seed evaluation exceeded {} bytes",
99                self.max_bytes
100            )));
101        }
102        Ok(())
103    }
104}
105
106/// Load and decode one node from the store by id, verifying content identity.
107pub fn load_node(store: &dyn SeedStore, id: &NodeId) -> Result<SeedNode> {
108    let bytes = store.get_node(id)?;
109    let node = SeedNode::decode_canonical(&bytes)?;
110    if node.content_id() != *id {
111        return Err(Error::integrity_mismatch(format!(
112            "seed node {id} decoded to a different content id"
113        )));
114    }
115    node.check_limits(&Limits::DEFAULT)?;
116    Ok(node)
117}
118
119/// The dependency ids of a node's already-decoded canonical bytes.
120pub fn deps_of_canonical(bytes: &[u8]) -> Result<Vec<NodeId>> {
121    Ok(SeedNode::decode_canonical(bytes)?.deps)
122}
123
124/// Serve an exact source byte range `[offset, offset+len)` using the descriptor's
125/// program via the shared partial-materialization path. Does **not** materialize
126/// the whole document.
127fn serve_source_range(
128    parsed: &ParsedDescriptor,
129    offset: u64,
130    len: u64,
131    limits: Limits,
132) -> Result<Vec<u8>> {
133    let d = &parsed.descriptor;
134    if d.objects
135        .iter()
136        .any(|o| matches!(o, ObjectSource::External { .. }))
137    {
138        return Err(Error::unsupported_feature(
139            "field v1 requires an inline descriptor (no external objects)",
140        ));
141    }
142    let end = offset
143        .checked_add(len)
144        .ok_or_else(|| Error::usage("source slice end overflows"))?;
145    if end > d.source_len {
146        return Err(Error::usage(format!(
147            "source slice {offset}..{end} exceeds source length {}",
148            d.source_len
149        )));
150    }
151    let objects: Vec<Vec<u8>> = d
152        .objects
153        .iter()
154        .map(|o| o.as_inline().unwrap_or(&[]).to_vec())
155        .collect();
156    let object_lens: Vec<u64> = objects.iter().map(|o| o.len() as u64).collect();
157    let channel_lens: Vec<u64> = d.channels.iter().map(|c| c.decoded_length).collect();
158    let window = select_ops(&d.program, &object_lens, &channel_lens, offset, end, limits)?;
159    let (objects_used, channels_used) =
160        selection_references(&window.ops, objects.len(), d.channels.len());
161    let served = serve_selection(
162        &objects,
163        &d.channels,
164        &d.models,
165        window,
166        &objects_used,
167        &channels_used,
168        offset,
169        end,
170        limits,
171    )?;
172    Ok(served.bytes)
173}
174
175/// A node-output cache keyed by [`NodeId`]. Because a node's id binds its full
176/// dependency closure, an unchanged closure hits and a changed dependency misses;
177/// there is no invalidation pass (ADR-0025).
178///
179/// Implementations are disposable: [`materialize_node_cached`] treats any `get`
180/// error as a miss and never trusts bytes it cannot validate, so a corrupt cache
181/// causes recomputation rather than wrong output.
182pub trait OutputCache {
183    /// Fetch a cached node output, or `None` on a miss. A `get` error is treated
184    /// as a miss by the caller (the cache is disposable, never authority).
185    fn get(&self, id: &NodeId) -> Result<Option<Vec<u8>>>;
186    /// Store a node output. Best-effort: a `put` error does not fail the
187    /// materialization.
188    fn put(&mut self, id: &NodeId, bytes: &[u8]) -> Result<()>;
189}
190
191/// A cache that stores nothing; used by the backward-compatible
192/// [`materialize_node`] wrapper and the `use_cache = false` cold court.
193#[derive(Debug, Clone, Copy, Default)]
194pub struct NoCache;
195
196impl OutputCache for NoCache {
197    fn get(&self, _id: &NodeId) -> Result<Option<Vec<u8>>> {
198        Ok(None)
199    }
200
201    fn put(&mut self, _id: &NodeId, _bytes: &[u8]) -> Result<()> {
202        Ok(())
203    }
204}
205
206/// Execution accounting for one materialization (ADR-0027). Reuse is claimed by
207/// an *execution counter*, never by wall-clock: `nodes_reused > 0` and a smaller
208/// `nodes_executed` are the evidence that persisted work was served from disk.
209#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
210pub struct ReuseStats {
211    /// Nodes actually evaluated (cache misses).
212    pub nodes_executed: u64,
213    /// Nodes served whole from the cache (their subtrees were not traversed).
214    pub nodes_reused: u64,
215    /// Output bytes written to the cache during this materialization.
216    pub cache_bytes_written: u64,
217}
218
219/// Materialize one node's output bytes, recursively resolving dependencies.
220pub fn materialize_node(
221    parsed: &ParsedDescriptor,
222    store: &dyn SeedStore,
223    node: &SeedNode,
224    limits: Limits,
225    budget: &mut EvalBudget,
226    depth: u16,
227) -> Result<Vec<u8>> {
228    let mut cache = NoCache;
229    let mut reuse = ReuseStats::default();
230    materialize_inner(
231        parsed, store, &mut cache, node, limits, budget, depth, &mut reuse,
232    )
233}
234
235/// Materialize one node's output bytes against an explicit [`SourceServer`].
236///
237/// This is the partial-reader entry point: the caller supplies a source that can
238/// serve ranges without being handed a fully parsed descriptor.
239pub fn materialize_node_with(
240    source: &dyn SourceServer,
241    store: &dyn SeedStore,
242    node: &SeedNode,
243    limits: Limits,
244    budget: &mut EvalBudget,
245    depth: u16,
246) -> Result<Vec<u8>> {
247    let mut cache = NoCache;
248    let mut reuse = ReuseStats::default();
249    materialize_inner(
250        source, store, &mut cache, node, limits, budget, depth, &mut reuse,
251    )
252}
253
254/// Materialize one node's output, consulting `cache` at **every** node (including
255/// dependencies).
256///
257/// A hit returns the cached bytes *without recursing into the node's dependency
258/// closure* and increments [`ReuseStats::nodes_reused`]; a miss evaluates the node
259/// (recursing through the same cache) and stores the output. Exact `Q_ref` nodes
260/// are cached too, since their output is equally a pure function of their id.
261#[allow(clippy::too_many_arguments)]
262pub fn materialize_node_cached(
263    parsed: &ParsedDescriptor,
264    store: &dyn SeedStore,
265    cache: &mut dyn OutputCache,
266    node: &SeedNode,
267    limits: Limits,
268    budget: &mut EvalBudget,
269    depth: u16,
270    reuse: &mut ReuseStats,
271) -> Result<Vec<u8>> {
272    materialize_inner(parsed, store, cache, node, limits, budget, depth, reuse)
273}
274
275/// The [`SourceServer`] counterpart of [`materialize_node_cached`].
276#[allow(clippy::too_many_arguments)]
277pub fn materialize_node_cached_with(
278    source: &dyn SourceServer,
279    store: &dyn SeedStore,
280    cache: &mut dyn OutputCache,
281    node: &SeedNode,
282    limits: Limits,
283    budget: &mut EvalBudget,
284    depth: u16,
285    reuse: &mut ReuseStats,
286) -> Result<Vec<u8>> {
287    materialize_inner(source, store, cache, node, limits, budget, depth, reuse)
288}
289
290#[allow(clippy::too_many_arguments)]
291fn materialize_inner(
292    source: &dyn SourceServer,
293    store: &dyn SeedStore,
294    cache: &mut dyn OutputCache,
295    node: &SeedNode,
296    limits: Limits,
297    budget: &mut EvalBudget,
298    depth: u16,
299    reuse: &mut ReuseStats,
300) -> Result<Vec<u8>> {
301    if depth == 0 {
302        return Err(Error::resource_limit("seed DAG exceeded its depth bound"));
303    }
304
305    let id = node.content_id();
306    // A hit is the whole subtree: return it without traversing dependencies. A
307    // cache error is a miss (the cache is disposable, never authority). Oversized
308    // cached bytes are likewise treated as a poisoned miss, not returned.
309    if let Ok(Some(bytes)) = cache.get(&id)
310        && bytes.len() as u64 <= node.limits.max_output_bytes
311    {
312        reuse.nodes_reused = reuse.nodes_reused.saturating_add(1);
313        budget.charge_bytes(bytes.len() as u64)?;
314        return Ok(bytes);
315    }
316
317    budget.charge_node()?;
318    reuse.nodes_executed = reuse.nodes_executed.saturating_add(1);
319
320    let out = match node.kind {
321        NodeKind::DocumentExact => source.serve_document(limits)?,
322        NodeKind::SourceSlice | NodeKind::ResourceRef => {
323            let (offset, len) = read_span_params(&node.params)?;
324            source.serve_range(offset, len, limits)?
325        }
326        NodeKind::PdfRevision | NodeKind::PdfObject | NodeKind::PdfStreamEncoded => {
327            let (_number, _generation, extra) = read_object_params(&node.params)?;
328            // `extra` packs `offset` in the high 32 bits and `len` in the low 32.
329            let offset = extra >> 32;
330            let len = extra & 0xFFFF_FFFF;
331            source.serve_range(offset, len, limits)?
332        }
333        NodeKind::Concat | NodeKind::PageContent => {
334            let mut out = Vec::new();
335            for dep in &node.deps {
336                let child = load_node(store, dep)?;
337                let bytes = materialize_inner(
338                    source,
339                    store,
340                    cache,
341                    &child,
342                    limits,
343                    budget,
344                    depth - 1,
345                    reuse,
346                )?;
347                budget.charge_bytes(bytes.len() as u64)?;
348                out.extend_from_slice(&bytes);
349            }
350            out
351        }
352        NodeKind::Literal => node.params.clone(),
353        NodeKind::PdfStreamDecoded => {
354            let dep = node
355                .deps
356                .first()
357                .ok_or_else(|| Error::usage("PdfStreamDecoded has no dependency"))?;
358            let child = load_node(store, dep)?;
359            let encoded = materialize_inner(
360                source,
361                store,
362                cache,
363                &child,
364                limits,
365                budget,
366                depth - 1,
367                reuse,
368            )?;
369            derive::inflate_zlib(&encoded, node.logical_output_len, limits)?
370        }
371        NodeKind::ContentOperators => {
372            let dep = node
373                .deps
374                .first()
375                .ok_or_else(|| Error::usage("ContentOperators has no dependency"))?;
376            let child = load_node(store, dep)?;
377            let decoded = materialize_inner(
378                source,
379                store,
380                cache,
381                &child,
382                limits,
383                budget,
384                depth - 1,
385                reuse,
386            )?;
387            derive::content_operators(&decoded, limits)?
388        }
389        NodeKind::TextRuns => {
390            let dep = node
391                .deps
392                .first()
393                .ok_or_else(|| Error::usage("TextRuns has no dependency"))?;
394            let child = load_node(store, dep)?;
395            let ops = materialize_inner(
396                source,
397                store,
398                cache,
399                &child,
400                limits,
401                budget,
402                depth - 1,
403                reuse,
404            )?;
405            derive::text_runs(&ops, limits)?
406        }
407        NodeKind::PagePreview => {
408            let dep = node
409                .deps
410                .first()
411                .ok_or_else(|| Error::usage("PagePreview has no dependency"))?;
412            let child = load_node(store, dep)?;
413            let content = materialize_inner(
414                source,
415                store,
416                cache,
417                &child,
418                limits,
419                budget,
420                depth - 1,
421                reuse,
422            )?;
423            let page = read_u32_params(&node.params)?;
424            derive::page_preview(page, &content, limits)?
425        }
426    };
427
428    if out.len() as u64 > node.limits.max_output_bytes {
429        return Err(Error::resource_limit(format!(
430            "node {} produced {} bytes > its cap {}",
431            node.kind.name(),
432            out.len(),
433            node.limits.max_output_bytes
434        )));
435    }
436    budget.charge_bytes(out.len() as u64)?;
437    // Best-effort persistence: a cache write failure never fails the observation.
438    if cache.put(&id, &out).is_ok() {
439        reuse.cache_bytes_written = reuse.cache_bytes_written.saturating_add(out.len() as u64);
440    }
441    Ok(out)
442}
443
444#[cfg(test)]
445mod tests {
446    use super::*;
447    use crate::container::ObjectSource;
448    use crate::dra::{Op, Program};
449    use crate::integrity::sha256;
450    use crate::{EXACTNESS_PROFILE_EXACT_BYTES, SOURCE_FORMAT_OPAQUE};
451
452    fn parsed_for(source: &[u8]) -> ParsedDescriptor {
453        let d = crate::container::Descriptor {
454            universe: crate::container::UNIVERSE.to_string(),
455            source_format: SOURCE_FORMAT_OPAQUE,
456            format_basis: "opaque:test".to_string(),
457            models: vec![],
458            channels: vec![],
459            objects: vec![ObjectSource::Inline(source.to_vec())],
460            program: Program::new(vec![Op::EmitObject { object_id: 0 }]),
461            observation_index: None,
462            seek_directory: false,
463            source_sha256: sha256(source),
464            source_len: source.len() as u64,
465        };
466        let (bytes, _cost) = d.serialize().unwrap();
467        let _ = EXACTNESS_PROFILE_EXACT_BYTES;
468        crate::container::Descriptor::parse(&bytes, Limits::DEFAULT).unwrap()
469    }
470
471    #[test]
472    fn source_slice_matches_exact_bytes() {
473        let parsed = parsed_for(b"hello field world");
474        let mut budget = EvalBudget::default();
475        let (off, len) = (6u64, 5u64);
476        let node = SeedNode::new(
477            NodeKind::SourceSlice,
478            len,
479            crate::field::node::span_params(off, len),
480            vec![],
481            "test",
482        );
483        let bytes = materialize_node(
484            &parsed,
485            &crate::store::FsSeedStore::open(
486                std::env::temp_dir().join(format!("vole-dag-{}", std::process::id())),
487            )
488            .unwrap(),
489            &node,
490            Limits::DEFAULT,
491            &mut budget,
492            8,
493        )
494        .unwrap();
495        assert_eq!(bytes, b"field");
496    }
497
498    #[test]
499    fn literal_roundtrips() {
500        let parsed = parsed_for(b"x");
501        let mut budget = EvalBudget::default();
502        let node = SeedNode::new(NodeKind::Literal, 3, b"abc".to_vec(), vec![], "test");
503        let bytes = materialize_node(
504            &parsed,
505            &crate::store::FsSeedStore::open(
506                std::env::temp_dir().join(format!("vole-lit-{}", std::process::id())),
507            )
508            .unwrap(),
509            &node,
510            Limits::DEFAULT,
511            &mut budget,
512            4,
513        )
514        .unwrap();
515        assert_eq!(bytes, b"abc");
516    }
517}