← run

py-13-windowed-aggregator

1.000
13/13 testsĀ· architecture
Challenge Ā· difficulty 5/5
# Windowed streaming aggregator

Implement a file **`solution.py`** containing a class `WindowedAggregator` that ingests
timestamped, grouped events and computes aggregate statistics over a sliding time window.

```python
class WindowedAggregator:
    def __init__(self, window: float):
        """window = length of the time window (same units as timestamps)."""

    def add(self, ts: float, group: str, value: float) -> None:
        """Record an event. Events may arrive OUT OF ORDER (ts not monotonic)."""

    def stats(self, group: str, at: float) -> dict:
        """Aggregate the events of `group` whose timestamp is in the half-open
        window (at - window, at] — i.e. at-window < ts <= at."""
```

## Semantics

- The window is **half-open**: an event whose timestamp is exactly `at - window` is **excluded**,
  and an event whose timestamp is exactly `at` is **included**. Formally, an event with timestamp
  `ts` is in the window iff `at - window < ts <= at`.
- **Groups are independent.** `stats` for a group must never observe events recorded under any other
  group.
- **Out-of-order ingestion is allowed.** `add` may be called with timestamps in any order. `stats`
  considers every event added so far regardless of insertion order; calling `add` with the same
  events in a different order yields identical `stats` results.
- **Duplicates count.** Adding the same `(ts, group)` pair (or even the same `(ts, group, value)`)
  more than once records multiple independent events, all of which contribute to the aggregates.
- `stats` returns a dict with **exactly** these keys:
  - `"count"`: `int` — number of events in the window.
  - `"sum"`: `float` — sum of their values (`0.0` if there are none).
  - `"min"`: `float` or `None` — minimum value, or `None` if there are no events.
  - `"max"`: `float` or `None` — maximum value, or `None` if there are no events.
  - `"mean"`: `float` or `None` — `sum / count`, or `None` if there are no events.
- An **unknown group**, or a window that **contains no events**, returns
  `{"count": 0, "sum": 0.0, "min": None, "max": None, "mean": None}`.

## Example

```python
agg = WindowedAggregator(window=10.0)
agg.add(100.0, "cpu", 5.0)
agg.add(95.0,  "cpu", 1.0)
agg.add(105.0, "cpu", 3.0)
agg.add(100.0, "mem", 50.0)   # different group, ignored by "cpu" stats

# Window for at=105, window=10 is (95, 105]:
#   ts=95  is excluded (boundary at-window)
#   ts=100 is included
#   ts=105 is included
s = agg.stats("cpu", at=105.0)
assert s == {"count": 2, "sum": 8.0, "min": 3.0, "max": 5.0, "mean": 4.0}

# Advance `at` so the window (100, 110] only contains ts=105:
s2 = agg.stats("cpu", at=110.0)
assert s2 == {"count": 1, "sum": 3.0, "min": 3.0, "max": 3.0, "mean": 3.0}

# Unknown group:
assert agg.stats("disk", at=105.0) == {
    "count": 0, "sum": 0.0, "min": None, "max": None, "mean": None,
}
```
tests/test_windowed_aggregator.py
import random

from solution import WindowedAggregator

ZERO = {"count": 0, "sum": 0.0, "min": None, "max": None, "mean": None}


def test_empty_and_unknown_group_returns_zero_dict():
    agg = WindowedAggregator(window=10.0)
    assert agg.stats("nope", at=0.0) == ZERO
    agg.add(5.0, "g", 1.0)
    # group exists, but window (90, 100] contains nothing
    assert agg.stats("g", at=100.0) == ZERO
    # still-unknown group
    assert agg.stats("other", at=5.0) == ZERO


def test_basic_single_event_stats():
    agg = WindowedAggregator(window=10.0)
    agg.add(100.0, "g", 4.0)
    s = agg.stats("g", at=100.0)
    assert s == {"count": 1, "sum": 4.0, "min": 4.0, "max": 4.0, "mean": 4.0}


