1use std::any::{Any, TypeId};
12use std::path::{Path, PathBuf};
13use std::sync::{Arc, Mutex};
14
15use polars::prelude::*;
16
17use crate::fixed_records::{Bytes, ColumnLayout};
18
19pub const MAX_RECORDS: usize = 64 << 20;
22
23#[derive(Debug, Clone)]
25pub enum Offsets {
26 Narrow(Vec<u32>),
27 Wide(Vec<u64>),
28}
29
30impl Default for Offsets {
31 fn default() -> Self {
32 Self::Narrow(Vec::new())
33 }
34}
35
36impl Offsets {
37 pub fn for_file(len: usize) -> Self {
39 if u32::try_from(len).is_ok() {
40 Self::Narrow(Vec::new())
41 } else {
42 Self::Wide(Vec::new())
43 }
44 }
45
46 pub fn push(&mut self, at: usize) {
47 match self {
48 Self::Narrow(v) => v.push(at as u32),
49 Self::Wide(v) => v.push(at as u64),
50 }
51 }
52
53 pub fn len(&self) -> usize {
54 match self {
55 Self::Narrow(v) => v.len(),
56 Self::Wide(v) => v.len(),
57 }
58 }
59
60 pub fn is_empty(&self) -> bool {
61 self.len() == 0
62 }
63
64 pub fn get(&self, i: usize) -> usize {
65 match self {
66 Self::Narrow(v) => v[i] as usize,
67 Self::Wide(v) => v[i] as usize,
68 }
69 }
70
71 pub fn append(&mut self, other: &Offsets) {
73 match (&mut *self, other) {
74 (Self::Narrow(v), Self::Narrow(o)) => v.extend_from_slice(o),
75 (Self::Wide(v), Self::Wide(o)) => v.extend_from_slice(o),
76 (Self::Wide(v), Self::Narrow(o)) => v.extend(o.iter().map(|&at| u64::from(at))),
77 (Self::Narrow(v), Self::Wide(o)) => {
78 let mut wide: Vec<u64> = v.iter().map(|&at| u64::from(at)).collect();
79 wide.extend_from_slice(o);
80 *self = Self::Wide(wide);
81 }
82 }
83 }
84
85 pub fn truncate(&mut self, len: usize) {
87 match self {
88 Self::Narrow(v) => v.truncate(len),
89 Self::Wide(v) => v.truncate(len),
90 }
91 }
92
93 pub fn widen_for(&mut self, len: usize) {
95 if let Self::Narrow(v) = self
96 && u32::try_from(len).is_err()
97 {
98 *self = Self::Wide(v.iter().map(|&at| u64::from(at)).collect());
99 }
100 }
101
102 pub fn shrink(&mut self) {
103 match self {
104 Self::Narrow(v) => v.shrink_to_fit(),
105 Self::Wide(v) => v.shrink_to_fit(),
106 }
107 }
108}
109
110pub struct IndexedRecords {
112 bytes: Arc<Bytes>,
113 offsets: Arc<Offsets>,
114 columns: Vec<ColumnLayout>,
115 schema: SchemaRef,
116}
117
118impl IndexedRecords {
119 pub fn new(
123 bytes: Arc<Bytes>,
124 offsets: Arc<Offsets>,
125 columns: Vec<ColumnLayout>,
126 ) -> PolarsResult<Self> {
127 let schema: Schema = columns
128 .iter()
129 .map(|c| Field::new(c.name.clone(), c.dtype()))
130 .collect();
131 polars_ensure!(
132 schema.len() == columns.len(),
133 Duplicate: "two columns have the same name"
134 );
135 Ok(Self {
136 bytes,
137 offsets,
138 columns,
139 schema: Arc::new(schema),
140 })
141 }
142
143 pub fn rows(&self) -> usize {
144 self.offsets.len().min(crate::row_index::MAX_ROWS)
145 }
146
147 pub fn lazy(self: &Arc<Self>) -> LazyFrame {
148 crate::row_index::lazy(self)
149 }
150
151 fn starts(&self, rows: impl Iterator<Item = usize>) -> Vec<usize> {
152 rows.map(|i| self.offsets.get(i)).collect()
153 }
154
155 pub fn collect_window(&self, start: usize, len: usize) -> PolarsResult<DataFrame> {
157 let start = start.min(self.rows());
158 let len = len.min(self.rows() - start);
159 self.bytes.still_whole()?;
160 let starts = self.starts(start..start + len);
161 let columns = self
162 .columns
163 .iter()
164 .map(|c| crate::fixed_records::decode_at(self.bytes.as_slice(), c, &starts))
165 .collect::<PolarsResult<Vec<_>>>()?;
166 DataFrame::new(len, columns)
167 }
168}
169
170impl crate::row_index::RowSource for IndexedRecords {
171 fn height(&self) -> usize {
172 self.rows()
173 }
174
175 fn schema(&self) -> SchemaRef {
176 self.schema.clone()
177 }
178
179 fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
180 let rows = crate::row_index::checked(index, self.rows())?;
181 self.bytes.still_whole()?;
182 let starts = self.starts(rows.iter().map(|&r| r as usize));
183 crate::fixed_records::decode_at(self.bytes.as_slice(), &self.columns[column], &starts)
184 }
185}
186
187impl crate::pushdown::Windowed for IndexedRecords {
188 fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
189 Ok(self.collect_window(start, len)?.lazy())
190 }
191}
192
193const KEPT: usize = 4;
197const KEPT_FILE_BYTES: u64 = 2 << 30;
201const MOST_KEPT: usize = 64;
203
204type Key = (PathBuf, u64, Option<std::time::SystemTime>, TypeId);
205
206static KEPT_INDEXES: Mutex<Vec<(Key, Arc<dyn Any + Send + Sync>)>> = Mutex::new(Vec::new());
207
208fn key<T: 'static>(path: &Path) -> Option<Key> {
209 let meta = std::fs::metadata(path).ok()?;
210 let path = crate::canonical::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
211 Some((path, meta.len(), meta.modified().ok(), TypeId::of::<T>()))
212}
213
214pub fn peek<T: Any + Send + Sync>(path: &Path) -> Option<Arc<T>> {
216 let key = key::<T>(path)?;
217 let kept = KEPT_INDEXES.lock().unwrap_or_else(|e| e.into_inner());
218 kept.iter()
219 .find(|(k, _)| *k == key)
220 .and_then(|(_, index)| index.clone().downcast::<T>().ok())
221}
222
223pub fn forget<T: Any + Send + Sync>(path: &Path) {
225 if let Some(key) = key::<T>(path) {
226 let mut kept = KEPT_INDEXES.lock().unwrap_or_else(|e| e.into_inner());
227 kept.retain(|(k, _)| *k != key);
228 }
229}
230
231pub fn cached<T: Any + Send + Sync, E>(
233 path: &Path,
234 build: impl FnOnce() -> Result<T, E>,
235) -> Result<Arc<T>, E> {
236 if let Some(index) = peek::<T>(path) {
237 return Ok(index);
238 }
239 let index = Arc::new(build()?);
240 keep(path, index.clone());
241 Ok(index)
242}
243
244pub fn keep<T: Any + Send + Sync>(path: &Path, index: Arc<T>) {
246 if let Some(key) = key::<T>(path) {
247 let mut kept = KEPT_INDEXES.lock().unwrap_or_else(|e| e.into_inner());
248 kept.retain(|(k, _)| *k != key);
249 kept.push((key, index as Arc<dyn Any + Send + Sync>));
250 let mut total: u64 = kept.iter().map(|((_, len, ..), _)| *len).sum();
252 let mut excess = 0;
253 while kept.len() - excess > KEPT
254 && (total > KEPT_FILE_BYTES || kept.len() - excess > MOST_KEPT)
255 {
256 total -= kept[excess].0.1;
257 excess += 1;
258 }
259 kept.drain(..excess);
260 }
261}
262
263#[cfg(test)]
264mod tests {
265 use super::*;
266 use crate::fixed_records::Physical;
267
268 #[test]
269 fn records_at_offsets_decode_and_window() {
270 let mut bytes = vec![0xEEu8; 20];
272 for (at, a, b) in [(0usize, 1u16, -1i32), (10, 2, -2)] {
273 bytes[at..at + 2].copy_from_slice(&a.to_le_bytes());
274 bytes[at + 2..at + 6].copy_from_slice(&b.to_le_bytes());
275 }
276 let mut offsets = Offsets::for_file(bytes.len());
277 for at in [0, 10] {
278 offsets.push(at);
279 }
280 let records = Arc::new(
281 IndexedRecords::new(
282 Arc::new(Bytes::Owned(bytes)),
283 Arc::new(offsets),
284 vec![
285 ColumnLayout::new("a", 0, 0, Physical::Unsigned(2), 2),
286 ColumnLayout::new("b", 2, 0, Physical::Signed(4), 4),
287 ],
288 )
289 .unwrap(),
290 );
291 let df = records.lazy().collect().unwrap();
292 assert_eq!(
293 df.column("b").unwrap().i32().unwrap().to_vec(),
294 [Some(-1), Some(-2)]
295 );
296 let w = records.collect_window(1, 5).unwrap();
297 assert_eq!(w.column("a").unwrap().u16().unwrap().to_vec(), [Some(2)]);
298 }
299
300 #[test]
301 fn an_index_is_kept_until_the_file_changes() {
302 let dir = tempfile::tempdir().unwrap();
303 let path = dir.path().join("log.bin");
304 std::fs::write(&path, b"one").unwrap();
305 let built = std::cell::Cell::new(0);
306 let build = || -> Result<usize, ()> {
307 built.set(built.get() + 1);
308 Ok(7)
309 };
310 assert_eq!(*cached(&path, build).unwrap(), 7);
311 assert_eq!(*cached(&path, build).unwrap(), 7);
312 assert_eq!(built.get(), 1);
313 std::fs::write(&path, b"longer").unwrap();
314 assert!(peek::<usize>(&path).is_none());
315 }
316}