From 7324b3285b9d339fed89f0b675f73bc17b62530e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Erik=20Bj=C3=A4reholt?= Date: Sat, 26 Sep 2026 18:44:53 +0200 Subject: [PATCH 1/2] fix(transform): port aw-core's flood() so both servers agree aw-core's flood() was fixed in ActivityWatch/aw-core#143 so that its output never double-counts time (ActivityWatch/activitywatch#1369), but the Rust flood() still let overlapping events with different data through. Port aw-core's algorithm: - a pairwise pass that merges same-data events and splits gaps between differing events at the midpoint, for gaps up to and including pulsetime (was: strictly less than) - a normalization pass in which the later of two overlapping events with different data wins, and zero-duration events are dropped - sort by (timestamp, duration) like aw-core The flood() query function now takes an optional pulsetime in seconds, defaulting to 5 like aw-core. Fixes #746 --- aw-query/src/functions.rs | 18 +- aw-query/tests/query.rs | 2 + aw-transform/src/flood.rs | 343 ++++++++++++++++++++++++-------------- 3 files changed, 237 insertions(+), 126 deletions(-) diff --git a/aw-query/src/functions.rs b/aw-query/src/functions.rs index 24605507..c45fd15c 100644 --- a/aw-query/src/functions.rs +++ b/aw-query/src/functions.rs @@ -272,10 +272,22 @@ mod qfunctions { _ds: &Datastore, ) -> Result { // typecheck - validate::args_length(&args, 1)?; - let events: Vec = 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 = 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(..) { diff --git a/aw-query/tests/query.rs b/aw-query/tests/query.rs index bec7c99e..2b7f607a 100644 --- a/aw-query/tests/query.rs +++ b/aw-query/tests/query.rs @@ -353,6 +353,8 @@ mod query_tests { r#" events = query_bucket(find_bucket("{}", "testhost")); events = flood(events); + events = flood(events, 10); + events = flood(events, 0.5); events = sort_by_duration(events); events = limit_events(events, 10000); events = sort_by_timestamp(events); diff --git a/aw-transform/src/flood.rs b/aw-transform/src/flood.rs index 7d6e7ba7..069b079f 100644 --- a/aw-transform/src/flood.rs +++ b/aw-transform/src/flood.rs @@ -1,20 +1,24 @@ -use aw_models::Event; +use std::cmp::{max, min}; -use crate::sort_by_timestamp; +use aw_models::Event; +use chrono::Duration; -/// Floods event to the nearest neighbouring event if within the specified pulsetime. +/// Fills short gaps between events and merges nearby events with the same data. /// -/// Also merges events if they have the same data and are within the pulsetime. +/// This is a port of aw-core's `flood()` and must stay in sync with it, since the Python and Rust +/// servers are expected to return the same results for the same query: +/// https://github.com/ActivityWatch/aw-core/blob/master/aw_transform/flood.py /// -/// "Flooding" actually performs two similar, but slightly different tasks, which we can describe as: -/// - Flooding: removes small gaps between events, by extending nearby events to fill the space. -/// - Merging: merging adjacent or overlapping events that have the same data. -/// - This is sometimes done in conjunction with flooding, if events with nearby events have the same data. +/// "Flooding" performs two similar, but slightly different tasks: +/// - Flooding: removes gaps of at most `pulsetime` between events with different data, by +/// extending both events to the middle of the gap. +/// - Merging: merges overlapping events, and events separated by a gap of at most `pulsetime`, +/// that have the same data. /// -/// Python implementation: -/// https://github.com/ActivityWatch/aw-core/blob/master/aw_transform/flood.py -/// Python re-implementation using generators (different spec): -/// https://github.com/ErikBjare/copilot-testing/blob/master/playground/flooding.py +/// Events are processed in (timestamp, duration) order. The output never contains overlapping +/// events: where events with different data overlap, the later event takes precedence and the +/// earlier one is cut off at its start (see ActivityWatch/activitywatch#1369). Zero-duration +/// events are dropped. /// /// # Example /// @@ -33,132 +37,100 @@ use crate::sort_by_timestamp; /// input: [a] [a] [b ][b] /// output: [a ] [b ] /// ``` -pub fn flood(events: Vec, pulsetime: chrono::Duration) -> Vec { - let mut new_events = Vec::new(); - let mut events_sorted = sort_by_timestamp(events); - let mut e1_iter = events_sorted.drain(..).peekable(); - - let mut gap_prev: Option = None; - let mut retry_e: Option = None; +pub fn flood(events: Vec, pulsetime: Duration) -> Vec { + let mut events = events; + // Sort by (timestamp, duration) so shorter same-timestamp events come first. + events.sort_by_key(|e| (e.timestamp, e.duration)); // If negative gaps are smaller than this, prune them to become zero - let negative_gap_trim_thres = chrono::Duration::milliseconds(100); + let negative_gap_trim_thres = Duration::milliseconds(100); + let zero = Duration::zero(); let mut warned_negative_gap_safe = false; let mut warned_negative_gap_unsafe = false; - while let Some(mut e1) = match retry_e { - Some(e) => { - retry_e = None; - Some(e) - } - None => e1_iter.next(), - } { - if let Some(gap) = gap_prev { - e1.timestamp -= gap / 2; - e1.duration += gap / 2; - gap_prev = None; - } - let e2 = match e1_iter.peek() { - Some(e) => e, - None => { - new_events.push(e1); - break; - } - }; - - // Ensure that ordering is preserved, at least in tests - debug_assert!(e1.timestamp <= e2.timestamp); + // Pairwise pass: each event is compared with the next one, which may be modified before it is + // compared with its own successor. + for i in 1..events.len() { + let (head, tail) = events.split_at_mut(i); + let e1 = &mut head[i - 1]; + let e2 = &mut tail[0]; - let gap = e2.timestamp - e1.calculate_endtime(); + let e1_end = e1.calculate_endtime(); + let e2_end = e2.calculate_endtime(); + let gap = e2.timestamp - e1_end; - // We split the program flow into 2 parts: positive and negative gaps - // First we check negative gaps (if events overlap) - if gap < chrono::Duration::seconds(0) { - // Python implementation: - // - // if gap < timedelta(0) and e1.data == e2.data: - // start = min(e1.timestamp, e2.timestamp) - // end = max(e1.timestamp + e1.duration, e2.timestamp + e2.duration) - // e1.timestamp, e1.duration = start, (end - start) - // e2.timestamp, e2.duration = end, timedelta(0) - // if not warned_about_negative_gap_safe: - // logger.warning( - // "Gap was of negative duration but could be safely merged ({}s). This message will only show once per batch.".format( - // gap.total_seconds() - // ) - // ) - // warned_about_negative_gap_safe = True - - // If data is same, we can safely merge them. - if e1.data == e2.data { - // Gap was negative and could be safely merged - if !warned_negative_gap_safe { - warn!("Gap was of negative duration ({}s), but could be safely merged. This error will only show once per batch.", gap); - warned_negative_gap_safe = true; - } - let start = std::cmp::min(e1.timestamp, e2.timestamp); - let end = std::cmp::max(e1.calculate_endtime(), e2.calculate_endtime()); // e2 isn't guaranteed to end last - e1.timestamp = start; - e1.duration = end - start; - // Drop next event since it is merged/flooded into e1 - e1_iter.next(); - // Retry this event again to give it a change to merge e1 - // with 'e3' - retry_e = Some(e1); - continue; - // If data differs and gap is negative beyond the trim thres, throw a warning - } else if gap < -negative_gap_trim_thres && !warned_negative_gap_unsafe { - // Events with negative gap but differing data cannot be merged safely + if gap < zero && e1.data == e2.data { + // Events with negative gap but same data can safely be merged + if !warned_negative_gap_safe { + warn!("Gap was of negative duration ({}s), but could be safely merged. This warning will only show once per batch.", gap); + warned_negative_gap_safe = true; + } + let start = min(e1.timestamp, e2.timestamp); + let end = max(e1_end, e2_end); + e1.timestamp = start; + e1.duration = end - start; + e2.timestamp = end; + e2.duration = zero; + } else if gap < -negative_gap_trim_thres { + // Events with negative gap but differing data cannot be merged here, they are + // resolved by the normalization pass below. + if !warned_negative_gap_unsafe { warn!("Gap was of negative duration and could NOT be safely merged ({}s). This warning will only show once per batch.", gap); warned_negative_gap_unsafe = true; } - - // By now, we've ensured gap is non-negative, now we want to know if the gap is smaller - // than pulsetime, in which case we should fill it. - } else if gap < pulsetime { - // Python implementation: - // - // elif gap < -negative_gap_trim_thres and not warned_about_negative_gap_unsafe: - // # Events with negative gap but differing data cannot be merged safely - // logger.warning( - // "Gap was of negative duration and could NOT be safely merged ({}s). This warning will only show once per batch.".format( - // gap.total_seconds() - // ) - // ) - // warned_about_negative_gap_unsafe = True - - // Ensure that gap is actually non-negative here, at least in tests - debug_assert!(gap >= chrono::Duration::seconds(0)); - - // If data is the same, we should merge them. + } else if gap > -negative_gap_trim_thres && gap <= pulsetime { if e1.data == e2.data { - // Choose the longest event and set the endtime to it - let start = std::cmp::min(e1.timestamp, e2.timestamp); // isn't e1 guaranteed to start? - let end = std::cmp::max(e1.calculate_endtime(), e2.calculate_endtime()); // e2 isn't guaranteed to end last - e1.timestamp = start; - e1.duration = end - start; - // Drop next event since it is merged/flooded into e1 - e1_iter.next(); - // Retry this event again to give it a change to merge e1 - // with 'e3' - retry_e = Some(e1); - // Since we are retrying on this event we don't want to push it - // to the new_events vec - continue; + // Keep the longer event and extend it over the gap and the other event. + if e1.duration >= e2.duration { + e1.duration = e2_end - e1.timestamp; + e2.timestamp = e2_end; + e2.duration = zero; + } else { + e2.timestamp = e1.timestamp; + e2.duration = e2_end - e2.timestamp; + e1.duration = zero; + } } else { - // Extend e1 to middle of the gap. - e1.duration += gap / 2; - - // Make sure next event (e2) is gets extended before it's processed - gap_prev = Some(gap); + // The gap is an interval of uncertainty: without evidence that either neighbour + // owns more of it, split it at the midpoint. + let midpoint = e1_end + gap / 2; + e1.duration = midpoint - e1.timestamp; + e2.timestamp = midpoint; + e2.duration = e2_end - midpoint; } } // else: nothing to do, events not near each other + } - new_events.push(e1); + // The pairwise pass can modify an event after it was compared with its predecessor, and does + // not resolve overlapping events with different data. Normalize the result so that it never + // contains overlapping events. For differing data, the later event wins. + let mut normalized: Vec = Vec::with_capacity(events.len()); + for event in events.into_iter().filter(|e| e.duration > zero) { + let mut merged = false; + while let Some(previous) = normalized.last_mut() { + let previous_end = previous.calculate_endtime(); + if previous_end <= event.timestamp { + break; + } + if previous.data == event.data { + let end = max(previous_end, event.calculate_endtime()); + previous.duration = end - previous.timestamp; + merged = true; + break; + } + previous.duration = event.timestamp - previous.timestamp; + if previous.duration > zero { + break; + } + normalized.pop(); + } + if !merged { + normalized.push(event); + } } - new_events + normalized } #[cfg(test)] @@ -280,8 +252,8 @@ mod tests { #[test] fn test_flood_containing_diff() { - // Tests flooding an event with different data contained within another event - // Events should pass unmodified. + // An event with different data contained within another event takes precedence from its + // start, so that the result has no overlap (ActivityWatch/activitywatch#1369). let e1 = Event { id: None, timestamp: DateTime::from_str("2000-01-01T00:00:00Z").unwrap(), @@ -296,7 +268,9 @@ mod tests { }; let res = flood(vec![e1.clone(), e2.clone()], Duration::seconds(5)); assert_eq!(2, res.len()); - assert_eq!(&res[0], &e1); + let mut e1_expected = e1; + e1_expected.duration = Duration::seconds(1); + assert_eq!(&res[0], &e1_expected); assert_eq!(&res[1], &e2); } @@ -386,4 +360,127 @@ mod tests { assert_eq!(&res[1], &e4); assert_eq!(&res[2], &e5); } + + /// Event starting `start` seconds after 10:00:00Z, lasting `duration` seconds, with data + /// `{"app": app}`. + fn ev(start: f64, duration: f64, app: &str) -> Event { + let base: DateTime = DateTime::from_str("2026-01-01T10:00:00Z").unwrap(); + Event { + id: None, + timestamp: base + Duration::milliseconds((start * 1000.0) as i64), + duration: Duration::milliseconds((duration * 1000.0) as i64), + data: json_map! {"app": json!(app)}, + } + } + + fn assert_no_overlap(events: &[Event]) { + for pair in events.windows(2) { + assert!( + pair[0].calculate_endtime() <= pair[1].timestamp, + "{:?} overlaps {:?}", + pair[0], + pair[1] + ); + } + } + + // The following tests are the repros from ActivityWatch/aw-server-rust#746, where Rust and + // aw-core disagreed. Expected results are aw-core's. + + #[test] + fn test_flood_overlap_diff_data() { + // activitywatch#1369: overlapping events with different data must not count time twice + let res = flood( + vec![ev(0., 10., "x"), ev(9., 10., "y")], + Duration::seconds(5), + ); + assert_eq!(res, vec![ev(0., 9., "x"), ev(9., 10., "y")]); + } + + #[test] + fn test_flood_drops_zero_duration() { + let res = flood( + vec![ev(0., 0., "x"), ev(60., 10., "y")], + Duration::seconds(5), + ); + assert_eq!(res, vec![ev(60., 10., "y")]); + } + + #[test] + fn test_flood_gap_equal_to_pulsetime() { + // A gap of exactly pulsetime is filled, meeting in the middle + let res = flood( + vec![ev(0., 10., "x"), ev(15., 10., "y")], + Duration::seconds(5), + ); + assert_eq!(res, vec![ev(0., 12.5, "x"), ev(12.5, 12.5, "y")]); + + // A gap just above pulsetime is left alone + let events = vec![ev(0., 10., "x"), ev(15.001, 10., "y")]; + assert_eq!(flood(events.clone(), Duration::seconds(5)), events); + } + + #[test] + fn test_flood_adjacent_same_data() { + let res = flood( + vec![ev(0., 10., "x"), ev(10., 10., "x")], + Duration::seconds(5), + ); + assert_eq!(res, vec![ev(0., 20., "x")]); + } + + #[test] + fn test_flood_pulsetime() { + let events = vec![ev(0., 10., "x"), ev(20., 10., "y")]; + assert_eq!(flood(events.clone(), Duration::seconds(5)), events); + assert_eq!( + flood(events, Duration::seconds(10)), + vec![ev(0., 15., "x"), ev(15., 15., "y")] + ); + } + + #[test] + fn test_flood_small_negative_gap_diff_data() { + // Overlaps smaller than 100ms between differing events are split in the middle + let res = flood( + vec![ev(0., 10., "x"), ev(9.95, 10., "y")], + Duration::seconds(5), + ); + assert_eq!(res, vec![ev(0., 9.975, "x"), ev(9.975, 9.975, "y")]); + } + + #[test] + fn test_flood_later_event_wins() { + // A later event overlapping several earlier ones truncates or removes them all + let res = flood( + vec![ev(0., 10., "x"), ev(2., 10., "y"), ev(2., 20., "z")], + Duration::seconds(5), + ); + assert_no_overlap(&res); + assert_eq!(res, vec![ev(0., 2., "x"), ev(2., 20., "z")]); + } + + #[test] + fn test_flood_no_overlap_randomized() { + // Deterministic pseudo-random inputs; the output must be sorted and free of overlap. + let mut seed: u64 = 0x2545_f491_4f6c_dd1d; + let mut next = |m: u64| { + seed ^= seed << 13; + seed ^= seed >> 7; + seed ^= seed << 17; + seed % m + }; + for _ in 0..500 { + let n = next(12) as usize; + let events: Vec = (0..n) + .map(|_| { + let app = ["a", "b", "c"][next(3) as usize]; + ev(next(60) as f64, next(15) as f64, app) + }) + .collect(); + let res = flood(events, Duration::seconds(5)); + assert_no_overlap(&res); + assert!(res.iter().all(|e| e.duration > Duration::zero())); + } + } } From e780b24a6742b1f3259379c463f26cb06805e530 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Erik=20Bj=C3=A4reholt?= Date: Sat, 26 Sep 2026 19:07:21 +0200 Subject: [PATCH 2/2] fix(transform): merge chains of same-data events in one flood() pass The pairwise pass leaves a merged-away event as a zero-duration placeholder at the end of the merged event, so a chain of same-data events (e.g. three adjacent 10 s events) came out of normalization as adjacent pieces. Merge touching same-data events in normalization too, like ActivityWatch/aw-core#165. Also check the pulsetime argument's effect in a query test. --- aw-query/tests/query.rs | 53 ++++++++++++++++++++++++++++++++++-- aw-transform/src/flood.rs | 56 ++++++++++++++++++++++++++++++++++++--- 2 files changed, 104 insertions(+), 5 deletions(-) diff --git a/aw-query/tests/query.rs b/aw-query/tests/query.rs index 2b7f607a..e151b429 100644 --- a/aw-query/tests/query.rs +++ b/aw-query/tests/query.rs @@ -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 = 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 { + 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(); @@ -353,8 +404,6 @@ mod query_tests { r#" events = query_bucket(find_bucket("{}", "testhost")); events = flood(events); - events = flood(events, 10); - events = flood(events, 0.5); events = sort_by_duration(events); events = limit_events(events, 10000); events = sort_by_timestamp(events); diff --git a/aw-transform/src/flood.rs b/aw-transform/src/flood.rs index 069b079f..bf1fec98 100644 --- a/aw-transform/src/flood.rs +++ b/aw-transform/src/flood.rs @@ -105,16 +105,21 @@ pub fn flood(events: Vec, pulsetime: Duration) -> Vec { // The pairwise pass can modify an event after it was compared with its predecessor, and does // not resolve overlapping events with different data. Normalize the result so that it never - // contains overlapping events. For differing data, the later event wins. + // contains overlapping events, and adjacent same-data events are merged. For differing data, + // the later event wins. let mut normalized: Vec = Vec::with_capacity(events.len()); for event in events.into_iter().filter(|e| e.duration > zero) { let mut merged = false; while let Some(previous) = normalized.last_mut() { let previous_end = previous.calculate_endtime(); - if previous_end <= event.timestamp { + let same_data = previous.data == event.data; + // Merge same-data events that touch, not only those that overlap: the pairwise pass + // leaves a merged-away event as a zero-duration placeholder at the end of the merged + // event, so a chain of same-data events can come out as adjacent pieces. + if previous_end < event.timestamp || (previous_end == event.timestamp && !same_data) { break; } - if previous.data == event.data { + if same_data { let end = max(previous_end, event.calculate_endtime()); previous.duration = end - previous.timestamp; merged = true; @@ -429,6 +434,51 @@ mod tests { assert_eq!(res, vec![ev(0., 20., "x")]); } + #[test] + fn test_flood_merge_chains() { + // Chains of same-data events merge into one event, whether adjacent or within pulsetime + let pt = Duration::seconds(5); + let adjacent = vec![ev(0., 10., "x"), ev(10., 10., "x"), ev(20., 10., "x")]; + assert_eq!(flood(adjacent, pt), vec![ev(0., 30., "x")]); + + let gaps = vec![ev(0., 10., "x"), ev(12., 1., "x"), ev(14., 1., "x")]; + assert_eq!(flood(gaps, pt), vec![ev(0., 15., "x")]); + + let then_other = vec![ev(0., 10., "x"), ev(12., 1., "x"), ev(14., 1., "y")]; + assert_eq!( + flood(then_other, pt), + vec![ev(0., 13.5, "x"), ev(13.5, 1.5, "y")] + ); + } + + #[test] + fn test_flood_merges_adjacent_randomized() { + let mut seed: u64 = 0x6a09_e667_f3bc_c908; + let mut next = |m: u64| { + seed ^= seed << 13; + seed ^= seed >> 7; + seed ^= seed << 17; + seed % m + }; + for _ in 0..500 { + let n = next(12) as usize; + let events: Vec = (0..n) + .map(|_| { + let app = ["a", "b"][next(2) as usize]; + ev(next(60) as f64, next(15) as f64, app) + }) + .collect(); + let once = flood(events.clone(), Duration::seconds(5)); + // No two adjacent same-data events are left unmerged + for pair in once.windows(2) { + assert!( + pair[0].data != pair[1].data || pair[0].calculate_endtime() < pair[1].timestamp, + "{events:?} -> {once:?}" + ); + } + } + } + #[test] fn test_flood_pulsetime() { let events = vec![ev(0., 10., "x"), ev(20., 10., "y")];