use std::any::{Any, TypeId};
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use polars::prelude::*;
use crate::fixed_records::{Bytes, ColumnLayout};
pub const MAX_RECORDS: usize = 64 << 20;
#[derive(Debug, Clone)]
pub enum Offsets {
Narrow(Vec<u32>),
Wide(Vec<u64>),
}
impl Default for Offsets {
fn default() -> Self {
Self::Narrow(Vec::new())
}
}
impl Offsets {
pub fn for_file(len: usize) -> Self {
if u32::try_from(len).is_ok() {
Self::Narrow(Vec::new())
} else {
Self::Wide(Vec::new())
}
}
pub fn push(&mut self, at: usize) {
match self {
Self::Narrow(v) => v.push(at as u32),
Self::Wide(v) => v.push(at as u64),
}
}
pub fn len(&self) -> usize {
match self {
Self::Narrow(v) => v.len(),
Self::Wide(v) => v.len(),
}
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn get(&self, i: usize) -> usize {
match self {
Self::Narrow(v) => v[i] as usize,
Self::Wide(v) => v[i] as usize,
}
}
pub fn append(&mut self, other: &Offsets) {
match (&mut *self, other) {
(Self::Narrow(v), Self::Narrow(o)) => v.extend_from_slice(o),
(Self::Wide(v), Self::Wide(o)) => v.extend_from_slice(o),
(Self::Wide(v), Self::Narrow(o)) => v.extend(o.iter().map(|&at| u64::from(at))),
(Self::Narrow(v), Self::Wide(o)) => {
let mut wide: Vec<u64> = v.iter().map(|&at| u64::from(at)).collect();
wide.extend_from_slice(o);
*self = Self::Wide(wide);
}
}
}
pub fn truncate(&mut self, len: usize) {
match self {
Self::Narrow(v) => v.truncate(len),
Self::Wide(v) => v.truncate(len),
}
}
pub fn widen_for(&mut self, len: usize) {
if let Self::Narrow(v) = self
&& u32::try_from(len).is_err()
{
*self = Self::Wide(v.iter().map(|&at| u64::from(at)).collect());
}
}
pub fn shrink(&mut self) {
match self {
Self::Narrow(v) => v.shrink_to_fit(),
Self::Wide(v) => v.shrink_to_fit(),
}
}
}
pub struct IndexedRecords {
bytes: Arc<Bytes>,
offsets: Arc<Offsets>,
columns: Vec<ColumnLayout>,
schema: SchemaRef,
}
impl IndexedRecords {
pub fn new(
bytes: Arc<Bytes>,
offsets: Arc<Offsets>,
columns: Vec<ColumnLayout>,
) -> PolarsResult<Self> {
let schema: Schema = columns
.iter()
.map(|c| Field::new(c.name.clone(), c.dtype()))
.collect();
polars_ensure!(
schema.len() == columns.len(),
Duplicate: "two columns have the same name"
);
Ok(Self {
bytes,
offsets,
columns,
schema: Arc::new(schema),
})
}
pub fn rows(&self) -> usize {
self.offsets.len().min(crate::row_index::MAX_ROWS)
}
pub fn lazy(self: &Arc<Self>) -> LazyFrame {
crate::row_index::lazy(self)
}
fn starts(&self, rows: impl Iterator<Item = usize>) -> Vec<usize> {
rows.map(|i| self.offsets.get(i)).collect()
}
pub fn collect_window(&self, start: usize, len: usize) -> PolarsResult<DataFrame> {
let start = start.min(self.rows());
let len = len.min(self.rows() - start);
self.bytes.still_whole()?;
let starts = self.starts(start..start + len);
let columns = self
.columns
.iter()
.map(|c| crate::fixed_records::decode_at(self.bytes.as_slice(), c, &starts))
.collect::<PolarsResult<Vec<_>>>()?;
DataFrame::new(len, columns)
}
}
impl crate::row_index::RowSource for IndexedRecords {
fn height(&self) -> usize {
self.rows()
}
fn schema(&self) -> SchemaRef {
self.schema.clone()
}
fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
let rows = crate::row_index::checked(index, self.rows())?;
self.bytes.still_whole()?;
let starts = self.starts(rows.iter().map(|&r| r as usize));
crate::fixed_records::decode_at(self.bytes.as_slice(), &self.columns[column], &starts)
}
}
impl crate::pushdown::Windowed for IndexedRecords {
fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
Ok(self.collect_window(start, len)?.lazy())
}
}
const KEPT: usize = 4;
const KEPT_FILE_BYTES: u64 = 2 << 30;
const MOST_KEPT: usize = 64;
type Key = (PathBuf, u64, Option<std::time::SystemTime>, TypeId);
static KEPT_INDEXES: Mutex<Vec<(Key, Arc<dyn Any + Send + Sync>)>> = Mutex::new(Vec::new());
fn key<T: 'static>(path: &Path) -> Option<Key> {
let meta = std::fs::metadata(path).ok()?;
let path = crate::canonical::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
Some((path, meta.len(), meta.modified().ok(), TypeId::of::<T>()))
}
pub fn peek<T: Any + Send + Sync>(path: &Path) -> Option<Arc<T>> {
let key = key::<T>(path)?;
let kept = KEPT_INDEXES.lock().unwrap_or_else(|e| e.into_inner());
kept.iter()
.find(|(k, _)| *k == key)
.and_then(|(_, index)| index.clone().downcast::<T>().ok())
}
pub fn forget<T: Any + Send + Sync>(path: &Path) {
if let Some(key) = key::<T>(path) {
let mut kept = KEPT_INDEXES.lock().unwrap_or_else(|e| e.into_inner());
kept.retain(|(k, _)| *k != key);
}
}
pub fn cached<T: Any + Send + Sync, E>(
path: &Path,
build: impl FnOnce() -> Result<T, E>,
) -> Result<Arc<T>, E> {
if let Some(index) = peek::<T>(path) {
return Ok(index);
}
let index = Arc::new(build()?);
keep(path, index.clone());
Ok(index)
}
pub fn keep<T: Any + Send + Sync>(path: &Path, index: Arc<T>) {
if let Some(key) = key::<T>(path) {
let mut kept = KEPT_INDEXES.lock().unwrap_or_else(|e| e.into_inner());
kept.retain(|(k, _)| *k != key);
kept.push((key, index as Arc<dyn Any + Send + Sync>));
let mut total: u64 = kept.iter().map(|((_, len, ..), _)| *len).sum();
let mut excess = 0;
while kept.len() - excess > KEPT
&& (total > KEPT_FILE_BYTES || kept.len() - excess > MOST_KEPT)
{
total -= kept[excess].0.1;
excess += 1;
}
kept.drain(..excess);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::fixed_records::Physical;
#[test]
fn records_at_offsets_decode_and_window() {
let mut bytes = vec![0xEEu8; 20];
for (at, a, b) in [(0usize, 1u16, -1i32), (10, 2, -2)] {
bytes[at..at + 2].copy_from_slice(&a.to_le_bytes());
bytes[at + 2..at + 6].copy_from_slice(&b.to_le_bytes());
}
let mut offsets = Offsets::for_file(bytes.len());
for at in [0, 10] {
offsets.push(at);
}
let records = Arc::new(
IndexedRecords::new(
Arc::new(Bytes::Owned(bytes)),
Arc::new(offsets),
vec![
ColumnLayout::new("a", 0, 0, Physical::Unsigned(2), 2),
ColumnLayout::new("b", 2, 0, Physical::Signed(4), 4),
],
)
.unwrap(),
);
let df = records.lazy().collect().unwrap();
assert_eq!(
df.column("b").unwrap().i32().unwrap().to_vec(),
[Some(-1), Some(-2)]
);
let w = records.collect_window(1, 5).unwrap();
assert_eq!(w.column("a").unwrap().u16().unwrap().to_vec(), [Some(2)]);
}
#[test]
fn an_index_is_kept_until_the_file_changes() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("log.bin");
std::fs::write(&path, b"one").unwrap();
let built = std::cell::Cell::new(0);
let build = || -> Result<usize, ()> {
built.set(built.get() + 1);
Ok(7)
};
assert_eq!(*cached(&path, build).unwrap(), 7);
assert_eq!(*cached(&path, build).unwrap(), 7);
assert_eq!(built.get(), 1);
std::fs::write(&path, b"longer").unwrap();
assert!(peek::<usize>(&path).is_none());
}
}