Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 15 additions & 3 deletions aw-query/src/functions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -272,10 +272,22 @@ mod qfunctions {
_ds: &Datastore,
) -> Result<DataType, QueryError> {
// typecheck
validate::args_length(&args, 1)?;
let events: Vec<Event> = args.into_iter().next().unwrap().try_into()?;
validate::args_length(&args, 1).or_else(|_| validate::args_length(&args, 2))?;
let mut args = args.into_iter();
let events: Vec<Event> = args.next().unwrap().try_into()?;
// Optional pulsetime in seconds, defaults to 5 like aw-core
let pulsetime_secs: f64 = match args.next() {
Some(arg) => arg.try_into()?,
None => 5.0,
};
if !pulsetime_secs.is_finite() || pulsetime_secs < 0.0 {
return Err(QueryError::InvalidFunctionParameters(format!(
"flood pulsetime must be a non-negative number of seconds, got {pulsetime_secs}"
)));
}
let pulsetime = chrono::Duration::nanoseconds((pulsetime_secs * 1e9).round() as i64);
// Run flood
let mut flooded_events = aw_transform::flood(events, chrono::Duration::seconds(5));
let mut flooded_events = aw_transform::flood(events, pulsetime);
// Put events back into DataType::Event container
let mut tagged_flooded_events = Vec::new();
for event in flooded_events.drain(..) {
Expand Down
51 changes: 51 additions & 0 deletions aw-query/tests/query.rs
Original file line number Diff line number Diff line change
Expand Up @@ -341,6 +341,57 @@ mod query_tests {
}
}

#[test]
fn test_flood_pulsetime() {
use chrono::{DateTime, Utc};
use std::str::FromStr;

let ds = setup_datastore_with_bucket();
let start: DateTime<Utc> = DateTime::from_str("2000-01-01T00:00:00Z").unwrap();
let e1 = Event {
id: None,
timestamp: start,
duration: Duration::seconds(10),
data: json_map! {"key": json!("a")},
};
let e2 = Event {
id: None,
timestamp: start + Duration::seconds(17),
duration: Duration::seconds(10),
data: json_map! {"key": json!("b")},
};
ds.insert_events(BUCKET_ID, &[e1, e2]).unwrap();
let interval = TimeInterval::new_from_string(TIME_INTERVAL).unwrap();

let durations = |args: &str| -> Vec<f64> {
let code = format!(r#"return flood(query_bucket("{BUCKET_ID}"){args});"#);
match aw_query::query(&code, &interval, &ds).unwrap() {
aw_query::DataType::List(l) => l
.into_iter()
.map(|e| match e {
aw_query::DataType::Event(e) => {
e.duration.num_milliseconds() as f64 / 1000.0
}
ref data => panic!("Wrong datatype, {data:?}"),
})
.collect(),
ref data => panic!("Wrong datatype, {data:?}"),
}
};
// The 7 s gap is left alone with the default pulsetime of 5 s
assert_eq!(durations(""), vec![10.0, 10.0]);
assert_eq!(durations(", 6.5"), vec![10.0, 10.0]);
// and filled in the middle when the gap is at most pulsetime
assert_eq!(durations(", 7"), vec![13.5, 13.5]);
assert_eq!(durations(", 10.5"), vec![13.5, 13.5]);

let code = format!(r#"return flood(query_bucket("{BUCKET_ID}"), 0 - 1);"#);
assert_err_type!(
aw_query::query(&code, &interval, &ds),
QueryError::InvalidFunctionParameters(_)
);
}

#[test]
fn test_all_functions() {
let ds = setup_datastore_populated();
Expand Down
Loading
Loading