def test_inclusion_at_upper_bound():
    # event exactly at `at` is INCLUDED
    agg = WindowedAggregator(window=10.0)
    agg.add(105.0, "g", 7.0)
    s = agg.stats("g", at=105.0)
    assert s["count"] == 1
    assert s["sum"] == 7.0


def test_exclusion_at_lower_bound():
    # event exactly at at-window is EXCLUDED (half-open lower bound)
    agg = WindowedAggregator(window=10.0)
    agg.add(95.0, "g", 7.0)
    s = agg.stats("g", at=105.0)  # window (95, 105]
    assert s == ZERO


def test_boundary_pair_together():
    agg = WindowedAggregator(window=10.0)
    agg.add(95.0, "g", 1.0)   # excluded
    agg.add(105.0, "g", 3.0)  # included
    s = agg.stats("g", at=105.0)
    assert s == {"count": 1, "sum": 3.0, "min": 3.0, "max": 3.0, "mean": 3.0}


def test_min_max_mean_sum_count_correctness():
    agg = WindowedAggregator(window=100.0)
    values = [3.0, -2.0, 10.5, 4.0, 0.0]
    for i, v in enumerate(values):
        agg.add(float(i), "g", v)
    # timestamps are 0..4; window (4-100, 4] = (-96, 4] covers all of them
    s = agg.stats("g", at=4.0)
    assert s["count"] == 5
    assert s["sum"] == sum(values)
    assert s["min"] == min(values)
    assert s["max"] == max(values)
    assert s["mean"] == sum(values) / len(values)


def test_group_isolation():
    agg = WindowedAggregator(window=10.0)
    agg.add(100.0, "a", 1.0)
    agg.add(100.0, "a", 2.0)
    agg.add(100.0, "b", 100.0)
    sa = agg.stats("a", at=100.0)
    sb = agg.stats("b", at=100.0)
    assert sa == {"count": 2, "sum": 3.0, "min": 1.0, "max": 2.0, "mean": 1.5}
    assert sb == {"count": 1, "sum": 100.0, "min": 100.0, "max": 100.0, "mean": 100.0}


def test_out_of_order_matches_in_order():
    events = [(100.0, 5.0), (95.0, 1.0), (105.0, 3.0), (102.0, 9.0), (98.0, 2.0)]

    ordered = WindowedAggregator(window=10.0)
    for ts, v in sorted(events):
        ordered.add(ts, "g", v)

    shuffled = WindowedAggregator(window=10.0)
    rnd = list(events)
    random.Random(1234).shuffle(rnd)
    for ts, v in rnd:
        shuffled.add(ts, "g", v)

    for at in (95.0, 100.0, 105.0, 110.0, 120.0):
        assert ordered.stats("g", at=at) == shuffled.stats("g", at=at)


def test_duplicate_ts_group_events_all_count():
    agg = WindowedAggregator(window=10.0)
    agg.add(100.0, "g", 2.0)
    agg.add(100.0, "g", 2.0)  # identical (ts, group, value)
    agg.add(100.0, "g", 5.0)  # same (ts, group), different value
    s = agg.stats("g", at=100.0)
    assert s["count"] == 3
    assert s["sum"] == 9.0
    assert s["min"] == 2.0
    assert s["max"] == 5.0
    assert s["mean"] == 3.0


def test_moving_window_includes_and_excludes():
    agg = WindowedAggregator(window=10.0)
    for ts in (90.0, 95.0, 100.0, 105.0, 110.0):
        agg.add(ts, "g", ts)  # value == ts for easy checking

    # at=100 -> window (90, 100]: ts in {95, 100}
    s = agg.stats("g", at=100.0)
    assert s["count"] == 2
    assert s["sum"] == 95.0 + 100.0
    assert s["min"] == 95.0 and s["max"] == 100.0

    # at=105 -> window (95, 105]: ts in {100, 105}
    s = agg.stats("g", at=105.0)
    assert s["count"] == 2
    assert s["min"] == 100.0 and s["max"] == 105.0

    # at=110 -> window (100, 110]: ts in {105, 110}
    s = agg.stats("g", at=110.0)
    assert s["count"] == 2
    assert s["min"] == 105.0 and s["max"] == 110.0

    # at far in the future -> window past all events
    assert agg.stats("g", at=1000.0) == ZERO


