use arrow::array::RecordBatch as ArrowRecordBatch;
use arrow::datatypes::Schema as ArrowSchema;
use duckdb::Connection;
use polars::datatypes::{AnyValue, PlSmallStr};
use polars::error::{PolarsError, PolarsResult};
use polars::frame::DataFrame;
use polars::prelude::*;
use std::collections::HashMap;
use std::fmt::Debug;
use std::sync::Arc;
#[cfg(debug_assertions)]
use std::time::Instant;
use tokio::sync::{Mutex, RwLock};
use nettan::DataFrameSchematic;
use crate::ipc::write_to_arrow;
use crate::persist::save_to_duckdb;
pub trait RawData: Debug {
fn validate(&self, df: &DataFrame) -> PolarsResult<bool>;
fn validate_series(
&self,
series_data: &HashMap<PlSmallStr, Series>,
df: &DataFrame,
) -> PolarsResult<bool>;
fn format(&self) -> PolarsResult<Vec<(PlSmallStr, AnyValue<'_>)>>;
}
type DataFrameMap = RwLock<HashMap<String, Arc<RwLock<HashMap<String, Option<DataFrame>>>>>>;
#[derive(Clone, Debug)]
pub struct DataFrameCollection {
df_map: Arc<DataFrameMap>,
df_schematics: Arc<RwLock<HashMap<String, DataFrameSchematic>>>,
}
impl DataFrameCollection {
pub fn new(df_schematics: HashMap<String, DataFrameSchematic>) -> DataFrameCollection {
let df_map = Arc::new(RwLock::new(HashMap::new()));
DataFrameCollection {
df_map: df_map,
df_schematics: Arc::new(RwLock::new(df_schematics)),
}
}
pub async fn set_minimum_rows(&self, schematic_key: &str, minimum_rows: u32) -> bool {
if let Some(df_schematic) = self.df_schematics.write().await.get_mut(schematic_key) {
df_schematic.with_minimum_rows(minimum_rows);
true
} else {
false
}
}
pub async fn insert_inner_map(&self, outer_key: String, inner_key: String) -> bool {
if !self.df_schematics.read().await.contains_key(&inner_key) {
return false;
}
let mut df_map = self.df_map.write().await;
let inner_map = df_map
.entry(outer_key)
.or_insert(Arc::new(RwLock::new(HashMap::new())));
let inner_map = Arc::clone(inner_map);
drop(df_map);
let mut inner_map = inner_map.write().await;
inner_map.entry(inner_key).or_insert(None);
true
}
async fn access_inner_map(
&self,
outer_key: &str,
) -> PolarsResult<Arc<RwLock<HashMap<String, Option<DataFrame>>>>> {
let df_map = self.df_map.read().await;
match df_map.get(outer_key) {
Some(inner_map) => {
let inner_map = Arc::clone(inner_map);
drop(df_map);
Ok(inner_map)
}
_ => Err(PolarsError::NoData(
format!("no DataFrame mapped for outer_key: {outer_key}").into(),
)),
}
}
pub async fn append_data_point<T>(&self, outer_key: &str, data_point: &T) -> PolarsResult<u32>
where
T: RawData,
{
let formatted_data_point = data_point.format()?;
let inner_map = self.access_inner_map(outer_key).await?;
let mut inner_map = inner_map.write().await;
#[cfg(debug_assertions)]
let start = Instant::now();
let mut num_appended = 0u32;
let df_schematics = self.df_schematics.read().await;
for (key, df) in inner_map.iter_mut() {
let df_schematic = df_schematics.get(key).ok_or(PolarsError::NoData(
format!("no DataFrameSchematic found for key: {key}").into(),
))?;
match df {
Some(df) => {
if !data_point.validate(df)? {
continue;
}
}
None => {
*df = Some(DataFrame::empty_with_schema(&df_schematic.schema));
}
}
if let Some(df) = df {
let prepared_df = df_schematic.prepare_data(&formatted_data_point)?;
df.vstack_mut(&prepared_df)?;
if let Some(minimum_rows) = df_schematic.minimum_rows {
*df = df.slice(-(minimum_rows as i64), minimum_rows as usize);
}
#[cfg(debug_assertions)]
dbg!(&df);
num_appended += 1;
}
}
#[cfg(debug_assertions)]
{
let elapsed = start.elapsed();
println!("Time elapsed: {elapsed:?}, outer_key: {outer_key}");
}
Ok(num_appended)
}
pub async fn append_series<T>(
&self,
outer_key: &str,
series_data: HashMap<PlSmallStr, Series>,
series_data_length: usize,
raw_data_type: &T,
) -> PolarsResult<u32>
where
T: RawData,
{
let inner_map = self.access_inner_map(outer_key).await?;
let mut inner_map = inner_map.write().await;
#[cfg(debug_assertions)]
let start = Instant::now();
let num_data_points = series_data_length as u32;
let mut num_appended = 0u32;
let df_schematics = self.df_schematics.read().await;
for (key, df) in inner_map.iter_mut() {
let df_schematic = df_schematics.get(key).ok_or(PolarsError::NoData(
format!("no DataFrameSchematic found for key: {key}").into(),
))?;
match df {
Some(df) => {
if !raw_data_type.validate_series(&series_data, df)? {
continue;
}
}
None => {
*df = Some(DataFrame::empty_with_schema(&df_schematic.schema));
}
}
if let Some(df) = df {
let prepared_df =
df_schematic.prepare_df_from_series(&series_data, series_data_length)?;
df.vstack_mut(&prepared_df)?;
#[cfg(debug_assertions)]
dbg!(&df);
num_appended += num_data_points;
}
}
#[cfg(debug_assertions)]
{
let elapsed = start.elapsed();
println!("Time elapsed: {elapsed:?}, outer_key: {outer_key}");
}
Ok(num_appended)
}
pub async fn df_schematic_keys(&self) -> Vec<String> {
self.df_schematics.read().await.keys().cloned().collect()
}
async fn outer_keys(&self) -> Vec<String> {
let df_map = self.df_map.read().await;
df_map.keys().cloned().collect()
}
pub async fn inner_keys_of_outer(&self, outer_key: &str) -> Option<Vec<String>> {
let df_map = self.df_map.read().await;
let map = df_map.get(outer_key)?;
let map = map.read().await;
Some(map.keys().cloned().collect())
}
pub async fn outer_keys_of_inner(&self, inner_key: &str) -> PolarsResult<Vec<String>> {
let mut inner_keys: Vec<String> = Vec::new();
let outer_keys = self.outer_keys().await;
for outer_key in outer_keys {
let inner_map = self.access_inner_map(&outer_key).await?;
let inner_map = inner_map.read().await;
if inner_map.contains_key(inner_key) {
inner_keys.push(outer_key);
}
}
Ok(inner_keys)
}
pub async fn df_height(&self, outer_key: &str, inner_key: &str) -> PolarsResult<usize> {
let inner_map = self.access_inner_map(outer_key).await?;
let inner_map = inner_map.read().await;
match inner_map.get(inner_key) {
Some(Some(df)) => Ok(df.height()),
_ => Err(PolarsError::NoData(
format!("no DataFrame found, outer_key: {outer_key}, inner_key: {inner_key}")
.into(),
)),
}
}
pub async fn evict_df(&self, outer_key: &str, inner_key: &str) -> PolarsResult<bool> {
let inner_map = self.access_inner_map(outer_key).await?;
let mut inner_map = inner_map.write().await;
if let Some(df) = inner_map.get_mut(inner_key) {
*df = None;
Ok(true)
} else {
Ok(false)
}
}
pub async fn evict_outer(&self, outer_key: &str) -> PolarsResult<bool> {
let mut df_map = self.df_map.write().await;
if let Some(outer_map) = df_map.get_mut(outer_key) {
*outer_map = Arc::new(RwLock::new(HashMap::<String, Option<DataFrame>>::new()));
Ok(true)
} else {
Ok(false)
}
}
pub async fn evict_inner(&self, inner_key: &str) -> PolarsResult<u32> {
let mut evicted_count = 0_u32;
let outer_keys = self.outer_keys().await;
for outer_key in outer_keys {
let inner_map = self.access_inner_map(&outer_key).await?;
let mut inner_map = inner_map.write().await;
if let Some(df) = inner_map.get_mut(inner_key) {
*df = None;
evicted_count += 1;
}
}
Ok(evicted_count)
}
pub async fn remove_df_map_entry(
&self,
outer_key: &str,
inner_key: &str,
) -> PolarsResult<bool> {
let inner_map = self.access_inner_map(outer_key).await?;
let mut inner_map = inner_map.write().await;
Ok(inner_map.remove(inner_key).is_some())
}
pub async fn df_to_arrow(
&self,
outer_key: &str,
inner_key: &str,
n_rows: Option<u32>,
exclude_columns: Option<Vec<String>>,
compression: Option<IpcCompression>,
) -> PolarsResult<(Arc<ArrowSchema>, Vec<ArrowRecordBatch>)> {
let inner_map = self.access_inner_map(outer_key).await?;
let inner_map = inner_map.read().await;
let df = inner_map.get(inner_key).cloned();
drop(inner_map);
match df {
Some(Some(df)) => {
let df_schematics = self.df_schematics.read().await;
let df_schematic = df_schematics.get(inner_key).ok_or(PolarsError::NoData(
format!("no DataFrameSchematic found for inner_key: {inner_key}").into(),
))?;
let mut df = df_schematic.feature_operator.collect_lazy_expressions(
df,
n_rows,
exclude_columns,
)?;
drop(df_schematics);
write_to_arrow(&mut df, compression)
}
Some(None) => Err(PolarsError::NoData(
format!("no DataFrame found, outer_key: {outer_key}, inner_key: {inner_key}")
.into(),
)),
None => Err(PolarsError::NoData(
format!("no value for inner_key: {inner_key}").into(),
)),
}
}
pub async fn persist_df(
&self,
duck_db_conn: Arc<Mutex<Connection>>,
outer_key: &str,
inner_key: &str,
compression: Option<IpcCompression>,
) -> Result<bool, String> {
let duck_db_conn = duck_db_conn.lock().await;
let inner_map = self
.access_inner_map(outer_key)
.await
.map_err(|e| e.to_string())?;
let inner_map = inner_map.read().await;
let df = inner_map.get(inner_key).cloned();
drop(inner_map);
match df {
Some(Some(df)) => {
let df_schematics = self.df_schematics.read().await;
let df_schematic = df_schematics.get(inner_key).ok_or(format!(
"no DataFrameSchematic found for inner_key: {inner_key}"
))?;
let mut df = df_schematic
.feature_operator
.collect_lazy_expressions(df, None, None)
.map_err(|e| e.to_string())?;
drop(df_schematics);
let (schema, record_batches) = write_to_arrow(&mut df, compression).unwrap();
save_to_duckdb(&duck_db_conn, inner_key, schema, record_batches)
.map_err(|e| e.to_string())?;
}
Some(None) => {
return Err(format!(
"no DataFrame found, outer_key: {outer_key}, inner_key: {inner_key}"
));
}
None => return Err(format!("no value for inner_key: {inner_key}")),
}
Ok(true)
}
}
#[cfg(test)]
mod tests {
use super::*;
use tokio::runtime::Runtime;
#[derive(Debug)]
struct DataPoint {
x: i32,
y: i32,
}
impl RawData for DataPoint {
fn validate(&self, _df: &DataFrame) -> PolarsResult<bool> {
Ok(true)
}
fn validate_series(
&self,
_series_data: &HashMap<PlSmallStr, Series>,
_df: &DataFrame,
) -> PolarsResult<bool> {
Ok(true)
}
fn format(&self) -> PolarsResult<Vec<(PlSmallStr, AnyValue<'_>)>> {
Ok(vec![
(PlSmallStr::from_static("x"), AnyValue::Int32(self.x)),
(PlSmallStr::from_static("y"), AnyValue::Int32(self.y)),
])
}
}
const OUTER_KEY: &str = "outer";
const ANOTHER_OUTER_KEY: &str = "another_outer";
const INNER_KEY: &str = "inner";
const ANOTHER_INNER_KEY: &str = "another_inner";
fn create_df_collection() -> DataFrameCollection {
let field_types: Vec<Field> = vec![
Field::new("x".into(), DataType::Int32),
Field::new("y".into(), DataType::Int32),
];
let mut df_schematic = DataFrameSchematic::new(field_types, vec![vec![]]);
df_schematic.with_minimum_rows(5);
let mut df_schematics = HashMap::new();
df_schematics.insert(String::from(INNER_KEY), df_schematic.clone());
df_schematics.insert(String::from(ANOTHER_INNER_KEY), df_schematic);
DataFrameCollection::new(df_schematics)
}
#[test]
fn append_data() {
let rt = Runtime::new().unwrap();
rt.block_on(async {
let df_collection = create_df_collection();
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(INNER_KEY))
.await;
let data_point = DataPoint { x: 1, y: 10 };
let num_appended = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await
.unwrap();
assert_eq!(num_appended, 1);
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(ANOTHER_INNER_KEY))
.await;
let data_point = DataPoint { x: 2, y: 10 };
let num_appended = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await
.unwrap();
assert_eq!(num_appended, 2);
let mut data_series_map: HashMap<PlSmallStr, Series> = HashMap::new();
data_series_map.insert(
PlSmallStr::from_static("x"),
Series::new(PlSmallStr::from_static("x"), vec![3, 4, 5, 6, 7]),
);
data_series_map.insert(
PlSmallStr::from_static("y"),
Series::new(PlSmallStr::from_static("y"), vec![10, 10, 10, 10, 10]),
);
let num_appended = df_collection
.append_series(OUTER_KEY, data_series_map, 5, &DataPoint { x: 3, y: 10 })
.await
.unwrap();
assert_eq!(num_appended, 10);
});
}
#[test]
fn minimum_rows() {
let rt = Runtime::new().unwrap();
rt.block_on(async {
let df_collection = create_df_collection();
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(INNER_KEY))
.await;
let data_point = DataPoint { x: 1, y: 10 };
let num_appended = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await
.unwrap();
assert_eq!(num_appended, 1);
let data_point = DataPoint { x: 2, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let data_point = DataPoint { x: 3, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let data_point = DataPoint { x: 4, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let data_point = DataPoint { x: 5, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 5);
let data_point = DataPoint { x: 6, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let df = df_collection
.df_map
.read()
.await
.get(OUTER_KEY)
.unwrap()
.read()
.await
.get(INNER_KEY)
.unwrap()
.clone()
.unwrap();
let select_last = df
.lazy()
.clone()
.select([col("x").last()])
.collect()
.unwrap();
let last_value = select_last.column("x").unwrap().get(0).unwrap();
assert_eq!(last_value, AnyValue::Int32(6));
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 5);
let data_point = DataPoint { x: 7, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 5);
let data_point = DataPoint { x: 8, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 5);
let data_point = DataPoint { x: 9, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let data_point = DataPoint { x: 10, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let data_point = DataPoint { x: 11, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 5);
let df = df_collection
.df_map
.read()
.await
.get(OUTER_KEY)
.unwrap()
.read()
.await
.get(INNER_KEY)
.unwrap()
.clone()
.unwrap();
let select_last = df
.lazy()
.clone()
.select([col("x").last()])
.collect()
.unwrap();
let last_value = select_last.column("x").unwrap().get(0).unwrap();
assert_eq!(last_value, AnyValue::Int32(11));
let set_operation_result = df_collection.set_minimum_rows(&INNER_KEY, 8).await;
assert!(set_operation_result);
let set_operation_result = df_collection.set_minimum_rows("bla", 8).await;
assert!(!set_operation_result);
let data_point = DataPoint { x: 12, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 6);
let data_point = DataPoint { x: 13, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 7);
let data_point = DataPoint { x: 14, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 8);
let data_point = DataPoint { x: 15, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 8);
let df = df_collection
.df_map
.read()
.await
.get(OUTER_KEY)
.unwrap()
.read()
.await
.get(INNER_KEY)
.unwrap()
.clone()
.unwrap();
let select_last = df
.lazy()
.clone()
.select([col("x").last()])
.collect()
.unwrap();
let last_value = select_last.column("x").unwrap().get(0).unwrap();
assert_eq!(last_value, AnyValue::Int32(15));
});
}
#[test]
fn key_presence() {
let rt = Runtime::new().unwrap();
rt.block_on(async {
let df_collection = create_df_collection();
let insert_successful = df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(INNER_KEY))
.await;
assert!(insert_successful);
let insert_successful = df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(ANOTHER_INNER_KEY))
.await;
assert!(insert_successful);
let insert_successful = df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from("third_inner"))
.await;
assert!(!insert_successful);
let outer_keys = df_collection.outer_keys().await;
assert!(outer_keys.contains(&String::from(OUTER_KEY)));
let inner_keys = df_collection.inner_keys_of_outer(OUTER_KEY).await.unwrap();
assert!(inner_keys.contains(&String::from(INNER_KEY)));
let outer_keys_of_inner = df_collection.outer_keys_of_inner(INNER_KEY).await.unwrap();
assert!(outer_keys_of_inner.contains(&String::from(OUTER_KEY)));
assert_eq!(outer_keys_of_inner, vec![OUTER_KEY]);
df_collection
.insert_inner_map(String::from(ANOTHER_OUTER_KEY), String::from(INNER_KEY))
.await;
df_collection
.insert_inner_map(String::from("third_outer"), String::from(INNER_KEY))
.await;
let outer_keys_of_inner = df_collection.outer_keys_of_inner(INNER_KEY).await.unwrap();
assert!(outer_keys_of_inner.contains(&String::from(OUTER_KEY)));
assert!(outer_keys_of_inner.contains(&String::from(ANOTHER_OUTER_KEY)));
assert!(outer_keys_of_inner.contains(&String::from("third_outer")));
assert_eq!(outer_keys_of_inner.len(), 3);
let outer_keys_of_inner = df_collection
.outer_keys_of_inner(ANOTHER_INNER_KEY)
.await
.unwrap();
assert!(outer_keys_of_inner.contains(&String::from(OUTER_KEY)));
assert!(outer_keys_of_inner.contains(&String::from(ANOTHER_OUTER_KEY)) == false);
assert_eq!(outer_keys_of_inner, vec![OUTER_KEY]);
});
}
#[test]
fn evict() {
let rt = Runtime::new().unwrap();
rt.block_on(async {
let df_collection = create_df_collection();
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(INNER_KEY))
.await;
let data_point = DataPoint { x: 5, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let outer_map = df_collection.df_map.read().await;
assert!(outer_map.contains_key(OUTER_KEY));
assert!(
df_collection
.outer_keys()
.await
.contains(&String::from(OUTER_KEY))
);
let inner_map = outer_map.get(OUTER_KEY).unwrap().read().await;
assert!(inner_map.contains_key(INNER_KEY));
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_some());
drop(inner_map);
drop(outer_map);
assert!(!df_collection.evict_df(OUTER_KEY, "nono").await.unwrap(),);
assert!(df_collection.evict_df(OUTER_KEY, INNER_KEY).await.unwrap(),);
let outer_map = df_collection.df_map.read().await;
let inner_map = outer_map.get(OUTER_KEY).unwrap().read().await;
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_none());
});
}
#[test]
fn evict_outer() {
let rt = Runtime::new().unwrap();
rt.block_on(async {
let df_collection = create_df_collection();
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(INNER_KEY))
.await;
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(ANOTHER_INNER_KEY))
.await;
df_collection
.insert_inner_map(String::from(ANOTHER_OUTER_KEY), String::from(INNER_KEY))
.await;
let data_point = DataPoint { x: 5, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let _ = df_collection
.append_data_point(ANOTHER_OUTER_KEY, &data_point)
.await;
let outer_map = df_collection.df_map.read().await;
assert!(outer_map.contains_key(OUTER_KEY));
assert!(outer_map.contains_key(ANOTHER_OUTER_KEY));
assert!(
df_collection
.outer_keys()
.await
.contains(&String::from(OUTER_KEY))
);
let inner_map = outer_map.get(OUTER_KEY).unwrap().read().await;
assert!(inner_map.contains_key(INNER_KEY));
assert!(inner_map.contains_key(ANOTHER_INNER_KEY));
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_some());
let df = inner_map.get(ANOTHER_INNER_KEY).unwrap();
assert!(df.is_some());
drop(inner_map);
drop(outer_map);
let outer_map = df_collection.df_map.read().await;
let inner_map = outer_map.get(ANOTHER_OUTER_KEY).unwrap().read().await;
assert!(inner_map.contains_key(INNER_KEY));
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_some());
drop(inner_map);
drop(outer_map);
assert!(!df_collection.evict_outer("nono").await.unwrap());
assert!(df_collection.evict_outer(OUTER_KEY).await.unwrap());
let outer_map = df_collection.df_map.read().await;
let inner_map = outer_map.get(OUTER_KEY).unwrap().read().await;
assert!(inner_map.is_empty());
let inner_map = outer_map.get(ANOTHER_OUTER_KEY).unwrap().read().await;
assert!(inner_map.contains_key(INNER_KEY));
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_some());
});
}
#[test]
fn evict_inner() {
let rt = Runtime::new().unwrap();
rt.block_on(async {
let df_collection = create_df_collection();
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(INNER_KEY))
.await;
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(ANOTHER_INNER_KEY))
.await;
df_collection
.insert_inner_map(String::from(ANOTHER_OUTER_KEY), String::from(INNER_KEY))
.await;
df_collection
.insert_inner_map(
String::from(ANOTHER_OUTER_KEY),
String::from(ANOTHER_INNER_KEY),
)
.await;
let data_point = DataPoint { x: 5, y: 10 };
let num_appended = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await
.unwrap();
assert_eq!(num_appended, 2);
let num_appended = df_collection
.append_data_point(ANOTHER_OUTER_KEY, &data_point)
.await
.unwrap();
assert_eq!(num_appended, 2);
let outer_map = df_collection.df_map.read().await;
let inner_map = outer_map.get(OUTER_KEY).unwrap().read().await;
assert!(outer_map.contains_key(OUTER_KEY));
assert!(outer_map.contains_key(ANOTHER_OUTER_KEY));
assert!(
df_collection
.outer_keys()
.await
.contains(&String::from(OUTER_KEY))
);
assert!(inner_map.contains_key(INNER_KEY));
assert!(inner_map.contains_key(ANOTHER_INNER_KEY));
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_some());
let df = inner_map.get(ANOTHER_INNER_KEY).unwrap();
assert!(df.is_some());
drop(inner_map);
drop(outer_map);
let outer_map = df_collection.df_map.read().await;
let inner_map = outer_map.get(ANOTHER_OUTER_KEY).unwrap().read().await;
assert!(inner_map.contains_key(INNER_KEY));
assert!(inner_map.contains_key(ANOTHER_INNER_KEY));
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_some());
let df = inner_map.get(ANOTHER_INNER_KEY).unwrap();
assert!(df.is_some());
drop(inner_map);
drop(outer_map);
df_collection.evict_inner(INNER_KEY).await.unwrap();
let outer_map = df_collection.df_map.read().await;
let inner_map = outer_map.get(OUTER_KEY).unwrap().read().await;
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_none());
let df = inner_map.get(ANOTHER_INNER_KEY).unwrap();
assert!(df.is_some());
drop(inner_map);
let inner_map = outer_map.get(ANOTHER_OUTER_KEY).unwrap().read().await;
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_none());
let df = inner_map.get(ANOTHER_INNER_KEY).unwrap();
assert!(df.is_some());
});
}
#[test]
fn remove_entry() {
let rt = Runtime::new().unwrap();
rt.block_on(async {
let df_collection = create_df_collection();
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(INNER_KEY))
.await;
let data_point = DataPoint { x: 5, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let outer_map = df_collection.df_map.read().await;
assert!(outer_map.contains_key(OUTER_KEY));
assert!(
df_collection
.outer_keys()
.await
.contains(&String::from(OUTER_KEY))
);
let inner_map = outer_map.get(OUTER_KEY).unwrap().read().await;
assert!(inner_map.contains_key(INNER_KEY));
let df = inner_map.get(INNER_KEY).unwrap();
assert!(df.is_some());
drop(inner_map);
drop(outer_map);
assert!(
!df_collection
.remove_df_map_entry(OUTER_KEY, "nono")
.await
.unwrap(),
);
assert!(
df_collection
.remove_df_map_entry(OUTER_KEY, INNER_KEY)
.await
.unwrap(),
);
let outer_map = df_collection.df_map.read().await;
let inner_map = outer_map.get(OUTER_KEY).unwrap().read().await;
let entry = inner_map.get(INNER_KEY);
assert!(entry.is_none())
});
}
#[test]
fn df_height() {
let rt = Runtime::new().unwrap();
rt.block_on(async {
let df_collection = create_df_collection();
df_collection
.insert_inner_map(String::from(OUTER_KEY), String::from(INNER_KEY))
.await;
let outer_keys = df_collection.outer_keys().await;
assert!(outer_keys.contains(&String::from(OUTER_KEY)));
let inner_keys = df_collection.inner_keys_of_outer(OUTER_KEY).await.unwrap();
assert!(inner_keys.contains(&String::from(INNER_KEY)));
let data_point = DataPoint { x: 5, y: 10 };
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 1);
let _ = df_collection
.append_data_point(OUTER_KEY, &data_point)
.await;
let height = df_collection.df_height(OUTER_KEY, INNER_KEY).await.unwrap();
assert_eq!(height, 2);
});
}
}