Skip to main content

vole_document/field/
ingest.rs

1//! Progressive early inverse-proceduralization of a PDF field (Phase 11.4).
2//!
3//! [`ingest_pdf`] runs Stage A (durable exact capture, delegated to
4//! [`FieldStore::ingest`]) followed by Stage B (cheap eager inversion): exact
5//! physical spans become `PdfObject` / `PdfRevision` / `PdfStreamEncoded` nodes,
6//! a lone-`/FlateDecode` stream is inflated once to learn its exact decoded
7//! length and becomes a `PdfStreamDecoded` node (an unfiltered stream's encoded
8//! bytes are used directly as content), and a bounded page-tree walk builds
9//! `PageContent` nodes. Page keys (`/Type`, `/Kids`, `/Contents`, `/Pages`) are
10//! read only at the leading dictionary's top level, so a nested sub-dictionary
11//! cannot shadow them. Every recovered observation is registered in the
12//! hierarchical index and the manifest is re-written with the richer root.
13//!
14//! Stage B is **best-effort and never fatal**: a heuristic step that declines
15//! (an unparseable scan, a non-inflatable stream, an ambiguous page) only lowers
16//! the recovered depth; it never fails the ingest or weakens exactness. Stage A
17//! remains the sole archival authority, so `materialize_exact` is byte-identical
18//! regardless of what Stage B recovered.
19//!
20//! Stage C ([`deepen_page`]) is demand-driven: it adds `ContentOperators`,
21//! `TextRuns`, and `PagePreview` nodes beneath an existing `PageContent` node and
22//! writes a *new* manifest (`DEC-7`: promotion is additive). The older field id
23//! keeps materializing exactly, so promotion is monotonic.
24//!
25//! Nothing here serializes known structure and reparses it (plan §13): the page
26//! tree and content references are read directly from the Phase-3 lexical cover,
27//! and a decoded stream is never re-DEFLATEd to be re-inflated.
28
29use std::collections::{BTreeMap, BTreeSet};
30
31use crate::adapter::pdf::cos::FilterClass;
32use crate::adapter::pdf::lexer::lex;
33use crate::adapter::pdf::physical::{PdfPhysical, PdfStreamSpan, RevisionInfo, scan};
34use crate::adapter::pdf::span::{Span, SpanKind};
35use crate::container::{Descriptor, ParsedDescriptor};
36use crate::error::{Error, Result};
37use crate::field::dag;
38use crate::field::index::{
39    FsIndexStore, IndexEntry, SEL_OBJECT, SEL_PAGE, SEL_REVISION, SEL_REVISION_LINEAGE,
40    SEL_REVISIONS, SEL_STREAM, SEL_STREAM_DECODED, SelectorKey, build, lookup, validate,
41};
42use crate::field::manifest::FieldRoot;
43use crate::field::node::{MAX_NODE_DEPS, NodeKind, SeedNode, object_params, u32_params};
44use crate::field::{Field, FieldId, FieldStore};
45use crate::limits::Limits;
46use crate::parallel::WorkerPool;
47use crate::store::NodeId;
48#[cfg(feature = "parallel")]
49use rayon::prelude::*;
50
51/// Hard cap on seed nodes one ingest may add.
52pub const MAX_INGEST_NODES: u64 = 1 << 20;
53/// Hard cap on index entries one ingest may add.
54pub const MAX_INGEST_INDEX_ENTRIES: usize = 1 << 20;
55/// Maximum decoded length admitted for a single stream (32 MiB).
56pub const MAX_STREAM_DECODE: u64 = 32 * 1024 * 1024;
57/// Maximum total decoded bytes admitted across one ingest (512 MiB).
58pub const MAX_TOTAL_DECODED: u64 = 512 * 1024 * 1024;
59/// Maximum number of page-tree objects visited during the Stage-B walk.
60const MAX_PAGE_TREE_NODES: usize = 1 << 16;
61/// Maximum number of references gathered from one `/Kids` or `/Contents`.
62const MAX_REFS: usize = 1 << 16;
63/// Maximum number of object streams indexed in one page-recovery pass.
64const MAX_OBJSTM: usize = 1 << 12;
65/// Maximum `/N` (object count) admitted from one object stream.
66const MAX_OBJSTM_OBJECTS: usize = 1 << 16;
67/// Maximum total decoded object-stream bytes retained for page recovery (64 MiB).
68const MAX_OBJSTM_BYTES: u64 = 64 * 1024 * 1024;
69/// Maximum pages recovered in one page-tree walk.
70const MAX_PAGES: usize = 1 << 16;
71
72/// What one ingest recovered.
73#[derive(Debug, Clone, PartialEq, Eq)]
74pub struct IngestReport {
75    /// The detected document format (byte-based, Phase 12.7).
76    pub format: crate::field::document_format::DocumentFormat,
77    /// The (richest) field id, whose manifest binds the recovered index.
78    pub field: FieldId,
79    /// The exact `DocumentExact` root node id.
80    pub root_node: NodeId,
81    /// The hierarchical index root, or `None` when no index was built.
82    pub index_root: Option<NodeId>,
83    /// Seed nodes in the manifest closure at ingest time (root + Stage-B nodes).
84    pub node_count: u64,
85    /// Number of index nodes in the built tree.
86    pub index_node_count: u64,
87    /// Exact reconstructed source length.
88    pub source_len: u64,
89    /// Exact per-object nodes created.
90    pub object_nodes: u64,
91    /// Encoded stream nodes created.
92    pub stream_nodes: u64,
93    /// Decoded stream nodes created.
94    pub decoded_stream_nodes: u64,
95    /// Page content nodes created.
96    pub page_nodes: u64,
97    /// Revision nodes created.
98    pub revision_nodes: u64,
99    /// Lone-`FlateDecode` streams we refused to inflate.
100    pub declined_streams: u64,
101    /// Content-addressed shared-resource blobs registered. Always `0`: the
102    /// Phase-11 PDF adapter does not extract embedded images/fonts as resource
103    /// blobs, so a PDF shares no resource with any other document in Phase 12
104    /// (a recorded limitation, not a failure). Package formats (DOCX/EPUB) do.
105    pub resource_blob_nodes: u64,
106    /// Resource blobs shared with an earlier document. Always `0` for PDF (see
107    /// [`Self::resource_blob_nodes`]).
108    pub shared_resource_ids: u64,
109    /// Resource bytes not rewritten because they were already present. Always
110    /// `0` for PDF (see [`Self::resource_blob_nodes`]).
111    pub shared_resource_bytes: u64,
112    /// Seed nodes whose content id already existed (nothing new written).
113    pub nodes_id_shared: u64,
114    /// Seed-node canonical bytes physically written by this ingest.
115    pub seed_bytes_written: u64,
116}
117
118/// Stage A (durable exact capture) + Stage B (cheap eager inversion).
119pub fn ingest_pdf(
120    store: &mut FieldStore,
121    descriptor_bytes: &[u8],
122    limits: Limits,
123) -> Result<IngestReport> {
124    ingest_pdf_with(store, descriptor_bytes, limits, None)
125}
126
127/// As [`ingest_pdf`], but with an optional bounded worker pool for the pure,
128/// order-independent Stage-B work (per-stream inflate, `/ObjStm` decode). `None`
129/// is byte-for-byte the serial path; the merge fold always runs serially in
130/// physical order, so the pool can never change a node id, a counter, or a limit.
131pub fn ingest_pdf_with(
132    store: &mut FieldStore,
133    descriptor_bytes: &[u8],
134    limits: Limits,
135    pool: Option<&WorkerPool>,
136) -> Result<IngestReport> {
137    // Stage A: the exact descriptor blob, the DocumentExact root, and a manifest.
138    // The stored blob gains a minimal advisory observation-index op table when it
139    // lacks one, so a later narrow observation can use the seek-based partial
140    // lane. Exactness is unchanged; only the ignorable record is added.
141    let observable = with_observation_index(descriptor_bytes, limits)?;
142    let base_id = store.ingest(&observable, limits)?;
143    let (manifest, source) = {
144        let field = Field::open(store, &base_id, limits)?;
145        (field.manifest().clone(), field.materialize_exact(limits)?)
146    };
147    ingest_pdf_stage_b(store, &source, manifest, limits, pool)
148}
149
150/// Direct-build variant of [`ingest_pdf_with`] for the fixed-profile
151/// [`crate::field::build`] path (Phase 18.2).
152///
153/// The caller already holds the exact original `source` and the already-enriched
154/// `observable` authority blob, so neither is reconstructed here: Stage A stores
155/// the authority and verifies it materializes to `source`
156/// ([`FieldStore::ingest_verified`]), and Stage B scans the **original** bytes.
157/// The observations, node ids, index, and manifest are byte-identical to
158/// [`ingest_pdf_with`]; only the redundant source → authority → source round
159/// trip is removed.
160pub(crate) fn ingest_pdf_direct(
161    store: &mut FieldStore,
162    observable: &[u8],
163    source: &[u8],
164    limits: Limits,
165    pool: Option<&WorkerPool>,
166) -> Result<IngestReport> {
167    let base_id = store.ingest_verified(observable, source, limits)?;
168    let manifest = Field::open(store, &base_id, limits)?.manifest().clone();
169    ingest_pdf_stage_b(store, source, manifest, limits, pool)
170}
171
172/// Shared Stage-B/C tail: scan `source` (the exact materialization of the
173/// authority `manifest` already binds) and write the richer manifest.
174fn ingest_pdf_stage_b(
175    store: &mut FieldStore,
176    source: &[u8],
177    manifest: FieldRoot,
178    limits: Limits,
179    pool: Option<&WorkerPool>,
180) -> Result<IngestReport> {
181    let source_len = source.len() as u64;
182    // Byte-based format detection, recorded in the manifest provenance so the
183    // universal observation API can dispatch common selectors without re-reading
184    // the source (Phase 12.7). Never derived from a file name.
185    let fmt = crate::field::document_format::detect_document_format(source, limits);
186
187    let mut acc = StageB::new(manifest.node_count);
188    let scanned = match scan(source, limits) {
189        Ok(physical) => {
190            run_stage_b(store, source, &physical, limits, pool, &mut acc)?;
191            true
192        }
193        Err(_) => false,
194    };
195
196    let (index_root, index_node_count, provenance) = if scanned && !acc.entries.is_empty() {
197        let mut istore = FsIndexStore::open(store.root())?;
198        let root = build(&mut istore, &acc.entries)?;
199        let (count, _depth) = validate(&istore, &root)?;
200        let prov = format!(
201            "field:ingest-b;objects={};streams={};decoded={};pages={};revs={}",
202            acc.object_nodes,
203            acc.stream_nodes,
204            acc.decoded_stream_nodes,
205            acc.page_nodes,
206            acc.revision_nodes
207        );
208        (Some(root), count, prov)
209    } else if scanned {
210        (None, 0, "field:ingest-b;entries=0".to_string())
211    } else {
212        (None, 0, "field:ingest-b;declined=scan".to_string())
213    };
214    // Prefix the machine-readable format token (idempotent if already present),
215    // then the Phase-12.8 sharing counters (representation facts, read back by
216    // observations so `nodes_id_shared` is reported alongside `nodes_reused`).
217    let provenance = format!(
218        "{}id_shared={};res_shared={};{}",
219        fmt.provenance_prefix(),
220        acc.nodes_id_shared,
221        acc.shared_resource_ids,
222        provenance
223    );
224
225    let mut new_manifest = manifest.clone();
226    if let Some(root) = &index_root {
227        new_manifest.index_root = *root.as_bytes();
228    }
229    new_manifest.node_count = acc.node_count;
230    new_manifest.index_node_count = index_node_count;
231    new_manifest.provenance = provenance;
232    let field = store.put_field(&new_manifest)?;
233
234    Ok(IngestReport {
235        format: fmt,
236        field,
237        root_node: manifest.root_node,
238        index_root,
239        node_count: acc.node_count,
240        index_node_count,
241        source_len,
242        object_nodes: acc.object_nodes,
243        stream_nodes: acc.stream_nodes,
244        decoded_stream_nodes: acc.decoded_stream_nodes,
245        page_nodes: acc.page_nodes,
246        revision_nodes: acc.revision_nodes,
247        declined_streams: acc.declined_streams,
248        resource_blob_nodes: acc.resource_blob_nodes,
249        shared_resource_ids: acc.shared_resource_ids,
250        shared_resource_bytes: acc.shared_resource_bytes,
251        nodes_id_shared: acc.nodes_id_shared,
252        seed_bytes_written: acc.seed_bytes_written,
253    })
254}
255
256/// A universal ingest outcome: which native inverse compiler ran.
257#[cfg(feature = "package")]
258#[derive(Debug, Clone, PartialEq, Eq)]
259pub enum IngestOutcome {
260    /// A PDF (or opaque non-ZIP) field, produced by [`ingest_pdf`].
261    Pdf(IngestReport),
262    /// A ZIP-based package (DOCX/EPUB/generic ZIP), produced by
263    /// [`crate::field::ingest_package::ingest_package`].
264    Package(crate::field::ingest_package::PackageIngestReport),
265}
266
267/// Detect the source format **from bytes** and invert it with the right adapter
268/// (Phase 12.7). A validated ZIP is inverted through the byte-authoritative
269/// package layer; everything else (PDF and the opaque floor) goes through
270/// [`ingest_pdf`]. Never consults a file name.
271#[cfg(feature = "package")]
272pub fn ingest(
273    store: &mut FieldStore,
274    descriptor_bytes: &[u8],
275    limits: Limits,
276) -> Result<IngestOutcome> {
277    ingest_with(store, descriptor_bytes, limits, None)
278}
279
280/// As [`ingest`], but with an optional bounded worker pool threaded into whichever
281/// adapter runs. `None` is byte-for-byte the serial path.
282#[cfg(feature = "package")]
283pub fn ingest_with(
284    store: &mut FieldStore,
285    descriptor_bytes: &[u8],
286    limits: Limits,
287    pool: Option<&WorkerPool>,
288) -> Result<IngestOutcome> {
289    let parsed = crate::container::Descriptor::parse(descriptor_bytes, limits)?;
290    let source = crate::materialize::materialize(&parsed, limits)?;
291    if crate::field::document_format::is_zip(&source, limits) {
292        Ok(IngestOutcome::Package(
293            crate::field::ingest_package::ingest_package_with(
294                store,
295                descriptor_bytes,
296                limits,
297                pool,
298            )?,
299        ))
300    } else {
301        Ok(IngestOutcome::Pdf(ingest_pdf_with(
302            store,
303            descriptor_bytes,
304            limits,
305            pool,
306        )?))
307    }
308}
309
310/// Add a minimal observation-index op table to a descriptor blob that lacks one,
311/// so the seek-based partial lane is available to narrow observations.
312///
313/// This never changes the reconstruction program, objects, channels, or declared
314/// source: it adds only the ignorable `OBSERVATION_INDEX` record whose op table is
315/// derived from [`Program::analyze_ops`] and re-validated by
316/// [`Descriptor::parse`] when the enriched blob is stored. When the descriptor
317/// already carries an index, or an op length does not fit the index's `u32`
318/// field, the input is returned unchanged (the observation then uses the full
319/// descriptor path).
320pub(crate) fn with_observation_index(bytes: &[u8], limits: Limits) -> Result<Vec<u8>> {
321    let parsed: ParsedDescriptor = Descriptor::parse(bytes, limits)?;
322    if parsed.descriptor.observation_index.is_some() {
323        return Ok(bytes.to_vec());
324    }
325    // The derivation is now a pure method on the descriptor (Phase 18.3), so the
326    // direct build can attach the same record before its single serialize pass.
327    // The fallbacks are preserved exactly: any decline (a limit breach, an op
328    // length that does not fit `u32`, a refused serialize) returns the input.
329    let enriched = parsed.descriptor.with_observation_index(limits);
330    if enriched.observation_index.is_none() {
331        return Ok(bytes.to_vec());
332    }
333    match enriched.serialize() {
334        Ok((out, _cost)) => Ok(out),
335        Err(_) => Ok(bytes.to_vec()),
336    }
337}
338
339/// Stage C: add decoded-stream operators/text for one page, returning a new
340/// field id. Idempotent: repeating on an already-promoted field returns it
341/// unchanged, and the original field id keeps materializing exactly.
342pub fn deepen_page(
343    store: &mut FieldStore,
344    field: &FieldId,
345    page: u32,
346    limits: Limits,
347) -> Result<FieldId> {
348    let manifest = Field::open(store, field, limits)?.manifest().clone();
349    deepen_page_with_manifest(store, &manifest, page)
350}
351
352/// Stage C against an already-loaded manifest, so an observation that already
353/// holds the parsed field does not re-read the descriptor blob just to promote a
354/// page (review fix #2: no hidden double descriptor load).
355pub fn deepen_page_with_manifest(
356    store: &mut FieldStore,
357    manifest: &FieldRoot,
358    page: u32,
359) -> Result<FieldId> {
360    let field = manifest.content_id();
361    if !manifest.has_index() {
362        return Ok(field);
363    }
364    // Already promoted for this page: idempotent no-op. The format token is
365    // preserved so a promoted field still serves common observations.
366    let format_token = manifest
367        .provenance
368        .split(';')
369        .next()
370        .filter(|t| t.starts_with("format="))
371        .map_or(String::new(), |t| format!("{t};"));
372    let target = format!("{format_token}field:deepen;page={page}");
373    if manifest.provenance == target {
374        return Ok(field);
375    }
376
377    let istore = FsIndexStore::open(store.root())?;
378    let root = NodeId::from_bytes(manifest.index_root);
379    let found = lookup(&istore, &root, &SelectorKey::new(SEL_PAGE, page))?;
380    let Some(entry) = found.first() else {
381        return Ok(field);
382    };
383    let page_content_id = entry.node_id;
384    let page_content = dag::load_node(store.seeds(), &page_content_id)?;
385    if page_content.kind != NodeKind::PageContent {
386        return Ok(field);
387    }
388
389    // Add only nodes: the page's SEL_PAGE entry keeps its exact semantics.
390    let (ops, text, preview) = derived_chain(page, page_content_id);
391    store.seeds_mut().put_node(&ops.encode_canonical())?;
392    store.seeds_mut().put_node(&text.encode_canonical())?;
393    store.seeds_mut().put_node(&preview.encode_canonical())?;
394
395    let mut new_manifest = manifest.clone();
396    new_manifest.node_count = manifest.node_count.saturating_add(3);
397    new_manifest.provenance = format!("{format_token}field:deepen;page={page}");
398    store.put_field(&new_manifest)
399}
400
401/// The deterministic Stage-C chain beneath a page's `PageContent` node.
402fn derived_chain(page: u32, page_content_id: NodeId) -> (SeedNode, SeedNode, SeedNode) {
403    let ops = SeedNode::new(
404        NodeKind::ContentOperators,
405        0,
406        u32_params(page),
407        vec![page_content_id],
408        "pdf:content-operators",
409    );
410    let text = SeedNode::new(
411        NodeKind::TextRuns,
412        0,
413        Vec::new(),
414        vec![ops.content_id()],
415        "pdf:text-runs",
416    );
417    let preview = SeedNode::new(
418        NodeKind::PagePreview,
419        0,
420        u32_params(page),
421        vec![page_content_id],
422        "pdf:page-preview",
423    );
424    (ops, text, preview)
425}
426
427/// The deterministic `PdfStreamDecoded` node for an encoded stream node.
428fn stream_decoded_node(object: u32, encoded_id: NodeId, decoded_len: u64) -> SeedNode {
429    SeedNode::new(
430        NodeKind::PdfStreamDecoded,
431        decoded_len,
432        u32_params(object),
433        vec![encoded_id],
434        "pdf:stream-decoded",
435    )
436}
437
438/// Mutable Stage-B accumulator.
439struct StageB {
440    entries: Vec<IndexEntry>,
441    decoded_by_object: BTreeMap<u32, (NodeId, u64)>,
442    /// Encoded node for streams with *no* `/Filter`, whose bytes are the raw
443    /// content (used to recover pages with unfiltered content streams).
444    plain_by_object: BTreeMap<u32, (NodeId, u64)>,
445    node_count: u64,
446    object_nodes: u64,
447    stream_nodes: u64,
448    decoded_stream_nodes: u64,
449    page_nodes: u64,
450    revision_nodes: u64,
451    declined_streams: u64,
452    total_decoded: u64,
453    /// Phase 12.8 cross-document sharing counters.
454    resource_blob_nodes: u64,
455    shared_resource_ids: u64,
456    shared_resource_bytes: u64,
457    nodes_id_shared: u64,
458    seed_bytes_written: u64,
459}
460
461impl StageB {
462    fn new(node_count: u64) -> Self {
463        StageB {
464            entries: Vec::new(),
465            decoded_by_object: BTreeMap::new(),
466            plain_by_object: BTreeMap::new(),
467            node_count,
468            object_nodes: 0,
469            stream_nodes: 0,
470            decoded_stream_nodes: 0,
471            page_nodes: 0,
472            revision_nodes: 0,
473            declined_streams: 0,
474            total_decoded: 0,
475            resource_blob_nodes: 0,
476            shared_resource_ids: 0,
477            shared_resource_bytes: 0,
478            nodes_id_shared: 0,
479            seed_bytes_written: 0,
480        }
481    }
482
483    /// Content-addressed `put_node` that records id-shared and newly-written
484    /// bytes (Phase 12.8). The `put_node` call is idempotent, so an id that
485    /// already existed writes nothing.
486    fn put(&mut self, store: &mut FieldStore, node: &SeedNode) -> Result<NodeId> {
487        if self.node_count >= MAX_INGEST_NODES {
488            return Err(Error::resource_limit(format!(
489                "ingest would exceed {MAX_INGEST_NODES} seed nodes"
490            )));
491        }
492        let id = node.content_id();
493        let preexisting = store.seeds().contains_node(&id)?;
494        let canonical = node.encode_canonical();
495        store.seeds_mut().put_node(&canonical)?;
496        self.node_count += 1;
497        if preexisting {
498            self.nodes_id_shared += 1;
499        } else {
500            self.seed_bytes_written = self
501                .seed_bytes_written
502                .saturating_add(canonical.len() as u64);
503        }
504        Ok(id)
505    }
506
507    fn add_entry(&mut self, entry: IndexEntry) -> Result<()> {
508        if self.entries.len() >= MAX_INGEST_INDEX_ENTRIES {
509            return Err(Error::resource_limit(format!(
510                "ingest would exceed {MAX_INGEST_INDEX_ENTRIES} index entries"
511            )));
512        }
513        self.entries.push(entry);
514        Ok(())
515    }
516}
517
518/// Stage B: exact spans, decoded streams, and the bounded page tree.
519fn run_stage_b(
520    store: &mut FieldStore,
521    source: &[u8],
522    physical: &PdfPhysical,
523    limits: Limits,
524    pool: Option<&WorkerPool>,
525    acc: &mut StageB,
526) -> Result<()> {
527    // Exact indirect-object spans.
528    for obj in &physical.objects {
529        if obj.number == 0 {
530            continue;
531        }
532        let (Some((number, generation)), Some((offset, len))) = (
533            identity32(obj.number, obj.generation),
534            span32(obj.start, obj.end.saturating_sub(obj.start)),
535        ) else {
536            continue;
537        };
538        let node = SeedNode::new(
539            NodeKind::PdfObject,
540            len,
541            object_params(number, generation, (offset << 32) | len),
542            Vec::new(),
543            "pdf:object",
544        );
545        let id = acc.put(store, &node)?;
546        acc.object_nodes += 1;
547        acc.add_entry(IndexEntry {
548            key: SelectorKey::new(SEL_OBJECT, number),
549            out_off: offset,
550            out_len: len,
551            node_id: id,
552        })?;
553    }
554
555    // Exact revision spans.
556    for rev in &physical.revisions {
557        let Some((offset, len)) = span32(rev.start, rev.end.saturating_sub(rev.start)) else {
558            continue;
559        };
560        let node = SeedNode::new(
561            NodeKind::PdfRevision,
562            len,
563            object_params(rev.index, 0, (offset << 32) | len),
564            Vec::new(),
565            "pdf:revision",
566        );
567        let id = acc.put(store, &node)?;
568        acc.revision_nodes += 1;
569        acc.add_entry(IndexEntry {
570            key: SelectorKey::new(SEL_REVISION, rev.index),
571            out_off: offset,
572            out_len: len,
573            node_id: id,
574        })?;
575    }
576
577    // Revision lineage (Phase 17): a compact derived projection of the physical
578    // revision chain, computed once here (the scan is already in hand) so a
579    // `Revisions`/`Revision(n) + Lineage` observation resolves through the index
580    // in O(depth) instead of re-scanning the source. One document-level node plus
581    // one node per revision; the per-revision entry carries the revision's exact
582    // source span so the scoped answer reports it without a second lookup.
583    if !physical.revisions.is_empty() {
584        let full = pdf_revision_lineage_json(source, physical);
585        let node = SeedNode::new(
586            NodeKind::PdfRevisionLineage,
587            full.len() as u64,
588            full.into_bytes(),
589            Vec::new(),
590            "pdf:revision-lineage",
591        );
592        let id = acc.put(store, &node)?;
593        acc.add_entry(IndexEntry {
594            key: SelectorKey::new(SEL_REVISIONS, 0),
595            out_off: 0,
596            out_len: 0,
597            node_id: id,
598        })?;
599        for rev in &physical.revisions {
600            let Some((offset, len)) = span32(rev.start, rev.end.saturating_sub(rev.start)) else {
601                continue;
602            };
603            let json = pdf_revision_json(rev, physical);
604            let node = SeedNode::new(
605                NodeKind::PdfRevisionLineage,
606                json.len() as u64,
607                json.into_bytes(),
608                Vec::new(),
609                "pdf:revision-lineage",
610            );
611            let id = acc.put(store, &node)?;
612            acc.add_entry(IndexEntry {
613                key: SelectorKey::new(SEL_REVISION_LINEAGE, rev.index),
614                out_off: offset,
615                out_len: len,
616                node_id: id,
617            })?;
618        }
619    }
620
621    // Encoded stream spans, plus one eager decode for a lone Flate stream. The
622    // decoded length of every candidate is learned up front (pure, order-free);
623    // the running `MAX_TOTAL_DECODED` gate below still runs serially in stream
624    // order, so the pool can never change which streams are admitted.
625    let lengths = decode_lengths(pool, source, physical, limits);
626    for (i, stream) in physical.streams.iter().enumerate() {
627        let (Some((number, generation)), Some((offset, len))) = (
628            identity32(stream.object, stream.generation),
629            span32(stream.data_start, stream.data_len),
630        ) else {
631            continue;
632        };
633        let encoded = SeedNode::new(
634            NodeKind::PdfStreamEncoded,
635            len,
636            object_params(number, generation, (offset << 32) | len),
637            Vec::new(),
638            "pdf:stream-encoded",
639        );
640        let encoded_id = acc.put(store, &encoded)?;
641        acc.stream_nodes += 1;
642        acc.add_entry(IndexEntry {
643            key: SelectorKey::new(SEL_STREAM, number),
644            out_off: offset,
645            out_len: len,
646            node_id: encoded_id,
647        })?;
648
649        // An unfiltered stream's encoded bytes *are* its content: remember the
650        // encoded node so a page whose `/Contents` is unfiltered can use it.
651        if stream.filter == FilterClass::Absent {
652            acc.plain_by_object.insert(number, (encoded_id, len));
653        }
654        if stream.filter != FilterClass::FlateDecode {
655            continue;
656        }
657        match lengths[i] {
658            Some(decoded_len)
659                if acc
660                    .total_decoded
661                    .checked_add(decoded_len)
662                    .is_some_and(|total| total <= MAX_TOTAL_DECODED) =>
663            {
664                let decoded = stream_decoded_node(number, encoded_id, decoded_len);
665                let decoded_id = acc.put(store, &decoded)?;
666                acc.decoded_stream_nodes += 1;
667                acc.total_decoded += decoded_len;
668                acc.decoded_by_object
669                    .insert(number, (decoded_id, decoded_len));
670                // Register the decoded node so a later `Stream(n) + DecodedBytes`
671                // observation resolves in O(depth) index reads, not by scanning
672                // the seed store (review fix #3). The span is the encoded
673                // stream's exact source span the node is derived from.
674                acc.add_entry(IndexEntry {
675                    key: SelectorKey::new(SEL_STREAM_DECODED, number),
676                    out_off: offset,
677                    out_len: len,
678                    node_id: decoded_id,
679                })?;
680            }
681            _ => acc.declined_streams += 1,
682        }
683    }
684
685    recover_pages(store, source, physical, limits, pool, acc)
686}
687
688/// The exact decoded length of every stream, in physical order.
689///
690/// Only a lone `/FlateDecode` stream is inflated, and only to learn its length;
691/// the decoded bytes are discarded. This is a pure function of the immutable
692/// inputs, so it may run on a worker pool. Indexed `par_iter().collect()`
693/// preserves order, so `lengths[i]` is stream `i`'s length whether the pool is
694/// used or not.
695///
696/// The running `MAX_TOTAL_DECODED` gate is deliberately **not** applied here: it
697/// is replayed serially by [`run_stage_b`], in stream order.
698fn decode_lengths(
699    pool: Option<&WorkerPool>,
700    source: &[u8],
701    physical: &PdfPhysical,
702    limits: Limits,
703) -> Vec<Option<u64>> {
704    let compute = |s: &PdfStreamSpan| {
705        if s.filter == FilterClass::FlateDecode {
706            try_decode_len(source, s.data_start, s.data_len, limits)
707        } else {
708            None
709        }
710    };
711    #[cfg(feature = "parallel")]
712    if let Some(p) = pool.filter(|p| p.workers() > 1) {
713        return p.install(|| physical.streams.par_iter().map(compute).collect());
714    }
715    #[cfg(not(feature = "parallel"))]
716    let _ = pool;
717    physical.streams.iter().map(compute).collect()
718}
719
720/// Best-effort page-tree recovery: `/Root` catalog → `/Pages` → `/Kids` → `/Page`.
721///
722/// Object bodies are read from the exact physical source *or* from a decoded
723/// `/ObjStm` buffer, so a producer that keeps its whole page tree inside an
724/// object stream (pdfTeX) still recovers its pages. Physical objects shadow
725/// object-stream objects of the same number, and every key read is
726/// top-level-dictionary-depth-1 correct so a nested `/Type` cannot shadow one.
727fn recover_pages(
728    store: &mut FieldStore,
729    source: &[u8],
730    physical: &PdfPhysical,
731    limits: Limits,
732    pool: Option<&WorkerPool>,
733    acc: &mut StageB,
734) -> Result<()> {
735    if physical.objects.is_empty() {
736        return Ok(());
737    }
738    let lexed = match lex(source, limits) {
739        Ok(lexed) => lexed,
740        Err(_) => return Ok(()),
741    };
742    let spans = lexed.spans.spans;
743
744    // Object number -> object index (a later revision shadows an earlier one).
745    let mut obj_index: BTreeMap<u64, usize> = BTreeMap::new();
746    for (i, obj) in physical.objects.iter().enumerate() {
747        obj_index.insert(obj.number, i);
748    }
749
750    // The leading dict/array range of every physical object.
751    let mut containers: Vec<Option<(bool, u64, u64)>> = Vec::with_capacity(physical.objects.len());
752    for obj in &physical.objects {
753        containers.push(leading_container(span_window(&spans, obj.start, obj.end)));
754    }
755
756    // Streams already decoded by Stage B (needed to read an `/ObjStm`'s bytes).
757    // Taken out of `acc` so the walk below can mutate it without a borrow clash.
758    let decoded_by_object = std::mem::take(&mut acc.decoded_by_object);
759    let plain_by_object = std::mem::take(&mut acc.plain_by_object);
760
761    // Index object streams: map each contained object number to its body range
762    // inside the retained decoded buffer. A stream with no decoded node is
763    // skipped outright -- the bytes are never guessed at.
764    //
765    // Candidate selection is pure and cheap (no inflate); the expensive pure work
766    // (inflate, header parse, and the per-buffer lex) runs on the pool. The result
767    // is index-aligned with `physical.streams`, so the serial fold below keeps the
768    // `MAX_OBJSTM`/`MAX_OBJSTM_BYTES` caps and the `buffers`/`objstm` insertion
769    // order identical to the serial path.
770    let mut inputs: Vec<Option<ObjStmInput>> = Vec::with_capacity(physical.streams.len());
771    for stream in &physical.streams {
772        inputs.push(objstm_input(
773            source,
774            stream,
775            &decoded_by_object,
776            &obj_index,
777            &containers,
778            &spans,
779        ));
780    }
781    let precomputed = precompute_objstm(pool, source, &inputs, limits);
782
783    let mut buffers: Vec<ObjStmBuf> = Vec::new();
784    let mut objstm: BTreeMap<u64, Resolved> = BTreeMap::new();
785    let mut total_objstm: u64 = 0;
786    for pre in precomputed.into_iter().flatten() {
787        if buffers.len() >= MAX_OBJSTM {
788            break;
789        }
790        let decoded = pre.decoded;
791        let Some(total) = total_objstm.checked_add(decoded.len() as u64) else {
792            continue;
793        };
794        if total > MAX_OBJSTM_BYTES {
795            continue;
796        }
797        let len = decoded.len() as u64;
798        // Per spec the pair offsets are relative to `/First` (the header end);
799        // when `/First` is absent the offsets are treated as absolute.
800        let base = pre.first.unwrap_or(0);
801        let buf_index = buffers.len();
802        let mut entries: Vec<(u64, Resolved)> = Vec::with_capacity(pre.pairs.len());
803        for (i, &(object, offset)) in pre.pairs.iter().enumerate() {
804            if object == 0 {
805                continue;
806            }
807            let Some(body_lo) = base.checked_add(offset) else {
808                continue;
809            };
810            let body_hi = pre
811                .pairs
812                .get(i + 1)
813                .and_then(|&(_, next)| base.checked_add(next))
814                .filter(|&next| next >= body_lo && next <= len)
815                .unwrap_or(len);
816            if body_lo > body_hi || body_hi > len {
817                continue;
818            }
819            let Some((body_is_array, clo, chi)) =
820                leading_container(span_window(&pre.spans, body_lo, body_hi))
821            else {
822                continue;
823            };
824            if clo < body_lo || chi > body_hi {
825                continue;
826            }
827            entries.push((
828                object,
829                Resolved {
830                    src: Src::ObjStm(buf_index),
831                    is_array: body_is_array,
832                    lo: clo,
833                    hi: chi,
834                },
835            ));
836        }
837        buffers.push(ObjStmBuf {
838            bytes: decoded,
839            spans: pre.spans,
840        });
841        total_objstm = total;
842        for (object, resolved) in entries {
843            objstm.insert(object, resolved);
844        }
845    }
846
847    // Unified resolver: physical objects shadow object-stream objects.
848    let physical_numbers: BTreeSet<u64> = physical
849        .objects
850        .iter()
851        .map(|o| o.number)
852        .filter(|&n| n != 0)
853        .collect();
854    let mut map: BTreeMap<u64, Resolved> = BTreeMap::new();
855    for (i, obj) in physical.objects.iter().enumerate() {
856        if obj.number == 0 {
857            continue;
858        }
859        if let Some((is_array, lo, hi)) = containers[i] {
860            map.insert(
861                obj.number,
862                Resolved {
863                    src: Src::Physical,
864                    is_array,
865                    lo,
866                    hi,
867                },
868            );
869        }
870    }
871    for (object, resolved) in objstm {
872        if !physical_numbers.contains(&object) {
873            map.insert(object, resolved);
874        }
875    }
876    let resolver = ObjResolver {
877        source,
878        source_spans: &spans,
879        buffers: &buffers,
880        map,
881    };
882
883    // The catalog is the lowest-numbered object whose *top-level* `/Type` is
884    // `/Catalog`, whether it lives in the physical source or an object stream.
885    let Some(catalog) = resolver
886        .map
887        .values()
888        .find(|&&res| resolver.is_type(res, b"Catalog"))
889        .copied()
890    else {
891        return Ok(());
892    };
893    let Some(pages_refs) = resolver.key_refs(catalog, b"Pages") else {
894        return Ok(());
895    };
896    let Some(&pages_root) = pages_refs.first() else {
897        return Ok(());
898    };
899
900    // Bounded depth-first walk in page order.
901    let mut stack = vec![pages_root];
902    let mut visited: BTreeSet<u64> = BTreeSet::new();
903    let mut pages: Vec<u64> = Vec::new();
904    while let Some(number) = stack.pop() {
905        if !visited.insert(number) {
906            continue;
907        }
908        if visited.len() > MAX_PAGE_TREE_NODES {
909            break;
910        }
911        let Some(res) = resolver.resolve(number) else {
912            continue;
913        };
914        if res.is_array {
915            continue;
916        }
917        if resolver.is_type(res, b"Page") {
918            if pages.len() < MAX_PAGES {
919                pages.push(number);
920            }
921        } else if let Some(kids) = resolver.key_refs(res, b"Kids") {
922            // A `/Pages` node, or a dict with no usable `/Type` that still
923            // carries `/Kids`: both are internal page-tree nodes. Best-effort.
924            let kids = resolver.expand(kids);
925            for kid in kids.into_iter().rev() {
926                if !visited.contains(&kid) {
927                    stack.push(kid);
928                }
929            }
930        }
931    }
932
933    let mut page_number: u32 = 0;
934    for page_obj in pages {
935        let Some(res) = resolver.resolve(page_obj) else {
936            continue;
937        };
938        let mut deps: Vec<NodeId> = Vec::new();
939        let mut total: u64 = 0;
940        let mut min_start: Option<u64> = None;
941        let mut max_end: u64 = 0;
942        if let Some(content_refs) = resolver.key_refs(res, b"Contents") {
943            let content_refs = resolver.expand(content_refs);
944            if content_refs.len() <= MAX_NODE_DEPS {
945                for content in &content_refs {
946                    let Ok(content_number) = u32::try_from(*content) else {
947                        continue;
948                    };
949                    // `/Contents` streams are physical (a stream cannot live in
950                    // an `/ObjStm`). Prefer the inflated node; otherwise the
951                    // encoded node of an unfiltered stream, whose bytes are the
952                    // content.
953                    let resolved = decoded_by_object
954                        .get(&content_number)
955                        .copied()
956                        .or_else(|| plain_by_object.get(&content_number).copied());
957                    let Some((node_id, content_len)) = resolved else {
958                        continue;
959                    };
960                    deps.push(node_id);
961                    total = total
962                        .checked_add(content_len)
963                        .ok_or_else(|| Error::resource_limit("page content length overflow"))?;
964                    if let Some(&ci) = obj_index.get(content) {
965                        let obj = &physical.objects[ci];
966                        min_start = Some(min_start.map_or(obj.start, |m| m.min(obj.start)));
967                        max_end = max_end.max(obj.end);
968                    }
969                }
970            }
971        }
972        // A page with no resolvable content is still a page: record an empty
973        // `PageContent` so `SEL_PAGE` numbering stays contiguous and meaningful.
974        page_number = page_number
975            .checked_add(1)
976            .ok_or_else(|| Error::resource_limit("page number overflow"))?;
977        let node = SeedNode::new(
978            NodeKind::PageContent,
979            total,
980            u32_params(page_number),
981            deps,
982            "pdf:page-content",
983        );
984        let node_id = acc.put(store, &node)?;
985        let (out_off, out_len) = match min_start {
986            Some(start) => (start, max_end.saturating_sub(start)),
987            None => match obj_index.get(&page_obj) {
988                Some(&pi) => {
989                    let obj = &physical.objects[pi];
990                    (obj.start, obj.end.saturating_sub(obj.start))
991                }
992                None => (0, 0),
993            },
994        };
995        acc.add_entry(IndexEntry {
996            key: SelectorKey::new(SEL_PAGE, page_number),
997            out_off,
998            out_len,
999            node_id,
1000        })?;
1001        acc.page_nodes += 1;
1002    }
1003
1004    Ok(())
1005}
1006
1007/// The cheap, store-free inputs needed to inflate and parse one candidate
1008/// `/ObjStm`, gathered during serial candidate selection.
1009struct ObjStmInput {
1010    /// Payload byte range in the source.
1011    start: usize,
1012    end: usize,
1013    /// The exact decoded length already learned by Stage B.
1014    decoded_len: u64,
1015    /// The object-stream header fields read from the source dictionary.
1016    first: Option<u64>,
1017    n: u64,
1018}
1019
1020/// The expensive, pure result for one `/ObjStm`: its inflated bytes, their
1021/// lexical cover, the parsed header pairs, and `/First`.
1022struct ObjStmPre {
1023    decoded: Vec<u8>,
1024    spans: Vec<Span>,
1025    pairs: Vec<(u64, u64)>,
1026    first: Option<u64>,
1027}
1028
1029/// Test whether `stream` is a candidate `/ObjStm` and, if so, gather the inputs
1030/// needed to inflate and parse it. Pure: no store, no counters.
1031///
1032/// The `/N` bound is checked here (before the inflate) rather than after it, as
1033/// the serial path did; that changes only how much work a rejected stream costs,
1034/// never the outcome, since a rejected stream is skipped either way.
1035fn objstm_input(
1036    source: &[u8],
1037    stream: &PdfStreamSpan,
1038    decoded_by_object: &BTreeMap<u32, (NodeId, u64)>,
1039    obj_index: &BTreeMap<u64, usize>,
1040    containers: &[Option<(bool, u64, u64)>],
1041    spans: &[Span],
1042) -> Option<ObjStmInput> {
1043    let number32 = u32::try_from(stream.object).ok()?;
1044    let &(_node, decoded_len) = decoded_by_object.get(&number32)?;
1045    let &idx = obj_index.get(&stream.object)?;
1046    let (is_array, lo, hi) = containers[idx]?;
1047    if is_array {
1048        return None;
1049    }
1050    let win = span_window(spans, lo, hi);
1051    if top_level_name_value(source, win, lo, hi, b"Type") != Some(&b"ObjStm"[..]) {
1052        return None;
1053    }
1054    let n = top_level_integer_value(source, win, lo, hi, b"N")?;
1055    if n == 0 || n > MAX_OBJSTM_OBJECTS as u64 {
1056        return None;
1057    }
1058    let start = usize::try_from(stream.data_start).ok()?;
1059    let end = stream
1060        .data_start
1061        .checked_add(stream.data_len)
1062        .and_then(|e| usize::try_from(e).ok())?;
1063    let first = top_level_integer_value(source, win, lo, hi, b"First");
1064    Some(ObjStmInput {
1065        start,
1066        end,
1067        decoded_len,
1068        first,
1069        n,
1070    })
1071}
1072
1073/// Inflate and lex every candidate `/ObjStm`, in parallel but index-aligned with
1074/// `inputs` (indexed `par_iter().collect()` preserves order). Pure: the inflated
1075/// bytes are owned and no store, counter, or cache is touched. The running
1076/// `MAX_OBJSTM*` caps are **not** applied here; they are replayed serially.
1077fn precompute_objstm(
1078    pool: Option<&WorkerPool>,
1079    source: &[u8],
1080    inputs: &[Option<ObjStmInput>],
1081    limits: Limits,
1082) -> Vec<Option<ObjStmPre>> {
1083    let compute = |input: &Option<ObjStmInput>| -> Option<ObjStmPre> {
1084        let input = input.as_ref()?;
1085        let encoded = source.get(input.start..input.end)?;
1086        let inflated = crate::field::derive::inflate_zlib(encoded, input.decoded_len, limits);
1087        let decoded = inflated.ok()?;
1088        if decoded.is_empty() {
1089            return None;
1090        }
1091        let pairs = parse_objstm_header(&decoded, input.first, input.n as usize)?;
1092        let lexed = lex(&decoded, limits).ok()?;
1093        Some(ObjStmPre {
1094            decoded,
1095            spans: lexed.spans.spans,
1096            pairs,
1097            first: input.first,
1098        })
1099    };
1100    #[cfg(feature = "parallel")]
1101    if let Some(p) = pool.filter(|p| p.workers() > 1) {
1102        return p.install(|| inputs.par_iter().map(compute).collect());
1103    }
1104    #[cfg(not(feature = "parallel"))]
1105    let _ = pool;
1106    inputs.iter().map(compute).collect()
1107}
1108
1109/// One decoded object stream retained for byte-level page recovery.
1110struct ObjStmBuf {
1111    bytes: Vec<u8>,
1112    spans: Vec<Span>,
1113}
1114
1115/// Where a resolved object's body bytes live.
1116#[derive(Clone, Copy, PartialEq, Eq)]
1117enum Src {
1118    /// The exact physical source (`N 0 obj ... endobj`).
1119    Physical,
1120    /// A decoded `/ObjStm` buffer, by index into the retained buffers.
1121    ObjStm(usize),
1122}
1123
1124/// A resolved object body: its source, leading-container kind, and byte range.
1125#[derive(Clone, Copy)]
1126struct Resolved {
1127    src: Src,
1128    is_array: bool,
1129    lo: u64,
1130    hi: u64,
1131}
1132
1133/// A read-only resolver over physical and object-stream object bodies.
1134struct ObjResolver<'a> {
1135    source: &'a [u8],
1136    source_spans: &'a [Span],
1137    buffers: &'a [ObjStmBuf],
1138    map: BTreeMap<u64, Resolved>,
1139}
1140
1141impl<'a> ObjResolver<'a> {
1142    /// The byte slice and lexical cover backing a resolved object body.
1143    fn view(&self, res: Resolved) -> (&'a [u8], &'a [Span]) {
1144        match res.src {
1145            Src::Physical => (self.source, self.source_spans),
1146            Src::ObjStm(i) => {
1147                let buf = &self.buffers[i];
1148                (&buf.bytes, &buf.spans)
1149            }
1150        }
1151    }
1152
1153    /// Resolve an object number to its body, if it has a leading container.
1154    fn resolve(&self, number: u64) -> Option<Resolved> {
1155        self.map.get(&number).copied()
1156    }
1157
1158    /// Whether the object's *top-level* `/Type` is `want`.
1159    fn is_type(&self, res: Resolved, want: &[u8]) -> bool {
1160        if res.is_array {
1161            return false;
1162        }
1163        let (bytes, spans) = self.view(res);
1164        let win = span_window(spans, res.lo, res.hi);
1165        top_level_name_value(bytes, win, res.lo, res.hi, b"Type") == Some(want)
1166    }
1167
1168    /// The indirect references following a top-level `key` in the object dict.
1169    fn key_refs(&self, res: Resolved, key: &[u8]) -> Option<Vec<u64>> {
1170        if res.is_array {
1171            return None;
1172        }
1173        let (bytes, spans) = self.view(res);
1174        let win = span_window(spans, res.lo, res.hi);
1175        collect_key_refs(bytes, win, res.lo, res.hi, key)
1176    }
1177
1178    /// Expand references that point at an array object into the refs inside it
1179    /// (one level, best-effort), across both sources.
1180    fn expand(&self, refs: Vec<u64>) -> Vec<u64> {
1181        let mut out: Vec<u64> = Vec::new();
1182        for reference in refs {
1183            if out.len() > MAX_REFS {
1184                break;
1185            }
1186            let expanded = self.resolve(reference).and_then(|res| {
1187                if !res.is_array {
1188                    return None;
1189                }
1190                let (bytes, spans) = self.view(res);
1191                let win = span_window(spans, res.lo, res.hi);
1192                parse_ref_array(bytes, win, 0, res.hi)
1193            });
1194            match expanded {
1195                Some(items) => out.extend(items),
1196                None => out.push(reference),
1197            }
1198        }
1199        out
1200    }
1201}
1202
1203/// Learn a lone zlib stream's exact decoded length, or decline.
1204///
1205/// The decoded length is unknown a priori, so it is learned by inflating under a
1206/// hard cap rather than invented. Exceeding the cap declines (whether the
1207/// inflater errors or would silently clamp). The compressed payload is never
1208/// re-baked.
1209fn try_decode_len(source: &[u8], data_start: u64, data_len: u64, limits: Limits) -> Option<u64> {
1210    if data_len > MAX_STREAM_DECODE {
1211        return None;
1212    }
1213    let start = usize::try_from(data_start).ok()?;
1214    let end = usize::try_from(data_start.checked_add(data_len)?).ok()?;
1215    let encoded = source.get(start..end)?;
1216    if !zlib_shape_ok(encoded) {
1217        return None;
1218    }
1219    // `+1` distinguishes "there are more bytes" from "exactly at the cap" for
1220    // implementations that clamp instead of erroring.
1221    let cap_u64 = MAX_STREAM_DECODE
1222        .min(limits.max_output_bytes)
1223        .checked_add(1)?;
1224    let cap = usize::try_from(cap_u64).ok()?;
1225    let decoded_len =
1226        super::inflate::inflate_len(encoded, cap, super::inflate::Wrapper::Zlib).ok()?;
1227    if decoded_len > MAX_STREAM_DECODE || decoded_len > limits.max_output_bytes {
1228        return None;
1229    }
1230    Some(decoded_len)
1231}
1232
1233/// Whether `bytes` begins with a structurally valid zlib (RFC 1950) header.
1234///
1235/// Mirrors the shape test in `codec::deflate::zlib_header_valid`; that function
1236/// is behind the opt-in `deflate-replay` feature, which the `field` feature does
1237/// not imply. This is a shape test only, never authority: the decode above
1238/// decides.
1239fn zlib_shape_ok(bytes: &[u8]) -> bool {
1240    if bytes.len() < 2 {
1241        return false;
1242    }
1243    let cmf = bytes[0];
1244    let flg = bytes[1];
1245    cmf & 0x0f == 8 && cmf >> 4 <= 7 && (u16::from(cmf) * 256 + u16::from(flg)).is_multiple_of(31)
1246}
1247
1248/// Pack an object number/generation into the canonical `(u32, u16)` identity.
1249fn identity32(number: u64, generation: u64) -> Option<(u32, u16)> {
1250    Some((u32::try_from(number).ok()?, u16::try_from(generation).ok()?))
1251}
1252
1253/// Require `offset` and `len` to fit in 32 bits (the `object_params` packing).
1254fn span32(offset: u64, len: u64) -> Option<(u64, u64)> {
1255    if offset >> 32 != 0 || len >> 32 != 0 {
1256        return None;
1257    }
1258    Some((offset, len))
1259}
1260
1261/// Spans whose start lies in `[lo, hi)`.
1262fn span_window(spans: &[Span], lo: u64, hi: u64) -> &[Span] {
1263    let a = spans.partition_point(|s| s.start < lo);
1264    let b = spans.partition_point(|s| s.start < hi);
1265    &spans[a..b]
1266}
1267
1268/// The leading `<< ... >>` or `[ ... ]` container byte range within a window.
1269///
1270/// Returns `(is_array, lo, hi)`. The first container opened in the window wins,
1271/// so an indirect-object header (`N G obj`) is skipped and a nested container
1272/// cannot be mistaken for the object's own body. `None` when the window holds no
1273/// balanced container (fail closed rather than guess).
1274fn leading_container(win: &[Span]) -> Option<(bool, u64, u64)> {
1275    for (i, span) in win.iter().enumerate() {
1276        let closer = match span.kind {
1277            SpanKind::DictOpen => SpanKind::DictClose,
1278            SpanKind::ArrayOpen => SpanKind::ArrayClose,
1279            _ => continue,
1280        };
1281        let opener = span.kind;
1282        let mut depth: u32 = 0;
1283        for s in &win[i..] {
1284            if s.kind == opener {
1285                depth = depth.checked_add(1)?;
1286            } else if s.kind == closer {
1287                if depth == 0 {
1288                    continue;
1289                }
1290                depth -= 1;
1291                if depth == 0 {
1292                    return Some((
1293                        opener == SpanKind::ArrayOpen,
1294                        span.start,
1295                        s.start.checked_add(s.len)?,
1296                    ));
1297                }
1298            }
1299        }
1300        return None;
1301    }
1302    None
1303}
1304
1305/// Collect the indirect references following a *top-level* `key` inside the
1306/// leading dictionary `[lo, hi)`.
1307///
1308/// Only a key at bracket-depth 1 and outside every array is considered, so a
1309/// nested sub-dictionary (e.g. Cairo's `/Group << ... >>` before `/Type /Page`)
1310/// can never shadow it. Accepts a single `N G R` or an inline array of them.
1311/// Returns `None` when the key is absent, the value is neither shape, or more
1312/// than [`MAX_REFS`] entries appear (fail closed rather than guess).
1313fn collect_key_refs(source: &[u8], win: &[Span], lo: u64, hi: u64, key: &[u8]) -> Option<Vec<u64>> {
1314    let name = find_top_level_name(win, source, lo, hi, key)?;
1315    let t0 = next_sig(win, name + 1, hi)?;
1316    if win[t0].kind == SpanKind::ArrayOpen {
1317        parse_ref_array(source, win, t0, hi)
1318    } else {
1319        let (number, _generation, _next) = read_ref(source, win, t0, hi)?;
1320        Some(vec![number])
1321    }
1322}
1323
1324/// The `/Name` value of a *top-level* `key` in the leading dictionary.
1325///
1326/// Depth-aware like [`collect_key_refs`], so a nested `/Type` (e.g. inside a
1327/// `/Group` sub-dictionary) is not mistaken for the object's own type.
1328fn top_level_name_value<'a>(
1329    source: &'a [u8],
1330    win: &[Span],
1331    lo: u64,
1332    hi: u64,
1333    key: &[u8],
1334) -> Option<&'a [u8]> {
1335    let name = find_top_level_name(win, source, lo, hi, key)?;
1336    let t = next_sig(win, name + 1, hi)?;
1337    let span = win[t];
1338    if span.kind != SpanKind::Name {
1339        return None;
1340    }
1341    span_bytes(source, span)?.strip_prefix(b"/")
1342}
1343
1344/// Parse an inline `[ N G R ... ]` array of references beginning at span index
1345/// `open_idx` (a `ArrayOpen`), bounded by `hi`.
1346fn parse_ref_array(source: &[u8], win: &[Span], open_idx: usize, hi: u64) -> Option<Vec<u64>> {
1347    let mut out: Vec<u64> = Vec::new();
1348    let mut i = open_idx + 1;
1349    loop {
1350        let t = next_sig(win, i, hi)?;
1351        if win[t].kind == SpanKind::ArrayClose {
1352            return Some(out);
1353        }
1354        let (number, _generation, next) = read_ref(source, win, t, hi)?;
1355        out.push(number);
1356        if out.len() > MAX_REFS {
1357            return None;
1358        }
1359        i = next;
1360    }
1361}
1362
1363/// The integer value of a *top-level* `key` in the given dictionary window.
1364fn top_level_integer_value(
1365    source: &[u8],
1366    win: &[Span],
1367    lo: u64,
1368    hi: u64,
1369    key: &[u8],
1370) -> Option<u64> {
1371    let name = find_top_level_name(win, source, lo, hi, key)?;
1372    let t = next_sig(win, name + 1, hi)?;
1373    integer_span(source, win[t])
1374}
1375
1376/// Parse an `/ObjStm` header of `n` `objnum offset` integer pairs.
1377///
1378/// When `/First` is present the header is bounded to the bytes before it and the
1379/// remainder must be blank; when it is absent the pairs are read from the start of
1380/// the buffer. Bounded by `n`, so a corrupted count cannot scan unboundedly.
1381fn parse_objstm_header(bytes: &[u8], first: Option<u64>, n: usize) -> Option<Vec<(u64, u64)>> {
1382    let mut pos = 0usize;
1383    let mut pairs = Vec::with_capacity(n.min(4096));
1384    for _ in 0..n {
1385        let object = read_uint_ws(bytes, &mut pos)?;
1386        let offset = read_uint_ws(bytes, &mut pos)?;
1387        pairs.push((object, offset));
1388    }
1389    if let Some(f) = first {
1390        let f = usize::try_from(f).ok()?;
1391        if f > bytes.len() || pos > f {
1392            return None;
1393        }
1394        if bytes[pos..f].iter().any(|&b| !is_pdf_ws(b)) {
1395            return None;
1396        }
1397    }
1398    Some(pairs)
1399}
1400
1401/// Read a whitespace-delimited unsigned decimal integer, advancing `pos`.
1402fn read_uint_ws(bytes: &[u8], pos: &mut usize) -> Option<u64> {
1403    while *pos < bytes.len() && is_pdf_ws(bytes[*pos]) {
1404        *pos += 1;
1405    }
1406    let start = *pos;
1407    let mut value: u64 = 0;
1408    while *pos < bytes.len() && bytes[*pos].is_ascii_digit() {
1409        value = value
1410            .checked_mul(10)?
1411            .checked_add(u64::from(bytes[*pos] - b'0'))?;
1412        *pos += 1;
1413    }
1414    if *pos == start {
1415        return None;
1416    }
1417    Some(value)
1418}
1419
1420/// Whether `b` is a PDF whitespace byte (PDF 32000-1 Table 1).
1421fn is_pdf_ws(b: u8) -> bool {
1422    matches!(b, 0x00 | 0x09 | 0x0A | 0x0C | 0x0D | 0x20)
1423}
1424
1425/// Read one `N G R` reference at a significant span index.
1426fn read_ref(source: &[u8], win: &[Span], at: usize, hi: u64) -> Option<(u64, u64, usize)> {
1427    let i = next_sig(win, at, hi)?;
1428    let number = integer_span(source, win[i])?;
1429    let j = next_sig(win, i + 1, hi)?;
1430    let generation = integer_span(source, win[j])?;
1431    let k = next_sig(win, j + 1, hi)?;
1432    if !is_regular(source, win[k], b"R") {
1433        return None;
1434    }
1435    Some((number, generation, k + 1))
1436}
1437
1438/// Index of a *top-level* (bracket-depth 1, outside any array) `Name` span equal
1439/// to `/<key>` fully inside `[lo, hi)`.
1440///
1441/// A nested dictionary or array increments the tracked depth, so a key that only
1442/// appears inside one is never returned.
1443fn find_top_level_name(win: &[Span], source: &[u8], lo: u64, hi: u64, key: &[u8]) -> Option<usize> {
1444    let mut dict_depth: i32 = 0;
1445    let mut array_depth: i32 = 0;
1446    for (i, span) in win.iter().enumerate() {
1447        if span.start < lo || span.start >= hi {
1448            continue;
1449        }
1450        match span.kind {
1451            SpanKind::DictOpen => dict_depth += 1,
1452            SpanKind::DictClose => dict_depth -= 1,
1453            SpanKind::ArrayOpen => array_depth += 1,
1454            SpanKind::ArrayClose => array_depth -= 1,
1455            SpanKind::Name
1456                if dict_depth == 1
1457                    && array_depth == 0
1458                    && span_bytes(source, *span).is_some_and(|b| {
1459                        b.len() == key.len() + 1 && b[0] == b'/' && &b[1..] == key
1460                    }) =>
1461            {
1462                return Some(i);
1463            }
1464            _ => {}
1465        }
1466    }
1467    None
1468}
1469
1470/// Index of the next non-whitespace, non-comment span starting before `hi`.
1471fn next_sig(win: &[Span], from: usize, hi: u64) -> Option<usize> {
1472    let mut j = from;
1473    while j < win.len() {
1474        let span = win[j];
1475        if span.start >= hi {
1476            return None;
1477        }
1478        match span.kind {
1479            SpanKind::Whitespace | SpanKind::Comment => j += 1,
1480            _ => return Some(j),
1481        }
1482    }
1483    None
1484}
1485
1486/// Parse a `Regular` decimal integer span.
1487fn integer_span(source: &[u8], span: Span) -> Option<u64> {
1488    if span.kind != SpanKind::Regular {
1489        return None;
1490    }
1491    let bytes = span_bytes(source, span)?;
1492    if bytes.is_empty() {
1493        return None;
1494    }
1495    let mut value: u64 = 0;
1496    for &b in bytes {
1497        if !b.is_ascii_digit() {
1498            return None;
1499        }
1500        value = value.checked_mul(10)?.checked_add(u64::from(b - b'0'))?;
1501    }
1502    Some(value)
1503}
1504
1505/// Whether `span` is the `Regular` keyword `keyword`.
1506fn is_regular(source: &[u8], span: Span, keyword: &[u8]) -> bool {
1507    span.kind == SpanKind::Regular && span_bytes(source, span) == Some(keyword)
1508}
1509
1510/// The `%PDF-x.y` header version string, from the scanned header span.
1511fn pdf_header_version(source: &[u8], physical: &PdfPhysical) -> Option<String> {
1512    let (off, len) = physical.header?;
1513    let start = usize::try_from(off).ok()?;
1514    let end = usize::try_from(off.checked_add(len)?).ok()?;
1515    Some(
1516        String::from_utf8_lossy(source.get(start..end)?)
1517            .trim()
1518            .to_string(),
1519    )
1520}
1521
1522fn pdf_u64_array(values: &[u64]) -> String {
1523    values
1524        .iter()
1525        .map(|v| v.to_string())
1526        .collect::<Vec<_>>()
1527        .join(",")
1528}
1529
1530/// One revision's lineage entry: its byte span, resolved `startxref`/`/Prev`
1531/// headers, and the indirect objects and streams whose bytes it defines.
1532fn pdf_revision_json(rev: &RevisionInfo, physical: &PdfPhysical) -> String {
1533    let objects: Vec<u64> = physical
1534        .objects
1535        .iter()
1536        .filter(|o| o.start >= rev.start && o.start < rev.end)
1537        .map(|o| o.number)
1538        .collect();
1539    let streams: Vec<u64> = physical
1540        .streams
1541        .iter()
1542        .filter(|s| s.data_start >= rev.start && s.data_start < rev.end)
1543        .map(|s| s.object)
1544        .collect();
1545    let startxref = rev
1546        .startxref
1547        .map_or_else(|| "null".to_string(), |v| v.to_string());
1548    let prev = rev
1549        .prev
1550        .map_or_else(|| "null".to_string(), |v| v.to_string());
1551    format!(
1552        "{{\"index\":{},\"start\":{},\"end\":{},\"len\":{},\"startxref\":{},\"prev\":{},\"objects\":[{}],\"streams\":[{}]}}",
1553        rev.index,
1554        rev.start,
1555        rev.end,
1556        rev.end.saturating_sub(rev.start),
1557        startxref,
1558        prev,
1559        pdf_u64_array(&objects),
1560        pdf_u64_array(&streams)
1561    )
1562}
1563
1564/// The whole-document revision lineage (Phase 17): a deterministic JSON object
1565/// with the `%PDF-` header, the revision count, and one entry per revision.
1566fn pdf_revision_lineage_json(source: &[u8], physical: &PdfPhysical) -> String {
1567    let header = match pdf_header_version(source, physical) {
1568        Some(h) => format!("\"{}\"", crate::field::provenance::json_escape(&h)),
1569        None => "null".to_string(),
1570    };
1571    let revs = physical
1572        .revisions
1573        .iter()
1574        .map(|r| pdf_revision_json(r, physical))
1575        .collect::<Vec<_>>()
1576        .join(",");
1577    format!(
1578        "{{\"format\":\"pdf\",\"header\":{},\"count\":{},\"revisions\":[{}]}}",
1579        header,
1580        physical.revisions.len(),
1581        revs
1582    )
1583}
1584
1585/// The bytes backing `span`, or `None` if the offset is out of range.
1586fn span_bytes(source: &[u8], span: Span) -> Option<&[u8]> {
1587    let start = usize::try_from(span.start).ok()?;
1588    let end = usize::try_from(span.start.checked_add(span.len)?).ok()?;
1589    source.get(start..end)
1590}
1591
1592#[cfg(test)]
1593mod tests {
1594    use super::*;
1595    use crate::container::{Descriptor, ObjectSource};
1596    use crate::dra::{Op, Program};
1597    use crate::field::dag::EvalBudget;
1598    use std::fs;
1599    use std::path::PathBuf;
1600
1601    fn temp_root(label: &str) -> PathBuf {
1602        let mut p = std::env::temp_dir();
1603        p.push(format!(
1604            "vole-ingest-{label}-{}-{}",
1605            std::process::id(),
1606            std::time::SystemTime::now()
1607                .duration_since(std::time::UNIX_EPOCH)
1608                .unwrap()
1609                .as_nanos()
1610        ));
1611        p
1612    }
1613
1614    /// Wrap raw bytes as an opaque exact `.voldoc` descriptor.
1615    fn opaque_descriptor(source: &[u8]) -> Vec<u8> {
1616        let d = Descriptor {
1617            universe: crate::container::UNIVERSE.to_string(),
1618            source_format: crate::SOURCE_FORMAT_PDF,
1619            format_basis: "pdf:ingest-test".to_string(),
1620            models: vec![],
1621            channels: vec![],
1622            objects: vec![ObjectSource::Inline(source.to_vec())],
1623            program: Program::new(vec![Op::EmitObject { object_id: 0 }]),
1624            observation_index: None,
1625            seek_directory: false,
1626            checkpoints: None,
1627            source_sha256: crate::integrity::sha256(source),
1628            source_len: source.len() as u64,
1629        };
1630        d.serialize().unwrap().0
1631    }
1632
1633    fn adler32(data: &[u8]) -> u32 {
1634        let mut a: u32 = 1;
1635        let mut b: u32 = 0;
1636        for &byte in data {
1637            a = (a + u32::from(byte)) % 65521;
1638            b = (b + a) % 65521;
1639        }
1640        (b << 16) | a
1641    }
1642
1643    /// A genuine zlib stream using stored (uncompressed) DEFLATE blocks, so the
1644    /// fixture needs no compressor dependency and is byte-deterministic.
1645    fn zlib_stored(data: &[u8]) -> Vec<u8> {
1646        assert!(!data.is_empty());
1647        let mut out = vec![0x78, 0x01];
1648        let chunks: Vec<&[u8]> = data.chunks(0xFFFF).collect();
1649        for (i, chunk) in chunks.iter().enumerate() {
1650            let final_block = u8::from(i + 1 == chunks.len());
1651            out.push(final_block); // BFINAL, BTYPE=00 (stored)
1652            let len = chunk.len() as u16;
1653            out.extend_from_slice(&len.to_le_bytes());
1654            out.extend_from_slice(&(!len).to_le_bytes());
1655            out.extend_from_slice(chunk);
1656        }
1657        out.extend_from_slice(&adler32(data).to_be_bytes());
1658        out
1659    }
1660
1661    struct PdfBuilder {
1662        buf: Vec<u8>,
1663        offsets: Vec<(u64, u64)>,
1664    }
1665
1666    impl PdfBuilder {
1667        fn new() -> Self {
1668            PdfBuilder {
1669                buf: Vec::new(),
1670                offsets: Vec::new(),
1671            }
1672        }
1673        fn text(&mut self, s: &str) {
1674            self.buf.extend_from_slice(s.as_bytes());
1675        }
1676        fn raw(&mut self, b: &[u8]) {
1677            self.buf.extend_from_slice(b);
1678        }
1679        fn obj(&mut self, number: u64, body: &[u8]) {
1680            self.offsets.push((number, self.buf.len() as u64));
1681            self.text(&format!("{number} 0 obj\n"));
1682            self.raw(body);
1683            self.text("\nendobj\n");
1684        }
1685        fn stream_obj(&mut self, number: u64, extra: &str, data: &[u8]) {
1686            self.offsets.push((number, self.buf.len() as u64));
1687            self.text(&format!(
1688                "{number} 0 obj\n<< /Length {}{extra} >>\nstream\n",
1689                data.len()
1690            ));
1691            self.raw(data);
1692            self.text("\nendstream\nendobj\n");
1693        }
1694        fn offset_of(&self, number: u64) -> u64 {
1695            self.offsets
1696                .iter()
1697                .find(|&&(n, _)| n == number)
1698                .map(|&(_, off)| off)
1699                .unwrap()
1700        }
1701        fn classic_trailer(&mut self, size: u64, extra: &str) {
1702            let xref = self.buf.len() as u64;
1703            self.text(&format!("xref\n0 {size}\n"));
1704            self.raw(b"0000000000 65535 f \n");
1705            for number in 1..size {
1706                let off = self.offset_of(number);
1707                self.text(&format!("{off:010} 00000 n \n"));
1708            }
1709            self.text(&format!(
1710                "trailer\n<< /Size {size}{extra} >>\nstartxref\n{xref}\n%%EOF\n"
1711            ));
1712        }
1713        /// A trailer (with no xref body) pointing at `/Root`; enough for the
1714        /// physical scan, which never consults the cross-reference table. Used by
1715        /// fixtures whose catalog is not a physical object (it lives in an
1716        /// `/ObjStm`, which has no physical offset to record).
1717        fn raw_trailer(&mut self, size: u64, root: u64) {
1718            self.text(&format!(
1719                "trailer\n<< /Size {size} /Root {root} 0 R >>\n%%EOF\n"
1720            ));
1721        }
1722    }
1723
1724    /// A classic-xref PDF with one page and one lone-Flate content stream whose
1725    /// plaintext contains `(Hello) Tj`.
1726    fn fixture_pdf() -> Vec<u8> {
1727        let content = b"BT /F1 12 Tf 72 720 Td (Hello) Tj ET\n";
1728        let encoded = zlib_stored(content);
1729        let mut w = PdfBuilder::new();
1730        w.text("%PDF-1.5\n");
1731        w.obj(1, b"<< /Type /Catalog /Pages 2 0 R >>");
1732        w.obj(2, b"<< /Type /Pages /Kids [3 0 R] /Count 1 >>");
1733        w.obj(
1734            3,
1735            b"<< /Type /Page /Parent 2 0 R /MediaBox [0 0 612 792] /Resources << /Font << /F1 5 0 R >> >> /Contents 4 0 R >>",
1736        );
1737        w.stream_obj(4, " /Filter /FlateDecode", &encoded);
1738        w.obj(5, b"<< /Type /Font /Subtype /Type1 /BaseFont /Helvetica >>");
1739        w.classic_trailer(6, " /Root 1 0 R");
1740        w.buf
1741    }
1742
1743    /// A one-page PDF whose page object carries a nested `/Group << ... /Type
1744    /// /Group ... >>` *before* its own `/Type /Page` (Cairo-style key order).
1745    fn fixture_pdf_nested_type() -> Vec<u8> {
1746        let content = b"BT /F1 12 Tf 72 720 Td (Nested) Tj ET\n";
1747        let encoded = zlib_stored(content);
1748        let mut w = PdfBuilder::new();
1749        w.text("%PDF-1.5\n");
1750        w.obj(1, b"<< /Type /Catalog /Pages 2 0 R >>");
1751        w.obj(2, b"<< /Type /Pages /Kids [3 0 R] /Count 1 >>");
1752        w.obj(
1753            3,
1754            b"<< /Contents 4 0 R /Group << /S /Transparency /Type /Group >> /MediaBox [0 0 612 792] /Parent 2 0 R /Resources << /Font << /F1 5 0 R >> >> /Type /Page >>",
1755        );
1756        w.stream_obj(4, " /Filter /FlateDecode", &encoded);
1757        w.obj(5, b"<< /Type /Font /Subtype /Type1 /BaseFont /Helvetica >>");
1758        w.classic_trailer(6, " /Root 1 0 R");
1759        w.buf
1760    }
1761
1762    /// A one-page PDF whose content stream has no `/Filter`, so its encoded
1763    /// bytes are the content.
1764    fn fixture_pdf_unfiltered() -> Vec<u8> {
1765        let content = b"BT /F1 12 Tf 72 720 Td (Plain) Tj ET\n";
1766        let mut w = PdfBuilder::new();
1767        w.text("%PDF-1.5\n");
1768        w.obj(1, b"<< /Type /Catalog /Pages 2 0 R >>");
1769        w.obj(2, b"<< /Type /Pages /Kids [3 0 R] /Count 1 >>");
1770        w.obj(
1771            3,
1772            b"<< /Type /Page /Parent 2 0 R /MediaBox [0 0 612 792] /Resources << /Font << /F1 5 0 R >> >> /Contents 4 0 R >>",
1773        );
1774        w.stream_obj(4, "", content);
1775        w.obj(5, b"<< /Type /Font /Subtype /Type1 /BaseFont /Helvetica >>");
1776        w.classic_trailer(6, " /Root 1 0 R");
1777        w.buf
1778    }
1779
1780    /// Build an `/ObjStm` payload: an `objnum offset` header plus the bodies.
1781    ///
1782    /// Returns the stored-block-zlib-compressed stream, the pair count (`/N`),
1783    /// and the header length (`/First`). Offsets are relative to `/First` and
1784    /// bodies are newline-separated, as PDF 32000-1 §7.5.7 requires.
1785    fn objstm_stream(objs: &[(u64, &[u8])]) -> (Vec<u8>, u64, u64) {
1786        let mut bodies = Vec::new();
1787        let mut offsets = Vec::new();
1788        for (_, body) in objs {
1789            offsets.push(bodies.len());
1790            bodies.extend_from_slice(body);
1791            bodies.push(b'\n');
1792        }
1793        let mut header = String::new();
1794        for (i, (number, _)) in objs.iter().enumerate() {
1795            if i > 0 {
1796                header.push(' ');
1797            }
1798            header.push_str(&format!("{number} {}", offsets[i]));
1799        }
1800        header.push('\n');
1801        let first = header.len() as u64;
1802        let mut decoded = header.into_bytes();
1803        decoded.extend_from_slice(&bodies);
1804        (zlib_stored(&decoded), objs.len() as u64, first)
1805    }
1806
1807    /// A PDF whose `/Catalog`, `/Pages`, and `/Page` all live inside `/ObjStm`
1808    /// object 1, with a physical (Flate) content stream. The plaintext says
1809    /// `(Streamed)`.
1810    fn fixture_pdf_objstm() -> Vec<u8> {
1811        let content = b"BT /F1 12 Tf 72 720 Td (Streamed) Tj ET\n";
1812        let encoded = zlib_stored(content);
1813        let page = b"<< /Type /Page /Parent 3 0 R /MediaBox [0 0 612 792] /Resources << /Font << /F1 6 0 R >> >> /Contents 5 0 R >>";
1814        let pages = b"<< /Type /Pages /Kids [2 0 R] /Count 1 >>";
1815        let catalog = b"<< /Type /Catalog /Pages 3 0 R >>";
1816        let (objstm, n, first) = objstm_stream(&[(2, page), (3, pages), (4, catalog)]);
1817        let mut w = PdfBuilder::new();
1818        w.text("%PDF-1.5\n");
1819        w.stream_obj(
1820            1,
1821            &format!(" /Filter /FlateDecode /Type /ObjStm /N {n} /First {first}"),
1822            &objstm,
1823        );
1824        w.stream_obj(5, " /Filter /FlateDecode", &encoded);
1825        w.obj(6, b"<< /Type /Font /Subtype /Type1 /BaseFont /Helvetica >>");
1826        w.raw_trailer(7, 4);
1827        w.buf
1828    }
1829
1830    /// Like [`fixture_pdf_objstm`] but the object stream's own `/Pages` tree is
1831    /// shadowed by a *later physical* catalog and page tree. The physical tree's
1832    /// plaintext says `(Physical)`; the object stream's says `(Streamed)`.
1833    fn fixture_pdf_objstm_shadowed() -> Vec<u8> {
1834        let streamed = zlib_stored(b"BT /F1 12 Tf 72 720 Td (Streamed) Tj ET\n");
1835        let physical = zlib_stored(b"BT /F1 12 Tf 72 720 Td (Physical) Tj ET\n");
1836        let page = b"<< /Type /Page /Parent 2 0 R /Contents 5 0 R >>";
1837        let pages = b"<< /Type /Pages /Kids [3 0 R] /Count 1 >>";
1838        let catalog = b"<< /Type /Catalog /Pages 2 0 R >>";
1839        let (objstm, n, first) = objstm_stream(&[(2, pages), (3, page), (4, catalog)]);
1840        let mut w = PdfBuilder::new();
1841        w.text("%PDF-1.5\n");
1842        w.stream_obj(
1843            1,
1844            &format!(" /Filter /FlateDecode /Type /ObjStm /N {n} /First {first}"),
1845            &objstm,
1846        );
1847        // A later physical object number 4 shadows the object-stream catalog.
1848        w.obj(4, b"<< /Type /Catalog /Pages 7 0 R >>");
1849        w.stream_obj(5, " /Filter /FlateDecode", &streamed);
1850        w.obj(7, b"<< /Type /Pages /Kids [8 0 R] /Count 1 >>");
1851        w.obj(
1852            8,
1853            b"<< /Type /Page /Parent 7 0 R /MediaBox [0 0 612 792] /Contents 9 0 R >>",
1854        );
1855        w.stream_obj(9, " /Filter /FlateDecode", &physical);
1856        w.raw_trailer(10, 4);
1857        w.buf
1858    }
1859
1860    fn ingest(store: &mut FieldStore, pdf: &[u8]) -> IngestReport {
1861        let descriptor = opaque_descriptor(pdf);
1862        ingest_pdf(store, &descriptor, Limits::DEFAULT).unwrap()
1863    }
1864
1865    /// Ingest `pdf` once serial and once with a `workers`-thread pool, and assert
1866    /// the two runs are byte-identical: the same report (node ids, counters, field
1867    /// id) and the same exact materialization as the source.
1868    #[cfg(feature = "parallel")]
1869    fn assert_parallel_matches_serial(label: &str, pdf: &[u8], workers: usize) {
1870        use crate::parallel::WorkerPool;
1871
1872        let descriptor = opaque_descriptor(pdf);
1873        let serial_root = temp_root(&format!("{label}-serial"));
1874        let par_root = temp_root(&format!("{label}-parallel"));
1875
1876        let mut serial_store = FieldStore::open(&serial_root).unwrap();
1877        let serial = ingest_pdf(&mut serial_store, &descriptor, Limits::DEFAULT).unwrap();
1878
1879        let mut par_store = FieldStore::open(&par_root).unwrap();
1880        let pool = WorkerPool::new(workers).unwrap();
1881        let parallel =
1882            ingest_pdf_with(&mut par_store, &descriptor, Limits::DEFAULT, Some(&pool)).unwrap();
1883
1884        assert!(
1885            serial == parallel,
1886            "{label}: parallel report must match serial"
1887        );
1888
1889        let serial_field = Field::open(&serial_store, &serial.field, Limits::DEFAULT).unwrap();
1890        let par_field = Field::open(&par_store, &parallel.field, Limits::DEFAULT).unwrap();
1891        let serial_bytes = serial_field.materialize_exact(Limits::DEFAULT).unwrap();
1892        let par_bytes = par_field.materialize_exact(Limits::DEFAULT).unwrap();
1893        assert_eq!(serial_bytes, par_bytes, "{label}: exact bytes must match");
1894        assert_eq!(par_bytes, pdf, "{label}: exact bytes must equal the source");
1895        assert_eq!(
1896            crate::integrity::sha256(&par_bytes),
1897            crate::integrity::sha256(pdf)
1898        );
1899
1900        fs::remove_dir_all(&serial_root).ok();
1901        fs::remove_dir_all(&par_root).ok();
1902    }
1903
1904    /// The parallel path must be byte-identical to serial, for the rank-1
1905    /// (per-stream inflate) and rank-2 (`/ObjStm` decode + lex) sites.
1906    #[cfg(feature = "parallel")]
1907    #[test]
1908    fn parallel_ingest_is_byte_identical_to_serial() {
1909        assert_parallel_matches_serial("pdf", &fixture_pdf(), 4);
1910        assert_parallel_matches_serial("pdf-nested", &fixture_pdf_nested_type(), 4);
1911        assert_parallel_matches_serial("pdf-unfiltered", &fixture_pdf_unfiltered(), 2);
1912        assert_parallel_matches_serial("pdf-objstm", &fixture_pdf_objstm(), 4);
1913        assert_parallel_matches_serial("pdf-objstm-shadowed", &fixture_pdf_objstm_shadowed(), 8);
1914    }
1915
1916    /// The rank-1 length table is exactly the serial one whatever the pool: the
1917    /// indexed collect preserves physical order, which is what makes the running
1918    /// `MAX_TOTAL_DECODED` gate (replayed serially in `run_stage_b`) order-safe.
1919    #[cfg(feature = "parallel")]
1920    #[test]
1921    fn decoded_lengths_are_identical_with_and_without_a_pool() {
1922        use crate::parallel::WorkerPool;
1923        let pdf = fixture_pdf_objstm();
1924        let physical = scan(&pdf, Limits::DEFAULT).unwrap();
1925        let serial = decode_lengths(None, &pdf, &physical, Limits::DEFAULT);
1926        let pool = WorkerPool::new(4).unwrap();
1927        let parallel = decode_lengths(Some(&pool), &pdf, &physical, Limits::DEFAULT);
1928        assert!(serial == parallel, "length table must be pool-independent");
1929    }
1930
1931    #[test]
1932    fn ingest_recovers_streams_and_index() {
1933        let root = temp_root("basic");
1934        let mut store = FieldStore::open(&root).unwrap();
1935        let pdf = fixture_pdf();
1936        let report = ingest(&mut store, &pdf);
1937
1938        assert!(report.stream_nodes >= 1);
1939        assert!(report.decoded_stream_nodes >= 1);
1940        assert!(report.object_nodes >= 3);
1941        assert!(report.revision_nodes >= 1);
1942        assert!(report.index_root.is_some());
1943        assert_eq!(report.index_node_count, 1);
1944
1945        // Every recovered node materializes through the DAG.
1946        let field = Field::open(&store, &report.field, Limits::DEFAULT).unwrap();
1947        let mut budget = EvalBudget::default();
1948        let root_bytes = field
1949            .materialize_node(&report.root_node, Limits::DEFAULT, &mut budget)
1950            .unwrap();
1951        assert_eq!(root_bytes, pdf);
1952
1953        let istore = FsIndexStore::open(store.root()).unwrap();
1954        let index_root = report.index_root.unwrap();
1955        let stream_entries =
1956            lookup(&istore, &index_root, &SelectorKey::new(SEL_STREAM, 4)).unwrap();
1957        assert!(!stream_entries.is_empty());
1958        for entry in &stream_entries {
1959            let bytes = field
1960                .materialize_node(&entry.node_id, Limits::DEFAULT, &mut budget)
1961                .unwrap();
1962            assert_eq!(
1963                bytes,
1964                pdf[entry.out_off as usize..(entry.out_off + entry.out_len) as usize]
1965            );
1966        }
1967        fs::remove_dir_all(&root).ok();
1968    }
1969
1970    #[test]
1971    fn exactness_is_preserved_after_ingest() {
1972        let root = temp_root("exact");
1973        let mut store = FieldStore::open(&root).unwrap();
1974        let pdf = fixture_pdf();
1975        let report = ingest(&mut store, &pdf);
1976        let field = Field::open(&store, &report.field, Limits::DEFAULT).unwrap();
1977        let materialized = field.materialize_exact(Limits::DEFAULT).unwrap();
1978        assert_eq!(materialized.len() as u64, report.source_len);
1979        assert_eq!(
1980            crate::integrity::sha256(&materialized),
1981            crate::integrity::sha256(&pdf)
1982        );
1983        assert_eq!(materialized, pdf);
1984        fs::remove_dir_all(&root).ok();
1985    }
1986
1987    #[test]
1988    fn every_exact_node_matches_its_declared_span() {
1989        let root = temp_root("spans");
1990        let mut store = FieldStore::open(&root).unwrap();
1991        let pdf = fixture_pdf();
1992        let report = ingest(&mut store, &pdf);
1993        let field = Field::open(&store, &report.field, Limits::DEFAULT).unwrap();
1994        let istore = FsIndexStore::open(store.root()).unwrap();
1995        let index_root = report.index_root.unwrap();
1996        let physical = scan(&pdf, Limits::DEFAULT).unwrap();
1997        let mut budget = EvalBudget::default();
1998
1999        let mut check = |key: SelectorKey, offset: u64, len: u64| {
2000            let entries = lookup(&istore, &index_root, &key).unwrap();
2001            assert!(!entries.is_empty(), "missing entry for {key:?}");
2002            let expected = &pdf[offset as usize..(offset + len) as usize];
2003            for entry in entries {
2004                let bytes = field
2005                    .materialize_node(&entry.node_id, Limits::DEFAULT, &mut budget)
2006                    .unwrap();
2007                assert_eq!(bytes, expected);
2008            }
2009        };
2010
2011        for obj in &physical.objects {
2012            if obj.number == 0 {
2013                continue;
2014            }
2015            check(
2016                SelectorKey::new(SEL_OBJECT, obj.number as u32),
2017                obj.start,
2018                obj.end - obj.start,
2019            );
2020        }
2021        for stream in &physical.streams {
2022            check(
2023                SelectorKey::new(SEL_STREAM, stream.object as u32),
2024                stream.data_start,
2025                stream.data_len,
2026            );
2027        }
2028        for rev in &physical.revisions {
2029            check(
2030                SelectorKey::new(SEL_REVISION, rev.index),
2031                rev.start,
2032                rev.end - rev.start,
2033            );
2034        }
2035        fs::remove_dir_all(&root).ok();
2036    }
2037
2038    #[test]
2039    fn decoded_stream_matches_inflated_bytes() {
2040        let root = temp_root("decoded");
2041        let mut store = FieldStore::open(&root).unwrap();
2042        let pdf = fixture_pdf();
2043        let report = ingest(&mut store, &pdf);
2044        let field = Field::open(&store, &report.field, Limits::DEFAULT).unwrap();
2045        let istore = FsIndexStore::open(store.root()).unwrap();
2046        let index_root = report.index_root.unwrap();
2047        let physical = scan(&pdf, Limits::DEFAULT).unwrap();
2048
2049        let stream = &physical.streams[0];
2050        let encoded =
2051            &pdf[stream.data_start as usize..(stream.data_start + stream.data_len) as usize];
2052        let expected = miniz_oxide::inflate::decompress_to_vec_zlib(encoded).unwrap();
2053
2054        let encoded_entry = lookup(
2055            &istore,
2056            &index_root,
2057            &SelectorKey::new(SEL_STREAM, stream.object as u32),
2058        )
2059        .unwrap()
2060        .pop()
2061        .unwrap();
2062        let node = stream_decoded_node(
2063            stream.object as u32,
2064            encoded_entry.node_id,
2065            expected.len() as u64,
2066        );
2067        let mut budget = EvalBudget::default();
2068        let bytes = field
2069            .materialize_node(&node.content_id(), Limits::DEFAULT, &mut budget)
2070            .unwrap();
2071        assert_eq!(bytes, expected);
2072        assert_eq!(bytes, b"BT /F1 12 Tf 72 720 Td (Hello) Tj ET\n");
2073        fs::remove_dir_all(&root).ok();
2074    }
2075
2076    #[test]
2077    fn object_and_stream_selectors_resolve() {
2078        let root = temp_root("selectors");
2079        let mut store = FieldStore::open(&root).unwrap();
2080        let pdf = fixture_pdf();
2081        let report = ingest(&mut store, &pdf);
2082        let istore = FsIndexStore::open(store.root()).unwrap();
2083        let index_root = report.index_root.unwrap();
2084        assert!(
2085            !lookup(&istore, &index_root, &SelectorKey::new(SEL_OBJECT, 1))
2086                .unwrap()
2087                .is_empty()
2088        );
2089        assert!(
2090            !lookup(&istore, &index_root, &SelectorKey::new(SEL_STREAM, 4))
2091                .unwrap()
2092                .is_empty()
2093        );
2094        assert!(
2095            lookup(&istore, &index_root, &SelectorKey::new(SEL_OBJECT, 999))
2096                .unwrap()
2097                .is_empty()
2098        );
2099        fs::remove_dir_all(&root).ok();
2100    }
2101
2102    #[test]
2103    fn hostile_input_never_panics() {
2104        let root = temp_root("hostile");
2105        let mut store = FieldStore::open(&root).unwrap();
2106        let mut state = 0x9E37_79B9_7F4A_7C15u64;
2107        let mut garbage = Vec::with_capacity(4096);
2108        while garbage.len() < 4096 {
2109            state ^= state << 13;
2110            state ^= state >> 7;
2111            state ^= state << 17;
2112            garbage.extend_from_slice(&state.to_le_bytes());
2113        }
2114        let descriptor = opaque_descriptor(&garbage);
2115        match ingest_pdf(&mut store, &descriptor, Limits::DEFAULT) {
2116            Ok(report) => {
2117                let field = Field::open(&store, &report.field, Limits::DEFAULT).unwrap();
2118                assert_eq!(field.materialize_exact(Limits::DEFAULT).unwrap(), garbage);
2119            }
2120            Err(e) => assert_eq!(e.class(), crate::ErrorClass::InvalidPdfStructure),
2121        }
2122        fs::remove_dir_all(&root).ok();
2123    }
2124
2125    #[test]
2126    fn deepen_is_idempotent_and_old_field_stays_exact() {
2127        let root = temp_root("deepen");
2128        let mut store = FieldStore::open(&root).unwrap();
2129        let pdf = fixture_pdf();
2130        let report = ingest(&mut store, &pdf);
2131
2132        let first = report.field;
2133        let deepened = deepen_page(&mut store, &first, 1, Limits::DEFAULT).unwrap();
2134        let deepened_again = deepen_page(&mut store, &first, 1, Limits::DEFAULT).unwrap();
2135        assert_eq!(deepened, deepened_again);
2136        assert_eq!(
2137            deepen_page(&mut store, &deepened, 1, Limits::DEFAULT).unwrap(),
2138            deepened
2139        );
2140        assert_ne!(deepened, first);
2141
2142        // The original field id still materializes the exact source.
2143        let old = Field::open(&store, &first, Limits::DEFAULT).unwrap();
2144        assert_eq!(old.materialize_exact(Limits::DEFAULT).unwrap(), pdf);
2145        let new = Field::open(&store, &deepened, Limits::DEFAULT).unwrap();
2146        assert_eq!(new.materialize_exact(Limits::DEFAULT).unwrap(), pdf);
2147        fs::remove_dir_all(&root).ok();
2148    }
2149
2150    #[test]
2151    fn deepened_page_yields_nonempty_text_runs() {
2152        let root = temp_root("text");
2153        let mut store = FieldStore::open(&root).unwrap();
2154        let pdf = fixture_pdf();
2155        let report = ingest(&mut store, &pdf);
2156        assert_eq!(report.page_nodes, 1);
2157
2158        let istore = FsIndexStore::open(store.root()).unwrap();
2159        let index_root = report.index_root.unwrap();
2160        let page_entry = lookup(&istore, &index_root, &SelectorKey::new(SEL_PAGE, 1))
2161            .unwrap()
2162            .pop()
2163            .unwrap();
2164        let page_content_id = page_entry.node_id;
2165
2166        let deepened = deepen_page(&mut store, &report.field, 1, Limits::DEFAULT).unwrap();
2167        let field = Field::open(&store, &deepened, Limits::DEFAULT).unwrap();
2168
2169        let (_ops, text, _preview) = derived_chain(1, page_content_id);
2170        let mut budget = EvalBudget::default();
2171        let bytes = field
2172            .materialize_node(&text.content_id(), Limits::DEFAULT, &mut budget)
2173            .unwrap();
2174        assert!(!bytes.is_empty());
2175        assert!(String::from_utf8_lossy(&bytes).contains("Hello"));
2176        fs::remove_dir_all(&root).ok();
2177    }
2178
2179    #[test]
2180    fn nested_type_dict_does_not_shadow_page() {
2181        let root = temp_root("nested-type");
2182        let mut store = FieldStore::open(&root).unwrap();
2183        let pdf = fixture_pdf_nested_type();
2184        let report = ingest(&mut store, &pdf);
2185        assert_eq!(
2186            report.page_nodes, 1,
2187            "a nested /Type /Group must not hide the page"
2188        );
2189
2190        let istore = FsIndexStore::open(store.root()).unwrap();
2191        let index_root = report.index_root.unwrap();
2192        let page_entry = lookup(&istore, &index_root, &SelectorKey::new(SEL_PAGE, 1))
2193            .unwrap()
2194            .pop()
2195            .unwrap();
2196
2197        let deepened = deepen_page(&mut store, &report.field, 1, Limits::DEFAULT).unwrap();
2198        let field = Field::open(&store, &deepened, Limits::DEFAULT).unwrap();
2199        let (_ops, text, _preview) = derived_chain(1, page_entry.node_id);
2200        let mut budget = EvalBudget::default();
2201        let bytes = field
2202            .materialize_node(&text.content_id(), Limits::DEFAULT, &mut budget)
2203            .unwrap();
2204        assert!(String::from_utf8_lossy(&bytes).contains("Nested"));
2205
2206        assert_eq!(field.materialize_exact(Limits::DEFAULT).unwrap(), pdf);
2207        fs::remove_dir_all(&root).ok();
2208    }
2209
2210    #[test]
2211    fn unfiltered_content_stream_recovers_page() {
2212        let root = temp_root("unfiltered");
2213        let mut store = FieldStore::open(&root).unwrap();
2214        let pdf = fixture_pdf_unfiltered();
2215        let report = ingest(&mut store, &pdf);
2216        assert_eq!(
2217            report.page_nodes, 1,
2218            "an unfiltered /Contents stream must not be dropped"
2219        );
2220
2221        let istore = FsIndexStore::open(store.root()).unwrap();
2222        let index_root = report.index_root.unwrap();
2223        let page_entry = lookup(&istore, &index_root, &SelectorKey::new(SEL_PAGE, 1))
2224            .unwrap()
2225            .pop()
2226            .unwrap();
2227
2228        let deepened = deepen_page(&mut store, &report.field, 1, Limits::DEFAULT).unwrap();
2229        let field = Field::open(&store, &deepened, Limits::DEFAULT).unwrap();
2230        let (_ops, text, _preview) = derived_chain(1, page_entry.node_id);
2231        let mut budget = EvalBudget::default();
2232        let bytes = field
2233            .materialize_node(&text.content_id(), Limits::DEFAULT, &mut budget)
2234            .unwrap();
2235        assert!(String::from_utf8_lossy(&bytes).contains("Plain"));
2236
2237        assert_eq!(field.materialize_exact(Limits::DEFAULT).unwrap(), pdf);
2238        fs::remove_dir_all(&root).ok();
2239    }
2240
2241    #[test]
2242    fn objstm_page_tree_is_recovered_and_deepened() {
2243        let root = temp_root("objstm-page");
2244        let mut store = FieldStore::open(&root).unwrap();
2245        let pdf = fixture_pdf_objstm();
2246        let report = ingest(&mut store, &pdf);
2247        assert_eq!(
2248            report.page_nodes, 1,
2249            "a page whose dictionary lives in an /ObjStm must be recovered"
2250        );
2251
2252        let istore = FsIndexStore::open(store.root()).unwrap();
2253        let index_root = report.index_root.unwrap();
2254        let page_entry = lookup(&istore, &index_root, &SelectorKey::new(SEL_PAGE, 1))
2255            .unwrap()
2256            .pop()
2257            .unwrap();
2258
2259        let deepened = deepen_page(&mut store, &report.field, 1, Limits::DEFAULT).unwrap();
2260        let field = Field::open(&store, &deepened, Limits::DEFAULT).unwrap();
2261        let (_ops, text, _preview) = derived_chain(1, page_entry.node_id);
2262        let mut budget = EvalBudget::default();
2263        let bytes = field
2264            .materialize_node(&text.content_id(), Limits::DEFAULT, &mut budget)
2265            .unwrap();
2266        assert!(String::from_utf8_lossy(&bytes).contains("Streamed"));
2267
2268        assert_eq!(field.materialize_exact(Limits::DEFAULT).unwrap(), pdf);
2269        fs::remove_dir_all(&root).ok();
2270    }
2271
2272    #[test]
2273    fn malformed_objstm_is_skipped_without_panic() {
2274        let root = temp_root("objstm-bad");
2275        let mut store = FieldStore::open(&root).unwrap();
2276        // (a) an absurd `/N` above the object cap, (b) a `/First` past the end of
2277        // the decoded buffer (a truncated header).
2278        for (n, first) in [(100_000u64, 3u64), (3, 4096)] {
2279            let encoded = zlib_stored(b"BT /F1 12 Tf 72 720 Td (X) Tj ET\n");
2280            let page = b"<< /Type /Page /Parent 3 0 R /Contents 5 0 R >>";
2281            let pages = b"<< /Type /Pages /Kids [2 0 R] /Count 1 >>";
2282            let catalog = b"<< /Type /Catalog /Pages 3 0 R >>";
2283            let (objstm, _, _) = objstm_stream(&[(2, page), (3, pages), (4, catalog)]);
2284            let mut w = PdfBuilder::new();
2285            w.text("%PDF-1.5\n");
2286            w.stream_obj(
2287                1,
2288                &format!(" /Filter /FlateDecode /Type /ObjStm /N {n} /First {first}"),
2289                &objstm,
2290            );
2291            w.stream_obj(5, " /Filter /FlateDecode", &encoded);
2292            w.obj(6, b"<< /Type /Font /Subtype /Type1 /BaseFont /Helvetica >>");
2293            w.raw_trailer(7, 4);
2294            let pdf = w.buf;
2295
2296            let report = ingest(&mut store, &pdf);
2297            assert_eq!(
2298                report.page_nodes, 0,
2299                "a malformed /ObjStm must yield no pages"
2300            );
2301            let field = Field::open(&store, &report.field, Limits::DEFAULT).unwrap();
2302            assert_eq!(field.materialize_exact(Limits::DEFAULT).unwrap(), pdf);
2303        }
2304        fs::remove_dir_all(&root).ok();
2305    }
2306
2307    #[test]
2308    fn physical_object_shadows_objstm_object() {
2309        let root = temp_root("objstm-shadow");
2310        let mut store = FieldStore::open(&root).unwrap();
2311        let pdf = fixture_pdf_objstm_shadowed();
2312        let report = ingest(&mut store, &pdf);
2313        assert_eq!(report.page_nodes, 1);
2314
2315        // The later physical catalog 4 and its page tree must win over the
2316        // object-stream catalog 4, so the recovered text says `Physical`.
2317        let istore = FsIndexStore::open(store.root()).unwrap();
2318        let index_root = report.index_root.unwrap();
2319        let page_entry = lookup(&istore, &index_root, &SelectorKey::new(SEL_PAGE, 1))
2320            .unwrap()
2321            .pop()
2322            .unwrap();
2323        let deepened = deepen_page(&mut store, &report.field, 1, Limits::DEFAULT).unwrap();
2324        let field = Field::open(&store, &deepened, Limits::DEFAULT).unwrap();
2325        let (_ops, text, _preview) = derived_chain(1, page_entry.node_id);
2326        let mut budget = EvalBudget::default();
2327        let bytes = field
2328            .materialize_node(&text.content_id(), Limits::DEFAULT, &mut budget)
2329            .unwrap();
2330        let rendered = String::from_utf8_lossy(&bytes);
2331        assert!(rendered.contains("Physical"), "got {rendered:?}");
2332        assert!(!rendered.contains("Streamed"), "got {rendered:?}");
2333
2334        assert_eq!(field.materialize_exact(Limits::DEFAULT).unwrap(), pdf);
2335        fs::remove_dir_all(&root).ok();
2336    }
2337}