def test_stats_is_non_destructive_and_repeatable():
    agg = WindowedAggregator(window=10.0)
    agg.add(100.0, "g", 1.0)
    agg.add(105.0, "g", 3.0)
    first = agg.stats("g", at=105.0)
    second = agg.stats("g", at=105.0)
    assert first == second  # querying must not consume events


def test_exact_keys_present():
    agg = WindowedAggregator(window=10.0)
    agg.add(100.0, "g", 1.0)
    s = agg.stats("g", at=100.0)
    assert set(s.keys()) == {"count", "sum", "min", "max", "mean"}
    assert isinstance(s["count"], int)
    # empty case keys too
    e = agg.stats("g", at=1000.0)
    assert set(e.keys()) == {"count", "sum", "min", "max", "mean"}


def test_float_window_and_timestamps():
    agg = WindowedAggregator(window=1.5)
    agg.add(10.0, "g", 1.0)
    agg.add(10.4, "g", 2.0)
    agg.add(10.5, "g", 3.0)
    # at=11.0 -> window (9.5, 11.0]: all three included
    s = agg.stats("g", at=11.0)
    assert s["count"] == 3
    # at=11.5 -> window (10.0, 11.5]: ts=10.0 excluded (lower bound)
    s2 = agg.stats("g", at=11.5)
    assert s2["count"] == 2
    assert s2["min"] == 2.0 and s2["max"] == 3.0
System prompt
You are an expert programmer. Solve the task exactly as specified. Output your solution as fenced code blocks using the required file name(s) and the exact function/type signatures requested. Prefer correctness; do not include prose outside code unless asked.
Proposed solution


To solve this problem, we need to implement a class `WindowedAggregator` that can efficiently compute aggregate statistics over a sliding time window for timestamped, grouped events. The class must handle out-of-order events and provide accurate statistics for each group independently.

### Approach
The solution involves maintaining a dictionary where each key is a group name, and the value is a tuple containing two lists: one for storing events (as tuples of timestamp and value) and another for storing timestamps. This allows us to efficiently manage and query the events using binary search.

1. **Initialization**: The constructor initializes the window size and an empty dictionary to store events for each group.
2. **Adding Events**: The `add` method inserts each event into the appropriate position in a sorted list of timestamps and events. This ensures that events are always stored in chronological order.
3. **Computing Statistics**: The `stats` method calculates the required statistics for a given group and time window. It uses binary search to quickly find the range of events that fall within the window, then computes the count, sum, min, max, and mean of the values in this range.

### Solution Code
```python
import bisect

class WindowedAggregator:
    def __init__(self, window: float):
        """window = length of the time window (same units as timestamps)."""
        self.window = window
        self.groups = {}  # Maps group names to (events, ts_list)

    def add(self, ts: float, group: str, value: float) -> None:
        """Record an event. Events may arrive OUT OF ORDER (ts not monotonic)."""
        if group not in self.groups:
            self.groups[group] = ([], [])
        events, ts_list = self.groups[group]
        index = bisect.bisect_left(ts_list, ts)
        ts_list.insert(index, ts)
        events.insert(index, (ts, value))

    def stats(self, group: str, at: float) -> dict:
        """Aggregate the events of `group` whose timestamp is in the half-open window (at - window, at]."""
        if group not in self.groups:
            return {
                "count": 0,
                "sum": 0.0,
                "min": None,
                "max": None,
                "mean": None
            }
        events, ts_list = self.groups[group]
        window_start = at - self.window
        left = bisect.bisect_right(ts_list, window_start)
        right = bisect.bisect_right(ts_list, at)
        window_events = events[left:right]
        values = [value for ts, value in window_events]
        count = len(values)
        if count == 0:
            return {
                "count": 0,
                "sum": 0.0,
                "min": None,
                "max": None,
                "mean": None
            }
        sum_val = sum(values)
        min_val = min(values)
        max_val = max(values)
        mean_val = sum_val / count
        return {
            "count": count,
            "sum": sum_val,
            "min": min_val,
            "max": max_val,
            "mean": mean_val
        }
```

