Skip to content
Open
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
14 changes: 14 additions & 0 deletions SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Empty file added errors.md
Empty file.
2 changes: 2 additions & 0 deletions src/activity_frames/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -39,6 +40,7 @@
"CommSurface",
"Coverage",
"Database",
"Event",
"PageRef",
"RecorderDBNotFound",
"SCHEMA_VERSION",
Expand Down
25 changes: 25 additions & 0 deletions src/activity_frames/events.py
Original file line number Diff line number Diff line change
@@ -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
6 changes: 6 additions & 0 deletions src/activity_frames/frames.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
140 changes: 139 additions & 1 deletion src/activity_frames/sessionize.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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.

Expand Down Expand Up @@ -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(
Expand Down
95 changes: 93 additions & 2 deletions tests/test_sessionize.py
Original file line number Diff line number Diff line change
@@ -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

Expand Down Expand Up @@ -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)
Expand All @@ -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