Skip to content
Merged
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
56 changes: 55 additions & 1 deletion src/slopometry/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import sys
import warnings
from importlib.metadata import version
from pathlib import Path

# REASON: analyzed repos may contain invalid escape sequences that emit SyntaxWarnings during AST parsing
warnings.filterwarnings("ignore", category=SyntaxWarning)
Expand Down Expand Up @@ -118,7 +119,7 @@ def hook_opencode(event_type: str) -> None:
"--source",
required=True,
type=click.Choice(["claude_code", "opencode"]),
help="Which harness produced the event on stdin.",
help="Which harness produced the event on stdin. Third-party collectors without an adapter should use 'slopometry ingest' instead.",
)
@click.option(
"--type",
Expand Down Expand Up @@ -153,6 +154,59 @@ def emit_event(source: str, event_type: str | None) -> None:
sys.exit(emit_event_from_stdin(abstract_source, event_type_override=abstract_type))


@cli.command("ingest")
@click.option("--source", required=True, help="Agent tool that produced the events, e.g. 'mmkr'.")
@click.option(
"--file",
"input_file",
type=click.Path(exists=True, dir_okay=False, path_type=Path),
default=None,
help="JSONL file of event envelopes; reads from stdin when omitted.",
)
@click.option(
"--working-directory",
type=click.Path(exists=True, file_okay=False, path_type=Path),
default=None,
help="Working directory recorded on ingested events; defaults to the current directory.",
)
def ingest(source: str, input_file: Path | None, working_directory: Path | None) -> None:
"""Ingest hook-protocol event envelopes (JSONL) from any agent tool.

Each line must be an EventEnvelope JSON object; see
slopometry.core.protocol.schema for the stable schema. Re-ingesting the
same (source, event_id) pairs is an idempotent no-op.
"""
import json

from pydantic import ValidationError

from slopometry.core.protocol.ingest import ingest_envelopes
from slopometry.core.protocol.schema import EventEnvelope

content = input_file.read_text() if input_file else sys.stdin.read()
envelopes: list[EventEnvelope] = []
for line_number, line in enumerate(content.splitlines(), start=1):
if not line.strip():
continue
try:
envelope = EventEnvelope.model_validate(json.loads(line))
except (json.JSONDecodeError, ValidationError) as e:
console.print(f"[red]Invalid envelope at line {line_number}: {e}[/red]")
sys.exit(2)
if envelope.source != source:
console.print(
f"[red]Envelope source '{envelope.source}' at line {line_number} does not match --source '{source}'[/red]"
)
sys.exit(2)
envelopes.append(envelope)

report = ingest_envelopes(envelopes, working_directory=str(working_directory) if working_directory else None)
console.print(
f"[green]Ingested {report.inserted} events from source '{source}'"
f" (skipped {report.skipped_duplicates} duplicates)[/green]"
)


@cli.command("shell-completion")
@click.argument("shell", type=click.Choice(["bash", "zsh", "fish"]))
def shell_completion(shell: str) -> None:
Expand Down
37 changes: 32 additions & 5 deletions src/slopometry/core/database.py
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,8 @@ def _create_tables(self) -> None:
project_source TEXT,
transcript_path TEXT,
source TEXT DEFAULT 'claude_code',
parent_session_id TEXT
parent_session_id TEXT,
event_id TEXT
)
""")
conn.execute("""
Expand All @@ -157,6 +158,10 @@ def _create_tables(self) -> None:
ON hook_events(session_id, event_type, source)
""")

# NOTE: the (source, event_id) unique partial index is created by
# Migration020EnvelopeIdempotency only; _create_tables runs before
# migrations on pre-existing databases that lack the event_id column.

conn.execute("""
CREATE TABLE IF NOT EXISTS experiment_runs (
id TEXT PRIMARY KEY,
Expand Down Expand Up @@ -398,6 +403,26 @@ def _create_tables(self) -> None:

conn.commit()

def has_event(self, source: str, event_id: str) -> bool:
"""Check whether an event with the given (source, event_id) is already stored.

Used by envelope ingestion to make backfills idempotent."""
with self._get_db_connection() as conn:
row = conn.execute(
"SELECT 1 FROM hook_events WHERE source = ? AND event_id = ? LIMIT 1",
(source, event_id),
).fetchone()
return row is not None

def get_max_sequence_number(self, session_id: str) -> int:
"""Get the highest stored sequence number for a session (0 when the session has no events)."""
with self._get_db_connection() as conn:
row = conn.execute(
"SELECT COALESCE(MAX(sequence_number), 0) FROM hook_events WHERE session_id = ?",
(session_id,),
).fetchone()
return row[0]

def save_event(self, event: AbstractHookEvent) -> int:
"""Save a hook event to the database."""
tool_name = event.tool_call.tool_name if event.tool_call else None
Expand All @@ -414,8 +439,8 @@ def save_event(self, event: AbstractHookEvent) -> int:
tool_name, tool_type, metadata, duration_ms,
exit_code, error_message, git_state, working_directory,
project_name, project_source, transcript_path,
source, parent_session_id
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
source, parent_session_id, event_id
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
event.session_id,
Expand All @@ -433,8 +458,9 @@ def save_event(self, event: AbstractHookEvent) -> int:
event.project.name if event.project else None,
event.project.source.value if event.project else None,
event.transcript_location,
event.source.value,
event.source,
event.parent_session_id,
event.event_id,
),
)
return cursor.lastrowid or 0
Expand Down Expand Up @@ -469,14 +495,15 @@ def get_session_events(self, session_id: str) -> list[AbstractHookEvent]:
)

