1use std::sync::Arc;
5
6use futures::{FutureExt, future::BoxFuture};
7use lance_encoding::EncodingsIo;
8use lance_io::scheduler::FileScheduler;
9
10use super::reader::DEFAULT_READ_CHUNK_SIZE;
11
12#[derive(Debug)]
13pub struct LanceEncodingsIo {
14 scheduler: FileScheduler,
15 read_chunk_size: u64,
17}
18
19impl LanceEncodingsIo {
20 pub fn new(scheduler: FileScheduler) -> Self {
21 Self {
22 scheduler,
23 read_chunk_size: DEFAULT_READ_CHUNK_SIZE,
24 }
25 }
26
27 pub fn with_read_chunk_size(mut self, read_chunk_size: u64) -> Self {
28 self.read_chunk_size = read_chunk_size;
29 self
30 }
31}
32
33impl EncodingsIo for LanceEncodingsIo {
34 fn with_bypass_backpressure(&self) -> Option<Arc<dyn EncodingsIo>> {
35 Some(Arc::new(Self {
36 scheduler: self.scheduler.with_bypass_backpressure(),
37 read_chunk_size: self.read_chunk_size,
38 }))
39 }
40
41 fn with_io_stats(
42 &self,
43 stats: Arc<dyn lance_core::utils::io_stats::IoStatsRecorder>,
44 ) -> Option<Arc<dyn EncodingsIo>> {
45 Some(Arc::new(Self {
46 scheduler: self.scheduler.with_io_stats(stats),
47 read_chunk_size: self.read_chunk_size,
48 }))
49 }
50
51 fn submit_request(
52 &self,
53 ranges: Vec<std::ops::Range<u64>>,
54 priority: u64,
55 ) -> BoxFuture<'static, lance_core::Result<Vec<bytes::Bytes>>> {
56 let mut split_ranges = Vec::new();
57 let mut split_indices = Vec::new(); let mut needs_reassembly = false;
63
64 for (idx, range) in ranges.iter().enumerate() {
68 if range.start == range.end {
69 needs_reassembly = true;
73 continue;
74 }
75 let range_size = range.end - range.start;
76
77 if range_size > self.read_chunk_size {
78 needs_reassembly = true;
79 let num_chunks = range_size.div_ceil(self.read_chunk_size);
80 let chunk_size = range_size / num_chunks;
81
82 for i in 0..num_chunks {
83 let start = range.start + i * chunk_size;
84 let end = if i == num_chunks - 1 {
85 range.end } else {
87 start + chunk_size
88 };
89 split_ranges.push(start..end);
90 split_indices.push(idx);
91 }
92 } else {
93 split_ranges.push(range.clone());
94 split_indices.push(idx);
95 }
96 }
97
98 let fut = self.scheduler.submit_request(split_ranges, priority);
99
100 async move {
101 let split_results = fut.await?;
102
103 if split_results.len() != split_indices.len() {
104 return Err(lance_core::Error::internal(format!(
105 "Encoding I/O returned {} results for {} requested range chunks",
106 split_results.len(),
107 split_indices.len()
108 )));
109 }
110 if !needs_reassembly {
111 return Ok(split_results);
112 }
113
114 let mut results = vec![Vec::new(); ranges.len()];
115
116 for (split_result, orig_idx) in split_results.into_iter().zip(split_indices) {
117 results[orig_idx].push(split_result);
118 }
119
120 let mut reassembled = Vec::with_capacity(ranges.len());
121 for (range, chunks) in ranges.iter().zip(results) {
122 if chunks.is_empty() {
123 if range.start == range.end {
124 reassembled.push(bytes::Bytes::new());
125 continue;
126 }
127 return Err(lance_core::Error::internal(format!(
128 "Encoding I/O returned no data for non-empty range {}..{}",
129 range.start, range.end
130 )));
131 }
132 if chunks.len() == 1 {
133 reassembled.push(chunks[0].clone());
134 continue;
135 }
136
137 let total_size: usize = chunks.iter().map(|c| c.len()).sum();
138 let mut combined = Vec::with_capacity(total_size);
139 for chunk in chunks {
140 combined.extend_from_slice(&chunk);
141 }
142 reassembled.push(bytes::Bytes::from(combined));
143 }
144 Ok(reassembled)
145 }
146 .boxed()
147 }
148}