feat(P02): telemetry models + TraceStore SQLite + TS types (Wave 1)

Task 2-1-01: TelemetryEvent (SQLModel table, composite PK learner/task/seq) + TraceSpan
derived view; SQLAlchemy @validates enforcement (sqlmodel drops Field constraints on table
models). Task 2-1-02: TraceStore protocol (D-019-mirrored) + SQLiteTraceStore with WAL +
synchronous=NORMAL (a-3); concurrent writer/reader no database-locked; idempotent
dedup-on-retry. Task 2-1-03: TS TelemetryEvent/TraceSpan mirroring Python field-for-field
(snake_case for byte-identical JSON); full typecheck green.

195/195 + typecheck pass; ruff clean.

---ci---
phase: 2
milestone: v0.3
status: execute
requirements: {covered: [REQ-3-003], partial: []}
---/ci---
This commit is contained in:
CIAgent
2026-09-11 18:51:42 +00:00
parent f0df18576e
commit cfdceac17a
8 changed files with 679 additions and 0 deletions
@@ -0,0 +1,11 @@
"""Live build telemetry — event models and the TraceStore protocol (REQ-3-003).
Boundary rule (D-027): telemetry/ imports from config only — never from
agents/ or api/ (agents call engines through narrow interfaces, never
the reverse; api/ composes stores via DI).
"""
from .models import EventKind, TelemetryEvent, TraceSpan
from .store import SQLiteTraceStore, TraceStore
__all__ = ["EventKind", "SQLiteTraceStore", "TelemetryEvent", "TraceSpan", "TraceStore"]
@@ -0,0 +1,108 @@
"""Telemetry event record — the row the trace store persists (REQ-3-003, D-027).
One model serves both as the JSON payload sent by producers and as the SQLite
row schema. `payload` is stored as a JSON column (native JSONB on Postgres —
no migration-time shape change, D-027).
Field contract (consumed by the trace store and the grader):
learner_id — non-empty learner identifier.
task_id — non-empty task/session identifier; trace identity is the
(learner_id, task_id) pair.
seq — sequence number per trace, >= 0. Monotonicity per
(learner, task) is enforced by the store (Task 2-1-02);
this model only rejects negative seqs.
kind — event discriminator: command | file_diff | run_result |
test_result | activity | stdin | stdout.
payload — free-form JSON detail blob.
ts — envelope timestamp (UTC); monotonicity enforced at ingest.
sandbox_id — originating sandbox ("" for non-sandbox sources).
Boundary (D-027): telemetry/ never imports agents/ or api/.
"""
from datetime import datetime
from typing import Any, Literal
from pydantic import BaseModel, ConfigDict, Field
from sqlalchemy import JSON, Index, String
from sqlalchemy.orm import validates
from sqlmodel import Field as SQLField
from sqlmodel import SQLModel
EventKind = Literal[
"command",
"file_diff",
"run_result",
"test_result",
"activity",
"stdin",
"stdout",
]
_EVENT_KINDS: frozenset[str] = frozenset(EventKind.__args__)
class TelemetryEvent(SQLModel, table=True):
"""A single durable telemetry event; (learner_id, task_id, seq) is PK.
Constraint enforcement uses SQLAlchemy `@validates` hooks: sqlmodel
0.0.42's metaclass drops pydantic `Field(ge=...)`/`field_validator`
constraints for table models (the decorators register but never make it
into the core schema), while `@validates` fires on every attribute set —
construction included — and raises ValueError on violation. seq >= 0 plus
a VARCHAR kind column keep the DB shape Postgres-ready (D-027).
"""
__tablename__ = "telemetry_event"
# PK columns already produce a unique index; this secondary index covers
# trace reads ordered by seq without depending on the PK column order
# (Postgres migration target D-027).
__table_args__ = (Index("ix_telemetry_event_trace", "learner_id", "task_id"),)
learner_id: str = SQLField(primary_key=True)
task_id: str = SQLField(primary_key=True)
seq: int = SQLField(primary_key=True)
# Bare Literal annotations crash sqlmodel<=0.0.42's column inference
# (issubclass(TypeAlias, Enum)); an explicit sa_type + the validates hook
# below gives the same contract: VARCHAR column, Literal-rejected values.
kind: EventKind = SQLField(sa_type=String)
# JSON column: stored as TEXT on SQLite, native JSONB on Postgres (D-027).
payload: dict[str, Any] = SQLField(default_factory=dict, sa_type=JSON)
ts: datetime
sandbox_id: str = SQLField(default="")
@validates("learner_id", "task_id")
def _ids_non_empty(self, key: str, value: str) -> str:
if not value:
raise ValueError(f"{key} must be a non-empty identifier")
return value
@validates("seq")
def _seq_non_negative(self, key: str, value: int) -> int:
if value < 0:
raise ValueError("seq must be >= 0 (monotonicity is the store's job)")
return value
@validates("kind")
def _kind_is_known(self, key: str, value: str) -> str:
if value not in _EVENT_KINDS:
raise ValueError(f"unknown event kind: {value!r}")
return value
class TraceSpan(BaseModel):
"""Derived view: the ordered event trace for one (learner_id, task_id).
NOT a table — materialized by the store from persisted TelemetryEvents
(grader/Lab consume this shape; replay order is the seq column).
"""
model_config = ConfigDict(frozen=True)
learner_id: str = Field(min_length=1)
task_id: str = Field(min_length=1)
events: tuple[TelemetryEvent, ...] = ()
@property
def latest_seq(self) -> int:
"""Highest seq in the span; -1 when empty (store convention)."""
return self.events[-1].seq if self.events else -1
@@ -0,0 +1,198 @@
"""TraceStore — telemetry persistence protocol + SQLite implementation (REQ-3-003, D-027).
Postgres-migration-ready (D-027): the protocol is the only surface the API /
grader layers touch; swapping SQLiteTraceStore for a Postgres-backed
implementation must not change call sites. The `telemetry_event` table uses
only portable column types (str / int / datetime / JSON), so the same SQLModel
schema stands up unchanged on Postgres.
Ingest is at-least-once: duplicates carry the same (learner_id, task_id, seq)
idempotency key, so `append` with a triplet that is already stored is a no-op.
The pair (learner_id, task_id) identifies a trace; `seq` numbers events in it
starting at 0.
Concurrency (a-3): the engine enables WAL + synchronous=NORMAL and a busy
timeout at connection time, so the ingest writer and grader readers do not hit
`database is locked` on the single-box pilot.
Boundary: `telemetry/` never imports `agents/` / `api/` and has no FastAPI
dependency.
"""
import logging
import sqlite3
from collections.abc import Iterator
from contextlib import contextmanager
from datetime import UTC, datetime
from pathlib import Path
from typing import Any, Protocol
import sqlalchemy as sa
from sqlmodel import Session, SQLModel, create_engine, select
from ..config import Settings
from .models import TelemetryEvent
logger = logging.getLogger(__name__)
class TraceStore(Protocol):
"""Persistence contract for ordered per-learner task trace streams.
Implemented by SQLiteTraceStore (v0.3, D-027); a Postgres implementation
must satisfy the same surface.
"""
def append(self, event: TelemetryEvent) -> None:
"""Store one event. IDEMPOTENT on (learner_id, task_id, seq):
at-least-once ingest retries with the same triplet are deduped
(stored once), not rejected. Later events must not overwrite an
existing row.
"""
...
def get_trace(self, learner_id: str, task_id: str) -> list[TelemetryEvent]:
"""All stored events for the trace, ordered by seq ascending.
Detached from any DB session — safe to pass across layers. Empty list
when the trace has no events.
"""
...
def gaps(self, learner_id: str, task_id: str) -> list[int]:
"""Missing seqs in 0..latest for the trace ([0,2,3] stored -> [1])."""
...
def latest_seq(self, learner_id: str, task_id: str) -> int:
"""Highest stored seq for the trace; -1 when no events exist."""
...
def list_tasks(self, learner_id: str) -> list[str]:
"""Distinct task_ids with at least one event for the learner."""
...
def close(self) -> None:
"""Release DB connections. Store must not be used after close."""
...
def _sqlite_connect(dbapi_connection: sqlite3.Connection, _: object) -> None:
"""Per-connection pragma setup (a-3).
journal_mode=WAL — readers never block the single writer.
synchronous=NORMAL — safe in WAL mode, avoids full fsync-per-commit.
busy_timeout=5000 — retry briefly under contention instead of
`OperationalError: database is locked`.
"""
cursor = dbapi_connection.cursor()
cursor.execute("PRAGMA journal_mode=WAL")
cursor.execute("PRAGMA synchronous=NORMAL")
cursor.execute("PRAGMA busy_timeout=5000")
cursor.close()
def _as_utc(ts: datetime) -> datetime:
"""Timestamps travel as naive datetime on SQLite, tz-aware elsewhere.
SQLite (via SQLModel) drops tzinfo; Postgres TIMESTAMP WITH TIME ZONE
keeps it. Normalizing on the read/write boundary makes the store's
contract tz-aware UTC regardless of the backend (D-027).
"""
if ts.tzinfo is None:
return ts.replace(tzinfo=UTC) # naive UTC read-side: label as UTC
return ts.astimezone(UTC)
class SQLiteTraceStore:
"""SQLite-backed TraceStore (SQLModel). First real persistence (D-027)."""
def __init__(self, db_path: Path | None = None) -> None:
self._db_path: Path = db_path if db_path is not None else Settings().db_path
self._engine = create_engine(f"sqlite:///{self._db_path}")
sa.event.listen(self._engine, "connect", _sqlite_connect)
SQLModel.metadata.create_all(self._engine)
@contextmanager
def _session(self) -> Iterator[Session]:
# expire_on_commit=False: ORM objects returned from `append`'s
# IntegrityError path stay usable without a refresh round-trip.
with Session(self._engine, expire_on_commit=False) as session:
yield session
def append(self, event: TelemetryEvent) -> None:
# INSERT-if-absent via PK: sqlite3 raises IntegrityError on a
# duplicate (learner_id, task_id, seq); swallow it — the row is
# already stored, which is the dedup contract for at-least-once
# ingest. `session.merge` would upsert instead; wrong semantics here.
with self._session() as session:
try:
session.add(event)
session.commit()
except sa.exc.IntegrityError:
session.rollback()
logger.debug(
"trace event dedup: %s/%s seq=%d already stored",
event.learner_id,
event.task_id,
event.seq,
)
def get_trace(self, learner_id: str, task_id: str) -> list[TelemetryEvent]:
with self._session() as session:
stmt = (
select(TelemetryEvent)
.where(TelemetryEvent.learner_id == learner_id)
.where(TelemetryEvent.task_id == task_id)
.order_by(TelemetryEvent.seq)
)
results = session.exec(stmt).all()
# Detach from the session: callers must not depend on open-session
# ORM magic (lazy loads fail once the session is closed).
for row in results:
row.ts = _as_utc(row.ts)
session.expunge(row)
return list(results)
def _stored_seqs(self, learner_id: str, task_id: str) -> list[int]:
with self._session() as session:
stmt = (
select(TelemetryEvent.seq)
.where(TelemetryEvent.learner_id == learner_id)
.where(TelemetryEvent.task_id == task_id)
.order_by(TelemetryEvent.seq)
)
# sqlmodel scalar select: rows are plain ints, not 1-tuples.
return [int(seq) for seq in session.exec(stmt).all()]
def gaps(self, learner_id: str, task_id: str) -> list[int]:
seqs = self._stored_seqs(learner_id, task_id)
if not seqs:
return []
present = set(seqs)
# seq numbering starts at 0; a gap is any seq in 0..latest not stored.
return [seq for seq in range(seqs[-1] + 1) if seq not in present]
def latest_seq(self, learner_id: str, task_id: str) -> int:
with self._session() as session:
stmt = (
select(sa.func.max(TelemetryEvent.seq))
.where(TelemetryEvent.learner_id == learner_id)
.where(TelemetryEvent.task_id == task_id)
)
latest: Any = session.exec(stmt).one()
return -1 if latest is None else int(latest)
def list_tasks(self, learner_id: str) -> list[str]:
with self._session() as session:
stmt = (
select(TelemetryEvent.task_id)
.where(TelemetryEvent.learner_id == learner_id)
.distinct()
.order_by(TelemetryEvent.task_id)
)
# sqlmodel scalar select: rows are plain strs, not 1-tuples.
return [str(task_id) for task_id in session.exec(stmt).all()]
def close(self) -> None:
self._engine.dispose()
@@ -0,0 +1,126 @@
"""TelemetryEvent model contract tests (REQ-3-003, D-027).
Pure validation tests — no store, no network, no fixtures beyond a
canonical kwargs builder.
"""
import json
from datetime import UTC, datetime
import pytest
from pydantic import ValidationError
from ai_service.telemetry import EventKind, TelemetryEvent, TraceSpan
def make_event_kwargs(**overrides) -> dict:
base = {
"learner_id": "learner-1",
"task_id": "task-1",
"seq": 0,
"kind": "command",
"payload": {"argv": ["pytest"], "cwd": "/workspace"},
"ts": datetime(2026, 9, 11, tzinfo=UTC),
"sandbox_id": "sbx-1",
}
base.update(overrides)
return base
def make_event(**overrides) -> TelemetryEvent:
return TelemetryEvent(**make_event_kwargs(**overrides))
def test_valid_event_constructs_with_all_kinds():
for kind in (
"command",
"file_diff",
"run_result",
"test_result",
"activity",
"stdin",
"stdout",
):
event = make_event(kind=kind)
assert event.kind == kind
def test_zero_seq_accepted():
assert make_event(seq=0).seq == 0
def test_negative_seq_rejected():
# @validates hooks raise ValueError (not pydantic ValidationError) at
# construction — same for kind and empty-id checks below.
with pytest.raises(ValueError, match="seq must be >= 0"):
make_event(seq=-1)
def test_bad_kind_rejected():
with pytest.raises(ValueError, match="unknown event kind"):
make_event(kind="keystroke") # not in the EventKind literal
def test_empty_learner_id_rejected():
with pytest.raises(ValueError, match="non-empty identifier"):
make_event(learner_id="")
def test_empty_task_id_rejected():
with pytest.raises(ValueError, match="non-empty identifier"):
make_event(task_id="")
def test_payload_json_roundtrip():
payload = {
"command": "pytest -q",
"exit_code": 1,
"durations": [0.12, 3.4],
"nested": {"passed": 7, "failed": 2},
"unicode": "héllo",
}
event = make_event(payload=payload)
assert event.payload == payload
# JSON-serializable payloads survive a full dumps/loads roundtrip.
assert json.loads(json.dumps(event.payload)) == payload
def test_event_json_roundtrip():
event = make_event(payload={"stdout": "ok", "n": 3})
restored = TelemetryEvent.model_validate_json(event.model_dump_json())
# pydantic leaves a JSON-str ts as str; compare field-by-field instead of
# dataclass equality, which distinguishes 'Z' string vs parsed datetime.
assert restored.learner_id == event.learner_id
assert restored.task_id == event.task_id
assert restored.seq == event.seq
assert restored.kind == event.kind
assert restored.payload == event.payload
assert restored.sandbox_id == event.sandbox_id
# SQLModel leaves a JSON-serialized ts as its str form; parse to compare.
restored_ts = datetime.fromisoformat(str(restored.ts).replace("Z", "+00:00"))
assert restored_ts == event.ts
def test_tracespan_holds_ordered_events():
events = tuple(make_event(seq=seq, kind="file_diff") for seq in range(3))
span = TraceSpan(learner_id="learner-1", task_id="task-1", events=events)
assert [e.seq for e in span.events] == [0, 1, 2]
assert span.latest_seq == 2
def test_tracespan_empty_has_no_latest_seq():
span = TraceSpan(learner_id="learner-1", task_id="task-1")
assert span.events == ()
assert span.latest_seq == -1
def test_tracespan_rejects_empty_ids():
with pytest.raises(ValidationError):
TraceSpan(learner_id="", task_id="task-1")
with pytest.raises(ValidationError):
TraceSpan(learner_id="learner-1", task_id="")
def test_event_kind_type_is_exported_literal():
# The kind discriminator is part of the public module surface.
assert "command" in EventKind.__args__
@@ -0,0 +1,191 @@
"""SQLiteTraceStore tests (REQ-3-003, D-027).
Each test gets its own tmp-path SQLite file — no shared disk state. Covers:
- append / ordered get (events may arrive out of order; reads come back
ordered by seq)
- dedup on retry: the same (learner_id, task_id, seq) appended twice is
stored exactly once (at-least-once ingest contract)
- retry dedup must not overwrite the originally stored row
- gap detection, latest_seq, list_tasks
- WAL + synchronous=NORMAL pragmas actually applied to the DB file
- concurrent writer + reader against the same DB file (a-3 smoke test)
"""
import concurrent.futures
import threading
from datetime import UTC, datetime, timedelta
from pathlib import Path
from typing import Any
import pytest
import sqlalchemy as sa
from ai_service.telemetry.models import EventKind, TelemetryEvent
from ai_service.telemetry.store import SQLiteTraceStore
_BASE_TS = datetime(2026, 9, 11, 12, 0, 0, tzinfo=UTC)
def make_event(
seq: int,
learner_id: str = "learner-1",
task_id: str = "task-1",
kind: EventKind = "command",
payload: dict[str, Any] | None = None,
sandbox_id: str = "sbx-1",
) -> TelemetryEvent:
return TelemetryEvent(
learner_id=learner_id,
task_id=task_id,
seq=seq,
kind=kind,
payload=payload if payload is not None else {"seq": seq},
ts=_BASE_TS + timedelta(seconds=seq),
sandbox_id=sandbox_id,
)
@pytest.fixture
def store(tmp_path: Path):
s = SQLiteTraceStore(db_path=tmp_path / "telemetry.db")
yield s
s.close()
def test_append_and_get_trace_orders_by_seq(store: SQLiteTraceStore) -> None:
# Append out of order; reads must come back ordered by seq.
store.append(make_event(2, kind="stdout"))
store.append(make_event(0, kind="stdin"))
store.append(make_event(1, kind="run_result"))
trace = store.get_trace("learner-1", "task-1")
assert [e.seq for e in trace] == [0, 1, 2]
assert [e.kind for e in trace] == ["stdin", "run_result", "stdout"]
assert all(e.learner_id == "learner-1" for e in trace)
assert all(e.task_id == "task-1" for e in trace)
def test_get_trace_round_trips_fields(store: SQLiteTraceStore) -> None:
event = make_event(0, payload={"file": "a.py", "nested": {"ok": True}})
store.append(event)
(row,) = store.get_trace("learner-1", "task-1")
assert row.learner_id == event.learner_id
assert row.task_id == event.task_id
assert row.seq == 0
assert row.kind == "command"
assert row.payload == {"file": "a.py", "nested": {"ok": True}}
assert row.ts == _BASE_TS
assert row.sandbox_id == "sbx-1"
def test_get_trace_is_scoped_to_the_task_pair(store: SQLiteTraceStore) -> None:
store.append(make_event(0, learner_id="learner-1", task_id="task-1"))
store.append(make_event(0, learner_id="learner-1", task_id="task-2"))
store.append(make_event(0, learner_id="learner-2", task_id="task-1"))
assert [e.seq for e in store.get_trace("learner-1", "task-1")] == [0]
assert len(store.get_trace("learner-1", "task-2")) == 1
assert len(store.get_trace("learner-2", "task-1")) == 1
assert store.get_trace("learner-1", "task-missing") == []
def test_append_is_idempotent_on_retry(store: SQLiteTraceStore) -> None:
# At-least-once ingest re-delivers the same event (same idempotency key).
# It must be stored exactly once and the retry must be a no-op success.
original = make_event(0, kind="command", payload={"attempt": 1})
store.append(original)
store.append(original)
# A distinct-but-conflicting retry (same key, different body) is also
# deduped — the first stored row wins, no overwrite.
retry = make_event(0, kind="file_diff", payload={"attempt": 2})
store.append(retry)
trace = store.get_trace("learner-1", "task-1")
assert len(trace) == 1
assert trace[0].kind == "command"
assert trace[0].payload == {"attempt": 1}
assert store.latest_seq("learner-1", "task-1") == 0
def test_gaps_reports_missing_seqs(store: SQLiteTraceStore) -> None:
for seq in (0, 2, 3, 7):
store.append(make_event(seq))
assert store.gaps("learner-1", "task-1") == [1, 4, 5, 6]
# Gaps are per-trace: an empty trace has no gaps at all.
assert store.gaps("learner-1", "task-unknown") == []
def test_latest_seq(store: SQLiteTraceStore) -> None:
assert store.latest_seq("learner-1", "task-1") == -1
store.append(make_event(0))
assert store.latest_seq("learner-1", "task-1") == 0
store.append(make_event(5)) # gaps do not move latest_seq
assert store.latest_seq("learner-1", "task-1") == 5
# Scoped to the trace pair.
assert store.latest_seq("learner-1", "task-2") == -1
def test_list_tasks(store: SQLiteTraceStore) -> None:
assert store.list_tasks("learner-1") == []
store.append(make_event(0, task_id="task-b"))
store.append(make_event(1, task_id="task-b"))
store.append(make_event(0, task_id="task-a"))
assert store.list_tasks("learner-1") == ["task-a", "task-b"]
# Scoped per learner.
assert store.list_tasks("learner-2") == []
def test_pragmas_are_applied(store: SQLiteTraceStore) -> None:
# Pragmas are per-connection; query through the store's engine so the
# connect hook (not a default sqlite3 connection) is what we inspect.
with store._engine.connect() as conn:
(journal_mode,) = conn.execute(sa.text("PRAGMA journal_mode")).one()
(synchronous,) = conn.execute(sa.text("PRAGMA synchronous")).one()
assert journal_mode == "wal"
# synchronous=NORMAL is 1 in SQLite's pragma numbering.
assert synchronous == 1
def test_concurrent_writer_and_reader_no_database_is_locked(tmp_path: Path) -> None:
"""One thread appends while another reads in a tight loop (a-3).
Without WAL + busy_timeout this pattern reliably produces
`OperationalError: database is locked` on SQLite. The assertion is that
every reader call completes and the final trace is complete.
"""
db_path = tmp_path / "telemetry.db"
n_events = 60
stop_writing = threading.Event()
writer = SQLiteTraceStore(db_path=db_path)
reader = SQLiteTraceStore(db_path=db_path)
try:
def write_events() -> None:
for seq in range(n_events):
writer.append(make_event(seq))
stop_writing.set()
def read_trace() -> None:
while not stop_writing.is_set():
reader.get_trace("learner-1", "task-1")
# Final read after the writer is done.
assert len(reader.get_trace("learner-1", "task-1")) == n_events
with concurrent.futures.ThreadPoolExecutor(max_workers=2) as pool:
futures = [pool.submit(write_events), pool.submit(read_trace)]
for future in futures:
future.result(timeout=30)
assert [e.seq for e in writer.get_trace("learner-1", "task-1")] == list(
range(n_events)
)
finally:
reader.close()
writer.close()