use std::cell::Cell;
use std::time::Instant;
use crate::error::{Error, Result};
use crate::field::cache::DerivedCache;
use crate::field::dag::{self, EvalBudget, ReuseStats, SourceServer};
use crate::field::index::{
FsIndexStore, IndexEntry, SEL_OBJECT, SEL_PAGE, SEL_REVISION, SEL_STREAM, SEL_STREAM_DECODED,
SelectorKey, lookup,
};
use crate::field::ingest;
use crate::field::manifest::FieldRoot;
use crate::field::node::{NodeKind, SeedNode, read_u32_params, span_params, u32_params};
use crate::field::partial::{PartialDescriptor, PartialLoad};
use crate::field::{Field, FieldId, FieldStore, SeedSubstrate};
use crate::limits::Limits;
use crate::store::{Id, IoSnapshot, NodeId, SeedStore};
use super::provenance::{AnswerValue, Basis, FieldAnswer, IntegrityScope, json_escape};
const MAX_TEXTMATCH_PAGES: u32 = 1 << 20;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Selector {
Document,
Page(u32),
Object(u32),
Stream(u32),
Revision(u32),
ByteRange {
offset: u64,
len: u64,
},
TextMatch(String),
}
impl Selector {
pub fn canonical(&self) -> String {
match self {
Selector::Document => "document".to_string(),
Selector::Page(n) => format!("page:{n}"),
Selector::Object(n) => format!("object:{n}"),
Selector::Stream(n) => format!("stream:{n}"),
Selector::Revision(n) => format!("revision:{n}"),
Selector::ByteRange { offset, len } => format!("byte-range:{offset}:{len}"),
Selector::TextMatch(p) => format!("text-match:{p}"),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Representation {
Metadata,
Text,
Structure,
Operators,
EncodedBytes,
DecodedBytes,
ExactBytes,
Preview,
FullDocument,
}
impl Representation {
pub const fn name(self) -> &'static str {
match self {
Representation::Metadata => "metadata",
Representation::Text => "text",
Representation::Structure => "structure",
Representation::Operators => "operators",
Representation::EncodedBytes => "encoded",
Representation::DecodedBytes => "decoded",
Representation::ExactBytes => "exact",
Representation::Preview => "preview",
Representation::FullDocument => "full",
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct ObserveBudget {
pub max_output_bytes: u64,
pub max_nodes: u64,
}
impl Default for ObserveBudget {
fn default() -> Self {
ObserveBudget {
max_output_bytes: 64 * 1024 * 1024,
max_nodes: 1 << 20,
}
}
}
#[derive(Debug, Clone)]
pub struct ObserveRequest {
pub selector: Selector,
pub representation: Representation,
pub budget: ObserveBudget,
pub use_cache: bool,
}
impl ObserveRequest {
pub fn new(selector: Selector, representation: Representation) -> Self {
ObserveRequest {
selector,
representation,
budget: ObserveBudget::default(),
use_cache: true,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum DescriptorReadMode {
#[default]
Full,
Partial,
}
impl DescriptorReadMode {
pub const fn name(self) -> &'static str {
match self {
DescriptorReadMode::Full => "full",
DescriptorReadMode::Partial => "partial",
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ObserveStats {
pub index_nodes_read: u64,
pub seed_nodes_fetched: u64,
pub seed_nodes_materialized: u64,
pub seed_nodes_executed: u64,
pub seed_nodes_reused: u64,
pub cache_bytes_written: u64,
pub descriptor_bytes_read: u64,
pub descriptor_read_mode: DescriptorReadMode,
pub manifest_bytes_read: u64,
pub index_bytes_read: u64,
pub seed_bytes_read: u64,
pub bytes_read: u64,
pub bytes_returned: u64,
pub deepened: bool,
pub wall_micros: u64,
}
struct CountingSeedStore<S: SeedStore> {
inner: S,
gets: Cell<u64>,
}
impl<S: SeedStore> CountingSeedStore<S> {
fn new(inner: S) -> Self {
CountingSeedStore {
inner,
gets: Cell::new(0),
}
}
fn gets(&self) -> u64 {
self.gets.get()
}
fn note(&self) {
self.gets.set(self.gets.get() + 1);
}
}
impl<S: SeedStore> SeedStore for CountingSeedStore<S> {
fn put_node(&mut self, canonical: &[u8]) -> Result<NodeId> {
self.inner.put_node(canonical)
}
fn get_node(&self, id: &NodeId) -> Result<Vec<u8>> {
self.note();
self.inner.get_node(id)
}
fn get_node_range(&self, id: &NodeId, offset: u64, len: u64) -> Result<Vec<u8>> {
self.note();
self.inner.get_node_range(id, offset, len)
}
fn contains_node(&self, id: &NodeId) -> Result<bool> {
self.inner.contains_node(id)
}
fn list_nodes(&self) -> Result<Vec<(NodeId, u64)>> {
self.inner.list_nodes()
}
}
pub fn observe(
store: &mut FieldStore,
id: &FieldId,
req: &ObserveRequest,
limits: Limits,
) -> Result<(FieldAnswer, ObserveStats, FieldId)> {
let started = Instant::now();
match narrow_probe(store, id, req)? {
NarrowProbe::Probed {
hit: true,
manifest,
carry,
} => {
let view = FieldView {
manifest: manifest.as_ref(),
id: *id,
open_io: IoSnapshot::default(),
source: &NO_SOURCE,
loader: None,
object_count: 0,
graph_ops: 0,
read_mode: DescriptorReadMode::Partial,
};
let (seeds, istore) = open_sub_stores(store)?;
observe_with_stores_pre(store, view, req, limits, started, seeds, istore, carry)
}
NarrowProbe::Probed {
hit: false,
manifest,
carry,
} => {
let opened = OpenedField::open_with_manifest(store, req, *manifest, limits)?;
observe_view_pre(store, opened.view(), req, limits, started, carry)
}
NarrowProbe::NotEligible => {
let opened = OpenedField::open(store, id, req, limits)?;
observe_view(store, opened.view(), req, limits, started)
}
}
}
struct NoSource;
static NO_SOURCE: NoSource = NoSource;
impl SourceServer for NoSource {
fn serve_range(&self, _offset: u64, _len: u64, _limits: Limits) -> Result<Vec<u8>> {
Err(Error::internal_invariant(
"a cache-served observation attempted a descriptor range read",
))
}
fn serve_document(&self, _limits: Limits) -> Result<Vec<u8>> {
Err(Error::internal_invariant(
"a cache-served observation attempted a descriptor document read",
))
}
}
#[derive(Default)]
struct PrefetchedIndex {
entries: Vec<(SelectorKey, Vec<IndexEntry>)>,
}
impl PrefetchedIndex {
fn insert(&mut self, key: SelectorKey, entries: Vec<IndexEntry>) {
self.entries.push((key, entries));
}
fn get(&self, key: &SelectorKey) -> Option<&Vec<IndexEntry>> {
self.entries.iter().find(|(k, _)| k == key).map(|(_, v)| v)
}
}
#[derive(Default)]
struct ProbeCarry {
base_io: IoSnapshot,
prefetched: PrefetchedIndex,
output: Option<(NodeId, Vec<u8>)>,
}
enum NarrowProbe {
NotEligible,
Probed {
hit: bool,
manifest: Box<FieldRoot>,
carry: ProbeCarry,
},
}
fn narrow_probe(store: &FieldStore, id: &FieldId, req: &ObserveRequest) -> Result<NarrowProbe> {
use Representation as R;
if !req.use_cache {
return Ok(NarrowProbe::NotEligible);
}
if !store.supports_partial_descriptor() {
return Ok(NarrowProbe::NotEligible);
}
let cacheable = matches!(
(&req.selector, req.representation),
(Selector::Page(_), R::Text | R::Preview | R::Structure)
| (Selector::Stream(_), R::DecodedBytes | R::Operators)
);
if !cacheable {
return Ok(NarrowProbe::NotEligible);
}
let io_before = store.io().snapshot();
let manifest = store.get_field(id)?;
let mut prefetched = PrefetchedIndex::default();
if !manifest.has_index() {
let base_io = io_before.delta(&store.io().snapshot());
return Ok(NarrowProbe::Probed {
hit: false,
manifest: Box::new(manifest),
carry: ProbeCarry {
base_io,
prefetched,
output: None,
},
});
}
let istore = FsIndexStore::open_with_io(store.root(), store.io().handle())?;
let root = NodeId::from_bytes(manifest.index_root);
let seeds = store.seed_substrate();
let target: Option<(NodeId, u64)> = match (&req.selector, req.representation) {
(Selector::Page(page), R::Text | R::Preview | R::Structure) => {
let key = SelectorKey::new(SEL_PAGE, *page);
let entries = lookup(&istore, &root, &key)?;
prefetched.insert(key, entries.clone());
match entries.first() {
Some(entry) => {
let (ops, text, preview) = derived_nodes(*page, entry.node_id);
if seeds.contains_node(&ops.content_id())?
&& seeds.contains_node(&text.content_id())?
&& seeds.contains_node(&preview.content_id())?
{
let node = match req.representation {
R::Preview | R::Structure => preview,
_ => text,
};
Some((node.content_id(), node.limits.max_output_bytes))
} else {
None
}
}
None => None,
}
}
(Selector::Stream(object), R::DecodedBytes | R::Operators) => {
let enc_key = SelectorKey::new(SEL_STREAM, *object);
let enc = lookup(&istore, &root, &enc_key)?;
prefetched.insert(enc_key, enc.clone());
let dec_key = SelectorKey::new(SEL_STREAM_DECODED, *object);
let dec = lookup(&istore, &root, &dec_key)?;
prefetched.insert(dec_key, dec.clone());
if enc.is_empty() {
None
} else {
match dec.first() {
Some(entry) if req.representation == R::DecodedBytes => Some((
entry.node_id,
crate::field::node::NodeLimits::DEFAULT.max_output_bytes,
)),
Some(entry) => {
let node = SeedNode::new(
NodeKind::ContentOperators,
0,
Vec::new(),
vec![entry.node_id],
"pdf:content-operators",
);
Some((node.content_id(), node.limits.max_output_bytes))
}
None => None,
}
}
}
_ => None,
};
let (hit, output) = match target {
Some((target_id, max_output_bytes)) => {
match DerivedCache::open(store.root().join("cache"))?.get(&target_id) {
Ok(Some(bytes)) if bytes.len() as u64 <= max_output_bytes => {
(true, Some((target_id, bytes)))
}
_ => (false, None),
}
}
None => (false, None),
};
let base_io = io_before.delta(&store.io().snapshot());
Ok(NarrowProbe::Probed {
hit,
manifest: Box::new(manifest),
carry: ProbeCarry {
base_io,
prefetched,
output,
},
})
}
pub(crate) fn observe_opened(
store: &mut FieldStore,
opened: &OpenedField,
req: &ObserveRequest,
limits: Limits,
) -> Result<(FieldAnswer, ObserveStats, FieldId)> {
let started = Instant::now();
observe_view(store, opened.view(), req, limits, started)
}
pub fn observe_with_field(
store: &mut FieldStore,
field: &Field,
req: &ObserveRequest,
limits: Limits,
) -> Result<(FieldAnswer, ObserveStats, FieldId)> {
let started = Instant::now();
observe_view(store, FieldView::from_field(field), req, limits, started)
}
pub(crate) enum OpenedField {
Full(Box<Field>),
Partial(Box<PartialField>),
}
impl OpenedField {
pub(crate) fn open(
store: &FieldStore,
id: &FieldId,
req: &ObserveRequest,
limits: Limits,
) -> Result<OpenedField> {
if store.supports_partial_descriptor()
&& partial_eligible(req)
&& let Some(pf) = PartialField::try_open(store, id, limits)?
{
return Ok(OpenedField::Partial(Box::new(pf)));
}
Ok(OpenedField::Full(Box::new(Field::open(store, id, limits)?)))
}
pub(crate) fn open_with_manifest(
store: &FieldStore,
req: &ObserveRequest,
manifest: FieldRoot,
limits: Limits,
) -> Result<OpenedField> {
if store.supports_partial_descriptor()
&& partial_eligible(req)
&& let Some(pf) =
PartialField::finish_open(store, manifest.clone(), store.io().snapshot(), limits)?
{
return Ok(OpenedField::Partial(Box::new(pf)));
}
Ok(OpenedField::Full(Box::new(Field::open_after_manifest(
store,
manifest,
store.io().snapshot(),
limits,
)?)))
}
pub(crate) fn view(&self) -> FieldView<'_> {
match self {
OpenedField::Full(f) => FieldView::from_field(f),
OpenedField::Partial(p) => p.view(),
}
}
pub(crate) fn manifest(&self) -> &FieldRoot {
match self {
OpenedField::Full(f) => f.manifest(),
OpenedField::Partial(p) => &p.manifest,
}
}
}
pub(crate) struct FieldView<'a> {
pub manifest: &'a FieldRoot,
pub id: FieldId,
pub open_io: IoSnapshot,
pub source: &'a dyn SourceServer,
pub loader: Option<&'a PartialDescriptor>,
pub object_count: usize,
pub graph_ops: usize,
pub read_mode: DescriptorReadMode,
}
impl<'a> FieldView<'a> {
pub(crate) fn from_field(field: &'a Field) -> FieldView<'a> {
let parsed = field.parsed();
FieldView {
manifest: field.manifest(),
id: field.id(),
open_io: field.open_io(),
source: parsed,
loader: None,
object_count: parsed.descriptor.objects.len(),
graph_ops: parsed.descriptor.program.ops.len(),
read_mode: DescriptorReadMode::Full,
}
}
}
pub(crate) struct PartialField {
pub(crate) manifest: FieldRoot,
pub(crate) id: FieldId,
pub(crate) open_io: IoSnapshot,
loader: PartialDescriptor,
}
impl PartialField {
pub(crate) fn try_open(
store: &FieldStore,
id: &FieldId,
limits: Limits,
) -> Result<Option<PartialField>> {
let io_before = store.io().snapshot();
let manifest = store.get_field(id)?;
PartialField::finish_open(store, manifest, io_before, limits)
}
pub(crate) fn finish_open(
store: &FieldStore,
manifest: FieldRoot,
io_before: IoSnapshot,
limits: Limits,
) -> Result<Option<PartialField>> {
let descriptor_id = Id::from_bytes(manifest.descriptor_id);
let Some(path) = store.descriptor_path(&descriptor_id) else {
return Ok(None);
};
let loader = match PartialDescriptor::open(&path, limits)? {
PartialLoad::Ready(l) => l,
PartialLoad::Ineligible { bytes_read } => {
store.io().add_descriptor(bytes_read);
return Ok(None);
}
};
if loader.source_len() != manifest.source_len
|| loader.source_sha256() != manifest.source_sha256
{
return Err(Error::integrity_mismatch(
"partial descriptor does not match its field manifest's declared source",
));
}
let open_io = io_before.delta(&store.io().snapshot());
let id = manifest.content_id();
Ok(Some(PartialField {
id,
manifest,
open_io,
loader: *loader,
}))
}
pub(crate) fn view(&self) -> FieldView<'_> {
FieldView {
manifest: &self.manifest,
id: self.id,
open_io: self.open_io,
source: &self.loader,
loader: Some(&self.loader),
object_count: self.loader.object_count(),
graph_ops: self.loader.graph_ops(),
read_mode: DescriptorReadMode::Partial,
}
}
}
fn partial_eligible(req: &ObserveRequest) -> bool {
use Representation as R;
matches!(
(&req.selector, req.representation),
(Selector::ByteRange { .. }, R::ExactBytes)
| (Selector::Object(_), R::ExactBytes | R::EncodedBytes)
| (Selector::Revision(_), R::ExactBytes)
| (Selector::Stream(_), R::EncodedBytes)
| (Selector::Page(_), R::Text | R::Preview | R::Structure)
)
}
fn open_sub_stores(store: &FieldStore) -> Result<(CountingSeedStore<SeedSubstrate>, FsIndexStore)> {
let io = store.io();
let seeds = CountingSeedStore::new(store.seed_substrate());
let istore = FsIndexStore::open_with_io(store.root(), io.handle())?;
Ok((seeds, istore))
}
fn observe_view<'a>(
store: &'a mut FieldStore,
view: FieldView<'a>,
req: &ObserveRequest,
limits: Limits,
started: Instant,
) -> Result<(FieldAnswer, ObserveStats, FieldId)> {
observe_view_pre(store, view, req, limits, started, ProbeCarry::default())
}
fn observe_view_pre<'a>(
store: &'a mut FieldStore,
view: FieldView<'a>,
req: &ObserveRequest,
limits: Limits,
started: Instant,
carry: ProbeCarry,
) -> Result<(FieldAnswer, ObserveStats, FieldId)> {
let (seeds, istore) = open_sub_stores(store)?;
observe_with_stores_pre(store, view, req, limits, started, seeds, istore, carry)
}
#[cfg(test)]
fn observe_with_stores<'a, S: SeedStore>(
store: &'a mut FieldStore,
view: FieldView<'a>,
req: &ObserveRequest,
limits: Limits,
started: Instant,
seeds: CountingSeedStore<S>,
istore: FsIndexStore,
) -> Result<(FieldAnswer, ObserveStats, FieldId)> {
observe_with_stores_pre(
store,
view,
req,
limits,
started,
seeds,
istore,
ProbeCarry::default(),
)
}
#[allow(clippy::too_many_arguments)]
fn observe_with_stores_pre<'a, S: SeedStore>(
store: &'a mut FieldStore,
view: FieldView<'a>,
req: &ObserveRequest,
limits: Limits,
started: Instant,
seeds: CountingSeedStore<S>,
istore: FsIndexStore,
carry: ProbeCarry,
) -> Result<(FieldAnswer, ObserveStats, FieldId)> {
let ProbeCarry {
base_io,
prefetched,
output,
} = carry;
let io_base = store.io().snapshot();
let budget = EvalBudget {
max_nodes: req.budget.max_nodes,
..EvalBudget::default()
};
let cache = DerivedCache::open(store.root().join("cache"))?;
let field_id = view.id;
let mut ctx = Ctx {
store,
manifest: view.manifest,
source: view.source,
loader: view.loader,
open_io: view.open_io,
object_count: view.object_count,
graph_ops: view.graph_ops,
read_mode: view.read_mode,
seeds,
istore,
prefetched,
prefetched_output: output,
limits,
budget,
stats: ObserveStats::default(),
use_cache: req.use_cache,
cache,
reuse: ReuseStats::default(),
current_id: field_id,
};
let answer = ctx.dispatch(req)?;
let produced = answer.value.byte_len();
if produced > req.budget.max_output_bytes {
return Err(Error::resource_limit(format!(
"observation produced {produced} bytes, exceeding the {}-byte budget",
req.budget.max_output_bytes
)));
}
let mut stats = ctx.stats;
if let Some(loader) = ctx.loader {
ctx.store.io().add_descriptor(loader.bytes_read());
}
let open = ctx.open_io;
let extra = io_base.delta(&ctx.store.io().snapshot());
stats.descriptor_bytes_read = base_io
.descriptor_bytes
.saturating_add(open.descriptor_bytes)
.saturating_add(extra.descriptor_bytes);
stats.descriptor_read_mode = ctx.read_mode;
stats.manifest_bytes_read = base_io
.manifest_bytes
.saturating_add(open.manifest_bytes)
.saturating_add(extra.manifest_bytes);
stats.index_bytes_read = base_io.index_bytes.saturating_add(extra.index_bytes);
stats.seed_bytes_read = base_io.seed_bytes.saturating_add(extra.seed_bytes);
stats.bytes_read = stats
.descriptor_bytes_read
.saturating_add(stats.manifest_bytes_read)
.saturating_add(stats.index_bytes_read)
.saturating_add(stats.seed_bytes_read);
stats.seed_nodes_fetched = ctx.seeds.gets();
stats.seed_nodes_materialized = ctx.budget.nodes;
stats.seed_nodes_executed = ctx.reuse.nodes_executed;
stats.seed_nodes_reused = ctx.reuse.nodes_reused;
stats.cache_bytes_written = ctx.reuse.cache_bytes_written;
stats.bytes_returned = produced;
stats.wall_micros = started.elapsed().as_micros().min(u128::from(u64::MAX)) as u64;
Ok((answer, stats, ctx.current_id))
}
struct Ctx<'a, S: SeedStore> {
store: &'a mut FieldStore,
manifest: &'a FieldRoot,
source: &'a dyn SourceServer,
loader: Option<&'a PartialDescriptor>,
open_io: IoSnapshot,
object_count: usize,
graph_ops: usize,
read_mode: DescriptorReadMode,
seeds: CountingSeedStore<S>,
istore: FsIndexStore,
prefetched: PrefetchedIndex,
prefetched_output: Option<(NodeId, Vec<u8>)>,
limits: Limits,
budget: EvalBudget,
stats: ObserveStats,
use_cache: bool,
cache: DerivedCache,
reuse: ReuseStats,
current_id: FieldId,
}
impl<S: SeedStore> Ctx<'_, S> {
fn materialize(&mut self, node: &SeedNode) -> Result<Vec<u8>> {
let depth = node.limits.max_depth;
if self.use_cache {
if let Some((id, bytes)) = self.prefetched_output.take() {
if id == node.content_id() {
self.reuse.nodes_reused = self.reuse.nodes_reused.saturating_add(1);
self.budget.charge_bytes(bytes.len() as u64)?;
return Ok(bytes);
}
self.prefetched_output = Some((id, bytes));
}
dag::materialize_node_cached_with(
self.source,
&self.seeds,
&mut self.cache,
node,
self.limits,
&mut self.budget,
depth,
&mut self.reuse,
)
} else {
let mut cache = dag::NoCache;
dag::materialize_node_cached_with(
self.source,
&self.seeds,
&mut cache,
node,
self.limits,
&mut self.budget,
depth,
&mut self.reuse,
)
}
}
fn load(&self, id: &NodeId) -> Result<SeedNode> {
dag::load_node(&self.seeds, id)
}
fn lookup(&mut self, key: SelectorKey) -> Result<Vec<IndexEntry>> {
if !self.manifest.has_index() {
return Ok(Vec::new());
}
let prefetched = self.prefetched.get(&key).cloned();
let entries = match prefetched {
Some(entries) => entries,
None => {
let root = NodeId::from_bytes(self.manifest.index_root);
lookup(&self.istore, &root, &key)?
}
};
self.stats.index_nodes_read += entries.len() as u64;
Ok(entries)
}
fn require_entry(&mut self, key: SelectorKey, what: &str) -> Result<IndexEntry> {
let entries = self.lookup(key)?;
entries.into_iter().next().ok_or_else(|| {
Error::unsupported_feature(format!(
"no {what} matching selector number {} in the observation index",
key.number
))
})
}
fn dispatch(&mut self, req: &ObserveRequest) -> Result<FieldAnswer> {
use Representation as R;
match (&req.selector, req.representation) {
(Selector::Document, R::FullDocument | R::ExactBytes) => self.document_full(req),
(Selector::Document, R::Metadata) => self.document_metadata(req),
(Selector::ByteRange { offset, len }, R::ExactBytes) => {
self.byte_range(req, *offset, *len)
}
(Selector::Object(n), R::ExactBytes | R::EncodedBytes) => {
self.indexed_exact(req, SelectorKey::new(SEL_OBJECT, *n), "object")
}
(Selector::Revision(n), R::ExactBytes) => {
self.indexed_exact(req, SelectorKey::new(SEL_REVISION, *n), "revision")
}
(Selector::Stream(n), R::EncodedBytes) => {
self.indexed_exact(req, SelectorKey::new(SEL_STREAM, *n), "stream")
}
(Selector::Stream(n), R::DecodedBytes) => self.stream_decoded(req, *n),
(Selector::Stream(n), R::Operators) => self.stream_operators(req, *n),
(Selector::Page(n), R::Text) => self.page_text(req, *n),
(Selector::Page(n), R::Preview) => self.page_preview(req, *n),
(Selector::Page(n), R::Structure) => self.page_structure(req, *n),
(Selector::TextMatch(p), R::Text) => self.text_match(req, p),
_ => Err(Error::unsupported_feature(format!(
"unsupported observation: selector {} with representation {}",
req.selector.canonical(),
req.representation.name()
))),
}
}
fn document_full(&mut self, req: &ObserveRequest) -> Result<FieldAnswer> {
let bytes = self.source.serve_document(self.limits)?;
Ok(FieldAnswer {
value: AnswerValue::Bytes(bytes),
basis: Basis::DirectlyObserved,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: Some((0, self.manifest.source_len)),
dependency_ids: vec![self.manifest.root_node],
integrity_scope: IntegrityScope::WholeSource,
exact: true,
})
}
fn document_metadata(&mut self, req: &ObserveRequest) -> Result<FieldAnswer> {
let json = format!(
concat!(
"{{",
"\"source_len\":{},",
"\"source_sha256\":\"{}\",",
"\"object_count\":{},",
"\"graph_ops\":{},",
"\"node_count\":{}",
"}}"
),
self.manifest.source_len,
crate::integrity::to_hex(&self.manifest.source_sha256),
self.object_count,
self.graph_ops,
self.manifest.node_count,
);
Ok(FieldAnswer {
value: AnswerValue::Json(json),
basis: Basis::DeterministicallyDerived,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: None,
dependency_ids: Vec::new(),
integrity_scope: IntegrityScope::None,
exact: false,
})
}
fn byte_range(&mut self, req: &ObserveRequest, offset: u64, len: u64) -> Result<FieldAnswer> {
let end = offset
.checked_add(len)
.ok_or_else(|| Error::usage("byte-range end overflows"))?;
let node = SeedNode::new(
NodeKind::SourceSlice,
len,
span_params(offset, len),
Vec::new(),
"field:observe;source-slice",
);
let bytes = self.materialize(&node)?;
Ok(FieldAnswer {
value: AnswerValue::Bytes(bytes),
basis: Basis::DirectlyObserved,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: Some((offset, end)),
dependency_ids: Vec::new(),
integrity_scope: IntegrityScope::Node,
exact: true,
})
}
fn indexed_exact(
&mut self,
req: &ObserveRequest,
key: SelectorKey,
what: &str,
) -> Result<FieldAnswer> {
let entry = self.require_entry(key, what)?;
let node = self.load(&entry.node_id)?;
let bytes = self.materialize(&node)?;
let end = entry.out_off.saturating_add(entry.out_len);
Ok(FieldAnswer {
value: AnswerValue::Bytes(bytes),
basis: Basis::DirectlyObserved,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: Some((entry.out_off, end)),
dependency_ids: vec![entry.node_id],
integrity_scope: IntegrityScope::Node,
exact: true,
})
}
fn stream_decoded(&mut self, req: &ObserveRequest, object: u32) -> Result<FieldAnswer> {
let entry = self.require_entry(SelectorKey::new(SEL_STREAM, object), "stream")?;
let node = self.decoded_node(object, &entry.node_id)?;
let id = node.content_id();
let bytes = self.materialize(&node)?;
Ok(FieldAnswer {
value: AnswerValue::Bytes(bytes),
basis: Basis::DeterministicallyDerived,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: None,
dependency_ids: vec![id],
integrity_scope: IntegrityScope::None,
exact: false,
})
}
fn stream_operators(&mut self, req: &ObserveRequest, object: u32) -> Result<FieldAnswer> {
let entry = self.require_entry(SelectorKey::new(SEL_STREAM, object), "stream")?;
let decoded = self.decoded_node(object, &entry.node_id)?;
let decoded_id = decoded.content_id();
let node = SeedNode::new(
NodeKind::ContentOperators,
0,
Vec::new(),
vec![decoded_id],
"pdf:content-operators",
);
let id = node.content_id();
let bytes = self.materialize(&node)?;
Ok(FieldAnswer {
value: AnswerValue::Bytes(bytes),
basis: Basis::DeterministicallyDerived,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: None,
dependency_ids: vec![id, decoded_id],
integrity_scope: IntegrityScope::None,
exact: false,
})
}
fn page_content_id(&mut self, what_number: u32) -> Result<NodeId> {
Ok(self
.require_entry(SelectorKey::new(SEL_PAGE, what_number), "page")?
.node_id)
}
fn ensure_page_derived(
&mut self,
page: u32,
page_content: NodeId,
) -> Result<(SeedNode, SeedNode, SeedNode)> {
let (ops, text, preview) = derived_nodes(page, page_content);
let present = self.seeds.contains_node(&ops.content_id())?
&& self.seeds.contains_node(&text.content_id())?
&& self.seeds.contains_node(&preview.content_id())?;
if !present {
let promoted = ingest::deepen_page_with_manifest(self.store, self.manifest, page)?;
self.stats.deepened = true;
self.current_id = promoted;
}
Ok((ops, text, preview))
}
fn page_text(&mut self, req: &ObserveRequest, page: u32) -> Result<FieldAnswer> {
let pc = self.page_content_id(page)?;
let (ops, text, _preview) = self.ensure_page_derived(page, pc)?;
let text_id = text.content_id();
let bytes = self.materialize(&text)?;
let value = AnswerValue::Text(String::from_utf8_lossy(&bytes).into_owned());
Ok(FieldAnswer {
value,
basis: Basis::Heuristic,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: None,
dependency_ids: vec![text_id, ops.content_id(), pc],
integrity_scope: IntegrityScope::None,
exact: false,
})
}
fn page_preview(&mut self, req: &ObserveRequest, page: u32) -> Result<FieldAnswer> {
let pc = self.page_content_id(page)?;
let (_ops, _text, preview) = self.ensure_page_derived(page, pc)?;
let preview_id = preview.content_id();
let bytes = self.materialize(&preview)?;
Ok(FieldAnswer {
value: AnswerValue::Bytes(bytes),
basis: Basis::Heuristic,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: None,
dependency_ids: vec![preview_id, pc],
integrity_scope: IntegrityScope::None,
exact: false,
})
}
fn page_structure(&mut self, req: &ObserveRequest, page: u32) -> Result<FieldAnswer> {
let pc = self.page_content_id(page)?;
let (_ops, _text, preview) = self.ensure_page_derived(page, pc)?;
let preview_id = preview.content_id();
let preview_bytes = self.materialize(&preview)?;
let (text_bytes, draw_ops, path_ops) = preview_stats(&preview_bytes);
let pc_node = self.load(&pc)?;
let mut streams: Vec<u32> = Vec::new();
for dep in &pc_node.deps {
let dep_node = self.load(dep)?;
if !matches!(
dep_node.kind,
NodeKind::PdfStreamDecoded | NodeKind::PdfStreamEncoded
) {
continue;
}
if let Ok(object) = read_u32_params(&dep_node.params) {
streams.push(object);
}
}
let streams_json = streams
.iter()
.map(u32::to_string)
.collect::<Vec<_>>()
.join(",");
let json = format!(
"{{\"page\":{page},\"text_bytes\":{text_bytes},\"draw_ops\":{draw_ops},\"path_ops\":{path_ops},\"content_streams\":[{streams_json}]}}"
);
Ok(FieldAnswer {
value: AnswerValue::Json(json),
basis: Basis::DeterministicallyDerived,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: None,
dependency_ids: vec![preview_id, pc],
integrity_scope: IntegrityScope::None,
exact: false,
})
}
fn text_match(&mut self, req: &ObserveRequest, pattern: &str) -> Result<FieldAnswer> {
let mut items: Vec<(u32, String)> = Vec::new();
let mut estimated: u64 = 0;
let mut page: u32 = 1;
while page <= MAX_TEXTMATCH_PAGES {
let entries = self.lookup(SelectorKey::new(SEL_PAGE, page))?;
let Some(entry) = entries.into_iter().next() else {
break;
};
let pc = entry.node_id;
let (_ops, text, _preview) = self.ensure_page_derived(page, pc)?;
let bytes = self.materialize(&text)?;
let rendered = String::from_utf8_lossy(&bytes).into_owned();
for line in rendered.split('\n') {
if line.contains(pattern) {
estimated = estimated.saturating_add(line.len() as u64 + 32);
if estimated > req.budget.max_output_bytes {
return Err(Error::resource_limit(format!(
"text match exceeded the {}-byte budget",
req.budget.max_output_bytes
)));
}
items.push((page, line.to_string()));
}
}
page += 1;
}
let body = items
.iter()
.map(|(p, line)| format!("{{\"page\":{p},\"line\":\"{}\"}}", json_escape(line)))
.collect::<Vec<_>>()
.join(",");
Ok(FieldAnswer {
value: AnswerValue::Json(format!("[{body}]")),
basis: Basis::Heuristic,
selector: req.selector.canonical(),
representation: req.representation.name().to_string(),
source_span: None,
dependency_ids: Vec::new(),
integrity_scope: IntegrityScope::None,
exact: false,
})
}
fn decoded_node(&mut self, object: u32, encoded_id: &NodeId) -> Result<SeedNode> {
let entries = self.lookup(SelectorKey::new(SEL_STREAM_DECODED, object))?;
if let Some(entry) = entries.into_iter().next() {
return self.load(&entry.node_id);
}
self.deepen_stream(object, encoded_id)
}
fn deepen_stream(&mut self, object: u32, encoded_id: &NodeId) -> Result<SeedNode> {
let encoded_node = self.load(encoded_id)?;
let encoded = self.materialize(&encoded_node)?;
let cap = usize::try_from(self.limits.max_output_bytes).unwrap_or(usize::MAX);
let decoded = miniz_oxide::inflate::decompress_to_vec_zlib_with_limit(&encoded, cap)
.map_err(|e| {
Error::unsupported_feature(format!(
"stream {object} has no recovered decoded representation: {:?}",
e.status
))
})?;
let node = SeedNode::new(
NodeKind::PdfStreamDecoded,
decoded.len() as u64,
u32_params(object),
vec![*encoded_id],
"pdf:stream-decoded",
);
self.store.seeds_mut().put_node(&node.encode_canonical())?;
self.stats.deepened = true;
Ok(node)
}
}
pub(crate) fn derived_nodes(page: u32, page_content: NodeId) -> (SeedNode, SeedNode, SeedNode) {
let ops = SeedNode::new(
NodeKind::ContentOperators,
0,
u32_params(page),
vec![page_content],
"pdf:content-operators",
);
let text = SeedNode::new(
NodeKind::TextRuns,
0,
Vec::new(),
vec![ops.content_id()],
"pdf:text-runs",
);
let preview = SeedNode::new(
NodeKind::PagePreview,
0,
u32_params(page),
vec![page_content],
"pdf:page-preview",
);
(ops, text, preview)
}
fn preview_stats(bytes: &[u8]) -> (u64, u64, u64) {
let rendered = String::from_utf8_lossy(bytes);
let mut text_bytes = 0u64;
let mut draw_ops = 0u64;
let mut path_ops = 0u64;
for line in rendered.lines() {
if let Some(v) = line.strip_prefix("text-bytes ") {
text_bytes = v.trim().parse().unwrap_or(0);
} else if let Some(v) = line.strip_prefix("draw-ops ") {
draw_ops = v.trim().parse().unwrap_or(0);
} else if let Some(v) = line.strip_prefix("path-ops ") {
path_ops = v.trim().parse().unwrap_or(0);
}
}
(text_bytes, draw_ops, path_ops)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::container::{Descriptor, ObjectSource};
use crate::dra::{Op, Program};
use crate::field::plan;
use crate::store::FsSeedStore;
use std::fs;
use std::path::PathBuf;
fn temp_root(label: &str) -> PathBuf {
let mut p = std::env::temp_dir();
p.push(format!(
"vole-observe-{label}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
p
}
fn opaque_descriptor(source: &[u8]) -> Vec<u8> {
let d = Descriptor {
universe: crate::container::UNIVERSE.to_string(),
source_format: crate::SOURCE_FORMAT_PDF,
format_basis: "pdf:observe-test".to_string(),
models: vec![],
channels: vec![],
objects: vec![ObjectSource::Inline(source.to_vec())],
program: Program::new(vec![Op::EmitObject { object_id: 0 }]),
observation_index: None,
seek_directory: false,
source_sha256: crate::integrity::sha256(source),
source_len: source.len() as u64,
};
d.serialize().unwrap().0
}
fn adler32(data: &[u8]) -> u32 {
let mut a: u32 = 1;
let mut b: u32 = 0;
for &byte in data {
a = (a + u32::from(byte)) % 65521;
b = (b + a) % 65521;
}
(b << 16) | a
}
fn zlib_stored(data: &[u8]) -> Vec<u8> {
assert!(!data.is_empty());
let mut out = vec![0x78, 0x01];
let chunks: Vec<&[u8]> = data.chunks(0xFFFF).collect();
for (i, chunk) in chunks.iter().enumerate() {
let final_block = u8::from(i + 1 == chunks.len());
out.push(final_block);
let len = chunk.len() as u16;
out.extend_from_slice(&len.to_le_bytes());
out.extend_from_slice(&(!len).to_le_bytes());
out.extend_from_slice(chunk);
}
out.extend_from_slice(&adler32(data).to_be_bytes());
out
}
struct PdfBuilder {
buf: Vec<u8>,
offsets: Vec<(u64, u64)>,
}
impl PdfBuilder {
fn new() -> Self {
PdfBuilder {
buf: Vec::new(),
offsets: Vec::new(),
}
}
fn text(&mut self, s: &str) {
self.buf.extend_from_slice(s.as_bytes());
}
fn raw(&mut self, b: &[u8]) {
self.buf.extend_from_slice(b);
}
fn obj(&mut self, number: u64, body: &[u8]) {
self.offsets.push((number, self.buf.len() as u64));
self.text(&format!("{number} 0 obj\n"));
self.raw(body);
self.text("\nendobj\n");
}
fn stream_obj(&mut self, number: u64, extra: &str, data: &[u8]) {
self.offsets.push((number, self.buf.len() as u64));
self.text(&format!(
"{number} 0 obj\n<< /Length {}{extra} >>\nstream\n",
data.len()
));
self.raw(data);
self.text("\nendstream\nendobj\n");
}
fn offset_of(&self, number: u64) -> u64 {
self.offsets
.iter()
.find(|&&(n, _)| n == number)
.map(|&(_, off)| off)
.unwrap()
}
fn classic_trailer(&mut self, size: u64, extra: &str) {
let xref = self.buf.len() as u64;
self.text(&format!("xref\n0 {size}\n"));
self.raw(b"0000000000 65535 f \n");
for number in 1..size {
let off = self.offset_of(number);
self.text(&format!("{off:010} 00000 n \n"));
}
self.text(&format!(
"trailer\n<< /Size {size}{extra} >>\nstartxref\n{xref}\n%%EOF\n"
));
}
}
fn fixture_pdf(with_image: bool) -> Vec<u8> {
let content = b"BT /F1 12 Tf 72 720 Td (Hello) Tj ET\n";
let encoded = zlib_stored(content);
let mut w = PdfBuilder::new();
w.text("%PDF-1.5\n");
w.obj(1, b"<< /Type /Catalog /Pages 2 0 R >>");
w.obj(2, b"<< /Type /Pages /Kids [3 0 R] /Count 1 >>");
if with_image {
w.obj(
3,
b"<< /Type /Page /Parent 2 0 R /MediaBox [0 0 612 792] /Resources << /Font << /F1 5 0 R >> /XObject << /Im0 6 0 R >> >> /Contents 4 0 R >>",
);
} else {
w.obj(
3,
b"<< /Type /Page /Parent 2 0 R /MediaBox [0 0 612 792] /Resources << /Font << /F1 5 0 R >> >> /Contents 4 0 R >>",
);
}
w.stream_obj(4, " /Filter /FlateDecode", &encoded);
w.obj(5, b"<< /Type /Font /Subtype /Type1 /BaseFont /Helvetica >>");
if with_image {
let image = vec![0x80u8; 256 * 256];
let image_encoded = zlib_stored(&image);
w.stream_obj(
6,
" /Type /XObject /Subtype /Image /Width 256 /Height 256 /ColorSpace /DeviceGray /BitsPerComponent 8 /Filter /FlateDecode",
&image_encoded,
);
w.classic_trailer(7, " /Root 1 0 R");
} else {
w.classic_trailer(6, " /Root 1 0 R");
}
w.buf
}
struct Fixture {
root: PathBuf,
store: FieldStore,
field: FieldId,
source: Vec<u8>,
}
impl Fixture {
fn new(label: &str, with_image: bool) -> Fixture {
let root = temp_root(label);
let mut store = FieldStore::open(&root).unwrap();
let source = fixture_pdf(with_image);
let descriptor = opaque_descriptor(&source);
let report = ingest::ingest_pdf(&mut store, &descriptor, Limits::DEFAULT).unwrap();
Fixture {
root,
store,
field: report.field,
source,
}
}
}
impl Drop for Fixture {
fn drop(&mut self) {
fs::remove_dir_all(&self.root).ok();
}
}
fn observe_req(
fx: &mut Fixture,
selector: Selector,
representation: Representation,
) -> (FieldAnswer, ObserveStats) {
let req = ObserveRequest::new(selector, representation);
let (answer, stats, _field) =
observe(&mut fx.store, &fx.field, &req, Limits::DEFAULT).unwrap();
(answer, stats)
}
#[test]
fn document_full_document_is_exact() {
let mut fx = Fixture::new("full", false);
let (answer, stats) =
observe_req(&mut fx, Selector::Document, Representation::FullDocument);
assert_eq!(answer.value, AnswerValue::Bytes(fx.source.clone()));
assert_eq!(answer.basis, Basis::DirectlyObserved);
assert_eq!(answer.integrity_scope, IntegrityScope::WholeSource);
assert!(answer.exact);
assert_eq!(answer.source_span, Some((0, fx.source.len() as u64)));
assert_eq!(stats.bytes_returned, fx.source.len() as u64);
}
#[test]
fn byte_range_returns_exact_bytes() {
let mut fx = Fixture::new("range", false);
let (answer, _) = observe_req(
&mut fx,
Selector::ByteRange { offset: 9, len: 8 },
Representation::ExactBytes,
);
assert_eq!(answer.value, AnswerValue::Bytes(fx.source[9..17].to_vec()));
assert_eq!(answer.source_span, Some((9, 17)));
assert!(answer.exact);
assert_eq!(answer.integrity_scope, IntegrityScope::Node);
}
#[test]
fn page_text_is_heuristic_and_nonempty() {
let mut fx = Fixture::new("text", false);
let (answer, _) = observe_req(&mut fx, Selector::Page(1), Representation::Text);
match &answer.value {
AnswerValue::Text(t) => assert!(t.contains("Hello"), "got {t:?}"),
other => panic!("expected text, got {other:?}"),
}
assert_eq!(answer.basis, Basis::Heuristic);
assert!(!answer.exact);
}
#[test]
fn page_structure_procedural_closure_excludes_image_seed_node() {
let mut fx = Fixture::new("structure", true);
let (answer, stats) = observe_req(&mut fx, Selector::Page(1), Representation::Structure);
match &answer.value {
AnswerValue::Json(j) => {
assert!(j.contains("\"page\":1"), "got {j}");
assert!(j.contains("\"content_streams\":[4]"), "got {j}");
}
other => panic!("expected json, got {other:?}"),
}
assert_eq!(answer.basis, Basis::DeterministicallyDerived);
assert!(!answer.exact);
assert!(
stats.seed_bytes_read < 16 * 1024,
"structure observation read {} seed bytes (expected < 16384)",
stats.seed_bytes_read
);
assert!(
stats.descriptor_bytes_read > 0,
"the descriptor read must be accounted, not hidden"
);
assert!(
stats.bytes_read >= stats.descriptor_bytes_read,
"bytes_read {} must include descriptor_bytes_read {}",
stats.bytes_read,
stats.descriptor_bytes_read
);
assert_eq!(
stats.bytes_read,
stats
.descriptor_bytes_read
.saturating_add(stats.manifest_bytes_read)
.saturating_add(stats.index_bytes_read)
.saturating_add(stats.seed_bytes_read),
"bytes_read must be the exact sum of the four physical classes"
);
}
#[test]
fn plan_is_pure_and_deterministic() {
let fx = Fixture::new("plan", false);
let req = ObserveRequest::new(Selector::Page(1), Representation::Text);
let field = Field::open(&fx.store, &fx.field, Limits::DEFAULT).unwrap();
let before = fx.store.seeds().list_nodes().unwrap().len();
let a = plan::plan(field.manifest(), &fx.store, &req).unwrap();
let b = plan::plan(field.manifest(), &fx.store, &req).unwrap();
assert_eq!(a, b);
let after = fx.store.seeds().list_nodes().unwrap().len();
assert_eq!(before, after, "plan must not add seed nodes");
}
#[test]
fn explain_json_keys_and_unsupported_pair() {
use crate::field::explain;
let mut fx = Fixture::new("explain", false);
let req = ObserveRequest::new(Selector::Document, Representation::FullDocument);
let (plan, actual) =
explain::explain_analyze(&mut fx.store, &fx.field, &req, Limits::DEFAULT).unwrap();
assert_eq!(
plan.json,
"{\"selector\":\"document\",\"representation\":\"full\",\"shape\":\"full_materialize\",\"index_reads\":0,\"required_nodes\":1,\"will_materialize\":[\"DocumentExact\"],\"will_not_materialize\":[]}"
);
let actual_json = actual.to_json();
let mut keys = top_level_keys(&actual_json);
keys.sort();
let mut expected = vec![
"basis",
"bytes_read",
"bytes_returned",
"deepened",
"descriptor_bytes_read",
"descriptor_read_mode",
"exact",
"index_bytes_read",
"index_nodes_read",
"manifest_bytes_read",
"seed_bytes_read",
"seed_nodes_fetched",
"seed_nodes_materialized",
"wall_micros",
];
expected.sort_unstable();
assert_eq!(keys, expected, "actual json keys: {actual_json}");
let bad = ObserveRequest::new(Selector::Document, Representation::Text);
let err = observe(&mut fx.store, &fx.field, &bad, Limits::DEFAULT).unwrap_err();
assert_eq!(err.class(), crate::ErrorClass::UnsupportedFeature);
let field = Field::open(&fx.store, &fx.field, Limits::DEFAULT).unwrap();
let perr = plan::plan(field.manifest(), &fx.store, &bad).unwrap_err();
assert_eq!(perr.class(), crate::ErrorClass::UnsupportedFeature);
}
#[test]
fn budget_yields_resource_limit_not_truncation() {
let mut fx = Fixture::new("budget", false);
let req = ObserveRequest {
selector: Selector::Document,
representation: Representation::FullDocument,
budget: ObserveBudget {
max_output_bytes: 4,
max_nodes: 1 << 20,
},
use_cache: true,
};
let err = observe(&mut fx.store, &fx.field, &req, Limits::DEFAULT).unwrap_err();
assert_eq!(err.class(), crate::ErrorClass::ResourceLimit);
}
#[test]
fn find_returns_matching_line() {
let mut fx = Fixture::new("find", false);
let (answer, _) = observe_req(
&mut fx,
Selector::TextMatch("Hello".to_string()),
Representation::Text,
);
match &answer.value {
AnswerValue::Json(j) => {
assert!(j.contains("\"page\":1"), "got {j}");
assert!(j.contains("Hello"), "got {j}");
}
other => panic!("expected json, got {other:?}"),
}
assert_eq!(answer.basis, Basis::Heuristic);
}
#[test]
fn preview_is_deterministic() {
let mut fx = Fixture::new("preview", false);
let (a, _) = observe_req(&mut fx, Selector::Page(1), Representation::Preview);
let (b, _) = observe_req(&mut fx, Selector::Page(1), Representation::Preview);
assert_eq!(a.value, b.value);
match &a.value {
AnswerValue::Bytes(bytes) => {
assert!(String::from_utf8_lossy(bytes).contains("VOLE-PREVIEW v1"))
}
other => panic!("expected preview bytes, got {other:?}"),
}
}
#[test]
fn decoded_and_operators_streams_resolve() {
let mut fx = Fixture::new("decoded", false);
let (decoded, _) = observe_req(&mut fx, Selector::Stream(4), Representation::DecodedBytes);
match &decoded.value {
AnswerValue::Bytes(b) => {
assert_eq!(b, b"BT /F1 12 Tf 72 720 Td (Hello) Tj ET\n");
}
other => panic!("expected bytes, got {other:?}"),
}
assert_eq!(decoded.basis, Basis::DeterministicallyDerived);
let (ops, _) = observe_req(&mut fx, Selector::Stream(4), Representation::Operators);
assert!(matches!(ops.value, AnswerValue::Bytes(ref b) if !b.is_empty()));
}
#[test]
fn observation_byte_classes_sum_and_are_all_charged() {
let mut fx = Fixture::new("io-sum", false);
let (_, stats) = observe_req(&mut fx, Selector::Page(1), Representation::Structure);
assert!(
stats.descriptor_bytes_read > 0,
"descriptor bytes: {stats:?}"
);
assert!(stats.manifest_bytes_read > 0, "manifest bytes: {stats:?}");
assert!(stats.index_bytes_read > 0, "index bytes: {stats:?}");
assert!(stats.seed_bytes_read > 0, "seed bytes: {stats:?}");
assert_eq!(
stats.bytes_read,
stats
.descriptor_bytes_read
.saturating_add(stats.manifest_bytes_read)
.saturating_add(stats.index_bytes_read)
.saturating_add(stats.seed_bytes_read),
"bytes_read must equal the class sum: {stats:?}"
);
}
#[test]
fn explain_analyze_opens_the_descriptor_once() {
use crate::field::explain;
let mut fx = Fixture::new("opens-once", false);
let req = ObserveRequest::new(Selector::Document, Representation::Metadata);
let before = fx.store.io().descriptor_reads();
let (_, actual) =
explain::explain_analyze(&mut fx.store, &fx.field, &req, Limits::DEFAULT).unwrap();
assert_eq!(
fx.store.io().descriptor_reads() - before,
1,
"explain_analyze must open the descriptor exactly once"
);
assert!(
actual.stats.descriptor_bytes_read > 0,
"the single open must still be accounted: {:?}",
actual.stats
);
let page = ObserveRequest::new(Selector::Page(1), Representation::Text);
let before = fx.store.io().descriptor_reads();
let (_, page_actual) =
explain::explain_analyze(&mut fx.store, &fx.field, &page, Limits::DEFAULT).unwrap();
assert_eq!(
fx.store.io().descriptor_reads() - before,
1,
"a cold promotion must not re-open the field for its manifest"
);
assert!(page_actual.stats.deepened, "the cold page must promote");
assert!(page_actual.stats.descriptor_bytes_read > 0);
}
struct BoundedSeedStore {
inner: FsSeedStore,
fetches: Cell<u64>,
limit: u64,
}
impl SeedStore for BoundedSeedStore {
fn put_node(&mut self, canonical: &[u8]) -> Result<NodeId> {
self.inner.put_node(canonical)
}
fn get_node(&self, id: &NodeId) -> Result<Vec<u8>> {
let n = self.fetches.get() + 1;
assert!(
n <= self.limit,
"get_node #{n} exceeds the {}-fetch bound: the store was enumerated",
self.limit
);
self.fetches.set(n);
self.inner.get_node(id)
}
fn get_node_range(&self, id: &NodeId, offset: u64, len: u64) -> Result<Vec<u8>> {
let n = self.fetches.get() + 1;
assert!(
n <= self.limit,
"get_node_range #{n} exceeds the {}-fetch bound: the store was enumerated",
self.limit
);
self.fetches.set(n);
self.inner.get_node_range(id, offset, len)
}
fn contains_node(&self, id: &NodeId) -> Result<bool> {
self.inner.contains_node(id)
}
fn list_nodes(&self) -> Result<Vec<(NodeId, u64)>> {
panic!("an observation must never enumerate the seed store")
}
}
#[test]
fn decoded_stream_resolves_without_enumerating_the_store() {
let mut fx = Fixture::new("no-scan", true);
for i in 0..256u32 {
let decoy = SeedNode::new(
NodeKind::PdfObject,
1,
u32_params(10_000 + i),
Vec::new(),
"decoy",
);
fx.store
.seeds_mut()
.put_node(&decoy.encode_canonical())
.unwrap();
}
let total = fx.store.seeds().list_nodes().unwrap().len() as u64;
assert!(total > 200, "expected many seed nodes, got {total}");
let field = Field::open(&fx.store, &fx.field, Limits::DEFAULT).unwrap();
let io = fx.store.io().handle();
let seeds = CountingSeedStore::new(BoundedSeedStore {
inner: FsSeedStore::open_with_io(fx.store.root(), io.handle()).unwrap(),
fetches: Cell::new(0),
limit: 8,
});
let istore = FsIndexStore::open_with_io(fx.store.root(), io.handle()).unwrap();
let req = ObserveRequest::new(Selector::Stream(4), Representation::DecodedBytes);
let (answer, stats, _) = observe_with_stores(
&mut fx.store,
FieldView::from_field(&field),
&req,
Limits::DEFAULT,
Instant::now(),
seeds,
istore,
)
.unwrap();
match &answer.value {
AnswerValue::Bytes(b) => {
assert_eq!(b, b"BT /F1 12 Tf 72 720 Td (Hello) Tj ET\n")
}
other => panic!("expected bytes, got {other:?}"),
}
assert!(
stats.seed_nodes_fetched <= 8,
"fetched {} seed nodes for one decoded stream",
stats.seed_nodes_fetched
);
assert!(
stats.seed_nodes_fetched < total,
"must not enumerate the {total}-node store; fetched {}",
stats.seed_nodes_fetched
);
}
pub(super) fn top_level_keys(json: &str) -> Vec<String> {
let b = json.as_bytes();
let mut keys = Vec::new();
let mut depth = 0i32;
let mut in_str = false;
let mut esc = false;
let mut i = 0;
while i < b.len() {
let c = b[i];
if in_str {
if esc {
esc = false;
} else if c == b'\\' {
esc = true;
} else if c == b'"' {
in_str = false;
if depth == 1 && b.get(i + 1) == Some(&b':') {
let mut j = i;
while j > 0 {
j -= 1;
if b[j] == b'"' {
keys.push(String::from_utf8_lossy(&b[j + 1..i]).into_owned());
break;
}
}
}
}
} else {
match c {
b'"' => in_str = true,
b'{' | b'[' => depth += 1,
b'}' | b']' => depth -= 1,
_ => {}
}
}
i += 1;
}
keys
}
}