### Explanation
- **Initialization**: The `__init__` method sets up the window size and initializes an empty dictionary to store events for each group.
- **Adding Events**: The `add` method uses binary search to maintain a sorted list of timestamps and corresponding events. This ensures efficient insertion and retrieval.
- **Computing Statistics**: The `stats` method uses binary search to quickly locate the range of timestamps within the specified window. It then computes the necessary statistics (count, sum, min, max, mean) for the events in this range, ensuring efficient and accurate results.

This approach efficiently handles out-of-order events and provides accurate statistics for each group independently, ensuring optimal performance even for large datasets.

Test output (stdout)

.............                                                            [100%]
13 passed in 0.02s
Reasoning (chain-of-thought)

Okay, I need to implement the WindowedAggregator class in Python. Let's see what the requirements are.

The class has an __init__ method that takes a window length. Then there's an add method to record events, and a stats method to compute aggregates over a sliding window for a specific group.

Hmm, the window is half-open. So for a given 'at' time, the window is (at - window, at], meaning events with ts exactly at at - window are excluded, but those at at are included.

Groups are independent, so each group's events are processed separately. Also, add can be called with any order of timestamps, and stats must consider all events added so far, regardless of insertion order.

Duplicates are allowed and count as separate events. So if the same (ts, group) is added multiple times, each contributes to the aggregate.

The stats method returns a dictionary with count, sum, min, max, mean. If no events, count is 0, sum 0.0, others None.

So, how to structure this?

I think for each group, we need to keep track of all events, and when stats is called, we filter those events that fall within the window (at - window, at].

But since events can come in any order, we can't assume they're added in order. So for each group, we need to store all events, and when querying, we process them each time.

Wait, but that might be inefficient if there are a lot of events. But given the problem statement, perhaps it's acceptable for the solution.

So, the plan is:

- The class will have a dictionary to map group names to their list of events. Each event is a tuple of (ts, value).

When add is called, we append the (ts, value) to the group's list.

When stats is called for a group and a specific 'at' time, we:

1. Check if the group exists. If not, return the default dict.
2. For the group's events, filter those where ts > (at - window) and ts <= at.
3. Collect all such values and compute the required statistics.

Wait, but for each stats call, we have to process all events of the group and filter them. That could be slow for large numbers of events, but perhaps it's manageable.

Alternatively, we could maintain for each group a sorted list of events, so that when querying, we can perform a binary search to find the lower and upper bounds. But since events can be added in any order, maintaining a sorted list would require inserting in the correct position each time, which is O(n) for each add. That might not be efficient either.

Hmm, but for the problem's constraints, maybe the straightforward approach is acceptable.

So, for each group, we'll store all events as a list. When add is called, we just append to the list. When stats is called, we iterate through all events of the group and check if they fall within the window.

Wait, but for each stats call, we have to process all events of the group. If the group has a million events, this could be slow. But perhaps for the problem's purposes, it's acceptable.

Alternatively, perhaps we can pre-sort the events when adding, but that would require maintaining a sorted list, which could be done with the bisect module. But each add would take O(log n) time for insertion. For small n, this is fine, but for large n, it's better.

But given that the problem allows any order of adds, and the stats function must consider all events, perhaps the initial approach is acceptable.

So, let's proceed with the initial approach.

Now, the steps:

In __init__, we'll initialize a dictionary, say self.groups, where each key is the group name, and the value is a list of (ts, value) tuples.

The add method will append to the group's list. If the group doesn't exist, we create it.

The stats method will:

- Check if the group exists. If not, return the default dict.
- For each event in the group's list, check if ts is > (at - window) and <= at.
- Collect all such values.
- Compute count, sum, min, max, mean.

But wait, the window is (at - window, at], so events with ts exactly at at - window are excluded. So the condition is ts > (at - window) and ts <= at.

So, for each event in the group, if ts is in that range, include it.

