use key_paths_core::KeyPaths;
use std::collections::HashMap;
use std::time::SystemTime;
#[cfg(feature = "datetime")]
use chrono::{DateTime, TimeZone};
pub struct Query<'a, T: 'static> {
data: &'a [T],
filters: Vec<Box<dyn Fn(&T) -> bool>>,
}
impl<'a, T: 'static> Query<'a, T> {
pub fn new(data: &'a [T]) -> Self {
Self {
data,
filters: Vec::new(),
}
}
pub fn where_<F>(mut self, path: KeyPaths<T, F>, predicate: impl Fn(&F) -> bool + 'static) -> Self
where
F: 'static,
{
self.filters.push(Box::new(move |item| {
path.get(item).map_or(false, |val| predicate(val))
}));
self
}
pub fn all(&self) -> Vec<&T> {
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.collect()
}
pub fn first(&self) -> Option<&T> {
self.data
.iter()
.find(|item| self.filters.iter().all(|f| f(item)))
}
pub fn count(&self) -> usize {
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.count()
}
pub fn limit(&self, n: usize) -> Vec<&T> {
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.take(n)
.collect()
}
pub fn skip<'b>(&'b self, offset: usize) -> QueryWithSkip<'a, 'b, T> {
QueryWithSkip {
query: self,
offset,
}
}
pub fn select<F>(&self, path: KeyPaths<T, F>) -> Vec<F>
where
F: Clone + 'static,
{
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.collect()
}
pub fn sum<F>(&self, path: KeyPaths<T, F>) -> F
where
F: Clone + std::ops::Add<Output = F> + Default + 'static,
{
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.fold(F::default(), |acc, val| acc + val)
}
pub fn avg(&self, path: KeyPaths<T, f64>) -> Option<f64> {
let items: Vec<f64> = self
.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.collect();
if items.is_empty() {
None
} else {
Some(items.iter().sum::<f64>() / items.len() as f64)
}
}
pub fn min<F>(&self, path: KeyPaths<T, F>) -> Option<F>
where
F: Ord + Clone + 'static,
{
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.min()
}
pub fn max<F>(&self, path: KeyPaths<T, F>) -> Option<F>
where
F: Ord + Clone + 'static,
{
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.max()
}
pub fn min_float(&self, path: KeyPaths<T, f64>) -> Option<f64> {
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.min_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal))
}
pub fn max_float(&self, path: KeyPaths<T, f64>) -> Option<f64> {
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.max_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal))
}
pub fn exists(&self) -> bool {
self.data
.iter()
.any(|item| self.filters.iter().all(|f| f(item)))
}
pub fn where_after_systemtime(self, path: KeyPaths<T, SystemTime>, reference: SystemTime) -> Self {
self.where_(path, move |time| time > &reference)
}
pub fn where_before_systemtime(self, path: KeyPaths<T, SystemTime>, reference: SystemTime) -> Self {
self.where_(path, move |time| time < &reference)
}
pub fn where_between_systemtime(
self,
path: KeyPaths<T, SystemTime>,
start: SystemTime,
end: SystemTime,
) -> Self {
self.where_(path, move |time| time >= &start && time <= &end)
}
}
#[cfg(feature = "datetime")]
impl<'a, T: 'static> Query<'a, T> {
pub fn where_after<Tz>(self, path: KeyPaths<T, DateTime<Tz>>, reference: DateTime<Tz>) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
self.where_(path, move |time| time > &reference)
}
pub fn where_before<Tz>(self, path: KeyPaths<T, DateTime<Tz>>, reference: DateTime<Tz>) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
self.where_(path, move |time| time < &reference)
}
pub fn where_between<Tz>(
self,
path: KeyPaths<T, DateTime<Tz>>,
start: DateTime<Tz>,
end: DateTime<Tz>,
) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
self.where_(path, move |time| time >= &start && time <= &end)
}
pub fn where_today<Tz>(self, path: KeyPaths<T, DateTime<Tz>>, now: DateTime<Tz>) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
self.where_(path, move |time| {
time.date_naive() == now.date_naive()
})
}
pub fn where_year<Tz>(self, path: KeyPaths<T, DateTime<Tz>>, year: i32) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
use chrono::Datelike;
self.where_(path, move |time| time.year() == year)
}
pub fn where_month<Tz>(self, path: KeyPaths<T, DateTime<Tz>>, month: u32) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
use chrono::Datelike;
self.where_(path, move |time| time.month() == month)
}
pub fn where_day<Tz>(self, path: KeyPaths<T, DateTime<Tz>>, day: u32) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
use chrono::Datelike;
self.where_(path, move |time| time.day() == day)
}
pub fn where_weekend<Tz>(self, path: KeyPaths<T, DateTime<Tz>>) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
use chrono::Datelike;
self.where_(path, |time| {
let weekday = time.weekday().num_days_from_monday();
weekday >= 5
})
}
pub fn where_weekday<Tz>(self, path: KeyPaths<T, DateTime<Tz>>) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
use chrono::Datelike;
self.where_(path, |time| {
let weekday = time.weekday().num_days_from_monday();
weekday < 5
})
}
pub fn where_business_hours<Tz>(self, path: KeyPaths<T, DateTime<Tz>>) -> Self
where
Tz: TimeZone + 'static,
Tz::Offset: std::fmt::Display,
{
use chrono::Timelike;
self.where_(path, |time| {
let hour = time.hour();
hour >= 9 && hour < 17
})
}
}
impl<'a, T: 'static + Clone> Query<'a, T> {
pub fn order_by<F>(&self, path: KeyPaths<T, F>) -> Vec<T>
where
F: Ord + Clone + 'static,
{
let mut results: Vec<T> = self
.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.cloned()
.collect();
results.sort_by_key(|item| path.get(item).cloned());
results
}
pub fn order_by_desc<F>(&self, path: KeyPaths<T, F>) -> Vec<T>
where
F: Ord + Clone + 'static,
{
let mut results: Vec<T> = self
.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.cloned()
.collect();
results.sort_by(|a, b| {
let a_val = path.get(a).cloned();
let b_val = path.get(b).cloned();
b_val.cmp(&a_val)
});
results
}
pub fn order_by_float(&self, path: KeyPaths<T, f64>) -> Vec<T> {
let mut results: Vec<T> = self
.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.cloned()
.collect();
results.sort_by(|a, b| {
let a_val = path.get(a).cloned().unwrap_or(0.0);
let b_val = path.get(b).cloned().unwrap_or(0.0);
a_val.partial_cmp(&b_val).unwrap_or(std::cmp::Ordering::Equal)
});
results
}
pub fn order_by_float_desc(&self, path: KeyPaths<T, f64>) -> Vec<T> {
let mut results: Vec<T> = self
.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.cloned()
.collect();
results.sort_by(|a, b| {
let a_val = path.get(a).cloned().unwrap_or(0.0);
let b_val = path.get(b).cloned().unwrap_or(0.0);
b_val.partial_cmp(&a_val).unwrap_or(std::cmp::Ordering::Equal)
});
results
}
pub fn group_by<F>(&self, path: KeyPaths<T, F>) -> HashMap<F, Vec<T>>
where
F: Eq + std::hash::Hash + Clone + 'static,
{
let mut groups: HashMap<F, Vec<T>> = HashMap::new();
for item in self.data.iter() {
if self.filters.iter().all(|f| f(item)) {
if let Some(key) = path.get(item).cloned() {
groups.entry(key).or_insert_with(Vec::new).push(item.clone());
}
}
}
groups
}
#[cfg(feature = "datetime")]
pub fn min_timestamp(&self, path: KeyPaths<T, i64>) -> Option<i64> {
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.min()
}
#[cfg(feature = "datetime")]
pub fn max_timestamp(&self, path: KeyPaths<T, i64>) -> Option<i64> {
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.max()
}
#[cfg(feature = "datetime")]
pub fn avg_timestamp(&self, path: KeyPaths<T, i64>) -> Option<i64> {
let items: Vec<i64> = self
.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.collect();
if items.is_empty() {
None
} else {
Some(items.iter().sum::<i64>() / items.len() as i64)
}
}
#[cfg(feature = "datetime")]
pub fn sum_timestamp(&self, path: KeyPaths<T, i64>) -> i64 {
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter_map(|item| path.get(item).cloned())
.sum()
}
#[cfg(feature = "datetime")]
pub fn count_timestamp(&self, path: KeyPaths<T, i64>) -> usize {
self.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.filter(|item| path.get(item).is_some())
.count()
}
#[cfg(feature = "datetime")]
pub fn where_after_timestamp(self, path: KeyPaths<T, i64>, reference: i64) -> Self {
self.where_(path, move |timestamp| timestamp > &reference)
}
#[cfg(feature = "datetime")]
pub fn where_before_timestamp(self, path: KeyPaths<T, i64>, reference: i64) -> Self {
self.where_(path, move |timestamp| timestamp < &reference)
}
#[cfg(feature = "datetime")]
pub fn where_between_timestamp(self, path: KeyPaths<T, i64>, start: i64, end: i64) -> Self {
self.where_(path, move |timestamp| timestamp >= &start && timestamp <= &end)
}
#[cfg(feature = "datetime")]
pub fn where_last_days_timestamp(self, path: KeyPaths<T, i64>, days: i64) -> Self {
let now = chrono::Utc::now().timestamp_millis();
let cutoff = now - (days * 24 * 60 * 60 * 1000); self.where_after_timestamp(path, cutoff)
}
#[cfg(feature = "datetime")]
pub fn where_next_days_timestamp(self, path: KeyPaths<T, i64>, days: i64) -> Self {
let now = chrono::Utc::now().timestamp_millis();
let cutoff = now + (days * 24 * 60 * 60 * 1000); self.where_before_timestamp(path, cutoff)
}
#[cfg(feature = "datetime")]
pub fn where_last_hours_timestamp(self, path: KeyPaths<T, i64>, hours: i64) -> Self {
let now = chrono::Utc::now().timestamp_millis();
let cutoff = now - (hours * 60 * 60 * 1000); self.where_after_timestamp(path, cutoff)
}
#[cfg(feature = "datetime")]
pub fn where_next_hours_timestamp(self, path: KeyPaths<T, i64>, hours: i64) -> Self {
let now = chrono::Utc::now().timestamp_millis();
let cutoff = now + (hours * 60 * 60 * 1000); self.where_before_timestamp(path, cutoff)
}
#[cfg(feature = "datetime")]
pub fn where_last_minutes_timestamp(self, path: KeyPaths<T, i64>, minutes: i64) -> Self {
let now = chrono::Utc::now().timestamp_millis();
let cutoff = now - (minutes * 60 * 1000); self.where_after_timestamp(path, cutoff)
}
#[cfg(feature = "datetime")]
pub fn where_next_minutes_timestamp(self, path: KeyPaths<T, i64>, minutes: i64) -> Self {
let now = chrono::Utc::now().timestamp_millis();
let cutoff = now + (minutes * 60 * 1000); self.where_before_timestamp(path, cutoff)
}
#[cfg(feature = "datetime")]
pub fn order_by_timestamp(&self, path: KeyPaths<T, i64>) -> Vec<T> {
let mut results: Vec<T> = self
.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.cloned()
.collect();
results.sort_by(|a, b| {
let a_val = path.get(a).cloned().unwrap_or(0);
let b_val = path.get(b).cloned().unwrap_or(0);
a_val.cmp(&b_val)
});
results
}
#[cfg(feature = "datetime")]
pub fn order_by_timestamp_desc(&self, path: KeyPaths<T, i64>) -> Vec<T> {
let mut results: Vec<T> = self
.data
.iter()
.filter(|item| self.filters.iter().all(|f| f(item)))
.cloned()
.collect();
results.sort_by(|a, b| {
let a_val = path.get(a).cloned().unwrap_or(0);
let b_val = path.get(b).cloned().unwrap_or(0);
b_val.cmp(&a_val)
});
results
}
}
pub struct QueryWithSkip<'a, 'b, T: 'static> {
query: &'b Query<'a, T>,
offset: usize,
}
impl<'a, 'b, T: 'static> QueryWithSkip<'a, 'b, T> {
pub fn limit(&self, n: usize) -> Vec<&'a T> {
self.query
.data
.iter()
.filter(|item| self.query.filters.iter().all(|f| f(item)))
.skip(self.offset)
.take(n)
.collect()
}
}
#[cfg(feature = "parallel")]
impl<'a, T: 'static + Send + Sync> Query<'a, T> {
pub fn all_parallel(&self) -> Vec<&'a T> {
use rayon::prelude::*;
self.data.par_iter().collect()
}
pub fn count_parallel(&self) -> usize {
use rayon::prelude::*;
self.data.par_iter().count()
}
pub fn exists_parallel(&self) -> bool {
use rayon::prelude::*;
self.data.par_iter().any(|_| true)
}
pub fn min_parallel<F>(&self, path: KeyPaths<T, F>) -> Option<F>
where
F: Ord + Clone + 'static + Send + Sync,
{
use rayon::prelude::*;
self.data
.par_iter()
.filter_map(|item| path.get(item).cloned())
.min()
}
pub fn max_parallel<F>(&self, path: KeyPaths<T, F>) -> Option<F>
where
F: Ord + Clone + 'static + Send + Sync,
{
use rayon::prelude::*;
self.data
.par_iter()
.filter_map(|item| path.get(item).cloned())
.max()
}
pub fn sum_parallel<F>(&self, path: KeyPaths<T, F>) -> F
where
F: Clone + std::ops::Add<Output = F> + Default + 'static + Send + Sync + std::iter::Sum,
{
use rayon::prelude::*;
self.data
.par_iter()
.filter_map(|item| path.get(item).cloned())
.sum()
}
pub fn avg_parallel(&self, path: KeyPaths<T, f64>) -> Option<f64> {
use rayon::prelude::*;
let items: Vec<f64> = self.data
.par_iter()
.filter_map(|item| path.get(item).cloned())
.collect();
if items.is_empty() {
None
} else {
Some(items.par_iter().sum::<f64>() / items.len() as f64)
}
}
pub fn min_timestamp_parallel(&self, path: KeyPaths<T, i64>) -> Option<i64> {
use rayon::prelude::*;
self.data
.par_iter()
.filter_map(|item| path.get(item).cloned())
.min()
}
pub fn max_timestamp_parallel(&self, path: KeyPaths<T, i64>) -> Option<i64> {
use rayon::prelude::*;
self.data
.par_iter()
.filter_map(|item| path.get(item).cloned())
.max()
}
pub fn avg_timestamp_parallel(&self, path: KeyPaths<T, i64>) -> Option<i64> {
use rayon::prelude::*;
let items: Vec<i64> = self.data
.par_iter()
.filter_map(|item| path.get(item).cloned())
.collect();
if items.is_empty() {
None
} else {
Some(items.par_iter().sum::<i64>() / items.len() as i64)
}
}
pub fn sum_timestamp_parallel(&self, path: KeyPaths<T, i64>) -> i64 {
use rayon::prelude::*;
self.data
.par_iter()
.filter_map(|item| path.get(item).cloned())
.sum()
}
pub fn count_timestamp_parallel(&self, path: KeyPaths<T, i64>) -> usize {
use rayon::prelude::*;
self.data
.par_iter()
.filter(|item| path.get(item).is_some())
.count()
}
}