Implementing Incremental Session Processing in ADR Sensor to Avoid Reprocessing
ADR Sensor uses timestamp‑encoded filenames and a three‑stage filter‑persist‑index pipeline to skip already‑processed sessions while capturing any updates.
The ADR Sensor from Uber's ADR repository ingests telemetry from multiple LLM‑based coding agents (Claude, Cursor, Warp, Codex, etc.) and stores events as AgentEvent objects. Without incremental handling, every sensor run would reparse and rewrite identical session data, wasting I/O and compute. This article explains how the sensor implements incremental session processing using tightly coupled components in Sensor/adr_sensor/observer.py.
Core Components for Incremental Processing
Three methods in AgentObserver work together to enable idempotent, incremental ingestion:
save_sessions_to_individual_files — Persistent Storage with Temporal Identity
Each AgentEvent is written to a uniquely named file that encodes both session identity and timestamp. This naming convention is the foundation of the incremental system.
Key implementation details from Sensor/adr_sensor/observer.py (lines 23‑34):
- Filename format:
adr.{clean_session_id}.{timestamp}.json - Session ID sanitization via
_clean_filenameremoves filesystem‑unsafe characters - Timestamp formatting via
format_timestamp_for_filenameensures lexicographic sortability - Output defaults to
_get_default_session_dir()or accepts a customoutput_dir - Returns a list of created file paths for downstream verification
from pathlib import Path
from adr_sensor.observer import AgentObserver
observer = AgentObserver()
events, _ = observer.ingest_all(source_filter="all")
# Persist with custom output directory
saved_paths = observer.save_sessions_to_individual_files(
events,
output_dir=Path("/tmp/adr_sessions")
)
_get_existing_session_files — Building the Session Index
Before writing new events, the sensor scans for existing files and builds an index of the newest entry per session.
Implementation at lines 82‑92 of observer.py:
- Glob pattern
adr.*.jsonlocates candidate files parse_timestamp_from_filenameextracts temporal metadata from each filename- Per‑session deduplication keeps only the newest file based on parsed timestamp
- Returns a dictionary mapping
session_id→{path, timestamp, filename}
This index enables O(1) lookup when filtering incoming events.
filter_entries_by_existing_files — Differential Event Selection
The filtering layer compares incoming events against the existing index, removing any that are already persisted and unchanged.
Implementation at lines 60‑71 of observer.py:
- Calls
_get_existing_session_filesto build the current index - Normalizes each incoming event's timestamp via
normalize_timestamp - Retains events where:
session_idis not present in the index (new session), or- Normalized timestamp is newer than the indexed file's timestamp (updated session)
- Returns the filtered list ready for
save_sessions_to_individual_files
# Complete incremental workflow
observer = AgentObserver()
events, _ = observer.ingest_all(source_filter="all")
# Step 2: Filter out already-processed sessions
filtered_events = observer.filter_entries_by_existing_files(events)
# Step 3: Persist only new/updated data
saved_paths = observer.save_sessions_to_individual_files(
filtered_events,
output_dir=Path("/tmp/adr_sessions")
)
print(f"New sessions saved: {len(saved_paths)}")
How the Incremental Pipeline Executes
The control flow follows a strict sequence:
| Step | Method | Purpose |
|---|---|---|
| 1 | ingest_all() |
Collect raw AgentEvent objects from configured parsers |
| 2 | filter_entries_by_existing_files(events) |
Prune events already persisted with identical timestamps |
| 3 | save_sessions_to_individual_files(filtered_events) |
Write incremental results with temporal filenames |
Running the same ingestion twice demonstrates the idempotent behavior:
$ python -m adr_sensor.cli ingest --source all
# First run: 12 new session files created
$ python -m adr_sensor.cli ingest --source all
# Second run: 0 new session files (incremental skip)
Data Model and Supporting Infrastructure
AgentEvent Schema
The AgentEvent class defined in Sensor/adr_sensor/schemas/agent_event_schema.py provides the structured data model. Each event carries:
session_id: Stable identifier for a coding sessiontimestamp: Event generation time (normalized for comparison)
Timestamp Utilities
Sensor/adr_sensor/utils/timestamp_utils.py contains the formatting and parsing logic that makes incremental detection possible:
format_timestamp_for_filename: Produces filesystem‑safe, sortable stringsparse_timestamp_from_filename: Reverses the encoding for index buildingnormalize_timestamp: Ensures consistent comparison across timezone or precision variations
Testing Incremental Behavior
Unit tests in Sensor/tests/test_observer.py verify:
- Correct filtering when timestamps match exactly
- Proper retention when an incoming event has a newer timestamp
- Handling of malformed or missing timestamp fields
- Directory traversal with non‑ADR files present
CLI Integration
The adr-sensor command‑line tool wires these components together in Sensor/adr_sensor/cli.py. The ingest command automatically applies incremental filtering when invoked repeatedly against the same output directory.
Summary
- ADR Sensor avoids reprocessing through timestamp‑encoded filenames and index‑based filtering
- The
AgentObserverclass inobserver.pyimplements three coordinated methods:_get_existing_session_files,filter_entries_by_existing_files, andsave_sessions_to_individual_files - Sessions are identified by sanitized
session_idwith per‑event timestamps enabling precise change detection - The design reduces I/O, eliminates redundant parsing, and supports downstream batch exporters processing only fresh data
Frequently Asked Questions
How does ADR Sensor detect if a session has already been processed?
The sensor extracts the session_id and timestamp from each filename in the output directory using _get_existing_session_files. When new events arrive, filter_entries_by_existing_files compares their normalized timestamps against the index. Only events with newer timestamps or unknown session IDs proceed to persistence.
What happens if two events have the same session ID but different timestamps?
The incremental system treats this as a session update. The newer timestamp passes the filter and save_sessions_to_individual_files writes a new file. The older file remains in place; downstream consumers should use _get_existing_session_files logic to select the newest version per session.
Can I use a custom output directory for incremental processing?
Yes. Both save_sessions_to_individual_files and filter_entries_by_existing_files respect the output_dir parameter. The filter method automatically scans whichever directory you specify, making the incremental behavior portable across storage locations.
Does the incremental system handle clock skew or timezone differences?
The normalize_timestamp utility in timestamp_utils.py standardizes timestamps before comparison. This normalization ensures that minor formatting variations or timezone representations do not cause false positives for "new" events when the underlying instant is identical.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →