From 9a2f0f6a673d240661d373ba72b7d1028215f55c Mon Sep 17 00:00:00 2001 From: Yuvaraj_G Date: Sat, 15 Aug 2026 18:13:14 +0530 Subject: [PATCH] Add changes for activity frames pull request --- SPEC.md | 14 +++ errors.md | 0 src/activity_frames/__init__.py | 2 + src/activity_frames/events.py | 25 ++++++ src/activity_frames/frames.py | 6 ++ src/activity_frames/sessionize.py | 140 +++++++++++++++++++++++++++++- tests/test_sessionize.py | 95 +++++++++++++++++++- 7 files changed, 279 insertions(+), 3 deletions(-) create mode 100644 errors.md create mode 100644 src/activity_frames/events.py diff --git a/SPEC.md b/SPEC.md index f13c79e..e651bfa 100644 --- a/SPEC.md +++ b/SPEC.md @@ -111,6 +111,20 @@ Producers may offer a `min_minutes` filter for consumer convenience. Filtered-ou Each monitor records its own frame stream (`device_name`). Segmentation runs per device, so simultaneous monitors do not fragment each other's frames; the document lists all devices' frames sorted by start time, which means frames MAY overlap in time. Input events are assigned to exactly one containing segment (ties across devices resolved by the nearest captured frame), so input volume is never double-counted. A consequence to disclose: the same app visible on two monitors at once earns active time on both. +### 5.6 Event-aware segmentation (opt-in) + +Consumers may optionally provide authoritative event metadata from external systems such as `git.merge`, `pr.approved`, or `deploy.success`. These are not inferred from screen activity and are supplied explicitly by callers. The event type is a simple string; the package does not classify, rank, or guess event meaning. + +An `Event` object is a minimal record: + +```python +Event(event_type="git.merge", timestamp=1700.0, source="ci", priority=10) +``` + +When a caller passes `events` alongside the normal segmentation call, the existing context/session segmentation remains the base pipeline. After those segments are formed, the optional event boundary step may split an existing segment only at an observed boundary between consecutive captured frames. It never invents a frame at the event timestamp. When an event falls between two captured frames, the implementation chooses the closer inter-frame midpoint as the deterministic split point; ties are resolved in timestamp order and then by event index. This keeps the result reproducible for identical inputs. + +The event-aware layer is intentionally opt-in. If `events=()` and `forced_event_types=()` and `force_priority is None`, the output is identical to the old implementation. `forced_event_types` allows splitting only on specific event types; `force_priority` allows splitting on any externally supplied event whose `priority` is at least the threshold. Both filters are additive: an event qualifies when it matches either a forced type or the priority threshold. Events outside all segments are ignored safely. + ## 6. Blind spots Every document carries a `blind_spots` list: plain-language statements of what the capture pipeline systematically cannot see (e.g. "browser URLs are only captured for browser apps"). Producers must not remove entries to make output look more complete. diff --git a/errors.md b/errors.md new file mode 100644 index 0000000..e69de29 diff --git a/src/activity_frames/__init__.py b/src/activity_frames/__init__.py index a505913..6ae2b11 100644 --- a/src/activity_frames/__init__.py +++ b/src/activity_frames/__init__.py @@ -18,6 +18,7 @@ from .communications import surfaces as comm_surfaces from .db import Database, RecorderDBNotFound, find_default_db from .emit import context_block, to_json, to_markdown, to_yaml +from .events import Event from .entities import PageRef, parse_url from .frames import ( SCHEMA_VERSION, @@ -39,6 +40,7 @@ "CommSurface", "Coverage", "Database", + "Event", "PageRef", "RecorderDBNotFound", "SCHEMA_VERSION", diff --git a/src/activity_frames/events.py b/src/activity_frames/events.py new file mode 100644 index 0000000..bbb7a6b --- /dev/null +++ b/src/activity_frames/events.py @@ -0,0 +1,25 @@ +"""Minimal external event abstraction for deterministic segmentation. + +Events are authoritative, externally supplied markers (e.g. git.merge, +pr.approved) that callers may pass to the sessionizer to force +deterministic boundaries. The class is intentionally minimal and +frozen to preserve immutability across the pipeline. +""" +from __future__ import annotations + +from dataclasses import dataclass + + + +@dataclass(frozen=True) +class Event: + """An externally supplied, authoritative boundary marker. + + The caller decides what is meaningful; this package does not infer or + classify such events from activity data. + """ + + event_type: str + timestamp: float + source: str + priority: int = 0 diff --git a/src/activity_frames/frames.py b/src/activity_frames/frames.py index 6afcbad..19becc7 100644 --- a/src/activity_frames/frames.py +++ b/src/activity_frames/frames.py @@ -188,11 +188,17 @@ def build_frames( session_gap: float = SESSION_GAP, merge_flicker: float = MERGE_FLICKER, debug: bool = False, + events=(), + forced_event_types=(), + force_priority: int | None = None, ) -> ActivityDocument: """Compile a UTC window of recorder data into an ActivityDocument.""" segs = compute_segments( db, start_utc, end_utc, dwell_cap=dwell_cap, session_gap=session_gap, merge_flicker=merge_flicker, + events=events, + forced_event_types=forced_event_types, + force_priority=force_priority, ) cov = compute_coverage(db, start_utc, end_utc, session_gap=session_gap) diff --git a/src/activity_frames/sessionize.py b/src/activity_frames/sessionize.py index fc2957a..65f90e9 100644 --- a/src/activity_frames/sessionize.py +++ b/src/activity_frames/sessionize.py @@ -17,11 +17,13 @@ """ from __future__ import annotations +from collections.abc import Collection, Sequence from dataclasses import dataclass, field from urllib.parse import urlsplit from ._time import parse_epoch from .db import Database +from .events import Event DWELL_CAP = 90.0 # seconds; max credit for one frame SESSION_GAP = 300.0 # seconds; larger gap = user away / new session @@ -127,6 +129,130 @@ def load_frames(db: Database, start_utc: str, end_utc: str) -> list[RawFrame]: return out +def _event_matches( + event: Event, + *, + forced_event_types: Collection[str], + force_priority: int | None, +) -> bool: + if event.event_type in forced_event_types: + return True + if force_priority is not None and event.priority >= force_priority: + return True + return False + + +def _segment_active_seconds(frames: list[RawFrame], *, dwell_cap: float, session_gap: float) -> float: + total = 0.0 + for i, f in enumerate(frames[:-1]): + gap = frames[i + 1].epoch - f.epoch + if gap <= 0: + continue + if gap > session_gap: + continue + total += min(gap, dwell_cap) + return total + + +def _split_segment_by_event( + seg: Segment, + event: Event, + *, + dwell_cap: float, + session_gap: float, +) -> list[Segment]: + if not seg.frames or len(seg.frames) < 2: + return [seg] + + # Select the nearest inter-frame boundary, not an invented event timestamp. + # The event is mapped to the midpoint between the two closest captured + # frames in the segment, which preserves the underlying capture stream. + boundary_index = min( + range(len(seg.frames) - 1), + key=lambda i: ( + abs(((seg.frames[i].epoch + seg.frames[i + 1].epoch) / 2.0) - event.timestamp), + i, + ), + ) + if boundary_index <= 0 or boundary_index >= len(seg.frames) - 1: + return [seg] + + left_frames = seg.frames[: boundary_index + 1] + right_frames = seg.frames[boundary_index + 1 :] + if not left_frames or not right_frames: + return [seg] + + left = Segment( + app=seg.app, + domain=seg.domain, + start_epoch=seg.start_epoch, + end_epoch=left_frames[-1].epoch, + active_seconds=_segment_active_seconds(left_frames, dwell_cap=dwell_cap, session_gap=session_gap), + frames=list(left_frames), + interruptions=list(seg.interruptions), + break_reason=seg.break_reason, + debug_notes=list(seg.debug_notes), + ) + right = Segment( + app=seg.app, + domain=seg.domain, + start_epoch=right_frames[0].epoch, + end_epoch=seg.end_epoch, + active_seconds=_segment_active_seconds(right_frames, dwell_cap=dwell_cap, session_gap=session_gap), + frames=list(right_frames), + interruptions=list(seg.interruptions), + break_reason=f"event_boundary: {event.event_type}", + debug_notes=list(seg.debug_notes), + ) + left.debug_notes.append(f"event boundary: {event.event_type} at {event.timestamp}") + right.debug_notes.append(f"event boundary: {event.event_type} at {event.timestamp}") + return [left, right] + + +def _apply_event_boundaries( + segs: Sequence[Segment], + events: Sequence[Event] = (), + *, + forced_event_types: Collection[str] = (), + force_priority: int | None = None, + dwell_cap: float = DWELL_CAP, + session_gap: float = SESSION_GAP, +) -> list[Segment]: + """Deterministically split existing segments at authoritative events. + + Events are opt-in and supplied by callers. The implementation never + invents a frame that is not in the capture stream; it only splits + around the nearest captured boundary inside the relevant segment. + """ + if not segs or not events: + return list(segs) + + relevant = [ + e for e in events + if _event_matches(e, forced_event_types=forced_event_types, force_priority=force_priority) + ] + if not relevant: + return list(segs) + + current = list(segs) + for event in sorted(relevant, key=lambda e: (e.timestamp, e.event_type, e.source, e.priority)): + next_segs: list[Segment] = [] + for seg in current: + if not seg.frames: + next_segs.append(seg) + continue + if event.timestamp < seg.start_epoch or event.timestamp > seg.end_epoch: + next_segs.append(seg) + continue + split = _split_segment_by_event(seg, event, dwell_cap=dwell_cap, session_gap=session_gap) + if len(split) == 2 and split[0] is not seg and split[1] is not seg: + next_segs.extend(split) + else: + next_segs.append(seg) + current = next_segs + return current + + def segments( db: Database, start_utc: str, @@ -135,6 +261,9 @@ def segments( dwell_cap: float = DWELL_CAP, session_gap: float = SESSION_GAP, merge_flicker: float = MERGE_FLICKER, + events: Sequence[Event] = (), + forced_event_types: Collection[str] = (), + force_priority: int | None = None, ) -> list[Segment]: """Chronological (app, site) segments for a UTC window. @@ -163,7 +292,16 @@ def segments( include_device_note=include_device_note) ) merged.sort(key=lambda s: s.start_epoch) - return merged + if not events and not forced_event_types and force_priority is None: + return merged + return _apply_event_boundaries( + merged, + events, + forced_event_types=forced_event_types, + force_priority=force_priority, + dwell_cap=dwell_cap, + session_gap=session_gap, + ) def _segment_stream( diff --git a/tests/test_sessionize.py b/tests/test_sessionize.py index 320f52a..7656839 100644 --- a/tests/test_sessionize.py +++ b/tests/test_sessionize.py @@ -1,4 +1,12 @@ -from activity_frames.sessionize import app_ledger, coverage, segments +from activity_frames.events import Event +from activity_frames.sessionize import ( + RawFrame, + Segment, + _apply_event_boundaries, + app_ledger, + coverage, + segments, +) from activity_frames import build_frames @@ -101,7 +109,7 @@ def test_break_reason_in_output_with_debug(fixture_db, day_window): def test_quantified_debug_reasons_flicker_and_dwell(): - from activity_frames.sessionize import RawFrame, _segment_stream + from activity_frames.sessionize import _segment_stream # Frame 1 & 2: gap of 150s (> dwell_cap of 90s, <= session_gap of 300s) -> dwell capped # Frame 3 & 4: brief flicker to Slack (start 260.0, end 265.0 => span 5s <= 20s) @@ -121,3 +129,86 @@ def test_quantified_debug_reasons_flicker_and_dwell(): assert any("merged flicker: Slack 5s <= 20s" in n for n in notes) +def _segment_with_frames(*epochs: float, app: str = "App") -> Segment: + frames = [ + RawFrame(id=i, epoch=e, app=app, window="main", url=None, domain=None) + for i, e in enumerate(epochs, start=1) + ] + start_epoch = epochs[0] + end_epoch = epochs[-1] + return Segment( + app=app, + domain=None, + start_epoch=start_epoch, + end_epoch=end_epoch, + active_seconds=float(end_epoch - start_epoch), + frames=frames, + ) + + +def test_forced_event_splits_segment(): + seg = _segment_with_frames(100.0, 103.0, 108.0, 115.0) + split = _apply_event_boundaries([seg], [Event("git.merge", 105.0, "external", priority=5)], forced_event_types={"git.merge"}) + assert len(split) == 2 + assert [s.start_epoch for s in split] == [100.0, 108.0] + assert [s.end_epoch for s in split] == [103.0, 115.0] + + +def test_non_forced_event_does_not_split_segment(): + seg = _segment_with_frames(100.0, 103.0, 108.0, 115.0) + split = _apply_event_boundaries([seg], [Event("deploy.success", 105.0, "external")]) + assert split == [seg] + + +def test_event_matching_forced_event_types_splits_correctly(): + seg = _segment_with_frames(10.0, 15.0, 20.0, 25.0) + split = _apply_event_boundaries([seg], [Event("pr.approved", 17.0, "ci")], forced_event_types={"pr.approved"}) + assert len(split) == 2 + assert [s.frame_ids for s in split] == [[1, 2], [3, 4]] + + +def test_event_meeting_force_priority_splits_correctly(): + seg = _segment_with_frames(10.0, 15.0, 20.0, 25.0) + split = _apply_event_boundaries([seg], [Event("custom.event", 17.0, "external", priority=3)], force_priority=3) + assert len(split) == 2 + assert [s.frame_ids for s in split] == [[1, 2], [3, 4]] + + +def test_event_below_force_priority_does_not_split(): + seg = _segment_with_frames(10.0, 15.0, 20.0, 25.0) + split = _apply_event_boundaries([seg], [Event("custom.event", 17.0, "external", priority=2)], force_priority=3) + assert split == [seg] + + +def test_event_outside_all_segments_is_ignored_safely(): + seg = _segment_with_frames(100.0, 105.0, 110.0) + split = _apply_event_boundaries([seg], [Event("git.merge", 50.0, "external")], forced_event_types={"git.merge"}) + assert split == [seg] + + +def test_multiple_events_produce_deterministic_boundaries(): + seg = _segment_with_frames(100.0, 103.0, 108.0, 115.0, 120.0) + events = [ + Event("git.merge", 105.0, "external", priority=5), + Event("deploy.success", 118.0, "external", priority=5), + ] + split = _apply_event_boundaries([seg], events, forced_event_types={"git.merge"}, force_priority=5) + assert [s.start_epoch for s in split] == [100.0, 108.0, 120.0] + assert [s.end_epoch for s in split] == [103.0, 115.0, 120.0] + + +def test_existing_behavior_without_events_remains_unchanged(): + seg = _segment_with_frames(100.0, 103.0, 108.0, 115.0) + baseline = _apply_event_boundaries([seg], ()) + assert baseline == [seg] + + +def test_boundary_behavior_between_captured_frames_uses_nearest_boundary(): + seg = _segment_with_frames(100.0, 103.0, 108.0, 115.0) + split = _apply_event_boundaries([seg], [Event("git.merge", 105.5, "external")], forced_event_types={"git.merge"}) + assert len(split) == 2 + assert [s.end_epoch for s in split] == [103.0, 115.0] + # event is nearer to 103 than 108, so the split should land on the earlier captured boundary + assert split[0].frames[-1].epoch == 103.0 + +