A streaming job counts events in tumbling windows of window_s seconds, by event time — when the event happened, not when it reached the job. Window [start, start + window_s) holds the events whose event_time falls in it, with start = event_time // window_s * window_s.
Events arrive in the order of the events list, each a dict with id, event_time and arrival_time (integer seconds). Process them one at a time:
1. The watermark is the largest event_time seen so far minus lateness_s. Before the first event there is no watermark. 2. If an event's window has already closed — its end is at or before the current watermark — the event is late: drop it and count it. 3. Otherwise add it to its window, then update the watermark with its event_time. 4. Every window whose end is at or before the watermark now closes: emit it as (start, count). A window closes once, in order of start.
Windows still open when the list runs out are not emitted — the job is still running. Windows that never received an event are never emitted.
Write window_counts(events, window_s, lateness_s) returning {"windows": [(start, count), ...], "late": <number of dropped events>}, with windows in the order they closed.
Python 3.13 in your browser — the standard library plus pandas and numpy; no pip installs.