1use std::io;
2use std::sync::Arc;
3
4use crate::geometry::{Box2D, Box3D, Overlaps2D, Overlaps3D};
5use crate::persistence::{
6 CHUNK_ENTRY_LEN, CHUNK_FLAG_CRITICAL, FORMAT_VERSION, LoadError, PFIX_DESC_LEN, PYLD_DESC_LEN,
7 PYLD_DESC_LEN_FIXED, SUPERBLOCK_LEN, TAG_PFIX, TAG_PYLD, TAG_TREE, TREE_DESC_LEN,
8 derive_level_bounds, expected_tree_shape, parse_pfix_chunk, parse_pyld_chunk, parse_tree_chunk,
9 read_u32_at, read_u64_at, read_u64_le_unchecked,
10};
11
12use super::core::{align8_u64, checked_directory_span};
13use super::directory::directory_start;
14use super::limits::{Budget, directory_node_budget};
15use super::payload::{
16 PayloadSection, PrefixSection, emit_run_payloads, emit_run_payloads_fixed, payload_blob_span,
17 payload_run_end, payload_run_end_fixed,
18};
19use super::planner::{apply_gather_run, expand_frontier, plan_gather};
20use super::{
21 PayloadPrefix, StreamCore, StreamError, StreamIndex2D, StreamIndex2DF32, StreamIndex3D,
22 StreamIndex3DF32, StreamLimits, parse_box2d, parse_box2d_f32, parse_box3d, parse_box3d_f32,
23 read_index,
24};
25
26#[cfg(feature = "async")]
44#[allow(async_fn_in_trait, clippy::len_without_is_empty)]
45pub trait AsyncRangeReader {
46 async fn read_exact_at(&self, offset: u64, buf: &mut [u8]) -> io::Result<()>;
48
49 fn len(&self) -> Option<u64> {
51 None
52 }
53}
54
55#[cfg(feature = "async")]
57#[derive(Clone, Copy, PartialEq, Eq)]
58enum Want {
59 Ids,
60 Payloads,
61}
62
63#[cfg(feature = "async")]
64impl<R: AsyncRangeReader> StreamCore<R> {
65 async fn open_async(
66 reader: R,
67 dimensions: usize,
68 coord_bytes: usize,
69 limits: StreamLimits,
70 ) -> Result<Self, StreamError> {
71 let mut head = [0u8; SUPERBLOCK_LEN];
72 reader.read_exact_at(0, &mut head).await?;
73 if &head[..8] != b"PSINDEX\0" {
74 return Err(StreamError::Format(LoadError::BadMagic));
75 }
76 if u64::from_le_bytes(head[8..16].try_into().unwrap()) != FORMAT_VERSION {
77 return Err(StreamError::Format(LoadError::UnsupportedVersion));
78 }
79 let chunk_count = read_u32_at(&head, 16)? as usize;
80 let file_len = reader.len();
81 let (dir_len, dir_end) = checked_directory_span(chunk_count, file_len)?;
82 let mut dir = vec![0u8; dir_len];
83 reader
84 .read_exact_at(SUPERBLOCK_LEN as u64, &mut dir)
85 .await?;
86
87 let mut max_end = dir_end;
88 let mut tree: Option<(u64, u64)> = None;
89 let mut pyld: Option<(u64, u64)> = None;
90 let mut pfix: Option<(u64, u64)> = None;
91 for i in 0..chunk_count {
92 let base = i * CHUNK_ENTRY_LEN;
93 let mut tag = [0u8; 4];
94 tag.copy_from_slice(&dir[base..base + 4]);
95 let flags = read_u32_at(&dir, base + 4)?;
96 let offset = read_u64_at(&dir, base + 8)?;
97 let len = read_u64_at(&dir, base + 16)?;
98 let end = offset.checked_add(len).ok_or(LoadError::IntegerOverflow)?;
99 if file_len.is_some_and(|fl| end > fl) {
100 return Err(StreamError::Format(LoadError::InvalidTree));
101 }
102 max_end = max_end.max(end);
103 if tag == TAG_TREE {
104 tree = Some((offset, len));
105 } else if tag == TAG_PYLD {
106 pyld = Some((offset, len));
107 } else if tag == TAG_PFIX {
108 pfix = Some((offset, len));
109 } else if flags & CHUNK_FLAG_CRITICAL != 0 {
110 return Err(StreamError::Format(LoadError::UnsupportedVersion));
111 }
112 }
113
114 let aligned_end = align8_u64(max_end)?;
117 if let Some(fl) = file_len
118 && fl > aligned_end
119 {
120 return Err(StreamError::Format(LoadError::LengthMismatch {
121 expected: max_end as usize,
122 actual: fl as usize,
123 }));
124 }
125 let (toff, tlen) = tree.ok_or(LoadError::InvalidTree)?;
126 if tlen < TREE_DESC_LEN as u64 {
127 return Err(StreamError::Format(LoadError::Truncated));
128 }
129 let mut desc = [0u8; TREE_DESC_LEN];
130 reader.read_exact_at(toff, &mut desc).await?;
131 let (td, _) = parse_tree_chunk(&desc)?;
132 if td.dimensions != dimensions || td.coord_bytes != coord_bytes {
133 return Err(StreamError::Format(LoadError::UnsupportedVersion));
134 }
135 let (num_nodes, level_count) = expected_tree_shape(td.num_items, td.node_size)?;
136 let record = dimensions
137 .checked_mul(2 * coord_bytes)
138 .ok_or(LoadError::IntegerOverflow)?;
139 let box_stride = if td.interleaved { record + 8 } else { record };
140 let box0 = toff + td.desc_len as u64;
141 let node_len = num_nodes
142 .checked_mul(box_stride + if td.interleaved { 0 } else { 8 })
143 .ok_or(LoadError::IntegerOverflow)?;
144 if tlen != td.desc_len as u64 + node_len as u64 {
145 return Err(StreamError::Format(LoadError::InvalidTree));
146 }
147 let idx0 = if td.interleaved {
148 box0
149 } else {
150 box0 + (num_nodes * record) as u64
151 };
152 let level_bounds = derive_level_bounds(td.num_items, td.node_size, level_count);
153
154 let payload = match pyld {
155 Some((poff, plen)) => {
156 if plen < PYLD_DESC_LEN as u64 {
157 return Err(StreamError::Format(LoadError::Truncated));
158 }
159 let dn = (PYLD_DESC_LEN_FIXED as u64).min(plen) as usize;
160 let mut pd = [0u8; PYLD_DESC_LEN_FIXED];
161 reader.read_exact_at(poff, &mut pd[..dn]).await?;
162 let (pdesc, _) = parse_pyld_chunk(&pd[..dn])?;
163 let body0 = poff + pdesc.desc_len as u64;
164 if pdesc.record_stride != 0 {
165 let stride = pdesc.record_stride as u64;
166 let blob_total = (td.num_items as u64)
167 .checked_mul(stride)
168 .ok_or(StreamError::Format(LoadError::IntegerOverflow))?;
169 let need = pdesc.desc_len as u64 + blob_total;
170 if plen != need {
171 return Err(StreamError::Format(LoadError::InvalidTree));
172 }
173 Some(PayloadSection {
174 offsets_start: 0,
175 blobs_start: body0,
176 blob_total,
177 stride,
178 })
179 } else {
180 let offsets_start = body0;
181 let last_at = offsets_start + (td.num_items as u64) * 8;
182 let mut last = [0u8; 8];
183 reader.read_exact_at(last_at, &mut last).await?;
184 let blob_total = u64::from_le_bytes(last);
185 let blobs_start = offsets_start + (td.num_items as u64 + 1) * 8;
186 let need = pdesc.desc_len as u64 + (td.num_items as u64 + 1) * 8 + blob_total;
187 if plen != need {
188 return Err(StreamError::Format(LoadError::InvalidTree));
189 }
190 Some(PayloadSection {
191 offsets_start,
192 blobs_start,
193 blob_total,
194 stride: 0,
195 })
196 }
197 }
198 None => None,
199 };
200
201 let prefix = match pfix {
205 Some((poff, plen)) if payload.is_some() => {
206 if plen < PFIX_DESC_LEN as u64 {
207 return Err(StreamError::Format(LoadError::Truncated));
208 }
209 let mut pd = [0u8; PFIX_DESC_LEN];
210 reader.read_exact_at(poff, &mut pd).await?;
211 let desc = parse_pfix_chunk(&pd)?;
212 let need = desc.desc_len as u64 + td.num_items as u64 * desc.record_stride as u64;
213 if plen != need {
214 return Err(StreamError::Format(LoadError::InvalidTree));
215 }
216 Some(PrefixSection {
217 start: poff + desc.desc_len as u64,
218 stride: desc.record_stride,
219 })
220 }
221 _ => None,
222 };
223
224 let budget = directory_node_budget(&limits, box_stride, td.interleaved);
226 let dir_node_start = directory_start(&level_bounds, level_count, budget);
227 let cached_nodes = num_nodes - dir_node_start;
228 let mut dir_boxes = vec![0u8; cached_nodes * box_stride];
229 if !dir_boxes.is_empty() {
230 let offset = box0 + (dir_node_start * box_stride) as u64;
231 reader.read_exact_at(offset, &mut dir_boxes).await?;
232 }
233 let mut dir_indices = if td.interleaved {
234 Vec::new()
235 } else {
236 vec![0u8; cached_nodes * 8]
237 };
238 if !dir_indices.is_empty() {
239 let offset = idx0 + (dir_node_start * 8) as u64;
240 reader.read_exact_at(offset, &mut dir_indices).await?;
241 }
242 let dir_boxes: Arc<[u8]> = dir_boxes.into();
243 let dir_indices: Arc<[u8]> = dir_indices.into();
244 Ok(StreamCore {
245 reader,
246 node_size: td.node_size,
247 num_items: td.num_items,
248 num_nodes,
249 level_count,
250 level_bounds,
251 record,
252 box_stride,
253 interleaved: td.interleaved,
254 box0,
255 idx0,
256 dir_node_start,
257 dir_boxes,
258 dir_indices,
259 payload,
260 prefix,
261 limits,
262 })
263 }
264
265 async fn gather_async(
270 &self,
271 positions: &[usize],
272 section0: u64,
273 stride: usize,
274 cache: &[u8],
275 out: &mut Vec<u8>,
276 budget: &mut Budget,
277 ) -> Result<(), StreamError> {
278 let runs = plan_gather(
279 positions,
280 section0,
281 stride,
282 self.dir_node_start,
283 cache,
284 out,
285 self.coalesce_gap(),
286 );
287 for run in &runs {
288 budget.charge_read(run.len)?;
289 }
290 let mut bufs: Vec<Vec<u8>> = runs.iter().map(|run| vec![0u8; run.len]).collect();
291 let reads = runs
292 .iter()
293 .zip(bufs.iter_mut())
294 .map(|(run, buf)| self.reader.read_exact_at(run.offset, buf.as_mut_slice()));
295 futures_util::future::try_join_all(reads).await?;
296 for (run, buf) in runs.iter().zip(&bufs) {
297 apply_gather_run(out, run, buf, stride);
298 }
299 Ok(())
300 }
301
302 async fn gather_payloads_async<F>(
307 &self,
308 section: &PayloadSection,
309 leaf_positions: &[usize],
310 indices: &[u8],
311 budget: &mut Budget,
312 sink: &mut F,
313 ) -> Result<(), StreamError>
314 where
315 F: FnMut(usize, &[u8]),
316 {
317 let mut runs: Vec<(usize, usize)> = Vec::new();
319 let mut j = 0;
320 while j < leaf_positions.len() {
321 let k = payload_run_end(leaf_positions, j, self.coalesce_gap());
322 runs.push((j, k));
323 j = k + 1;
324 }
325
326 let mut off_bufs: Vec<Vec<u8>> = runs
328 .iter()
329 .map(|&(j, k)| vec![0u8; (leaf_positions[k] + 2 - leaf_positions[j]) * 8])
330 .collect();
331 for buf in &off_bufs {
332 budget.charge_read(buf.len())?;
333 }
334 let off_reads = runs.iter().zip(off_bufs.iter_mut()).map(|(&(j, _), buf)| {
335 let lo = leaf_positions[j];
336 self.reader
337 .read_exact_at(section.offsets_start + (lo * 8) as u64, buf.as_mut_slice())
338 });
339 futures_util::future::try_join_all(off_reads).await?;
340
341 let mut spans = Vec::with_capacity(runs.len());
343 for (&(j, k), off_buf) in runs.iter().zip(&off_bufs) {
344 spans.push(payload_blob_span(
345 off_buf,
346 leaf_positions[j],
347 leaf_positions[k],
348 section.blob_total,
349 )?);
350 }
351
352 let mut blob_bufs: Vec<Vec<u8>> = spans
354 .iter()
355 .map(|&(lo, hi)| vec![0u8; (hi - lo) as usize])
356 .collect();
357 for buf in &blob_bufs {
358 if !buf.is_empty() {
359 budget.charge_read(buf.len())?;
360 }
361 }
362 let blob_reads = spans
363 .iter()
364 .zip(blob_bufs.iter_mut())
365 .map(|(&(lo, _), buf)| {
366 self.reader
367 .read_exact_at(section.blobs_start + lo, buf.as_mut_slice())
368 });
369 futures_util::future::try_join_all(blob_reads).await?;
370
371 for ((&(j, k), off_buf), (&(blob_lo, blob_hi), blob_buf)) in
373 runs.iter().zip(&off_bufs).zip(spans.iter().zip(&blob_bufs))
374 {
375 emit_run_payloads(
376 leaf_positions,
377 indices,
378 j,
379 k,
380 leaf_positions[j],
381 off_buf,
382 blob_lo,
383 blob_hi,
384 blob_buf,
385 self.num_items,
386 budget,
387 sink,
388 )?;
389 }
390 Ok(())
391 }
392
393 async fn gather_payloads_fixed_async<F>(
397 &self,
398 section: &PayloadSection,
399 leaf_positions: &[usize],
400 indices: &[u8],
401 budget: &mut Budget,
402 sink: &mut F,
403 ) -> Result<(), StreamError>
404 where
405 F: FnMut(usize, &[u8]),
406 {
407 let stride = section.stride as usize;
408 let mut runs: Vec<(usize, usize)> = Vec::new();
409 let mut j = 0;
410 while j < leaf_positions.len() {
411 let k = payload_run_end_fixed(leaf_positions, j, stride, self.coalesce_gap());
412 runs.push((j, k));
413 j = k + 1;
414 }
415
416 let mut blob_bufs: Vec<Vec<u8>> = runs
417 .iter()
418 .map(|&(j, k)| vec![0u8; (leaf_positions[k] + 1 - leaf_positions[j]) * stride])
419 .collect();
420 for buf in &blob_bufs {
421 budget.charge_read(buf.len())?;
422 }
423 let reads = runs.iter().zip(blob_bufs.iter_mut()).map(|(&(j, _), buf)| {
424 let lo = leaf_positions[j];
425 self.reader.read_exact_at(
426 section.blobs_start + (lo * stride) as u64,
427 buf.as_mut_slice(),
428 )
429 });
430 futures_util::future::try_join_all(reads).await?;
431
432 for (&(j, k), blob_buf) in runs.iter().zip(&blob_bufs) {
433 emit_run_payloads_fixed(
434 leaf_positions,
435 indices,
436 j,
437 k,
438 leaf_positions[j],
439 stride,
440 blob_buf,
441 self.num_items,
442 budget,
443 sink,
444 )?;
445 }
446 Ok(())
447 }
448
449 async fn visit_payload_prefixes_async<O, F>(
450 &self,
451 overlaps: O,
452 prefix_len: usize,
453 mut emit: F,
454 ) -> Result<(), StreamError>
455 where
456 O: Fn(&[u8]) -> bool,
457 F: FnMut(PayloadPrefix<'_>),
458 {
459 let section = self.payload.as_ref().ok_or(StreamError::NoPayload)?;
460 if self.num_items == 0 {
461 return Ok(());
462 }
463
464 let mut budget = Budget::new(self.limits);
465 let mut frontier = vec![self.num_nodes - 1];
466 let mut level = self.level_count - 1;
467 let mut boxes = Vec::new();
468 let mut indices = Vec::new();
469 let mut survivors: Vec<usize> = Vec::new();
470
471 loop {
472 self.gather_async(
473 &frontier,
474 self.box0,
475 self.box_stride,
476 &self.dir_boxes,
477 &mut boxes,
478 &mut budget,
479 )
480 .await?;
481 survivors.clear();
482 indices.clear();
483 for (i, &pos) in frontier.iter().enumerate() {
484 let slot = i * self.box_stride;
485 if overlaps(&boxes[slot..slot + self.record]) {
486 survivors.push(pos);
487 if self.interleaved {
488 indices
489 .extend_from_slice(&boxes[slot + self.record..slot + self.record + 8]);
490 }
491 }
492 }
493 if survivors.is_empty() {
494 return Ok(());
495 }
496
497 if !self.interleaved {
498 self.gather_async(
499 &survivors,
500 self.idx0,
501 8,
502 &self.dir_indices,
503 &mut indices,
504 &mut budget,
505 )
506 .await?;
507 }
508
509 if level == 0 {
510 self.gather_payload_prefixes_async(
511 section,
512 &survivors,
513 &indices,
514 prefix_len,
515 &mut budget,
516 &mut emit,
517 )
518 .await?;
519 return Ok(());
520 }
521
522 frontier = expand_frontier(
523 &self.level_bounds,
524 self.node_size,
525 level,
526 survivors.len(),
527 &indices,
528 )?;
529 level -= 1;
530 }
531 }
532
533 async fn gather_payload_prefixes_async<F>(
534 &self,
535 section: &PayloadSection,
536 leaf_positions: &[usize],
537 indices: &[u8],
538 prefix_len: usize,
539 budget: &mut Budget,
540 emit: &mut F,
541 ) -> Result<(), StreamError>
542 where
543 F: FnMut(PayloadPrefix<'_>),
544 {
545 let mut spans = Vec::with_capacity(leaf_positions.len());
546 if section.stride != 0 {
547 let stride = section.stride as usize;
548 for (i, &p) in leaf_positions.iter().enumerate() {
549 spans.push(AsyncPrefixSpan {
550 run_index: i,
551 leaf_rank: p,
552 blob_start: (p * stride) as u64,
553 payload_len: stride,
554 });
555 }
556 } else {
557 let mut runs: Vec<(usize, usize)> = Vec::new();
558 let mut j = 0;
559 while j < leaf_positions.len() {
560 let k = payload_run_end(leaf_positions, j, self.coalesce_gap());
561 runs.push((j, k));
562 j = k + 1;
563 }
564
565 let mut off_bufs: Vec<Vec<u8>> = runs
566 .iter()
567 .map(|&(j, k)| vec![0u8; (leaf_positions[k] + 2 - leaf_positions[j]) * 8])
568 .collect();
569 for buf in &off_bufs {
570 budget.charge_read(buf.len())?;
571 }
572 let off_reads = runs.iter().zip(off_bufs.iter_mut()).map(|(&(j, _), buf)| {
573 let lo = leaf_positions[j];
574 self.reader
575 .read_exact_at(section.offsets_start + (lo * 8) as u64, buf.as_mut_slice())
576 });
577 futures_util::future::try_join_all(off_reads).await?;
578
579 for (&(j, k), off_buf) in runs.iter().zip(&off_bufs) {
580 let lo = leaf_positions[j];
581 for (offset, &p) in leaf_positions[j..=k].iter().enumerate() {
582 let o0 = read_u64_le_unchecked(off_buf, (p - lo) * 8);
583 let o1 = read_u64_le_unchecked(off_buf, (p + 1 - lo) * 8);
584 if o1 < o0 || o1 > section.blob_total {
585 return Err(StreamError::Format(LoadError::InvalidTree));
586 }
587 spans.push(AsyncPrefixSpan {
588 run_index: j + offset,
589 leaf_rank: p,
590 blob_start: o0,
591 payload_len: (o1 - o0) as usize,
592 });
593 }
594 }
595 }
596
597 if prefix_len == 0 {
598 for span in &spans {
599 let id = read_index(indices, span.run_index)?;
600 if id >= self.num_items {
601 return Err(StreamError::Format(LoadError::InvalidTree));
602 }
603 budget.charge_item()?;
604 emit(PayloadPrefix {
605 id,
606 leaf_rank: span.leaf_rank,
607 prefix: &[],
608 payload_len: span.payload_len,
609 });
610 }
611 return Ok(());
612 }
613
614 if let Some(pfix) = self.prefix.as_ref().filter(|p| prefix_len <= p.stride) {
619 let stride = pfix.stride;
620 let ranks: Vec<usize> = spans.iter().map(|s| s.leaf_rank).collect();
621 let mut runs: Vec<(usize, usize)> = Vec::new();
622 let mut j = 0;
623 while j < ranks.len() {
624 let k = payload_run_end_fixed(&ranks, j, stride, self.coalesce_gap());
625 runs.push((j, k));
626 j = k + 1;
627 }
628 let mut bufs: Vec<Vec<u8>> = runs
629 .iter()
630 .map(|&(j, k)| vec![0u8; (ranks[k] + 1 - ranks[j]) * stride])
631 .collect();
632 for buf in &bufs {
633 budget.charge_read(buf.len())?;
634 }
635 let reads = runs.iter().zip(bufs.iter_mut()).map(|(&(j, _), buf)| {
636 self.reader
637 .read_exact_at(pfix.start + (ranks[j] * stride) as u64, buf.as_mut_slice())
638 });
639 futures_util::future::try_join_all(reads).await?;
640
641 for (&(j, k), read_buf) in runs.iter().zip(&bufs) {
642 let lo = ranks[j];
643 for span in &spans[j..=k] {
644 let id = read_index(indices, span.run_index)?;
645 if id >= self.num_items {
646 return Err(StreamError::Format(LoadError::InvalidTree));
647 }
648 budget.charge_item()?;
649 let at = (span.leaf_rank - lo) * stride;
650 let take = span.payload_len.min(prefix_len);
651 emit(PayloadPrefix {
652 id,
653 leaf_rank: span.leaf_rank,
654 prefix: &read_buf[at..at + take],
655 payload_len: span.payload_len,
656 });
657 }
658 }
659 return Ok(());
660 }
661
662 let mut prefix_runs: Vec<(usize, usize, u64, u64)> = Vec::new();
663 let gap = self.prefix_coalesce_gap(prefix_len);
664 let mut j = 0;
665 while j < spans.len() {
666 let run_start = spans[j].blob_start;
667 let mut run_end = spans[j].prefix_end(prefix_len);
668 let mut k = j;
669 while k + 1 < spans.len() {
670 let next = &spans[k + 1];
671 if next.blob_start < run_start || next.blob_start.saturating_sub(run_end) > gap {
672 break;
673 }
674 run_end = run_end.max(next.prefix_end(prefix_len));
675 k += 1;
676 }
677 prefix_runs.push((j, k, run_start, run_end));
678 j = k + 1;
679 }
680
681 let mut bufs: Vec<Vec<u8>> = prefix_runs
682 .iter()
683 .map(|&(_, _, start, end)| vec![0u8; (end - start) as usize])
684 .collect();
685 for buf in &bufs {
686 if !buf.is_empty() {
687 budget.charge_read(buf.len())?;
688 }
689 }
690 let reads = prefix_runs
691 .iter()
692 .zip(bufs.iter_mut())
693 .filter(|(_, buf)| !buf.is_empty())
694 .map(|(&(_, _, start, _), buf)| {
695 self.reader
696 .read_exact_at(section.blobs_start + start, buf.as_mut_slice())
697 });
698 futures_util::future::try_join_all(reads).await?;
699
700 for (&(j, k, run_start, _), read_buf) in prefix_runs.iter().zip(&bufs) {
701 for span in &spans[j..=k] {
702 let id = read_index(indices, span.run_index)?;
703 if id >= self.num_items {
704 return Err(StreamError::Format(LoadError::InvalidTree));
705 }
706 budget.charge_item()?;
707 let at = (span.blob_start - run_start) as usize;
708 let take = span.payload_len.min(prefix_len);
709 emit(PayloadPrefix {
710 id,
711 leaf_rank: span.leaf_rank,
712 prefix: &read_buf[at..at + take],
713 payload_len: span.payload_len,
714 });
715 }
716 }
717 Ok(())
718 }
719
720 async fn visit_payloads_at_ranks_async<F>(
721 &self,
722 leaf_ranks: &[usize],
723 mut emit: F,
724 ) -> Result<(), StreamError>
725 where
726 F: FnMut(usize, &[u8]),
727 {
728 let section = self.payload.as_ref().ok_or(StreamError::NoPayload)?;
729 let mut ranks = leaf_ranks.to_vec();
730 ranks.sort_unstable();
731 ranks.dedup();
732 if ranks.last().is_some_and(|&max| max >= self.num_items) {
733 return Err(StreamError::InvalidRank);
734 }
735 let mut budget = Budget::new(self.limits);
736
737 if section.stride != 0 {
738 let stride = section.stride as usize;
739 let mut runs: Vec<(usize, usize)> = Vec::new();
740 let mut j = 0;
741 while j < ranks.len() {
742 let k = payload_run_end_fixed(&ranks, j, stride, self.coalesce_gap());
743 runs.push((j, k));
744 j = k + 1;
745 }
746 let mut bufs: Vec<Vec<u8>> = runs
747 .iter()
748 .map(|&(j, k)| vec![0u8; (ranks[k] + 1 - ranks[j]) * stride])
749 .collect();
750 for buf in &bufs {
751 budget.charge_read(buf.len())?;
752 }
753 let reads = runs.iter().zip(bufs.iter_mut()).map(|(&(j, _), buf)| {
754 let lo = ranks[j];
755 self.reader.read_exact_at(
756 section.blobs_start + (lo * stride) as u64,
757 buf.as_mut_slice(),
758 )
759 });
760 futures_util::future::try_join_all(reads).await?;
761 for (&(j, k), buf) in runs.iter().zip(&bufs) {
762 let lo = ranks[j];
763 for &p in &ranks[j..=k] {
764 budget.charge_item()?;
765 let within = (p - lo) * stride;
766 emit(p, &buf[within..within + stride]);
767 }
768 }
769 return Ok(());
770 }
771
772 let mut runs: Vec<(usize, usize)> = Vec::new();
773 let mut j = 0;
774 while j < ranks.len() {
775 let k = payload_run_end(&ranks, j, self.coalesce_gap());
776 runs.push((j, k));
777 j = k + 1;
778 }
779
780 let mut off_bufs: Vec<Vec<u8>> = runs
781 .iter()
782 .map(|&(j, k)| vec![0u8; (ranks[k] + 2 - ranks[j]) * 8])
783 .collect();
784 for buf in &off_bufs {
785 budget.charge_read(buf.len())?;
786 }
787 let off_reads = runs.iter().zip(off_bufs.iter_mut()).map(|(&(j, _), buf)| {
788 let lo = ranks[j];
789 self.reader
790 .read_exact_at(section.offsets_start + (lo * 8) as u64, buf.as_mut_slice())
791 });
792 futures_util::future::try_join_all(off_reads).await?;
793
794 let mut blob_spans = Vec::with_capacity(ranks.len());
795 for (&(j, k), off_buf) in runs.iter().zip(&off_bufs) {
796 let lo = ranks[j];
797 for &rank in &ranks[j..=k] {
798 let o0 = read_u64_le_unchecked(off_buf, (rank - lo) * 8);
799 let o1 = read_u64_le_unchecked(off_buf, (rank + 1 - lo) * 8);
800 if o1 < o0 || o1 > section.blob_total {
801 return Err(StreamError::Format(LoadError::InvalidTree));
802 }
803 blob_spans.push(AsyncRankBlobSpan {
804 rank,
805 blob_start: o0,
806 blob_end: o1,
807 });
808 }
809 }
810
811 let mut blob_runs: Vec<(usize, usize, u64, u64)> = Vec::new();
812 let gap = self.coalesce_gap();
813 let mut j = 0;
814 while j < blob_spans.len() {
815 let run_start = blob_spans[j].blob_start;
816 let mut run_end = blob_spans[j].blob_end;
817 let mut k = j;
818 while k + 1 < blob_spans.len() {
819 let next = &blob_spans[k + 1];
820 if next.blob_start < run_start || next.blob_start.saturating_sub(run_end) > gap {
821 break;
822 }
823 run_end = run_end.max(next.blob_end);
824 k += 1;
825 }
826 blob_runs.push((j, k, run_start, run_end));
827 j = k + 1;
828 }
829
830 let mut blob_bufs: Vec<Vec<u8>> = blob_runs
831 .iter()
832 .map(|&(_, _, lo, hi)| vec![0u8; (hi - lo) as usize])
833 .collect();
834 for buf in &blob_bufs {
835 if !buf.is_empty() {
836 budget.charge_read(buf.len())?;
837 }
838 }
839 let blob_reads = blob_runs
840 .iter()
841 .zip(blob_bufs.iter_mut())
842 .map(|(&(_, _, lo, _), buf)| {
843 self.reader
844 .read_exact_at(section.blobs_start + lo, buf.as_mut_slice())
845 });
846 futures_util::future::try_join_all(blob_reads).await?;
847
848 for (&(j, k, blob_lo, _blob_hi), blob_buf) in blob_runs.iter().zip(&blob_bufs) {
849 for span in &blob_spans[j..=k] {
850 budget.charge_item()?;
851 emit(
852 span.rank,
853 &blob_buf
854 [(span.blob_start - blob_lo) as usize..(span.blob_end - blob_lo) as usize],
855 );
856 }
857 }
858 Ok(())
859 }
860
861 async fn traverse_async<O, F>(
864 &self,
865 overlaps: O,
866 want: Want,
867 mut sink: F,
868 ) -> Result<(), StreamError>
869 where
870 O: Fn(&[u8]) -> bool,
871 F: FnMut(usize, &[u8]),
872 {
873 let section = if want == Want::Payloads {
874 Some(self.payload.as_ref().ok_or(StreamError::NoPayload)?)
875 } else {
876 None
877 };
878 if self.num_items == 0 {
879 return Ok(());
880 }
881
882 let mut budget = Budget::new(self.limits);
883 let mut frontier = vec![self.num_nodes - 1];
884 let mut level = self.level_count - 1;
885 let mut boxes = Vec::new();
886 let mut indices = Vec::new();
887 let mut survivors: Vec<usize> = Vec::new();
888
889 loop {
890 self.gather_async(
893 &frontier,
894 self.box0,
895 self.box_stride,
896 &self.dir_boxes,
897 &mut boxes,
898 &mut budget,
899 )
900 .await?;
901 survivors.clear();
902 indices.clear();
903 for (i, &pos) in frontier.iter().enumerate() {
904 let slot = i * self.box_stride;
905 if overlaps(&boxes[slot..slot + self.record]) {
906 survivors.push(pos);
907 if self.interleaved {
908 indices
909 .extend_from_slice(&boxes[slot + self.record..slot + self.record + 8]);
910 }
911 }
912 }
913 if survivors.is_empty() {
914 return Ok(());
915 }
916
917 if !self.interleaved {
918 self.gather_async(
919 &survivors,
920 self.idx0,
921 8,
922 &self.dir_indices,
923 &mut indices,
924 &mut budget,
925 )
926 .await?;
927 }
928
929 if level == 0 {
930 match section {
931 Some(section) if section.stride != 0 => {
932 self.gather_payloads_fixed_async(
933 section,
934 &survivors,
935 &indices,
936 &mut budget,
937 &mut sink,
938 )
939 .await?;
940 }
941 Some(section) => {
942 self.gather_payloads_async(
943 section,
944 &survivors,
945 &indices,
946 &mut budget,
947 &mut sink,
948 )
949 .await?;
950 }
951 None => {
952 for i in 0..survivors.len() {
953 let id = read_index(&indices, i)?;
954 if id >= self.num_items {
955 return Err(StreamError::Format(LoadError::InvalidTree));
956 }
957 budget.charge_item()?;
958 sink(id, &[]);
959 }
960 }
961 }
962 return Ok(());
963 }
964
965 frontier = expand_frontier(
966 &self.level_bounds,
967 self.node_size,
968 level,
969 survivors.len(),
970 &indices,
971 )?;
972 level -= 1;
973 }
974 }
975}
976
977struct AsyncPrefixSpan {
978 run_index: usize,
979 leaf_rank: usize,
980 blob_start: u64,
981 payload_len: usize,
982}
983
984struct AsyncRankBlobSpan {
985 rank: usize,
986 blob_start: u64,
987 blob_end: u64,
988}
989
990impl AsyncPrefixSpan {
991 fn prefix_end(&self, prefix_len: usize) -> u64 {
992 self.blob_start + self.payload_len.min(prefix_len) as u64
993 }
994}
995
996#[cfg(feature = "async")]
1000impl<R: AsyncRangeReader> StreamIndex2D<R> {
1001 pub async fn open_async(reader: R) -> Result<Self, StreamError> {
1003 Self::open_with_limits_async(reader, StreamLimits::default()).await
1004 }
1005
1006 pub async fn open_with_limits_async(
1009 reader: R,
1010 limits: StreamLimits,
1011 ) -> Result<Self, StreamError> {
1012 Ok(Self {
1013 core: StreamCore::open_async(reader, 2, 8, limits).await?,
1014 })
1015 }
1016
1017 pub async fn search_async(&self, query: Box2D) -> Result<Vec<usize>, StreamError> {
1019 let mut out = Vec::new();
1020 self.core
1021 .traverse_async(
1022 |r| parse_box2d(r).overlaps(query),
1023 Want::Ids,
1024 |id, _| out.push(id),
1025 )
1026 .await?;
1027 Ok(out)
1028 }
1029
1030 pub async fn search_payloads_async(
1032 &self,
1033 query: Box2D,
1034 ) -> Result<Vec<(usize, Vec<u8>)>, StreamError> {
1035 let mut out = Vec::new();
1036 self.core
1037 .traverse_async(
1038 |r| parse_box2d(r).overlaps(query),
1039 Want::Payloads,
1040 |id, blob| out.push((id, blob.to_vec())),
1041 )
1042 .await?;
1043 Ok(out)
1044 }
1045
1046 pub async fn visit_region_async<Q, F>(
1049 &self,
1050 query: &Q,
1051 mut visitor: F,
1052 ) -> Result<(), StreamError>
1053 where
1054 Q: Overlaps2D,
1055 F: FnMut(usize),
1056 {
1057 self.core
1058 .traverse_async(
1059 |r| query.overlaps_box(parse_box2d(r)),
1060 Want::Ids,
1061 |id, _| visitor(id),
1062 )
1063 .await
1064 }
1065
1066 pub async fn search_region_async<Q: Overlaps2D>(
1068 &self,
1069 query: &Q,
1070 ) -> Result<Vec<usize>, StreamError> {
1071 let mut out = Vec::new();
1072 self.visit_region_async(query, |index| out.push(index))
1073 .await?;
1074 Ok(out)
1075 }
1076
1077 pub async fn visit_payloads_region_async<Q, F>(
1080 &self,
1081 query: &Q,
1082 visitor: F,
1083 ) -> Result<(), StreamError>
1084 where
1085 Q: Overlaps2D,
1086 F: FnMut(usize, &[u8]),
1087 {
1088 self.core
1089 .traverse_async(
1090 |r| query.overlaps_box(parse_box2d(r)),
1091 Want::Payloads,
1092 visitor,
1093 )
1094 .await
1095 }
1096
1097 pub async fn search_payloads_region_async<Q: Overlaps2D>(
1100 &self,
1101 query: &Q,
1102 ) -> Result<Vec<(usize, Vec<u8>)>, StreamError> {
1103 let mut out = Vec::new();
1104 self.visit_payloads_region_async(query, |id, blob| out.push((id, blob.to_vec())))
1105 .await?;
1106 Ok(out)
1107 }
1108
1109 pub async fn visit_payload_prefixes_async<F: FnMut(PayloadPrefix<'_>)>(
1111 &self,
1112 query: Box2D,
1113 prefix_len: usize,
1114 visitor: F,
1115 ) -> Result<(), StreamError> {
1116 self.core
1117 .visit_payload_prefixes_async(
1118 |record| parse_box2d(record).overlaps(query),
1119 prefix_len,
1120 visitor,
1121 )
1122 .await
1123 }
1124
1125 pub async fn visit_payload_prefixes_region_async<Q, F>(
1127 &self,
1128 query: &Q,
1129 prefix_len: usize,
1130 visitor: F,
1131 ) -> Result<(), StreamError>
1132 where
1133 Q: Overlaps2D,
1134 F: FnMut(PayloadPrefix<'_>),
1135 {
1136 self.core
1137 .visit_payload_prefixes_async(
1138 |record| query.overlaps_box(parse_box2d(record)),
1139 prefix_len,
1140 visitor,
1141 )
1142 .await
1143 }
1144
1145 pub async fn visit_payloads_at_ranks_async<F: FnMut(usize, &[u8])>(
1147 &self,
1148 leaf_ranks: &[usize],
1149 visitor: F,
1150 ) -> Result<(), StreamError> {
1151 self.core
1152 .visit_payloads_at_ranks_async(leaf_ranks, visitor)
1153 .await
1154 }
1155
1156 pub fn has_payload_async(&self) -> bool {
1158 self.core.has_payload()
1159 }
1160}
1161
1162#[cfg(feature = "async")]
1165impl<R: AsyncRangeReader> StreamIndex3D<R> {
1166 pub async fn open_async(reader: R) -> Result<Self, StreamError> {
1168 Self::open_with_limits_async(reader, StreamLimits::default()).await
1169 }
1170
1171 pub async fn open_with_limits_async(
1173 reader: R,
1174 limits: StreamLimits,
1175 ) -> Result<Self, StreamError> {
1176 Ok(Self {
1177 core: StreamCore::open_async(reader, 3, 8, limits).await?,
1178 })
1179 }
1180
1181 pub async fn search_async(&self, query: Box3D) -> Result<Vec<usize>, StreamError> {
1183 let mut out = Vec::new();
1184 self.core
1185 .traverse_async(
1186 |r| parse_box3d(r).overlaps(query),
1187 Want::Ids,
1188 |id, _| out.push(id),
1189 )
1190 .await?;
1191 Ok(out)
1192 }
1193
1194 pub async fn search_payloads_async(
1196 &self,
1197 query: Box3D,
1198 ) -> Result<Vec<(usize, Vec<u8>)>, StreamError> {
1199 let mut out = Vec::new();
1200 self.core
1201 .traverse_async(
1202 |r| parse_box3d(r).overlaps(query),
1203 Want::Payloads,
1204 |id, blob| out.push((id, blob.to_vec())),
1205 )
1206 .await?;
1207 Ok(out)
1208 }
1209
1210 pub async fn visit_region_async<Q, F>(
1213 &self,
1214 query: &Q,
1215 mut visitor: F,
1216 ) -> Result<(), StreamError>
1217 where
1218 Q: Overlaps3D,
1219 F: FnMut(usize),
1220 {
1221 self.core
1222 .traverse_async(
1223 |r| query.overlaps_box(parse_box3d(r)),
1224 Want::Ids,
1225 |id, _| visitor(id),
1226 )
1227 .await
1228 }
1229
1230 pub async fn search_region_async<Q: Overlaps3D>(
1232 &self,
1233 query: &Q,
1234 ) -> Result<Vec<usize>, StreamError> {
1235 let mut out = Vec::new();
1236 self.visit_region_async(query, |index| out.push(index))
1237 .await?;
1238 Ok(out)
1239 }
1240
1241 pub async fn visit_payloads_region_async<Q, F>(
1244 &self,
1245 query: &Q,
1246 visitor: F,
1247 ) -> Result<(), StreamError>
1248 where
1249 Q: Overlaps3D,
1250 F: FnMut(usize, &[u8]),
1251 {
1252 self.core
1253 .traverse_async(
1254 |r| query.overlaps_box(parse_box3d(r)),
1255 Want::Payloads,
1256 visitor,
1257 )
1258 .await
1259 }
1260
1261 pub async fn search_payloads_region_async<Q: Overlaps3D>(
1264 &self,
1265 query: &Q,
1266 ) -> Result<Vec<(usize, Vec<u8>)>, StreamError> {
1267 let mut out = Vec::new();
1268 self.visit_payloads_region_async(query, |id, blob| out.push((id, blob.to_vec())))
1269 .await?;
1270 Ok(out)
1271 }
1272
1273 pub async fn visit_payload_prefixes_async<F: FnMut(PayloadPrefix<'_>)>(
1275 &self,
1276 query: Box3D,
1277 prefix_len: usize,
1278 visitor: F,
1279 ) -> Result<(), StreamError> {
1280 self.core
1281 .visit_payload_prefixes_async(
1282 |record| parse_box3d(record).overlaps(query),
1283 prefix_len,
1284 visitor,
1285 )
1286 .await
1287 }
1288
1289 pub async fn visit_payload_prefixes_region_async<Q, F>(
1291 &self,
1292 query: &Q,
1293 prefix_len: usize,
1294 visitor: F,
1295 ) -> Result<(), StreamError>
1296 where
1297 Q: Overlaps3D,
1298 F: FnMut(PayloadPrefix<'_>),
1299 {
1300 self.core
1301 .visit_payload_prefixes_async(
1302 |record| query.overlaps_box(parse_box3d(record)),
1303 prefix_len,
1304 visitor,
1305 )
1306 .await
1307 }
1308
1309 pub async fn visit_payloads_at_ranks_async<F: FnMut(usize, &[u8])>(
1311 &self,
1312 leaf_ranks: &[usize],
1313 visitor: F,
1314 ) -> Result<(), StreamError> {
1315 self.core
1316 .visit_payloads_at_ranks_async(leaf_ranks, visitor)
1317 .await
1318 }
1319
1320 pub fn has_payload_async(&self) -> bool {
1322 self.core.has_payload()
1323 }
1324}
1325
1326#[cfg(feature = "async")]
1329impl<R: AsyncRangeReader> StreamIndex2DF32<R> {
1330 pub async fn open_async(reader: R) -> Result<Self, StreamError> {
1332 Self::open_with_limits_async(reader, StreamLimits::default()).await
1333 }
1334
1335 pub async fn open_with_limits_async(
1337 reader: R,
1338 limits: StreamLimits,
1339 ) -> Result<Self, StreamError> {
1340 Ok(Self {
1341 core: StreamCore::open_async(reader, 2, 4, limits).await?,
1342 })
1343 }
1344
1345 pub async fn search_async(&self, query: Box2D) -> Result<Vec<usize>, StreamError> {
1347 let mut out = Vec::new();
1348 self.core
1349 .traverse_async(
1350 |r| parse_box2d_f32(r).overlaps(query),
1351 Want::Ids,
1352 |id, _| out.push(id),
1353 )
1354 .await?;
1355 Ok(out)
1356 }
1357
1358 pub async fn search_payloads_async(
1360 &self,
1361 query: Box2D,
1362 ) -> Result<Vec<(usize, Vec<u8>)>, StreamError> {
1363 let mut out = Vec::new();
1364 self.core
1365 .traverse_async(
1366 |r| parse_box2d_f32(r).overlaps(query),
1367 Want::Payloads,
1368 |id, blob| out.push((id, blob.to_vec())),
1369 )
1370 .await?;
1371 Ok(out)
1372 }
1373
1374 pub async fn visit_region_async<Q, F>(
1377 &self,
1378 query: &Q,
1379 mut visitor: F,
1380 ) -> Result<(), StreamError>
1381 where
1382 Q: Overlaps2D,
1383 F: FnMut(usize),
1384 {
1385 self.core
1386 .traverse_async(
1387 |r| query.overlaps_box(parse_box2d_f32(r)),
1388 Want::Ids,
1389 |id, _| visitor(id),
1390 )
1391 .await
1392 }
1393
1394 pub async fn search_region_async<Q: Overlaps2D>(
1396 &self,
1397 query: &Q,
1398 ) -> Result<Vec<usize>, StreamError> {
1399 let mut out = Vec::new();
1400 self.visit_region_async(query, |index| out.push(index))
1401 .await?;
1402 Ok(out)
1403 }
1404
1405 pub async fn visit_payloads_region_async<Q, F>(
1408 &self,
1409 query: &Q,
1410 visitor: F,
1411 ) -> Result<(), StreamError>
1412 where
1413 Q: Overlaps2D,
1414 F: FnMut(usize, &[u8]),
1415 {
1416 self.core
1417 .traverse_async(
1418 |r| query.overlaps_box(parse_box2d_f32(r)),
1419 Want::Payloads,
1420 visitor,
1421 )
1422 .await
1423 }
1424
1425 pub async fn search_payloads_region_async<Q: Overlaps2D>(
1428 &self,
1429 query: &Q,
1430 ) -> Result<Vec<(usize, Vec<u8>)>, StreamError> {
1431 let mut out = Vec::new();
1432 self.visit_payloads_region_async(query, |id, blob| out.push((id, blob.to_vec())))
1433 .await?;
1434 Ok(out)
1435 }
1436
1437 pub async fn visit_payload_prefixes_async<F: FnMut(PayloadPrefix<'_>)>(
1439 &self,
1440 query: Box2D,
1441 prefix_len: usize,
1442 visitor: F,
1443 ) -> Result<(), StreamError> {
1444 self.core
1445 .visit_payload_prefixes_async(
1446 |record| parse_box2d_f32(record).overlaps(query),
1447 prefix_len,
1448 visitor,
1449 )
1450 .await
1451 }
1452
1453 pub async fn visit_payload_prefixes_region_async<Q, F>(
1455 &self,
1456 query: &Q,
1457 prefix_len: usize,
1458 visitor: F,
1459 ) -> Result<(), StreamError>
1460 where
1461 Q: Overlaps2D,
1462 F: FnMut(PayloadPrefix<'_>),
1463 {
1464 self.core
1465 .visit_payload_prefixes_async(
1466 |record| query.overlaps_box(parse_box2d_f32(record)),
1467 prefix_len,
1468 visitor,
1469 )
1470 .await
1471 }
1472
1473 pub async fn visit_payloads_at_ranks_async<F: FnMut(usize, &[u8])>(
1475 &self,
1476 leaf_ranks: &[usize],
1477 visitor: F,
1478 ) -> Result<(), StreamError> {
1479 self.core
1480 .visit_payloads_at_ranks_async(leaf_ranks, visitor)
1481 .await
1482 }
1483
1484 pub fn has_payload_async(&self) -> bool {
1486 self.core.has_payload()
1487 }
1488}
1489
1490#[cfg(feature = "async")]
1493impl<R: AsyncRangeReader> StreamIndex3DF32<R> {
1494 pub async fn open_async(reader: R) -> Result<Self, StreamError> {
1496 Self::open_with_limits_async(reader, StreamLimits::default()).await
1497 }
1498
1499 pub async fn open_with_limits_async(
1501 reader: R,
1502 limits: StreamLimits,
1503 ) -> Result<Self, StreamError> {
1504 Ok(Self {
1505 core: StreamCore::open_async(reader, 3, 4, limits).await?,
1506 })
1507 }
1508
1509 pub async fn search_async(&self, query: Box3D) -> Result<Vec<usize>, StreamError> {
1511 let mut out = Vec::new();
1512 self.core
1513 .traverse_async(
1514 |r| parse_box3d_f32(r).overlaps(query),
1515 Want::Ids,
1516 |id, _| out.push(id),
1517 )
1518 .await?;
1519 Ok(out)
1520 }
1521
1522 pub async fn search_payloads_async(
1524 &self,
1525 query: Box3D,
1526 ) -> Result<Vec<(usize, Vec<u8>)>, StreamError> {
1527 let mut out = Vec::new();
1528 self.core
1529 .traverse_async(
1530 |r| parse_box3d_f32(r).overlaps(query),
1531 Want::Payloads,
1532 |id, blob| out.push((id, blob.to_vec())),
1533 )
1534 .await?;
1535 Ok(out)
1536 }
1537
1538 pub async fn visit_region_async<Q, F>(
1541 &self,
1542 query: &Q,
1543 mut visitor: F,
1544 ) -> Result<(), StreamError>
1545 where
1546 Q: Overlaps3D,
1547 F: FnMut(usize),
1548 {
1549 self.core
1550 .traverse_async(
1551 |r| query.overlaps_box(parse_box3d_f32(r)),
1552 Want::Ids,
1553 |id, _| visitor(id),
1554 )
1555 .await
1556 }
1557
1558 pub async fn search_region_async<Q: Overlaps3D>(
1560 &self,
1561 query: &Q,
1562 ) -> Result<Vec<usize>, StreamError> {
1563 let mut out = Vec::new();
1564 self.visit_region_async(query, |index| out.push(index))
1565 .await?;
1566 Ok(out)
1567 }
1568
1569 pub async fn visit_payloads_region_async<Q, F>(
1572 &self,
1573 query: &Q,
1574 visitor: F,
1575 ) -> Result<(), StreamError>
1576 where
1577 Q: Overlaps3D,
1578 F: FnMut(usize, &[u8]),
1579 {
1580 self.core
1581 .traverse_async(
1582 |r| query.overlaps_box(parse_box3d_f32(r)),
1583 Want::Payloads,
1584 visitor,
1585 )
1586 .await
1587 }
1588
1589 pub async fn search_payloads_region_async<Q: Overlaps3D>(
1592 &self,
1593 query: &Q,
1594 ) -> Result<Vec<(usize, Vec<u8>)>, StreamError> {
1595 let mut out = Vec::new();
1596 self.visit_payloads_region_async(query, |id, blob| out.push((id, blob.to_vec())))
1597 .await?;
1598 Ok(out)
1599 }
1600
1601 pub async fn visit_payload_prefixes_async<F: FnMut(PayloadPrefix<'_>)>(
1603 &self,
1604 query: Box3D,
1605 prefix_len: usize,
1606 visitor: F,
1607 ) -> Result<(), StreamError> {
1608 self.core
1609 .visit_payload_prefixes_async(
1610 |record| parse_box3d_f32(record).overlaps(query),
1611 prefix_len,
1612 visitor,
1613 )
1614 .await
1615 }
1616
1617 pub async fn visit_payload_prefixes_region_async<Q, F>(
1619 &self,
1620 query: &Q,
1621 prefix_len: usize,
1622 visitor: F,
1623 ) -> Result<(), StreamError>
1624 where
1625 Q: Overlaps3D,
1626 F: FnMut(PayloadPrefix<'_>),
1627 {
1628 self.core
1629 .visit_payload_prefixes_async(
1630 |record| query.overlaps_box(parse_box3d_f32(record)),
1631 prefix_len,
1632 visitor,
1633 )
1634 .await
1635 }
1636
1637 pub async fn visit_payloads_at_ranks_async<F: FnMut(usize, &[u8])>(
1639 &self,
1640 leaf_ranks: &[usize],
1641 visitor: F,
1642 ) -> Result<(), StreamError> {
1643 self.core
1644 .visit_payloads_at_ranks_async(leaf_ranks, visitor)
1645 .await
1646 }
1647
1648 pub fn has_payload_async(&self) -> bool {
1650 self.core.has_payload()
1651 }
1652}