Once we have the list of values, compute the stats.

Now, what about the case where no events are in the window? Then count is 0, sum is 0.0, and min, max, mean are None.

So, in code:

In the stats method:

group_events = self.groups.get(group, [])
window_start = at - window
filtered = [ (ts, val) for ts, val in group_events if window_start < ts <= at ]

Then, collect all the vals:

values = [val for ts, val in filtered]

count = len(values)
sum_val = sum(values) if count else 0.0

if count == 0:
    return {
        "count": 0,
        "sum": 0.0,
        "min": None,
        "max": None,
        "mean": None
    }
else:
    min_val = min(values)
    max_val = max(values)
    mean_val = sum_val / count
    return {
        "count": count,
        "sum": sum_val,
        "min": min_val,
        "max": max_val,
        "mean": mean_val
    }

Wait, but in the example given:

In the first example, after adding 100, 95, 105 for group 'cpu', when at=105, window=10, the window is (95, 105]. So ts=95 is excluded, 100 and 105 are included. So the count is 2, sum 8, min 3, max 5, mean 4.

Wait, but in the events, 100 is 5, 105 is 3. So the sum is 5+3=8, count 2. Min is 3, max is 5.

Yes, that's correct.

So, the code seems to handle that.

But wait, in the code I wrote, the filtered list includes all events in the group where ts is in the window. So for each event, regardless of when it was added, it's considered.

Yes.

Now, what about the case where the same (ts, group) is added multiple times? Like adding (100, 'cpu', 5) twice. Then, in the filtered list, both events are included, each contributing 5 to the sum, etc.

Yes, which is correct as per the problem statement.

Now, what about the stats method's parameters: group and at.

So, the code for stats is:

def stats(self, group: str, at: float) -> dict:

So, the code for the stats function is as I outlined.

Now, let's think about the data structures.

Each group's events are stored as a list of tuples. For each add, we append to the list.

But when we pro
…[8938 chars elided]…
, the events are added as:

First add 100.0: ts_list becomes [100.0], events [(100.0,5.0)].

Then add 95.0: bisect_left finds index 0, so ts_list becomes [95.0, 100.0], events [(95.0,1.0), (100.0,5.0)].

Then add 105.0: bisect_left finds index 2, so ts_list becomes [95.0, 100.0, 105.0], events [(95.0,1.0), (100.0,5.0), (105.0,3.0)].

When stats is called for 'cpu' at 105.0:

window_start = 105 -10 =95.0.

left = bisect.bisect_right(ts_list, 95.0) → since ts_list is [95, 100, 105], bisect_right(95) is 1.

right = bisect.bisect_right(ts_list, 105) → returns 3.

So, window_events is events[1:3], which are (100,5) and (105,3).

So, values are [5,3], count 2, sum 8, min 3, max 5, mean 4.

Which matches the example.

Another test case: at=110.0.

window_start = 110-10=100.

left = bisect_right(ts_list, 100) → ts_list is [95,100,105]. bisect_right(100) is 2.

right = bisect_right(ts_list, 110) → 3.

window_events is events[2:3], which is (105,3). So count 1, sum 3, min 3, max 3, mean 3.

Which is correct.

Another test case: group 'disk' is unknown, returns the default.

So, the code seems to handle that.

What about when the same (ts, group) is added multiple times?

For example, adding (100, 'cpu', 5) twice.

After adding, the ts_list for 'cpu' would have two entries at 100.0, and events would have two tuples (100,5) each.

When stats is called, both would be included in the window if applicable.

Yes, which is correct.

What about when events are added out of order?

Like adding 105, then 95, then 100.

The ts_list would be [95,100,105], regardless of insertion order.

Yes, because each add uses bisect to insert in the correct position.

So, the code should handle that.

Now, what about the performance? For each add, it's O(log n) time for the bisect, and O(n) for the insertion (since inserting into a list is O(n) time). For each stats call, it's O(log n) for the bisects and O(k) for processing the k events in the window.

But for the problem's constraints, this should be acceptable.

So, the code should be correct.