Skip to main content

packed_spatial_index/stream/
async_io.rs

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// ---- Async streaming (behind the `async` feature) ----
27//
28// Mirror of the synchronous traversal for sources whose reads are async (browser
29// / edge worker over HTTP range or object storage). The descent logic is the
30// same — only the reads are awaited; the overlap test and the result sink stay
31// synchronous closures so no async closures are needed. (The sync and async
32// paths are kept in lockstep by an equivalence test; a future sans-io refactor
33// could share one core.)
34
35/// Async counterpart of [`RangeReader`](super::RangeReader): read a byte range,
36/// returning a future.
37///
38/// Implement this to query an index that lives behind async I/O — an HTTP range
39/// request from WebAssembly, an object-storage `get(range)` in an edge worker.
40/// The returned futures need not be `Send` (edge/browser executors are
41/// single-threaded). See [`RangeReader`](super::RangeReader) for the sync
42/// analogue and an HTTP implementation sketch.
43#[cfg(feature = "async")]
44#[allow(async_fn_in_trait, clippy::len_without_is_empty)]
45pub trait AsyncRangeReader {
46    /// Read exactly `buf.len()` bytes starting at `offset`.
47    async fn read_exact_at(&self, offset: u64, buf: &mut [u8]) -> io::Result<()>;
48
49    /// Total length in bytes, if known.
50    fn len(&self) -> Option<u64> {
51        None
52    }
53}
54
55/// What a traversal collects at the leaves.
56#[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        // Reject a file longer than the last chunk plus its alignment pad — a
115        // stray trailing byte the directory does not account for.
116        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        // The optional prefix section: a dense copy of the blob heads. Only
202        // useful next to a payload, and only when its stride can satisfy the
203        // requested prefix, both of which the scan re-checks per query.
204        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        // Directory prefetch (mirror of the sync `open` epilogue).
225        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 mirror of [`gather`](StreamCore::gather), but issues all of a
266    /// level's coalesced runs concurrently (one buffer each). On a
267    /// single-threaded async executor this puts several range fetches in flight
268    /// at once, so the level's latency is one round trip rather than the sum.
269    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 mirror of [`gather_payloads`](StreamCore::gather_payloads). Reads
303    /// every run's offset table concurrently, then every run's blobs
304    /// concurrently — two round trips for the whole leaf frontier rather than two
305    /// per run.
306    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        // Group leaf positions into coalesced runs.
318        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        // Phase 1: read every run's offset table concurrently.
327        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        // Validate each run's blob span.
342        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        // Phase 2: read every run's blobs concurrently (empty spans are no-ops).
353        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        // Emit every run.
372        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    /// Fixed-width async payload gather: one contiguous blob read per coalesced
394    /// run, all runs issued concurrently. No offset-table phase (the variable
395    /// `gather_payloads_async` needs two round trips; this needs one).
396    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        // With a prefix section the same bytes sit in a dense array, so the scan
615        // reads rank runs instead of one strided range per match. The ordinary
616        // record gap applies: what lies between two prefixes here is other
617        // prefixes, not bodies, so over-reading is cheap.
618        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 mirror of the synchronous traversal, parameterized by `want` (ids or
862    /// id+payload). `overlaps` and `sink` are synchronous; only reads are awaited.
863    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            // One gather per level fetches each frontier node's box (interleaved:
891            // box + index together; SoA: box only).
892            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/// Streaming reader for a 2D `f64` index over async I/O. Mirrors
997/// [`StreamIndex2D`]; use it when reads return futures (e.g. browser / edge
998/// worker). Behind the `async` feature.
999#[cfg(feature = "async")]
1000impl<R: AsyncRangeReader> StreamIndex2D<R> {
1001    /// Open and validate a 2D `f64` index from an async `reader`.
1002    pub async fn open_async(reader: R) -> Result<Self, StreamError> {
1003        Self::open_with_limits_async(reader, StreamLimits::default()).await
1004    }
1005
1006    /// Open from an async `reader` with per-query [`StreamLimits`]. See
1007    /// [`StreamIndex2D::open_with_limits`].
1008    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    /// Stream the indices of every item whose box intersects `query`.
1018    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    /// Stream `(item index, payload blob)` for every item intersecting `query`.
1031    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    /// Stream the indices of every item whose box overlaps the region `query` —
1047    /// any [`Overlaps2D`] shape, not just a box.
1048    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    /// Collect the indices of every item whose box overlaps the region `query`.
1067    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    /// Visit `(item index, payload blob)` for every item whose box overlaps the
1078    /// region `query`.
1079    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    /// Collect `(item index, payload blob)` for every item whose box overlaps the
1098    /// region `query`.
1099    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    /// Async counterpart of [`StreamIndex2D::visit_payload_prefixes`].
1110    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    /// Async counterpart of [`StreamIndex2D::visit_payload_prefixes_region`].
1126    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    /// Async counterpart of [`StreamIndex2D::visit_payloads_at_ranks`].
1146    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    /// Whether this index was written with a payload section.
1157    pub fn has_payload_async(&self) -> bool {
1158        self.core.has_payload()
1159    }
1160}
1161
1162/// Streaming reader for a 3D `f64` index over async I/O. See [`StreamIndex2D`]'s
1163/// async methods. Behind the `async` feature.
1164#[cfg(feature = "async")]
1165impl<R: AsyncRangeReader> StreamIndex3D<R> {
1166    /// Open and validate a 3D `f64` index from an async `reader`.
1167    pub async fn open_async(reader: R) -> Result<Self, StreamError> {
1168        Self::open_with_limits_async(reader, StreamLimits::default()).await
1169    }
1170
1171    /// Open from an async `reader` with per-query [`StreamLimits`].
1172    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    /// Stream the indices of every item whose box intersects `query`.
1182    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    /// Stream `(item index, payload blob)` for every item intersecting `query`.
1195    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    /// Stream the indices of every item whose box overlaps the region `query` —
1211    /// any [`Overlaps3D`] shape, not just a box.
1212    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    /// Collect the indices of every item whose box overlaps the region `query`.
1231    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    /// Visit `(item index, payload blob)` for every item whose box overlaps the
1242    /// region `query`.
1243    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    /// Collect `(item index, payload blob)` for every item whose box overlaps the
1262    /// region `query`.
1263    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    /// Async counterpart of [`StreamIndex3D::visit_payload_prefixes`].
1274    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    /// Async counterpart of [`StreamIndex3D::visit_payload_prefixes_region`].
1290    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    /// Async counterpart of [`StreamIndex3D::visit_payloads_at_ranks`].
1310    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    /// Whether this index was written with a payload section.
1321    pub fn has_payload_async(&self) -> bool {
1322        self.core.has_payload()
1323    }
1324}
1325
1326/// Async streaming reader for a compact `f32` 2D index. Mirrors
1327/// [`StreamIndex2DF32`]'s sync methods over async I/O. Behind the `async` feature.
1328#[cfg(feature = "async")]
1329impl<R: AsyncRangeReader> StreamIndex2DF32<R> {
1330    /// Open and validate a 2D `f32` index from an async `reader`.
1331    pub async fn open_async(reader: R) -> Result<Self, StreamError> {
1332        Self::open_with_limits_async(reader, StreamLimits::default()).await
1333    }
1334
1335    /// Open from an async `reader` with per-query [`StreamLimits`].
1336    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    /// Stream the indices of every item whose (rounded) box intersects `query`.
1346    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    /// Stream `(item index, payload blob)` for every item intersecting `query`.
1359    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    /// Stream the indices of every item whose (rounded) box overlaps the region
1375    /// `query` — any [`Overlaps2D`] shape.
1376    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    /// Collect the indices of every item whose box overlaps the region `query`.
1395    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    /// Visit `(item index, payload blob)` for every item whose (rounded) box
1406    /// overlaps the region `query`.
1407    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    /// Collect `(item index, payload blob)` for every item whose box overlaps the
1426    /// region `query`.
1427    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    /// Async counterpart of [`StreamIndex2DF32::visit_payload_prefixes`].
1438    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    /// Async counterpart of [`StreamIndex2DF32::visit_payload_prefixes_region`].
1454    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    /// Async counterpart of [`StreamIndex2DF32::visit_payloads_at_ranks`].
1474    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    /// Whether this index was written with a payload section.
1485    pub fn has_payload_async(&self) -> bool {
1486        self.core.has_payload()
1487    }
1488}
1489
1490/// Async streaming reader for a compact `f32` 3D index. See
1491/// [`StreamIndex2DF32`]'s async methods. Behind the `async` feature.
1492#[cfg(feature = "async")]
1493impl<R: AsyncRangeReader> StreamIndex3DF32<R> {
1494    /// Open and validate a 3D `f32` index from an async `reader`.
1495    pub async fn open_async(reader: R) -> Result<Self, StreamError> {
1496        Self::open_with_limits_async(reader, StreamLimits::default()).await
1497    }
1498
1499    /// Open from an async `reader` with per-query [`StreamLimits`].
1500    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    /// Stream the indices of every item whose (rounded) box intersects `query`.
1510    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    /// Stream `(item index, payload blob)` for every item intersecting `query`.
1523    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    /// Stream the indices of every item whose (rounded) box overlaps the region
1539    /// `query` — any [`Overlaps3D`] shape.
1540    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    /// Collect the indices of every item whose box overlaps the region `query`.
1559    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    /// Visit `(item index, payload blob)` for every item whose (rounded) box
1570    /// overlaps the region `query`.
1571    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    /// Collect `(item index, payload blob)` for every item whose box overlaps the
1590    /// region `query`.
1591    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    /// Async counterpart of [`StreamIndex3DF32::visit_payload_prefixes`].
1602    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    /// Async counterpart of [`StreamIndex3DF32::visit_payload_prefixes_region`].
1618    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    /// Async counterpart of [`StreamIndex3DF32::visit_payloads_at_ranks`].
1638    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    /// Whether this index was written with a payload section.
1649    pub fn has_payload_async(&self) -> bool {
1650        self.core.has_payload()
1651    }
1652}