source_val = row["source"] if "source" in row.keys() else None
source = AbstractEventSource(source_val) if source_val else AbstractEventSource.CLAUDE_CODE
source = source_val or str(AbstractEventSource.CLAUDE_CODE)
parent_session_id = row["parent_session_id"] if "parent_session_id" in row.keys() else None

tool_name = row["tool_name"]
tool_type_value = row["tool_type"]
events.append(
AbstractHookEvent(
id=row["id"],
event_id=row["event_id"],
session_id=row["session_id"],
event_type=AbstractEventType(row["event_type"]),
timestamp=datetime.fromisoformat(row["timestamp"]),
Expand Down
51 changes: 51 additions & 0 deletions src/slopometry/core/migrations.py
Original file line number Diff line number Diff line change
Expand Up @@ -569,6 +569,56 @@ def up(self, conn: sqlite3.Connection) -> None:
raise


class Migration020EnvelopeIdempotency(Migration):
"""Enable idempotent envelope ingestion on hook_events.

- Adds the event_id column (if missing) with a partial unique index on
(source, event_id) so re-ingesting a trace is a no-op.
- Repairs event_type values written by the abandoned 'hook-protocol'
branch, which used a divergent taxonomy (tool_call/tool_result/stop/
subagent_stop/subagent_start). Every value maps 1:1 back onto this
taxonomy; databases that never ran that branch are unaffected.

Numbered 020: 017-019 were consumed by diverged local trees.
"""

_WRONG_DIALECT_RENAMES: dict[str, str] = {
"tool_call": "tool_call_started",
"tool_result": "tool_call_completed",
"stop": "turn_completed",
"subagent_stop": "subagent_completed",
"subagent_start": "subagent_started",
}

@property
def version(self) -> str:
return "020"

@property
def description(self) -> str:
return "Add event_id with unique partial index for envelope backfills; repair wrong-dialect event_type values"

def up(self, conn: sqlite3.Connection) -> None:
cursor = conn.execute("SELECT name FROM sqlite_master WHERE type='table' AND name='hook_events'")
if not cursor.fetchone():
return

columns = {row[1] for row in conn.execute("PRAGMA table_info(hook_events)").fetchall()}

if "event_type" in columns:
for wrong, canonical in self._WRONG_DIALECT_RENAMES.items():
conn.execute("UPDATE hook_events SET event_type = ? WHERE event_type = ?", (canonical, wrong))

if "event_id" not in columns:
conn.execute("ALTER TABLE hook_events ADD COLUMN event_id TEXT")

conn.execute("""
CREATE UNIQUE INDEX IF NOT EXISTS idx_hook_events_source_event_id
ON hook_events(source, event_id)
WHERE event_id IS NOT NULL
""")


class MigrationRunner:
"""Manages database migrations."""

Expand All @@ -591,6 +641,7 @@ def __init__(self, db_path: Path):
Migration014AddBehavioralPatternHistory(),
Migration015AbstractEventTypeValues(),
Migration016AddRetiredReasonToMemories(),
Migration020EnvelopeIdempotency(),
]

@contextmanager
Expand Down
16 changes: 14 additions & 2 deletions src/slopometry/core/models/protocol/events.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,13 @@


class AbstractEventSource(StrEnum):
"""Identity of the agent harness or collector that produced the event."""
"""Vocabulary for the harnesses slopometry ships adapters for.

This enum is NOT a gate: `AbstractHookEvent.source` is an open string so
third-party collectors can ingest events under their own source name
(see `slopometry ingest`). Membership here only means "has a wire-format
adapter registered in `protocol/adapters`".
"""

CLAUDE_CODE = "claude_code"
OPENCODE = "opencode"
Expand Down Expand Up @@ -86,13 +92,19 @@ class AbstractHookEvent(BaseModel):
default=None,
description="Database autoincrement id; None for in-memory events not yet persisted",
)
event_id: str | None = Field(
default=None,
description="Source-unique id from envelope ingestion; (source, event_id) is unique so backfills are idempotent. Live harness paths have no source-side event id.",
)
session_id: str
parent_session_id: str | None = Field(
default=None,
description="Parent session ID for subagent/child sessions; None for top-level",
)
event_type: AbstractEventType
source: AbstractEventSource
source: str = Field(
description="Source name of the producing agent; built-ins in AbstractEventSource, open for third-party collectors",
)
timestamp: datetime = Field(default_factory=datetime.now)
tool_call: ToolCallPayload | None = None
metadata: dict[str, Any] = Field(
Expand Down
97 changes: 97 additions & 0 deletions src/slopometry/core/protocol/ingest.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
"""Idempotent ingestion of `EventEnvelope` batches into slopometry storage.

This is the third-party counterpart to `protocol/dispatch.py`: dispatch
runs a wire payload through a registered adapter (live harness paths), while
ingest accepts already-canonical envelopes from collectors without adapters.
Session context (git state, project detection) is deliberately not captured:
the repository the ingester runs in is unrelated to the traced session.
"""

import sqlite3
from pathlib import Path

from pydantic import BaseModel

from slopometry.core.database import EventDatabase
from slopometry.core.models.protocol.events import AbstractHookEvent, ToolCallPayload
from slopometry.core.protocol.schema import EventEnvelope


class IngestReport(BaseModel):
"""Outcome of ingesting a batch of envelopes."""

inserted: int = 0
skipped_duplicates: int = 0


def envelope_to_event(envelope: EventEnvelope, working_directory: str, sequence_number: int) -> AbstractHookEvent:
"""Expand a validated envelope into the canonical stored event model."""
tool_call = (
ToolCallPayload(
tool_name=envelope.tool_name,
input=envelope.raw,
duration_ms=envelope.duration_ms,
exit_code=envelope.exit_code,
error_message=envelope.error_message,
)
if envelope.tool_name
else None
)
return AbstractHookEvent(
session_id=envelope.session_id,
parent_session_id=envelope.parent_session_id,
event_type=envelope.kind,
source=envelope.source,
timestamp=envelope.occurred_at,
tool_call=tool_call,
metadata=envelope.raw,
working_directory=working_directory,
sequence_number=sequence_number,
event_id=envelope.event_id,
)


def ingest_envelopes(
envelopes: list[EventEnvelope],
db: EventDatabase | None = None,
working_directory: str | None = None,
) -> IngestReport:
"""Store envelopes, skipping duplicates identified by (source, event_id).

Sequence numbers come from the envelope when provided; otherwise a
per-session counter continuing after the highest sequence number already
used (stored or earlier in the batch) is assigned. The live SessionManager
state files are not touched: backfills derive ordering from storage.
"""
db = db or EventDatabase()
working_directory = working_directory or str(Path.cwd())
report = IngestReport()
next_sequence: dict[str, int] = {}

for envelope in envelopes:
if db.has_event(envelope.source, envelope.event_id):
report.skipped_duplicates += 1
continue

if envelope.seq is not None:
sequence_number = envelope.seq
next_sequence[envelope.session_id] = max(
envelope.seq + 1,
next_sequence.get(envelope.session_id, 1 + db.get_max_sequence_number(envelope.session_id)),
)
else:
next_sequence.setdefault(envelope.session_id, 1 + db.get_max_sequence_number(envelope.session_id))
sequence_number = next_sequence[envelope.session_id]
next_sequence[envelope.session_id] = sequence_number + 1

try:
event = envelope_to_event(envelope, working_directory, sequence_number)
event.id = db.save_event(event)
except sqlite3.IntegrityError:
# Concurrent ingest of the same (source, event_id); the unique
# partial index makes this safe to treat as a duplicate.
report.skipped_duplicates += 1
continue
report.inserted += 1

return report
Loading
Loading