1use std::any::{Any, TypeId};
13use std::path::{Path, PathBuf};
14use std::sync::{Arc, Mutex};
15
16use polars::prelude::*;
17
18use crate::formats::fixed_records::{Bytes, ColumnLayout};
19use crate::formats::members::{Opened, Pick, Table};
20use crate::formats::readers::ScanIn;
21use crate::formats::text_formats::Detail;
22use crate::loading::scan::Scan;
23
24#[derive(Debug, Clone)]
26pub enum Offsets {
27 Narrow(Vec<u32>),
28 Wide(Vec<u64>),
29}
30
31impl Default for Offsets {
32 fn default() -> Self {
33 Self::Narrow(Vec::new())
34 }
35}
36
37impl Offsets {
38 pub fn for_file(len: usize) -> Self {
40 if u32::try_from(len).is_ok() {
41 Self::Narrow(Vec::new())
42 } else {
43 Self::Wide(Vec::new())
44 }
45 }
46
47 pub fn push(&mut self, at: usize) {
48 match self {
49 Self::Narrow(v) => v.push(at as u32),
50 Self::Wide(v) => v.push(at as u64),
51 }
52 }
53
54 pub fn len(&self) -> usize {
55 match self {
56 Self::Narrow(v) => v.len(),
57 Self::Wide(v) => v.len(),
58 }
59 }
60
61 pub fn is_empty(&self) -> bool {
62 self.len() == 0
63 }
64
65 pub fn get(&self, i: usize) -> usize {
66 match self {
67 Self::Narrow(v) => v[i] as usize,
68 Self::Wide(v) => v[i] as usize,
69 }
70 }
71
72 pub fn append(&mut self, other: &Offsets) {
74 match (&mut *self, other) {
75 (Self::Narrow(v), Self::Narrow(o)) => v.extend_from_slice(o),
76 (Self::Wide(v), Self::Wide(o)) => v.extend_from_slice(o),
77 (Self::Wide(v), Self::Narrow(o)) => v.extend(o.iter().map(|&at| u64::from(at))),
78 (Self::Narrow(v), Self::Wide(o)) => {
79 let mut wide: Vec<u64> = v.iter().map(|&at| u64::from(at)).collect();
80 wide.extend_from_slice(o);
81 *self = Self::Wide(wide);
82 }
83 }
84 }
85
86 pub fn truncate(&mut self, len: usize) {
88 match self {
89 Self::Narrow(v) => v.truncate(len),
90 Self::Wide(v) => v.truncate(len),
91 }
92 }
93
94 pub fn widen_for(&mut self, len: usize) {
96 if let Self::Narrow(v) = self
97 && u32::try_from(len).is_err()
98 {
99 *self = Self::Wide(v.iter().map(|&at| u64::from(at)).collect());
100 }
101 }
102
103 pub fn shrink(&mut self) {
104 match self {
105 Self::Narrow(v) => v.shrink_to_fit(),
106 Self::Wide(v) => v.shrink_to_fit(),
107 }
108 }
109}
110
111pub struct IndexedRecords {
113 bytes: Arc<Bytes>,
114 offsets: Arc<Offsets>,
115 columns: Vec<ColumnLayout>,
116 schema: SchemaRef,
117}
118
119impl IndexedRecords {
120 pub fn new(
124 bytes: Arc<Bytes>,
125 offsets: Arc<Offsets>,
126 columns: Vec<ColumnLayout>,
127 ) -> PolarsResult<Self> {
128 let schema: Schema = columns
129 .iter()
130 .map(|c| Field::new(c.name.clone(), c.dtype()))
131 .collect();
132 polars_ensure!(
133 schema.len() == columns.len(),
134 Duplicate: "two columns have the same name"
135 );
136 Ok(Self {
137 bytes,
138 offsets,
139 columns,
140 schema: Arc::new(schema),
141 })
142 }
143
144 pub fn rows(&self) -> usize {
145 self.offsets.len().min(crate::formats::row_index::MAX_ROWS)
146 }
147
148 pub fn lazy(self: &Arc<Self>) -> LazyFrame {
149 crate::formats::row_index::lazy(self)
150 }
151
152 fn starts(&self, rows: impl Iterator<Item = usize>) -> Vec<usize> {
153 rows.map(|i| self.offsets.get(i)).collect()
154 }
155
156 pub fn collect_window(&self, start: usize, len: usize) -> PolarsResult<DataFrame> {
158 let start = start.min(self.rows());
159 let len = len.min(self.rows() - start);
160 self.bytes.still_whole()?;
161 let starts = self.starts(start..start + len);
162 let columns = self
163 .columns
164 .iter()
165 .map(|c| crate::formats::fixed_records::decode_at(self.bytes.as_slice(), c, &starts))
166 .collect::<PolarsResult<Vec<_>>>()?;
167 DataFrame::new(len, columns)
168 }
169}
170
171impl crate::formats::row_index::RowSource for IndexedRecords {
172 fn height(&self) -> usize {
173 self.rows()
174 }
175
176 fn schema(&self) -> SchemaRef {
177 self.schema.clone()
178 }
179
180 fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
181 let rows = crate::formats::row_index::checked(index, self.rows())?;
182 self.bytes.still_whole()?;
183 let starts = self.starts(rows.iter().map(|&r| r as usize));
184 crate::formats::fixed_records::decode_at(
185 self.bytes.as_slice(),
186 &self.columns[column],
187 &starts,
188 )
189 }
190}
191
192impl crate::formats::pushdown::Windowed for IndexedRecords {
193 fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
194 Ok(self.collect_window(start, len)?.lazy())
195 }
196}
197
198const KEPT: usize = 4;
202const KEPT_FILE_BYTES: u64 = 2 << 30;
206const MOST_KEPT: usize = 64;
208
209type Key = (PathBuf, u64, Option<std::time::SystemTime>, TypeId);
210
211static KEPT_INDEXES: Mutex<Vec<(Key, Arc<dyn Any + Send + Sync>)>> = Mutex::new(Vec::new());
212
213fn key<T: 'static>(path: &Path) -> Option<Key> {
214 let meta = std::fs::metadata(path).ok()?;
215 let path = crate::canonical::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
216 Some((path, meta.len(), meta.modified().ok(), TypeId::of::<T>()))
217}
218
219pub fn peek<T: Any + Send + Sync>(path: &Path) -> Option<Arc<T>> {
221 let key = key::<T>(path)?;
222 let kept = KEPT_INDEXES.lock().unwrap_or_else(|e| e.into_inner());
223 kept.iter()
224 .find(|(k, _)| *k == key)
225 .and_then(|(_, index)| index.clone().downcast::<T>().ok())
226}
227
228pub fn forget<T: Any + Send + Sync>(path: &Path) {
230 if let Some(key) = key::<T>(path) {
231 let mut kept = KEPT_INDEXES.lock().unwrap_or_else(|e| e.into_inner());
232 kept.retain(|(k, _)| *k != key);
233 }
234}
235
236pub fn cached<T: Any + Send + Sync, E>(
238 path: &Path,
239 build: impl FnOnce() -> Result<T, E>,
240) -> Result<Arc<T>, E> {
241 if let Some(index) = peek::<T>(path) {
242 return Ok(index);
243 }
244 let index = Arc::new(build()?);
245 keep(path, index.clone());
246 Ok(index)
247}
248
249pub fn keep<T: Any + Send + Sync>(path: &Path, index: Arc<T>) {
251 if let Some(key) = key::<T>(path) {
252 let mut kept = KEPT_INDEXES.lock().unwrap_or_else(|e| e.into_inner());
253 kept.retain(|(k, _)| *k != key);
254 kept.push((key, index as Arc<dyn Any + Send + Sync>));
255 let mut total: u64 = kept.iter().map(|((_, len, ..), _)| *len).sum();
257 let mut excess = 0;
258 while kept.len() - excess > KEPT
259 && (total > KEPT_FILE_BYTES || kept.len() - excess > MOST_KEPT)
260 {
261 total -= kept[excess].0.1;
262 excess += 1;
263 }
264 kept.drain(..excess);
265 }
266}
267
268pub(crate) fn indexed<T: Any + Send + Sync>(
271 path: &Path,
272 index: impl FnOnce(&[u8]) -> Result<T, String>,
273) -> color_eyre::Result<(Arc<Bytes>, Arc<T>)> {
274 let bytes = Arc::new(Bytes::map(path)?);
275 let index = cached(path, || index(bytes.as_slice())).map_err(|e| color_eyre::eyre::eyre!(e))?;
276 Ok((bytes, index))
277}
278
279pub trait Log: Any + Send + Sync + Sized {
282 const EMPTY: &'static str;
284 fn index(data: &[u8]) -> Result<Self, String>;
285 fn tables(&self) -> Vec<Table>;
287 fn detail(&self) -> Detail;
289 fn notes(&self) -> Vec<String>;
291 fn table(
294 &self,
295 bytes: Arc<Bytes>,
296 name: &str,
297 opened: &mut Opened,
298 ) -> Result<LazyFrame, String>;
299}
300
301pub(crate) fn listed<L: Log>(path: &Path) -> color_eyre::Result<Vec<Table>> {
304 peek::<L>(path)
305 .map(|log| log.tables())
306 .ok_or_else(|| color_eyre::eyre::eyre!("Open the log to list its tables."))
307}
308
309pub(crate) fn scan<L: Log>(input: ScanIn<'_>) -> color_eyre::Result<Scan> {
313 let path = input.path().to_path_buf();
314 let (bytes, log) = indexed(&path, L::index)?;
315 let tables = log.tables();
316 let picked = match crate::formats::members::pick(
317 tables.clone(),
318 input.options.table.as_deref(),
319 &path,
320 L::EMPTY,
321 )? {
322 Pick::One(table) => table.name,
323 Pick::Several(tables) => return Ok(crate::formats::members::several(&input, tables)),
324 };
325 let mut opened = Opened::for_table(log.detail(), &tables, &picked, log.notes(), "the log");
326 let lf = log
327 .table(bytes, &picked, &mut opened)
328 .map_err(|e| color_eyre::eyre::eyre!(e))?;
329 Ok(opened.scan(input, lf))
330}
331
332#[cfg(test)]
333mod tests {
334 use super::*;
335 use crate::formats::fixed_records::Physical;
336
337 #[test]
338 fn records_at_offsets_decode_and_window() {
339 let mut bytes = vec![0xEEu8; 20];
341 for (at, a, b) in [(0usize, 1u16, -1i32), (10, 2, -2)] {
342 bytes[at..at + 2].copy_from_slice(&a.to_le_bytes());
343 bytes[at + 2..at + 6].copy_from_slice(&b.to_le_bytes());
344 }
345 let mut offsets = Offsets::for_file(bytes.len());
346 for at in [0, 10] {
347 offsets.push(at);
348 }
349 let records = Arc::new(
350 IndexedRecords::new(
351 Arc::new(Bytes::Owned(bytes)),
352 Arc::new(offsets),
353 vec![
354 ColumnLayout::new("a", 0, 0, Physical::Unsigned(2), 2),
355 ColumnLayout::new("b", 2, 0, Physical::Signed(4), 4),
356 ],
357 )
358 .unwrap(),
359 );
360 let df = records.lazy().collect().unwrap();
361 assert_eq!(
362 df.column("b").unwrap().i32().unwrap().to_vec(),
363 [Some(-1), Some(-2)]
364 );
365 let w = records.collect_window(1, 5).unwrap();
366 assert_eq!(w.column("a").unwrap().u16().unwrap().to_vec(), [Some(2)]);
367 }
368
369 #[test]
370 fn an_index_is_kept_until_the_file_changes() {
371 let dir = tempfile::tempdir().unwrap();
372 let path = dir.path().join("log.bin");
373 std::fs::write(&path, b"one").unwrap();
374 let built = std::cell::Cell::new(0);
375 let build = || -> Result<usize, ()> {
376 built.set(built.get() + 1);
377 Ok(7)
378 };
379 assert_eq!(*cached(&path, build).unwrap(), 7);
380 assert_eq!(*cached(&path, build).unwrap(), 7);
381 assert_eq!(built.get(), 1);
382 std::fs::write(&path, b"longer").unwrap();
383 assert!(peek::<usize>(&path).is_none());
384 }
385}