Skip to content

Commit 6cab89e

Browse files
evansenterclaude
andcommitted
fix: First bus event ingestion now captures full history (#108)
First run no longer applies a timestamp cutoff — it ingests ALL events from the event-bus database. Subsequent runs continue using the high-water mark for incremental updates. This prevents the scenario where a days=7 first run permanently skips older events. Closes #108 Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent 30eceef commit 6cab89e

2 files changed

Lines changed: 16 additions & 22 deletions

File tree

‎src/agent_session_analytics/bus_ingest.py‎

Lines changed: 5 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,6 @@
88
import json
99
import logging
1010
import sqlite3
11-
from datetime import datetime, timedelta
1211
from pathlib import Path
1312

1413
from agent_session_analytics.storage import SQLiteStorage
@@ -29,12 +28,12 @@ def ingest_bus_events(storage: SQLiteStorage, days: int = 7) -> dict:
2928
"""Ingest events from event-bus database.
3029
3130
Performs incremental ingestion by tracking the last ingested event ID.
32-
Events are read from the event-bus database in read-only mode.
31+
First run ingests all events; subsequent runs only pick up new events.
3332
Raw event JSON is stored alongside parsed data for future re-parsing.
3433
3534
Args:
3635
storage: Session analytics storage instance
37-
days: Number of days to look back for initial ingestion
36+
days: Unused (kept for backward compatibility)
3837
3938
Returns:
4039
Dict with ingestion stats including events_ingested count
@@ -50,9 +49,6 @@ def ingest_bus_events(storage: SQLiteStorage, days: int = 7) -> dict:
5049
last_event = storage.execute_query("SELECT MAX(event_id) as last_id FROM bus_events")
5150
last_id = last_event[0]["last_id"] if last_event and last_event[0]["last_id"] else 0
5251

53-
# Calculate cutoff for first-run ingestion
54-
cutoff = datetime.now() - timedelta(days=days)
55-
5652
# Read from event-bus DB (read-only mode)
5753
try:
5854
conn = sqlite3.connect(f"file:{EVENT_BUS_DB}?mode=ro", uri=True)
@@ -65,9 +61,8 @@ def ingest_bus_events(storage: SQLiteStorage, days: int = 7) -> dict:
6561
}
6662

6763
try:
68-
# Query events newer than last ingested ID, or from cutoff on first run
6964
if last_id > 0:
70-
# Incremental: get events after last ID
65+
# Incremental: get events after last ingested ID
7166
rows = conn.execute(
7267
"""
7368
SELECT id, event_type, channel, session_id, timestamp, payload
@@ -78,15 +73,13 @@ def ingest_bus_events(storage: SQLiteStorage, days: int = 7) -> dict:
7873
(last_id,),
7974
).fetchall()
8075
else:
81-
# First run: get events from cutoff
76+
# First run: ingest ALL events (no timestamp cutoff)
8277
rows = conn.execute(
8378
"""
8479
SELECT id, event_type, channel, session_id, timestamp, payload
8580
FROM events
86-
WHERE timestamp >= ?
8781
ORDER BY id
88-
""",
89-
(cutoff.isoformat(),),
82+
"""
9083
).fetchall()
9184

9285
if not rows:

‎tests/test_bus_ingest.py‎

Lines changed: 11 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -106,10 +106,10 @@ class TestIngestBusEvents:
106106
def test_ingest_from_bus_db(self, storage, bus_db):
107107
"""Test basic ingestion from an event-bus database."""
108108
with patch("agent_session_analytics.bus_ingest.EVENT_BUS_DB", bus_db):
109-
result = ingest_bus_events(storage, days=7)
109+
result = ingest_bus_events(storage)
110110

111111
assert result["status"] == "ok"
112-
assert result["events_ingested"] == 4 # 4 within 7 days
112+
assert result["events_ingested"] == 5 # First run gets ALL events
113113

114114
def test_incremental_ingestion(self, storage, bus_db):
115115
"""Test that second ingestion only picks up new events."""
@@ -160,11 +160,11 @@ def test_missing_db_skips(self, storage):
160160
def test_raw_events_stored(self, storage, bus_db):
161161
"""Test that raw event JSON is stored in raw_bus_events table."""
162162
with patch("agent_session_analytics.bus_ingest.EVENT_BUS_DB", bus_db):
163-
ingest_bus_events(storage, days=7)
163+
ingest_bus_events(storage)
164164

165165
# Check raw_bus_events table
166166
rows = storage.execute_query("SELECT * FROM raw_bus_events ORDER BY event_id")
167-
assert len(rows) == 4 # 4 within 7 days
167+
assert len(rows) == 5 # First run gets ALL events
168168

169169
# Verify raw JSON is parseable and contains original fields
170170
raw = json.loads(rows[0]["entry_json"])
@@ -187,19 +187,20 @@ def test_raw_events_dedup(self, storage, bus_db):
187187
def test_repo_extraction(self, storage, bus_db):
188188
"""Test that repo is correctly extracted from channel."""
189189
with patch("agent_session_analytics.bus_ingest.EVENT_BUS_DB", bus_db):
190-
ingest_bus_events(storage, days=7)
190+
ingest_bus_events(storage)
191191

192192
rows = storage.execute_query(
193193
"SELECT repo FROM bus_events WHERE event_type = 'gotcha_discovered' ORDER BY event_id"
194194
)
195195
assert rows[0]["repo"] == "dotfiles"
196196

197-
def test_full_history_ingestion(self, storage, bus_db):
198-
"""Test ingestion with large days window gets all events."""
197+
def test_first_run_ignores_days_param(self, storage, bus_db):
198+
"""Test that first run ingests ALL events regardless of days param."""
199199
with patch("agent_session_analytics.bus_ingest.EVENT_BUS_DB", bus_db):
200-
result = ingest_bus_events(storage, days=30)
200+
# Even with days=1, first run should get all 5 events (including 10-day-old one)
201+
result = ingest_bus_events(storage, days=1)
201202

202-
assert result["events_ingested"] == 5 # All events including old one
203+
assert result["events_ingested"] == 5
203204

204205

205206
class TestQueryBusEvents:
@@ -209,7 +210,7 @@ class TestQueryBusEvents:
209210
def storage_with_bus_events(self, storage, bus_db):
210211
"""Storage with bus events ingested."""
211212
with patch("agent_session_analytics.bus_ingest.EVENT_BUS_DB", bus_db):
212-
ingest_bus_events(storage, days=30)
213+
ingest_bus_events(storage)
213214
return storage
214215

215216
def test_query_all(self, storage_with_bus_events):

0 commit comments

Comments
 (0)