vole_document/materialize/
seek.rs1use std::io::{self, Read, Seek, SeekFrom};
41
42use crate::container::directory::{
43 DirectoryEntry, SECTION_CHANNEL_LENGTHS, SECTION_LOCATORS, SeekDirectory,
44};
45use crate::container::header::{FEATURE_SEEK_DIRECTORY, HEADER_LEN, Header};
46use crate::container::observation::ObservationIndex;
47use crate::container::record::{FLAG_OPTIONAL, RECORD_OVERHEAD, Record, RecordTag, read_record_at};
48use crate::dra::Program;
49use crate::entropy::codec::EntropyChannelDescriptor;
50use crate::entropy::model::EntropyModel;
51use crate::entropy::{CODER_ORDER0_BYTE_RANS, CODER_VERSION_1};
52use crate::error::{Error, Result};
53use crate::limits::Limits;
54use crate::materialize::observation::{
55 ObservationReport, ObservationSelector, ObservationStats, resolve_selector, select_ops,
56 selection_references, serve_selection,
57};
58
59#[derive(Debug)]
68pub struct CountingReader<R> {
69 inner: R,
70 bytes_read: u64,
71 read_calls: u32,
72 seeks: u32,
73}
74
75impl<R> CountingReader<R> {
76 pub fn new(inner: R) -> Self {
78 CountingReader {
79 inner,
80 bytes_read: 0,
81 read_calls: 0,
82 seeks: 0,
83 }
84 }
85
86 pub fn bytes_read(&self) -> u64 {
88 self.bytes_read
89 }
90
91 pub fn read_calls(&self) -> u32 {
93 self.read_calls
94 }
95
96 pub fn seeks(&self) -> u32 {
98 self.seeks
99 }
100
101 pub fn into_inner(self) -> R {
103 self.inner
104 }
105}
106
107impl<R: Read> Read for CountingReader<R> {
108 fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
109 let n = self.inner.read(buf)?;
110 self.bytes_read += n as u64;
111 self.read_calls += 1;
112 Ok(n)
113 }
114}
115
116impl<R: Seek> Seek for CountingReader<R> {
117 fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
118 self.seeks += 1;
119 self.inner.seek(pos)
120 }
121}
122
123pub fn materialize_observation_seeked<R: Read + Seek>(
130 reader: R,
131 selector: ObservationSelector,
132 limits: Limits,
133) -> Result<ObservationReport> {
134 let mut reader = CountingReader::new(reader);
135 let file_len = reader
136 .seek(SeekFrom::End(0))
137 .map_err(|e| Error::io(format!("seek to end failed: {e}")))?;
138
139 let header = read_header(&mut reader)?;
141 if header.optional_features & FEATURE_SEEK_DIRECTORY == 0 {
142 return Err(Error::unsupported_feature(
143 "descriptor has no seek directory (the header does not declare the seek feature)",
144 ));
145 }
146 let dir_rec = read_record_at(&mut reader, HEADER_LEN as u64, limits)?;
147 if dir_rec.tag != RecordTag::Directory as u8 {
148 return Err(Error::unsupported_feature(
149 "descriptor has no seek directory record at offset 64",
150 ));
151 }
152 if dir_rec.flags & FLAG_OPTIONAL == 0 {
153 return Err(Error::invalid_container(
154 "DIRECTORY record must carry FLAG_OPTIONAL",
155 ));
156 }
157 let dir = SeekDirectory::decode(&dir_rec.payload, limits)?;
158 dir.validate_structural(file_len, limits)?;
159 if dir.section_flags & SECTION_LOCATORS == 0 {
160 return Err(Error::unsupported_feature(
161 "seek directory omits the locator section",
162 ));
163 }
164 if dir.entries[0].payload_len != dir_rec.payload.len() as u32 {
165 return Err(Error::invalid_container(
166 "seek directory locator 0 payload length disagrees with the record framing",
167 ));
168 }
169
170 let object_entries = class_entries(&dir, RecordTag::Object)?;
171 let channel_entries = class_entries(&dir, RecordTag::EntropyChannel)?;
172 let model_entries = class_entries(&dir, RecordTag::Model)?;
173 if !class_entries(&dir, RecordTag::ExternalRef)?.is_empty() {
174 return Err(Error::unsupported_feature(
175 "seek-based partial read cannot resolve external objects",
176 ));
177 }
178 if !channel_entries.is_empty() && dir.section_flags & SECTION_CHANNEL_LENGTHS == 0 {
179 return Err(Error::unsupported_feature(
180 "seek directory omits the channel-lengths section",
181 ));
182 }
183 if dir.channel_lengths.len() != channel_entries.len() {
184 return Err(Error::invalid_container(
185 "seek directory channel-length table disagrees with the channel locators",
186 ));
187 }
188 let object_lens: Vec<u64> = object_entries
189 .iter()
190 .map(|e| u64::from(e.payload_len))
191 .collect();
192 let channel_lens: Vec<u64> = dir.channel_lengths.clone();
193
194 let graph_site = class_entries(&dir, RecordTag::Graph)?
197 .first()
198 .ok_or_else(|| {
199 Error::unsupported_feature("seek directory does not locate a GRAPH record")
200 })?;
201 let index_site = class_entries(&dir, RecordTag::ObservationIndex)?
202 .first()
203 .ok_or_else(|| {
204 Error::unsupported_feature("seek directory does not locate an OBSERVATION_INDEX record")
205 })?;
206 let integrity_site = class_entries(&dir, RecordTag::Integrity)?
207 .first()
208 .ok_or_else(|| {
209 Error::unsupported_feature("seek directory does not locate an INTEGRITY record")
210 })?;
211
212 let graph_rec = read_checked(&mut reader, graph_site, limits)?;
213 let program = Program::decode(&graph_rec.payload, limits)?;
214 let index_rec = read_checked(&mut reader, index_site, limits)?;
215 let index = ObservationIndex::decode(&index_rec.payload, limits)?;
216 let integrity_rec = read_checked(&mut reader, integrity_site, limits)?;
217 if integrity_rec.payload.len() != 40 {
218 return Err(Error::invalid_container(
219 "INTEGRITY payload must be 40 bytes",
220 ));
221 }
222 let declared_len = u64::from_le_bytes([
223 integrity_rec.payload[32],
224 integrity_rec.payload[33],
225 integrity_rec.payload[34],
226 integrity_rec.payload[35],
227 integrity_rec.payload[36],
228 integrity_rec.payload[37],
229 integrity_rec.payload[38],
230 integrity_rec.payload[39],
231 ]);
232 if declared_len != header.declared_source_len {
233 return Err(Error::integrity_mismatch(format!(
234 "INTEGRITY length {declared_len} disagrees with header {}",
235 header.declared_source_len
236 )));
237 }
238
239 index.validate(&program, &object_lens, &channel_lens, limits)?;
242
243 let (a, b) = resolve_selector(&index, selector, declared_len)?;
245 let window = select_ops(&program, &object_lens, &channel_lens, a, b, limits)?;
246 let (objects_used, channels_used) =
247 selection_references(&window.ops, object_entries.len(), channel_entries.len());
248
249 let mut objects: Vec<Vec<u8>> = vec![Vec::new(); object_entries.len()];
251 for (id, used) in objects_used.iter().enumerate() {
252 if *used {
253 objects[id] = read_checked(&mut reader, &object_entries[id], limits)?.payload;
254 }
255 }
256
257 let mut channels: Vec<EntropyChannelDescriptor> =
260 vec![placeholder_channel(); channel_entries.len()];
261 let mut models_needed = vec![false; model_entries.len()];
262 for (id, used) in channels_used.iter().enumerate() {
263 if *used {
264 let rec = read_checked(&mut reader, &channel_entries[id], limits)?;
265 let channel = EntropyChannelDescriptor::decode(&rec.payload, limits)?;
266 if channel.decoded_length != channel_lens[id] {
267 return Err(Error::invalid_container(format!(
268 "seek directory channel-length {id} disagrees with the channel record"
269 )));
270 }
271 if channel.model_id as usize >= model_entries.len() {
272 return Err(Error::invalid_model(format!(
273 "entropy channel {id} references missing model {}",
274 channel.model_id
275 )));
276 }
277 models_needed[channel.model_id as usize] = true;
278 channels[id] = channel;
279 }
280 }
281 let mut models: Vec<EntropyModel> = vec![placeholder_model(); model_entries.len()];
282 for (id, used) in models_needed.iter().enumerate() {
283 if *used {
284 let rec = read_checked(&mut reader, &model_entries[id], limits)?;
285 models[id] = EntropyModel::decode(&rec.payload)?;
286 }
287 }
288
289 let served = serve_selection(
291 &objects,
292 &channels,
293 &models,
294 window,
295 &objects_used,
296 &channels_used,
297 a,
298 b,
299 limits,
300 )?;
301
302 let descriptor_bytes_traversed = graph_rec.payload.len() as u64
305 + index_rec.payload.len() as u64
306 + RECORD_OVERHEAD as u64
307 + served.referenced_object_bytes
308 + served.referenced_channel_bytes;
309
310 let stats = ObservationStats {
311 ops_evaluated: served.ops_evaluated,
312 ops_total: served.ops_total,
313 objects_fetched: served.objects_fetched,
314 objects_total: object_entries.len(),
315 channels_decoded: served.channels_decoded,
316 channels_total: channel_entries.len(),
317 entropy_bytes_decoded: served.entropy_bytes_decoded,
318 descriptor_bytes_traversed,
319 output_bytes: served.bytes.len() as u64,
320 bytes_read: reader.bytes_read(),
322 integrity_verified: false,
323 };
324
325 Ok(ObservationReport {
326 range: (a, b),
327 bytes: served.bytes,
328 stats,
329 })
330}
331
332fn read_header<R: Read + Seek>(reader: &mut R) -> Result<Header> {
334 reader
335 .seek(SeekFrom::Start(0))
336 .map_err(|e| Error::io(format!("seek to header failed: {e}")))?;
337 let mut buf = [0u8; HEADER_LEN];
338 reader
339 .read_exact(&mut buf)
340 .map_err(|_| Error::invalid_container("truncated header"))?;
341 Header::decode(&buf)
342}
343
344fn read_checked<R: Read + Seek>(
346 reader: &mut R,
347 site: &DirectoryEntry,
348 limits: Limits,
349) -> Result<Record> {
350 let rec = read_record_at(reader, site.offset, limits)?;
351 if rec.tag != site.tag || rec.payload.len() as u64 != u64::from(site.payload_len) {
352 return Err(Error::invalid_container(format!(
353 "DIRECTORY locator for tag {:#04x} disagrees with the record framing at offset {}",
354 site.tag, site.offset
355 )));
356 }
357 Ok(rec)
358}
359
360fn class_entries(dir: &SeekDirectory, tag: RecordTag) -> Result<&[DirectoryEntry]> {
364 let (first, count) = match dir.classes.iter().find(|c| c.tag == tag as u8) {
365 Some(c) => (c.first as usize, c.count as usize),
366 None => (0, 0),
367 };
368 let scan_first = dir.entries.iter().position(|e| e.tag == tag as u8);
369 let scan_count = dir.entries.iter().filter(|e| e.tag == tag as u8).count();
370 let expected_first = if count == 0 { None } else { Some(first) };
371 if scan_first != expected_first || scan_count != count {
372 return Err(Error::invalid_container(
373 "seek directory class index disagrees with the locator table",
374 ));
375 }
376 if count == 0 {
377 return Ok(&[]);
378 }
379 let end = first
380 .checked_add(count)
381 .ok_or_else(|| Error::invalid_container("seek directory class range overflow"))?;
382 dir.entries
383 .get(first..end)
384 .ok_or_else(|| Error::invalid_container("seek directory class range out of bounds"))
385}
386
387fn placeholder_channel() -> EntropyChannelDescriptor {
391 EntropyChannelDescriptor {
392 coder: CODER_ORDER0_BYTE_RANS,
393 coder_version: CODER_VERSION_1,
394 scale_bits: 0,
395 lane_count: 1,
396 model_id: 0,
397 symbol_count: 0,
398 decoded_length: 0,
399 initial_state: 0,
400 payload: Vec::new(),
401 }
402}
403
404fn placeholder_model() -> EntropyModel {
406 EntropyModel {
407 scale_bits: 0,
408 frequencies: Vec::new(),
409 }
410}