From b5ef512fcdd9c659a2d6aca0918941c9239777bc Mon Sep 17 00:00:00 2001 From: Will Washburn Date: Thu, 24 Sep 2026 16:17:52 -0700 Subject: [PATCH 1/2] feat(sdk): change feed carries every kind and verbatim columns The change feed reports all twelve evidence tables the export journal captures -- adding history, presence, commit_link, trajectory, source_observation and observation_evidence -- and every upsert carries the row as stored (StoredRow: every column but revision, in table order, raw SQLite values, read from the live table) beside the typed row. Every change carries the journal's record key, tombstones included, and a row from an unknown source is carried with source None and source_name set instead of failing the drain. A database the six-kind feed reached stamps the new tables once on open, above its head. The session re-stamp triggers move to their own names and skip a presence's own stamp; the history FTS update trigger fires only on the columns it indexes. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 17 + crates/ai-hist/README.md | 2 +- crates/ai-hist/public-api.txt | 23 +- crates/ai-hist/src/change_feed.rs | 1107 +++++++++++++++++++++---- crates/ai-hist/src/export/schema.rs | 15 +- crates/ai-hist/src/lib.rs | 2 +- crates/ai-hist/src/source_evidence.rs | 2 + crates/ai-hist/src/store.rs | 12 +- crates/ai-hist/tests/change_feed.rs | 432 ++++++++-- docs/sourcing-sdk.md | 59 +- 10 files changed, 1399 insertions(+), 272 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 87958742..df0b1722 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -184,6 +184,23 @@ Notable changes to the native `ai-hist` CLI are documented here. ### Rust API +- The change feed reports every evidence table and each row exactly as + stored. `ChangeKind` gains `History`, `Presence`, `CommitLink`, + `Trajectory`, `SourceObservation` and `ObservationEvidence` (`ALL` lists + twelve kinds); their rows are stamped and tombstoned by the same triggers, + and a database the six-kind feed reached stamps them once on open, above its + head, so a cursor bound to every kind resumes into all of them. + `EvidenceRow::History(HistoryEntry)` is the typed prompt row and + `EvidenceRow::Untyped` marks a kind with none. `Change` gains `columns: + Option` — every column but `revision`, in table order, values as + SQLite holds them, read from the live table so a new column is carried + without a code change; the export journal's payload is built from the same + column list — and `key: Vec`, the record's identity as + the journal keys it (`["history", source, timestamp_ms, prompt]`), on + upserts and tombstones alike. `Change::source` is `Option` and + `Change::source_name` holds the stored name, so a row from a source this + build does not know is carried instead of failing the drain. The `history` + FTS update trigger fires only on the columns it indexes. - `SessionStore::discover(DiscoveryOptions)` is the shallow catalog sweep: every local provider's sessions from metadata, as `Shallow` rows, hydrating none and taking no `SyncRunLock`. `discover`, `sync` and `hydrate` take an diff --git a/crates/ai-hist/README.md b/crates/ai-hist/README.md index 5a635703..d0641591 100644 --- a/crates/ai-hist/README.md +++ b/crates/ai-hist/README.md @@ -22,7 +22,7 @@ fn main() -> Result<(), ai_hist::Error> { let query = ai_hist::ChangeQuery::default().consumer("my-ingest"); let mut changes = store.changes_since(ai_hist::Watermark::CONSUMER, query)?; for change in changes.by_ref() { - let _change = change?; // `Upsert(EvidenceRow)` or `Delete`, keyed by kind + record_key + let _change = change?; // `Upsert` or `Delete` of `change.key`; `columns` is the stored row } changes.commit()?; // the cursor moves only here Ok(()) diff --git a/crates/ai-hist/public-api.txt b/crates/ai-hist/public-api.txt index e8e5f026..c0f6312f 100644 --- a/crates/ai-hist/public-api.txt +++ b/crates/ai-hist/public-api.txt @@ -35,12 +35,18 @@ pub ai_hist::Capability::Full pub ai_hist::Capability::Partial pub ai_hist::Capability::ShallowOnly #[non_exhaustive] pub enum ai_hist::ChangeKind +pub ai_hist::ChangeKind::CommitLink pub ai_hist::ChangeKind::FileEdit +pub ai_hist::ChangeKind::History +pub ai_hist::ChangeKind::ObservationEvidence +pub ai_hist::ChangeKind::Presence pub ai_hist::ChangeKind::Relationship pub ai_hist::ChangeKind::Session pub ai_hist::ChangeKind::SessionEvent pub ai_hist::ChangeKind::SessionMarker +pub ai_hist::ChangeKind::SourceObservation pub ai_hist::ChangeKind::ToolCall +pub ai_hist::ChangeKind::Trajectory impl ai_hist::ChangeKind pub const ai_hist::ChangeKind::ALL: &'static [ai_hist::ChangeKind] pub fn ai_hist::ChangeKind::as_str(self) -> &'static str @@ -107,11 +113,13 @@ impl ai_hist::EvidenceKind pub fn ai_hist::EvidenceKind::as_str(self) -> &'static str #[non_exhaustive] pub enum ai_hist::EvidenceRow pub ai_hist::EvidenceRow::FileEdit(ai_hist::SessionFileEdit) +pub ai_hist::EvidenceRow::History(ai_hist::HistoryEntry) pub ai_hist::EvidenceRow::Relationship(ai_hist::SessionRelationship) pub ai_hist::EvidenceRow::Session(ai_hist::ShallowSession) pub ai_hist::EvidenceRow::SessionEvent(ai_hist::SessionEvent) pub ai_hist::EvidenceRow::SessionMarker(ai_hist::SessionMarker) pub ai_hist::EvidenceRow::ToolCall(ai_hist::SessionToolCall) +pub ai_hist::EvidenceRow::Untyped #[non_exhaustive] pub enum ai_hist::HydrateStatus pub ai_hist::HydrateStatus::CapabilityLimited pub ai_hist::HydrateStatus::Hydrated @@ -294,12 +302,15 @@ pub ai_hist::CatalogSession::workspace_roots: alloc::vec::Vec ai_hist::SessionRef #[non_exhaustive] pub struct ai_hist::Change +pub ai_hist::Change::columns: core::option::Option +pub ai_hist::Change::key: alloc::vec::Vec pub ai_hist::Change::kind: ai_hist::ChangeKind pub ai_hist::Change::op: ai_hist::ChangeOp pub ai_hist::Change::record_key: alloc::string::String pub ai_hist::Change::revision: u64 pub ai_hist::Change::session_id: alloc::string::String -pub ai_hist::Change::source: ai_hist::Source +pub ai_hist::Change::source: core::option::Option +pub ai_hist::Change::source_name: alloc::string::String #[non_exhaustive] pub struct ai_hist::ChangeQuery pub ai_hist::ChangeQuery::batch: usize pub ai_hist::ChangeQuery::consumer: core::option::Option @@ -705,6 +716,16 @@ pub ai_hist::StoreOptions::db_path: core::option::Option pub ai_hist::StoreOptions::home: core::option::Option pub ai_hist::StoreOptions::read_only: bool pub ai_hist::StoreOptions::roots: core::option::Option +pub struct ai_hist::StoredRow +impl ai_hist::StoredRow +pub fn ai_hist::StoredRow::get(&self, &str) -> core::option::Option<&serde_json::value::Value> +pub fn ai_hist::StoredRow::is_empty(&self) -> bool +pub fn ai_hist::StoredRow::iter(&self) -> impl core::iter::traits::iterator::Iterator + '_ +pub fn ai_hist::StoredRow::len(&self) -> usize +impl serde_core::ser::Serialize for ai_hist::StoredRow +pub fn ai_hist::StoredRow::serialize(&self, S) -> core::result::Result<::Ok, ::Error> +impl<'de> serde_core::de::Deserialize<'de> for ai_hist::StoredRow +pub fn ai_hist::StoredRow::deserialize>(D) -> core::result::Result::Error> #[non_exhaustive] pub struct ai_hist::SyncOptions pub ai_hist::SyncOptions::force: bool pub ai_hist::SyncOptions::lock_timeout_ms: u64 diff --git a/crates/ai-hist/src/change_feed.rs b/crates/ai-hist/src/change_feed.rs index 9bd375d1..6c3cc37c 100644 --- a/crates/ai-hist/src/change_feed.rs +++ b/crates/ai-hist/src/change_feed.rs @@ -7,10 +7,18 @@ //! and a pull cursor over it: //! //! - Every row of `sessions`, `session_events`, `tool_calls`, `file_edits`, -//! `session_markers` and `session_relationships` carries a `revision`. A -//! trigger stamps the current clock on every insert and every update, so a -//! re-parse that upserts a row it already holds re-stamps it: consumers must -//! treat a re-seen `record_key` as a replace, never as a duplicate. +//! `session_markers`, `session_relationships`, `history`, +//! `session_presences`, `session_commit_links`, `trajectories`, +//! `session_observations` and `observation_evidence` carries a `revision`. +//! A trigger stamps the current clock on every insert and every update, so +//! a re-parse that upserts a row it already holds re-stamps it: consumers +//! must treat a re-seen key as a replace, never as a duplicate. +//! - An upsert carries the row twice over: typed, where the kind has a typed +//! row, and as stored ([`StoredRow`]) -- every column but `revision`, read +//! from the live table, values as SQLite holds them -- so an embedder can +//! rebuild the row exactly without knowing the table's shape in advance. +//! Every change, a delete included, carries the record's identity as the +//! export journal keys it ([`Change::key`]). //! - A deleted row — a sidechain heal moving records onto the child, a session //! dropped from the catalog — leaves a tombstone in `evidence_tombstones` //! carrying its own revision, so a consumer learns about the removal in the @@ -39,16 +47,19 @@ use crate::relationship_graph::{map_relationship, SessionRelationship, RELATIONS use crate::session_store::{Error, SessionStore, Source}; use crate::store::{ ensure_columns, migration_applied, open_db, open_db_readonly, row_to_file_edit, - row_to_session_event, row_to_session_marker, row_to_tool_call, SessionEvent, SessionFileEdit, - SessionMarker, SessionToolCall, FILE_EDIT_COLUMNS, SESSION_EVENT_COLUMNS, + row_to_session_event, row_to_session_marker, row_to_tool_call, HistoryEntry, SessionEvent, + SessionFileEdit, SessionMarker, SessionToolCall, FILE_EDIT_COLUMNS, SESSION_EVENT_COLUMNS, SESSION_MARKER_COLUMNS, TOOL_CALL_COLUMNS, }; use crate::EvidenceKind; use anyhow::{Context, Result}; +use rusqlite::types::ValueRef; use rusqlite::{params, Connection, OptionalExtension}; use serde::{Deserialize, Serialize}; +use serde_json::Value; use std::collections::VecDeque; use std::path::{Path, PathBuf}; +use std::sync::Arc; /// The column every fed table carries. Named once so the delivery capture /// triggers can leave it out of their payloads. @@ -60,7 +71,10 @@ pub const MAX_CHANGE_BATCH: usize = 10_000; /// Page size when [`ChangeQuery::batch`] is zero. pub const DEFAULT_CHANGE_BATCH: usize = 1_000; -const MIGRATION: &str = "change_feed_v1"; +/// The feed schema's marker. `v2` is every kind the feed reports; a database +/// stamped by `v1` gains the kinds it lacked, backfilled above its head, in +/// the one pass the missing marker triggers. +const MIGRATION: &str = "change_feed_v2"; /// The column a named cursor records its kind set in. const KINDS_COLUMN: &str = "kinds"; @@ -149,19 +163,40 @@ impl Watermark { /// Which table a change is about. /// /// This is not [`EvidenceKind`]: that enum names the record kinds a source -/// adapter can supply, and the catalog row is not one of them, while the feed -/// has to report a session's catalog row changing. [`ChangeKind::evidence_kind`] -/// maps the overlap. +/// adapter can supply, and the feed also reports tables no adapter writes -- +/// the catalog row, where a session was seen, what a connector observed. +/// [`ChangeKind::evidence_kind`] maps the overlap. The wire names are the +/// export journal's kinds, so a record keeps one kind name whichever of the +/// two an embedder read it from. #[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] #[non_exhaustive] pub enum ChangeKind { + /// `sessions`: the catalog row. Session, + /// `session_events`. SessionEvent, + /// `tool_calls`. ToolCall, + /// `file_edits`. FileEdit, + /// `session_markers`. SessionMarker, + /// `session_relationships`. Relationship, + /// `history`: one prompt from a provider's prompt log. + History, + /// `session_presences`: where a session was seen, local or remote. + Presence, + /// `session_commit_links`: a commit a session is linked to. + CommitLink, + /// `trajectories`. + Trajectory, + /// `session_observations`: one connector's observation of a session. + SourceObservation, + /// `observation_evidence`: a record a connector supplied with an + /// observation. + ObservationEvidence, } impl ChangeKind { @@ -173,10 +208,16 @@ impl ChangeKind { ChangeKind::FileEdit, ChangeKind::SessionMarker, ChangeKind::Relationship, + ChangeKind::History, + ChangeKind::Presence, + ChangeKind::CommitLink, + ChangeKind::Trajectory, + ChangeKind::SourceObservation, + ChangeKind::ObservationEvidence, ]; - /// The wire name, identical to the serde representation and to the `kind` - /// stored on a tombstone. + /// The wire name, identical to the serde representation, to the `kind` + /// stored on a tombstone and to the first element of [`Change::key`]. pub fn as_str(self) -> &'static str { match self { Self::Session => "session", @@ -185,83 +226,398 @@ impl ChangeKind { Self::FileEdit => "file_edit", Self::SessionMarker => "session_marker", Self::Relationship => "relationship", + Self::History => "history", + Self::Presence => "presence", + Self::CommitLink => "commit_link", + Self::Trajectory => "trajectory", + Self::SourceObservation => "source_observation", + Self::ObservationEvidence => "observation_evidence", } } - /// The source-evidence kind this change carries, or `None` for the - /// catalog row, which no adapter supplies. + /// The source-evidence kind this change carries, or `None` for a table + /// no adapter supplies. pub fn evidence_kind(self) -> Option { match self { - Self::Session => None, Self::SessionEvent => Some(EvidenceKind::SessionEvent), Self::ToolCall => Some(EvidenceKind::ToolCall), Self::FileEdit => Some(EvidenceKind::FileEdit), Self::SessionMarker => Some(EvidenceKind::SessionMarker), Self::Relationship => Some(EvidenceKind::Relationship), + Self::History => Some(EvidenceKind::History), + Self::CommitLink => Some(EvidenceKind::CommitLink), + Self::Session + | Self::Presence + | Self::Trajectory + | Self::SourceObservation + | Self::ObservationEvidence => None, } } fn table(self) -> FedTable { + let table = |name, session, key, record| FedTable { + name, + source: "source", + session, + optional_session: false, + key, + record, + }; match self { - Self::Session => FedTable { - name: "sessions", - session: "session_id", - key: "session_id", - }, - Self::SessionEvent => FedTable { - name: "session_events", - session: "session_id", - key: "event_uid", - }, - Self::ToolCall => FedTable { - name: "tool_calls", - session: "session_id", - key: "tool_use_id", - }, - Self::FileEdit => FedTable { - name: "file_edits", - session: "session_id", - key: "tool_use_id", - }, - Self::SessionMarker => FedTable { - name: "session_markers", - session: "session_id", - key: "marker_uid", + Self::Session => table( + "sessions", + "session_id", + &["source", "session_id"], + &["session_id"], + ), + Self::SessionEvent => table( + "session_events", + "session_id", + &["source", "session_id", "event_uid"], + &["event_uid"], + ), + Self::ToolCall => table( + "tool_calls", + "session_id", + &["source", "session_id", "tool_use_id"], + &["tool_use_id"], + ), + Self::FileEdit => table( + "file_edits", + "session_id", + &["source", "session_id", "tool_use_id"], + &["tool_use_id"], + ), + Self::SessionMarker => table( + "session_markers", + "session_id", + &["source", "session_id", "marker_uid"], + &["marker_uid"], + ), + Self::Relationship => table( + "session_relationships", + "parent_session_id", + &["source", "parent_session_id", "relationship_uid"], + &["relationship_uid"], + ), + // A prompt may name no session; its identity is the prompt log's + // own `UNIQUE(source, timestamp_ms, prompt)`. + Self::History => FedTable { + optional_session: true, + ..table( + "history", + "session_id", + &["source", "timestamp_ms", "prompt"], + &["timestamp_ms", "prompt"], + ) }, - Self::Relationship => FedTable { - name: "session_relationships", - session: "parent_session_id", - key: "relationship_uid", + Self::Presence => table( + "session_presences", + "session_id", + &["source", "session_id", "location"], + &["location"], + ), + Self::CommitLink => table( + "session_commit_links", + "session_id", + &["source", "session_id", "commit_sha", "match_method"], + &["commit_sha", "match_method"], + ), + // A trajectory row carries no source column: every one is the + // `trajectory` source, and its id is its session. + Self::Trajectory => FedTable { + source: "'trajectory'", + ..table("trajectories", "id", &["id"], &["id"]) }, + Self::SourceObservation => table( + "session_observations", + "session_id", + &[ + "source", + "session_id", + "location", + "connector_id", + "connector_instance", + ], + &["location", "connector_id", "connector_instance"], + ), + Self::ObservationEvidence => table( + "observation_evidence", + "session_id", + &[ + "source", + "session_id", + "location", + "connector_id", + "connector_instance", + "evidence_uid", + ], + &[ + "location", + "connector_id", + "connector_instance", + "evidence_uid", + ], + ), } } - fn columns(self) -> &'static str { + /// The typed row's column list, for the kinds that have one. + fn columns(self) -> Option<&'static str> { match self { - Self::Session => SESSION_COLUMNS, - Self::SessionEvent => SESSION_EVENT_COLUMNS, - Self::ToolCall => TOOL_CALL_COLUMNS, - Self::FileEdit => FILE_EDIT_COLUMNS, - Self::SessionMarker => SESSION_MARKER_COLUMNS, - Self::Relationship => RELATIONSHIP_COLUMNS, + Self::Session => Some(SESSION_COLUMNS), + Self::SessionEvent => Some(SESSION_EVENT_COLUMNS), + Self::ToolCall => Some(TOOL_CALL_COLUMNS), + Self::FileEdit => Some(FILE_EDIT_COLUMNS), + Self::SessionMarker => Some(SESSION_MARKER_COLUMNS), + Self::Relationship => Some(RELATIONSHIP_COLUMNS), + Self::History => Some(HISTORY_COLUMNS), + Self::Presence + | Self::CommitLink + | Self::Trajectory + | Self::SourceObservation + | Self::ObservationEvidence => None, } } } -/// One stamped table: where its source, session and record key live. +/// The typed [`HistoryEntry`] columns, in the order [`row_to_history`] reads. +const HISTORY_COLUMNS: &str = "id, source, session_id, project, prompt, prompt_hash, timestamp_ms"; + +fn row_to_history(row: &rusqlite::Row<'_>) -> rusqlite::Result { + Ok(HistoryEntry { + id: row.get(0)?, + source: row.get(1)?, + session_id: row.get(2)?, + project: row.get(3)?, + prompt: row.get(4)?, + prompt_hash: row.get(5)?, + timestamp_ms: row.get(6)?, + }) +} + +/// One stamped table: where a record's source, session and identity live. +/// +/// A record is named two ways. `key` is its identity as the export journal +/// keys it -- the columns of the table's own uniqueness constraint -- and is +/// what [`Change::key`] carries. `record` is the part of that identity the +/// tombstone's `record_key` holds beside the source and session: one column's +/// text, or a JSON array of several, so a tombstone keeps every key column's +/// stored type and [`Change::key`] can be rebuilt from it. struct FedTable { name: &'static str, + /// A column, or a quoted literal for a table that stores no source. + source: &'static str, session: &'static str, - key: &'static str, + /// Whether `session` may be NULL; a tombstone stores it as `''`. + optional_session: bool, + key: &'static [&'static str], + record: &'static [&'static str], } impl FedTable { fn revision_index(&self) -> String { format!("idx_{}_revision", self.name) } + + /// The source as SQL over `row` (`NEW`, `OLD` or the table name). + fn source_sql(&self, row: &str) -> String { + if self.source.starts_with('\'') { + self.source.to_string() + } else { + format!("{row}.{}", self.source) + } + } + + /// The session as SQL over `row`, as a tombstone stores it. + fn session_sql(&self, row: &str) -> String { + if self.optional_session { + format!("COALESCE({row}.{}, '')", self.session) + } else { + format!("{row}.{}", self.session) + } + } + + /// The `record_key` as SQL over `row`. + fn record_key_sql(&self, row: &str) -> String { + match self.record { + [column] => format!("{row}.{column}"), + columns => format!( + "json_array({})", + columns + .iter() + .map(|column| format!("{row}.{column}")) + .collect::>() + .join(", ") + ), + } + } + + /// True when a write moved the row to another identity. + fn identity_changed_sql(&self) -> String { + let mut changed = Vec::new(); + if !self.source.starts_with('\'') { + changed.push(format!("OLD.{0} IS NOT NEW.{0}", self.source)); + } + changed.push(format!("OLD.{0} IS NOT NEW.{0}", self.session)); + changed.push(format!( + "{} IS NOT {}", + self.record_key_sql("OLD"), + self.record_key_sql("NEW") + )); + changed.join(" OR ") + } + + /// [`Change::key`] from a tombstone's identity columns. + fn key_from_identity( + &self, + kind: ChangeKind, + source: &str, + session: &str, + record_key: &str, + ) -> Result> { + let record: Vec = match self.record { + [_] => vec![Value::String(record_key.to_string())], + _ => serde_json::from_str(record_key).with_context(|| { + format!( + "change feed: {} tombstone record key {record_key:?} is not a JSON array", + kind.as_str() + ) + })?, + }; + let mut key = Vec::with_capacity(self.key.len() + 1); + key.push(Value::String(kind.as_str().to_string())); + for column in self.key { + if let Some(index) = self.record.iter().position(|part| part == column) { + key.push(record.get(index).cloned().unwrap_or(Value::Null)); + } else if *column == self.source { + key.push(Value::String(source.to_string())); + } else { + key.push(Value::String(session.to_string())); + } + } + Ok(key) + } + + /// [`Change::key`] from a stored row. + fn key_from_row(&self, kind: ChangeKind, row: &StoredRow) -> Vec { + let mut key = Vec::with_capacity(self.key.len() + 1); + key.push(Value::String(kind.as_str().to_string())); + key.extend( + self.key + .iter() + .map(|column| row.get(column).cloned().unwrap_or(Value::Null)), + ); + key + } +} + +/// The columns of `table` a stored row carries: every column but the feed's +/// own [`REVISION_COLUMN`], in table order, read from the live schema so a +/// column a migration adds is carried without a code change. +/// +/// The export journal's capture payload is built from the same list. +pub(crate) fn stored_columns(conn: &Connection, table: &str) -> Result> { + let columns = conn + .prepare_cached("SELECT name FROM pragma_table_info(?1) ORDER BY cid")? + .query_map([table], |row| row.get::<_, String>(0))? + .collect::>>()?; + Ok(columns + .into_iter() + .filter(|column| column != REVISION_COLUMN) + .collect()) +} + +/// A record exactly as its table stores it: every column but `revision`, in +/// table order, each value as SQLite holds it. +/// +/// Text stays text even when it holds JSON, an integer stays an integer, a +/// real stays a real and NULL is `null`; nothing is parsed, defaulted or +/// derived. A column the table gains is carried as soon as it exists. A BLOB, +/// which no stamped table declares, is carried as an array of its bytes. +/// +/// Serializes as one JSON object whose keys are in table order. +#[derive(Debug, Clone, Default, PartialEq)] +pub struct StoredRow { + columns: Vec<(Arc, Value)>, +} + +impl StoredRow { + /// The stored value of `column`, or `None` when the table has no such + /// column. + pub fn get(&self, column: &str) -> Option<&Value> { + self.columns + .iter() + .find(|(name, _)| &**name == column) + .map(|(_, value)| value) + } + + /// Every column and its stored value, in table order. + pub fn iter(&self) -> impl Iterator + '_ { + self.columns.iter().map(|(name, value)| (&**name, value)) + } + + /// How many columns the row carries. + pub fn len(&self) -> usize { + self.columns.len() + } + + /// Whether the row carries no column. + pub fn is_empty(&self) -> bool { + self.columns.is_empty() + } +} + +impl Serialize for StoredRow { + fn serialize( + &self, + serializer: S, + ) -> std::result::Result { + use serde::ser::SerializeMap; + let mut map = serializer.serialize_map(Some(self.columns.len()))?; + for (name, value) in &self.columns { + map.serialize_entry(&**name, value)?; + } + map.end() + } +} + +impl<'de> Deserialize<'de> for StoredRow { + fn deserialize>( + deserializer: D, + ) -> std::result::Result { + struct Columns; + impl<'de> serde::de::Visitor<'de> for Columns { + type Value = StoredRow; + fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("an object of column names to stored values") + } + fn visit_map>( + self, + mut map: A, + ) -> std::result::Result { + let mut columns = Vec::with_capacity(map.size_hint().unwrap_or(0)); + while let Some((name, value)) = map.next_entry::()? { + columns.push((Arc::from(name), value)); + } + Ok(StoredRow { columns }) + } + } + deserializer.deserialize_map(Columns) + } +} + +fn stored_value(value: ValueRef<'_>) -> Value { + match value { + ValueRef::Null => Value::Null, + ValueRef::Integer(integer) => Value::from(integer), + ValueRef::Real(real) => Value::from(real), + ValueRef::Text(text) => Value::String(String::from_utf8_lossy(text).into_owned()), + ValueRef::Blob(bytes) => Value::from(bytes.to_vec()), + } } -/// The row an upsert carries, typed per kind so no second read is needed. +/// The typed row an upsert carries, per kind, so no second read is needed. /// /// The variants differ in size because the rows do; an event row is several /// times a tool call. Boxing the large ones would put an allocation between @@ -278,6 +634,10 @@ pub enum EvidenceRow { FileEdit(SessionFileEdit), SessionMarker(SessionMarker), Relationship(SessionRelationship), + History(HistoryEntry), + /// A kind with no typed row: presences, commit links, trajectories and + /// connector observations. [`Change::columns`] is the row. + Untyped, } /// What happened to the record. @@ -298,15 +658,35 @@ pub enum ChangeOp { #[non_exhaustive] pub struct Change { pub kind: ChangeKind, - pub source: Source, - /// For a relationship, the parent session. + /// The record's source, or `None` when the stored name is one this build + /// does not know -- a row written by a newer release. The drain carries + /// such a row rather than failing on it; `source_name` names it. + pub source: Option, + /// The source exactly as stored. + #[serde(default)] + pub source_name: String, + /// For a relationship, the parent session; for a trajectory, its id; for + /// a prompt that names no session, empty. pub session_id: String, - /// The record's provider-native identity within its session and kind: - /// `event_uid`, `tool_use_id`, `marker_uid`, `relationship_uid`, or the - /// session id itself for a catalog row. + /// The record's identity within its source, session and kind: the one + /// identity column's text (`event_uid`, `tool_use_id`, `marker_uid`, + /// `relationship_uid`, `location`, a trajectory's id, or the session id + /// itself for a catalog row), or for a kind whose identity spans several + /// columns, those columns' stored values as a JSON array. pub record_key: String, + /// The record's identity as the export journal keys it: the kind's wire + /// name, then the stored value of each column of the table's uniqueness + /// constraint, in order -- `["history", source, timestamp_ms, prompt]`, + /// `["trajectory", id]`. An upsert and a delete of one record carry the + /// same key. + #[serde(default)] + pub key: Vec, pub revision: u64, pub op: ChangeOp, + /// The row as stored, for an upsert; `None` for a delete. See + /// [`StoredRow`]. + #[serde(default)] + pub columns: Option, } /// How to read the feed. @@ -874,28 +1254,33 @@ fn commit_cursor( Err(kinds_mismatch(name, &stored, &offered)) } -fn parse_source(name: &str) -> Result { - match name { - "claude" => Ok(Source::Claude), - "codex" => Ok(Source::Codex), - "cursor" => Ok(Source::Cursor), - "grok" => Ok(Source::Grok), - "relay" => Ok(Source::Relay), - "trajectory" => Ok(Source::Trajectory), - "opencode" => Ok(Source::OpenCode), - other => anyhow::bail!("change feed: unknown source {other:?}"), - } -} - /// The page query for one kind: an indexed range read, oldest first. -fn upsert_sql(kind: ChangeKind) -> String { +/// +/// Four groups of columns, in order: the typed row's columns, for the kinds +/// that have one, read by position from zero; the tombstone identity -- +/// source, session, record key -- computed by the same SQL the triggers use, +/// so an upsert and a delete of one record name it identically; the stored +/// columns; and the revision. +fn upsert_sql(kind: ChangeKind, stored: &[String]) -> String { let table = kind.table(); + let name = table.name; + let mut select = Vec::new(); + if let Some(columns) = kind.columns() { + select.push(columns.to_string()); + } + select.push(table.source_sql(name)); + select.push(table.session_sql(name)); + select.push(table.record_key_sql(name)); + select.extend( + stored + .iter() + .map(|column| format!("{name}.\"{}\"", column.replace('"', "\"\""))), + ); format!( - "SELECT {columns}, {REVISION_COLUMN} FROM {name} \ + "SELECT {select}, {name}.{REVISION_COLUMN} FROM {name} \ WHERE {REVISION_COLUMN} > ?1 AND {REVISION_COLUMN} <= ?2 \ ORDER BY {REVISION_COLUMN} ASC LIMIT ?3", - columns = kind.columns(), - name = table.name, + select = select.join(", "), ) } @@ -906,77 +1291,63 @@ fn read_upserts( hi: u64, batch: usize, ) -> Result> { - let mut statement = conn.prepare_cached(&upsert_sql(kind))?; + let table = kind.table(); + let stored = stored_columns(conn, table.name)?; + let names: Vec> = stored + .iter() + .map(|column| Arc::from(column.as_str())) + .collect(); + let mut statement = conn.prepare_cached(&upsert_sql(kind, &stored))?; + // Everything before the identity is the typed row. + let identity = statement.column_count() - stored.len() - 4; let rows = statement.query_map(params![lo as i64, hi as i64, batch as i64], |row| { - let revision: i64 = row.get(REVISION_COLUMN)?; - let (source, session_id, record_key, evidence) = match kind { - ChangeKind::Session => { - let session = row_to_session(row)?; - ( - session.source.clone(), - session.session_id.clone(), - session.session_id.clone(), - EvidenceRow::Session(session), - ) - } - ChangeKind::SessionEvent => { - let event = row_to_session_event(row)?; - ( - event.source.clone(), - event.session_id.clone(), - event.event_uid.clone(), - EvidenceRow::SessionEvent(event), - ) - } - ChangeKind::ToolCall => { - let call = row_to_tool_call(row)?; - ( - call.source.clone(), - call.session_id.clone(), - call.tool_use_id.clone(), - EvidenceRow::ToolCall(call), - ) - } - ChangeKind::FileEdit => { - let edit = row_to_file_edit(row)?; - ( - edit.source.clone(), - edit.session_id.clone(), - edit.tool_use_id.clone(), - EvidenceRow::FileEdit(edit), - ) - } - ChangeKind::SessionMarker => { - let marker = row_to_session_marker(row)?; - ( - marker.source.clone(), - marker.session_id.clone(), - marker.marker_uid.clone(), - EvidenceRow::SessionMarker(marker), - ) - } - ChangeKind::Relationship => { - let relationship = map_relationship(row)?; - ( - relationship.source.clone(), - relationship.parent_session_id.clone(), - relationship.relationship_uid.clone(), - EvidenceRow::Relationship(relationship), - ) - } + let evidence = match kind { + ChangeKind::Session => EvidenceRow::Session(row_to_session(row)?), + ChangeKind::SessionEvent => EvidenceRow::SessionEvent(row_to_session_event(row)?), + ChangeKind::ToolCall => EvidenceRow::ToolCall(row_to_tool_call(row)?), + ChangeKind::FileEdit => EvidenceRow::FileEdit(row_to_file_edit(row)?), + ChangeKind::SessionMarker => EvidenceRow::SessionMarker(row_to_session_marker(row)?), + ChangeKind::Relationship => EvidenceRow::Relationship(map_relationship(row)?), + ChangeKind::History => EvidenceRow::History(row_to_history(row)?), + ChangeKind::Presence + | ChangeKind::CommitLink + | ChangeKind::Trajectory + | ChangeKind::SourceObservation + | ChangeKind::ObservationEvidence => EvidenceRow::Untyped, }; - Ok((source, session_id, record_key, revision, evidence)) + let source: String = row.get(identity)?; + let session_id: String = row.get(identity + 1)?; + let record_key: String = row.get(identity + 2)?; + let mut columns = Vec::with_capacity(names.len()); + for (offset, name) in names.iter().enumerate() { + columns.push(( + Arc::clone(name), + stored_value(row.get_ref(identity + 3 + offset)?), + )); + } + let revision: i64 = row.get(identity + 3 + names.len())?; + Ok(( + source, + session_id, + record_key, + revision, + evidence, + StoredRow { columns }, + )) })?; let mut changes = Vec::new(); for row in rows { - let (source, session_id, record_key, revision, evidence) = row?; + let (source, session_id, record_key, revision, evidence, columns) = row?; changes.push(Change { kind, - source: parse_source(&source)?, + source: Source::parse(&source), + source_name: source, session_id, record_key, + key: table.key_from_row(kind, &columns), revision: revision.max(0) as u64, op: ChangeOp::Upsert(evidence), + columns: Some(columns), }); } Ok(changes) @@ -1004,6 +1375,7 @@ fn read_tombstones( let mut changes = Vec::new(); let mut statement = conn.prepare_cached(&tombstone_sql())?; for kind in kinds { + let table = kind.table(); let rows = statement.query_map( params![kind.as_str(), lo as i64, hi as i64, batch as i64], |row| { @@ -1017,13 +1389,17 @@ fn read_tombstones( )?; for row in rows { let (source, session_id, record_key, revision) = row?; + let key = table.key_from_identity(*kind, &source, &session_id, &record_key)?; changes.push(Change { kind: *kind, - source: parse_source(&source)?, + source: Source::parse(&source), + source_name: source, session_id, record_key, + key, revision: revision.max(0) as u64, op: ChangeOp::Delete, + columns: None, }); } } @@ -1041,6 +1417,15 @@ fn read_tombstones( /// not written — a subagent cleanup that drops the local presence and keeps /// the remote one changes nothing else. const PRESENCE_TRIGGERS: &[&str] = &[ + "change_feed_session_locations_insert", + "change_feed_session_locations_update", + "change_feed_session_locations_delete", +]; + +/// The names the session re-stamp triggers had before a presence was a kind +/// of its own. They are the presence kind's stamping triggers' names now, so +/// the migration to [`MIGRATION`] drops them before either set is created. +const RETIRED_PRESENCE_TRIGGERS: &[&str] = &[ "change_feed_session_presences_insert", "change_feed_session_presences_update", "change_feed_session_presences_delete", @@ -1162,6 +1547,15 @@ pub(crate) fn init_schema(conn: &Connection) -> Result<()> { &[(KINDS_COLUMN, "TEXT NOT NULL DEFAULT '*'")], )?; let backfill = !migration_applied(conn, MIGRATION)?; + if backfill { + // A database from before presences were fed: its session re-stamp + // triggers hold the names the presence kind's own triggers take + // below, and fire on the backfill's stamp. Both sets are created + // afresh after the backfill. + for trigger in RETIRED_PRESENCE_TRIGGERS { + conn.execute_batch(&format!("DROP TRIGGER IF EXISTS {trigger};"))?; + } + } for kind in ChangeKind::ALL { let table = kind.table(); ensure_columns( @@ -1170,10 +1564,13 @@ pub(crate) fn init_schema(conn: &Connection) -> Result<()> { &[(REVISION_COLUMN, "INTEGER NOT NULL DEFAULT 0")], )?; if backfill { - // Rows written before the feed existed are stamped once, in + // Rows written before their table was fed are stamped once, in // rowid order, each above everything stamped before it, so a - // replay from START reports the whole store. The triggers do not - // exist yet, so this UPDATE stamps exactly what it names. + // replay from START reports the whole store and a cursor that + // predates the kind still receives every one of its rows. The + // table's triggers do not exist yet, so this UPDATE stamps + // exactly what it names; a table fed already holds no unstamped + // row and is left alone. conn.execute( &format!( "UPDATE {name} SET {REVISION_COLUMN} = rowid + \ @@ -1204,14 +1601,19 @@ pub(crate) fn init_schema(conn: &Connection) -> Result<()> { let [insert, update, delete] = trigger_names(*kind); let name = table.name; let kind = kind.as_str(); - let session = table.session; - let key = table.key; + let new_source = table.source_sql("NEW"); + let old_source = table.source_sql("OLD"); + let new_session = table.session_sql("NEW"); + let old_session = table.session_sql("OLD"); + let new_record = table.record_key_sql("NEW"); + let old_record = table.record_key_sql("OLD"); + let moved = table.identity_changed_sql(); // The update trigger's own stamp changes `revision`, and only that, // so `NEW.revision = OLD.revision` is what stops it re-firing under // `recursive_triggers` — and what makes an external write that leaves - // the stamp alone (every upsert in this crate) take a new one. A key - // change is a delete of the old key and an upsert of the new, each at - // its own revision. + // the stamp alone (every upsert in this crate) take a new one. A + // change of identity is a delete of the old one and an upsert of the + // new, each at its own revision. conn.execute_batch(&format!( "CREATE TRIGGER IF NOT EXISTS {insert} AFTER INSERT ON {name} BEGIN UPDATE observation_clock SET version = version + 1 WHERE singleton = 1; @@ -1219,41 +1621,41 @@ pub(crate) fn init_schema(conn: &Connection) -> Result<()> { (SELECT version FROM observation_clock WHERE singleton = 1) \ WHERE rowid = NEW.rowid; DELETE FROM evidence_tombstones WHERE kind = '{kind}' \ - AND source = NEW.source AND session_id = NEW.{session} \ - AND record_key = NEW.{key}; + AND source = {new_source} AND session_id = {new_session} \ + AND record_key = {new_record}; END; CREATE TRIGGER IF NOT EXISTS {update} AFTER UPDATE ON {name} WHEN NEW.{REVISION_COLUMN} = OLD.{REVISION_COLUMN} BEGIN UPDATE observation_clock SET version = version + 1 WHERE singleton = 1 \ - AND (OLD.source IS NOT NEW.source OR OLD.{session} IS NOT NEW.{session} \ - OR OLD.{key} IS NOT NEW.{key}); + AND ({moved}); INSERT OR REPLACE INTO evidence_tombstones \ (kind, source, session_id, record_key, revision) \ - SELECT '{kind}', OLD.source, OLD.{session}, OLD.{key}, \ + SELECT '{kind}', {old_source}, {old_session}, {old_record}, \ (SELECT version FROM observation_clock WHERE singleton = 1) \ - WHERE OLD.source IS NOT NEW.source OR OLD.{session} IS NOT NEW.{session} \ - OR OLD.{key} IS NOT NEW.{key}; + WHERE {moved}; UPDATE observation_clock SET version = version + 1 WHERE singleton = 1; UPDATE {name} SET {REVISION_COLUMN} = \ (SELECT version FROM observation_clock WHERE singleton = 1) \ WHERE rowid = NEW.rowid; DELETE FROM evidence_tombstones WHERE kind = '{kind}' \ - AND source = NEW.source AND session_id = NEW.{session} \ - AND record_key = NEW.{key}; + AND source = {new_source} AND session_id = {new_session} \ + AND record_key = {new_record}; END; CREATE TRIGGER IF NOT EXISTS {delete} AFTER DELETE ON {name} BEGIN UPDATE observation_clock SET version = version + 1 WHERE singleton = 1; INSERT OR REPLACE INTO evidence_tombstones \ (kind, source, session_id, record_key, revision) \ - VALUES ('{kind}', OLD.source, OLD.{session}, OLD.{key}, \ + VALUES ('{kind}', {old_source}, {old_session}, {old_record}, \ (SELECT version FROM observation_clock WHERE singleton = 1)); END;" ))?; } // A direct write of `revision` does not re-fire the sessions update // trigger (its guard is `NEW.revision = OLD.revision`), so this is one - // stamp, not two. A key change is a presence leaving one session and - // arriving at another; both rows are stamped. + // stamp, not two. The update trigger carries the same guard, so a + // presence's own stamp does not re-stamp its session a second time. A + // key change is a presence leaving one session and arriving at another; + // both rows are stamped. let stamp_session = |row: &str| { format!( "UPDATE observation_clock SET version = version + 1 WHERE singleton = 1; @@ -1264,13 +1666,17 @@ pub(crate) fn init_schema(conn: &Connection) -> Result<()> { }; let stamp_new = stamp_session("NEW"); let stamp_old = stamp_session("OLD"); + let [insert, update, delete] = PRESENCE_TRIGGERS else { + unreachable!("three presence triggers") + }; conn.execute_batch(&format!( - "CREATE TRIGGER IF NOT EXISTS change_feed_session_presences_insert \ + "CREATE TRIGGER IF NOT EXISTS {insert} \ AFTER INSERT ON session_presences BEGIN {stamp_new} END; - CREATE TRIGGER IF NOT EXISTS change_feed_session_presences_update \ - AFTER UPDATE ON session_presences BEGIN + CREATE TRIGGER IF NOT EXISTS {update} \ + AFTER UPDATE ON session_presences \ + WHEN NEW.{REVISION_COLUMN} = OLD.{REVISION_COLUMN} BEGIN {stamp_new} UPDATE observation_clock SET version = version + 1 WHERE singleton = 1 \ AND (OLD.source IS NOT NEW.source OR OLD.session_id IS NOT NEW.session_id); @@ -1279,7 +1685,7 @@ pub(crate) fn init_schema(conn: &Connection) -> Result<()> { WHERE source = OLD.source AND session_id = OLD.session_id \ AND (OLD.source IS NOT NEW.source OR OLD.session_id IS NOT NEW.session_id); END; - CREATE TRIGGER IF NOT EXISTS change_feed_session_presences_delete \ + CREATE TRIGGER IF NOT EXISTS {delete} \ AFTER DELETE ON session_presences BEGIN {stamp_old} END;" @@ -1398,7 +1804,7 @@ mod tests { assert_eq!(keys, vec!["s1", "e1", "e2", "t1", "t2", "mk1", "r1"]); assert!(changes .iter() - .all(|change| change.source == Source::Claude && change.session_id == "s1")); + .all(|change| change.source == Some(Source::Claude) && change.session_id == "s1")); match &changes[1].op { ChangeOp::Upsert(EvidenceRow::SessionEvent(event)) => { assert_eq!(event.text.as_deref(), Some("one")); @@ -1669,7 +2075,12 @@ mod tests { [], ) .unwrap(); - let query = || ChangeQuery::default().consumer("c"); + // The catalog row alone: the presence rows are a kind of their own. + let query = || { + ChangeQuery::default() + .kinds([ChangeKind::Session]) + .consumer("c") + }; let mut seen = store.changes_since(Watermark::CONSUMER, query()).unwrap(); let mut last = None; for change in seen.by_ref() { @@ -1710,7 +2121,11 @@ mod tests { [], ) .unwrap(); - let delta = drain(store.changes_since(head, ChangeQuery::default()).unwrap()); + let delta = drain( + store + .changes_since(head, ChangeQuery::default().kinds([ChangeKind::Session])) + .unwrap(), + ); assert_eq!(delta.len(), 1, "{delta:?}"); match &delta[0].op { ChangeOp::Upsert(EvidenceRow::Session(session)) => { @@ -2124,7 +2539,8 @@ mod tests { let page: Vec = vec![0i64.into(), 1_000i64.into(), 100i64.into()]; let mut plans = Vec::new(); for kind in ChangeKind::ALL { - plans.push((kind.table().name, upsert_sql(*kind), page.clone())); + let stored = stored_columns(&conn, kind.table().name).unwrap(); + plans.push((kind.table().name, upsert_sql(*kind, &stored), page.clone())); } let mut tombstone_page = vec![rusqlite::types::Value::from("session_event".to_string())]; tombstone_page.extend(page.iter().cloned()); @@ -2561,7 +2977,7 @@ mod tests { let conn = open_db(&db).unwrap(); insert_event(&conn, "s1", "e1", "x"); conn.execute_batch( - "DELETE FROM schema_migrations WHERE name = 'change_feed_v1'; \ + "DELETE FROM schema_migrations WHERE name = 'change_feed_v2'; \ UPDATE observation_clock SET version = 42;", ) .unwrap(); @@ -2606,7 +3022,7 @@ mod tests { conn.execute_batch( "DROP INDEX idx_session_events_revision; \ ALTER TABLE session_events DROP COLUMN revision; \ - DELETE FROM schema_migrations WHERE name = 'change_feed_v1'; \ + DELETE FROM schema_migrations WHERE name = 'change_feed_v2'; \ UPDATE observation_clock SET version = 0;", ) .unwrap(); @@ -2633,6 +3049,183 @@ mod tests { assert_eq!(all(&store), changes); } + /// A store the feed reached before it reported every kind: the kinds it + /// lacked are stamped on migration above its head, so a cursor committed + /// for every kind resumes into all of their rows and none of the rows it + /// already had; the session re-stamp triggers move to their own names; + /// and from then on a presence write stamps its session once. + #[test] + fn a_store_fed_before_every_kind_gains_the_rest_above_its_head() { + const NEW_KINDS: [ChangeKind; 6] = [ + ChangeKind::History, + ChangeKind::Presence, + ChangeKind::CommitLink, + ChangeKind::Trajectory, + ChangeKind::SourceObservation, + ChangeKind::ObservationEvidence, + ]; + let dir = tempfile::tempdir().unwrap(); + let db = dir.path().join("fed-v1.db"); + let committed = { + let (store, conn) = { + let store = SessionStore::open(StoreOptions { + db_path: Some(db.clone()), + ..StoreOptions::default() + }) + .unwrap(); + (store, open_db(&db).unwrap()) + }; + conn.execute( + "INSERT INTO sessions (session_id, source) VALUES ('s1', 'claude')", + [], + ) + .unwrap(); + insert_event(&conn, "s1", "e1", "x"); + // Take the database back to the six-kind feed: no stamp, index or + // triggers on the other tables, and the session re-stamp triggers + // under the names they had. + for trigger in PRESENCE_TRIGGERS { + conn.execute_batch(&format!("DROP TRIGGER {trigger};")) + .unwrap(); + } + for kind in NEW_KINDS { + for trigger in trigger_names(kind) { + conn.execute_batch(&format!("DROP TRIGGER {trigger};")) + .unwrap(); + } + let table = kind.table(); + conn.execute_batch(&format!( + "DROP INDEX {index}; ALTER TABLE {name} DROP COLUMN {REVISION_COLUMN};", + index = table.revision_index(), + name = table.name + )) + .unwrap(); + } + let stamp = |row: &str| { + format!( + "UPDATE observation_clock SET version = version + 1 WHERE singleton = 1; \ + UPDATE sessions SET revision = \ + (SELECT version FROM observation_clock WHERE singleton = 1) \ + WHERE source = {row}.source AND session_id = {row}.session_id;" + ) + }; + conn.execute_batch(&format!( + "CREATE TRIGGER change_feed_session_presences_insert \ + AFTER INSERT ON session_presences BEGIN {new} END; \ + CREATE TRIGGER change_feed_session_presences_update \ + AFTER UPDATE ON session_presences BEGIN {new} END; \ + CREATE TRIGGER change_feed_session_presences_delete \ + AFTER DELETE ON session_presences BEGIN {old} END; \ + DELETE FROM schema_migrations WHERE name = 'change_feed_v2'; \ + INSERT OR IGNORE INTO schema_migrations (name) VALUES ('change_feed_v1');", + new = stamp("NEW"), + old = stamp("OLD") + )) + .unwrap(); + // Rows the six-kind feed never stamped. + conn.execute_batch( + "INSERT INTO history (source, session_id, prompt, timestamp_ms) \ + VALUES ('claude', 's1', 'hello', 1000); \ + INSERT INTO session_presences (source, session_id, location) \ + VALUES ('claude', 's1', 'local'); \ + INSERT INTO trajectories (id, decisions_json, retrospective_json, \ + search_text, updated_ms, timestamp_ms) \ + VALUES ('traj-1', '[]', '{}', 'x', 1, 1);", + ) + .unwrap(); + assert!(!schema_is_current(&conn).unwrap()); + // A consumer of every kind has read the six-kind store to its + // head. This build refuses to drain a schema it has not migrated, + // so the cursor is written as the six-kind build committed it. + assert!(store + .changes_since(Watermark::START, ChangeQuery::default()) + .is_err()); + let head: i64 = conn + .query_row( + "SELECT version FROM observation_clock WHERE singleton = 1", + [], + |row| row.get(0), + ) + .unwrap(); + conn.execute( + "INSERT INTO consumer_cursors (name, revision, updated_ms, kinds) \ + VALUES ('all', ?, 0, '*')", + [head], + ) + .unwrap(); + head as u64 + }; + + let store = SessionStore::open(StoreOptions { + db_path: Some(db.clone()), + ..StoreOptions::default() + }) + .unwrap(); + let conn = open_db(&db).unwrap(); + assert!(schema_is_current(&conn).unwrap()); + let resumed = drain( + store + .changes_since(Watermark::CONSUMER, ChangeQuery::default().consumer("all")) + .unwrap(), + ); + let kinds: Vec = resumed.iter().map(|change| change.kind).collect(); + assert_eq!( + kinds, + vec![ + ChangeKind::History, + ChangeKind::Presence, + ChangeKind::Trajectory + ], + "the new kinds' rows, and nothing the cursor already accounted for: {resumed:?}" + ); + assert!(resumed.iter().all(|change| change.revision > committed)); + assert_eq!( + resumed[0].key, + vec![ + Value::from("history"), + Value::from("claude"), + Value::from(1000), + Value::from("hello") + ] + ); + + // The retired names now stamp the presence kind; the re-stamp has its + // own. + let body: String = conn + .query_row( + "SELECT sql FROM sqlite_master WHERE name = 'change_feed_session_presences_insert'", + [], + |row| row.get(0), + ) + .unwrap(); + assert!(body.contains("evidence_tombstones"), "{body}"); + let head = store.head_revision().unwrap(); + conn.execute( + "INSERT INTO session_presences (source, session_id, location) \ + VALUES ('claude', 's1', 'remote')", + [], + ) + .unwrap(); + let delta = drain(store.changes_since(head, ChangeQuery::default()).unwrap()); + let kinds: Vec = delta.iter().map(|change| change.kind).collect(); + assert_eq!(kinds.len(), 2, "{delta:?}"); + assert!(kinds.contains(&ChangeKind::Session) && kinds.contains(&ChangeKind::Presence)); + assert_eq!( + store.head_revision().unwrap().revision, + head.revision + 2, + "one stamp for the presence and one for its session" + ); + + // Re-opening does not stamp again. + let before = all(&store); + SessionStore::open(StoreOptions { + db_path: Some(db), + ..StoreOptions::default() + }) + .unwrap(); + assert_eq!(all(&store), before); + } + #[test] fn a_read_only_store_reads_the_feed_but_cannot_commit() { let dir = tempfile::tempdir().unwrap(); @@ -2657,6 +3250,188 @@ mod tests { assert_eq!(drain(changes).len(), 1); } + /// Every kind's stored row and key, as the feed reports them, are what + /// the export journal records for the same writes -- inserts, updates + /// that change a value, one that moves a row to another identity, and + /// deletes, direct and cascaded. Compared as the state each stream + /// replays to, so the feed's netting of a re-stamped row is not a + /// difference, and column order is part of the comparison. + #[cfg(feature = "export")] + #[test] + fn every_kind_matches_the_export_journal_for_the_same_writes() { + use crate::export::capture; + use std::collections::{BTreeMap, BTreeSet}; + + type State = BTreeMap>; + + fn journal(conn: &Connection) -> State { + let mut state = State::new(); + let mut position = 0; + while let Some(record) = capture::next_change(conn, position, None).unwrap() { + position = record.position; + let key: Vec = serde_json::from_str(&record.key).unwrap(); + let row = match record.operation.as_str() { + "upsert" => Some(serde_json::from_str::(&record.payload).unwrap()), + "delete" => None, + other => panic!("journal operation {other}"), + }; + state.insert(serde_json::to_string(&key).unwrap(), row); + } + state + } + + fn feed(store: &SessionStore) -> (State, BTreeSet) { + let mut state = State::new(); + let mut kinds = BTreeSet::new(); + for change in all(store) { + kinds.insert(change.kind); + assert_eq!(change.key[0], Value::from(change.kind.as_str())); + match &change.op { + ChangeOp::Upsert(_) => { + let columns = change.columns.clone().expect("an upsert carries its row"); + assert!(columns.get(REVISION_COLUMN).is_none()); + assert_eq!( + change.kind.table().key_from_row(change.kind, &columns), + change.key + ); + state.insert(serde_json::to_string(&change.key).unwrap(), Some(columns)); + } + ChangeOp::Delete => { + assert!(change.columns.is_none()); + state.insert(serde_json::to_string(&change.key).unwrap(), None); + } + } + } + (state, kinds) + } + + fn assert_parity( + store: &SessionStore, + conn: &Connection, + step: &str, + ) -> BTreeSet { + let journal = journal(conn); + let (feed, kinds) = feed(store); + let keys = |state: &State| state.keys().cloned().collect::>(); + assert_eq!(keys(&feed), keys(&journal), "{step}: the same records"); + for (key, row) in &journal { + assert_eq!(&feed[key], row, "{step}: {key}"); + } + kinds + } + + let dir = tempfile::tempdir().unwrap(); + let (store, conn) = store(dir.path()); + // One subscription over everything, its snapshot already read: every + // write from here on is journaled. + { + let tx = conn.unchecked_transaction().unwrap(); + capture::save_subscription( + &tx, + &capture::Subscription { + id: "parity", + session: None, + cursor: 0, + kind: capture::kind_count(), + rowid: 0, + complete: true, + }, + ) + .unwrap(); + tx.commit().unwrap(); + } + + conn.execute_batch( + r#" +INSERT INTO sessions (session_id, source, cwd, models_json, workspace_roots_json) + VALUES ('s1', 'claude', '/p', '["opus"]', 'not json'); +INSERT INTO sessions (session_id, source) VALUES ('s2', 'claude'); +INSERT INTO session_events (source, session_id, message_id, ts_ms, role, kind, text, event_uid, + token_json, project_key_method, raw_facts_version) + VALUES ('claude', 's1', 'm1', 10, 'assistant', 'text', 'one', 'e1', '{"input":3}', 'git', 2); +INSERT INTO tool_calls (source, session_id, tool_use_id, name, args_json, is_error) + VALUES ('claude', 's1', 't1', 'Bash', '{"command":"ls"}', 0); +INSERT INTO file_edits (source, session_id, tool_use_id, file_path, tool_name, lines_added) + VALUES ('claude', 's1', 't2', '/p/a.rs', 'Edit', 3); +INSERT INTO session_markers (source, session_id, marker_uid, kind, payload_json) + VALUES ('claude', 's1', 'mk1', 'compaction', '{"trigger":"auto"}'); +INSERT INTO session_relationships (source, parent_session_id, relationship_uid, + child_session_id, relationship, identity_status, evidence_kind, child_has_events, + created_ms, updated_ms) + VALUES ('claude', 's1', 'r1', 's2', 'delegated', 'observed', 'sidecar', 1, 1, 2); +INSERT INTO history (source, session_id, project, prompt, timestamp_ms) + VALUES ('claude', 's1', '/p', 'hello', 1000); +INSERT INTO history (source, session_id, prompt, timestamp_ms) + VALUES ('codex', NULL, 'no session yet', 2000); +INSERT INTO session_presences (source, session_id, location, raw_locator) + VALUES ('claude', 's1', 'local', '/p/s1.jsonl'); +INSERT INTO session_commit_links (source, session_id, repo, commit_sha, match_method, + confidence, files_json, created_at_ms) + VALUES ('claude', 's1', 'repo', 'abc123', 'trailer', 0.75, '["a.rs"]', 1); +INSERT INTO trajectories (id, version, status, decisions_json, retrospective_json, + search_text, updated_ms, timestamp_ms) + VALUES ('traj-1', 1, 'active', '[]', '{}', 'x', 1, 1); +INSERT INTO session_observations (source, session_id, location, connector_id, + connector_instance, updated_ms) + VALUES ('claude', 's1', 'remote', 'conn', 'default', 1); +INSERT INTO observation_evidence (source, session_id, location, connector_id, + connector_instance, evidence_uid, payload_json) + VALUES ('claude', 's1', 'remote', 'conn', 'default', 'ev1', '{"a":1}'); +"#, + ) + .unwrap(); + let kinds = assert_parity(&store, &conn, "inserts"); + let every: BTreeSet = ChangeKind::ALL.iter().copied().collect(); + assert_eq!(kinds, every, "the writes cover every kind"); + + conn.execute_batch( + r#" +UPDATE sessions SET cwd = '/q' WHERE session_id = 's1'; +UPDATE session_events SET text = 'one, edited' WHERE event_uid = 'e1'; +UPDATE tool_calls SET tool_use_id = 't1b' WHERE tool_use_id = 't1'; +UPDATE file_edits SET lines_removed = 1 WHERE tool_use_id = 't2'; +UPDATE session_markers SET text = 'summary' WHERE marker_uid = 'mk1'; +UPDATE session_relationships SET updated_ms = 9 WHERE relationship_uid = 'r1'; +UPDATE history SET project = '/q' WHERE prompt = 'hello'; +UPDATE history SET session_id = 'c1' WHERE prompt = 'no session yet'; +UPDATE session_presences SET raw_locator = '/q/s1.jsonl'; +UPDATE session_commit_links SET confidence = 0.5; +UPDATE trajectories SET status = 'completed', completed_at = '2026-09-24'; +UPDATE session_observations SET access_state = 'unavailable'; +UPDATE observation_evidence SET payload_json = '{"a":2}'; +"#, + ) + .unwrap(); + assert_parity(&store, &conn, "updates"); + + conn.execute_batch( + "DELETE FROM history WHERE prompt = 'hello'; + DELETE FROM trajectories; + DELETE FROM observation_evidence; + DELETE FROM session_commit_links; + DELETE FROM sessions WHERE session_id = 's1';", + ) + .unwrap(); + assert_parity(&store, &conn, "deletes"); + let tombstoned: BTreeSet = all(&store) + .into_iter() + .filter(|change| change.op == ChangeOp::Delete) + .map(|change| change.kind) + .collect(); + for kind in [ + ChangeKind::Session, + ChangeKind::ToolCall, + ChangeKind::History, + ChangeKind::Presence, + ChangeKind::CommitLink, + ChangeKind::Trajectory, + ChangeKind::SourceObservation, + ChangeKind::ObservationEvidence, + ] { + assert!(tombstoned.contains(&kind), "{kind:?} left a tombstone"); + } + } + /// A read-only handle over a database the feed schema has not reached is /// told what to do, not served `no such table`. #[test] diff --git a/crates/ai-hist/src/export/schema.rs b/crates/ai-hist/src/export/schema.rs index d5683c3f..ba5856c5 100644 --- a/crates/ai-hist/src/export/schema.rs +++ b/crates/ai-hist/src/export/schema.rs @@ -172,7 +172,8 @@ impl Table { own } } - /// The columns a capture trigger carries, in table order. + /// The columns a capture trigger carries, in table order: the change + /// feed's stored row, so the two cannot disagree about a record's shape. /// /// The change feed's `revision` stamp is excluded on purpose. It is this /// database's bookkeeping, not a fact about the record, and it changes on @@ -180,17 +181,7 @@ impl Table { /// each insert read as a second change of the row, and journal every /// record twice. pub(crate) fn captured_columns(&self, conn: &Connection) -> Result> { - let columns = conn - .prepare(&format!( - "SELECT name FROM pragma_table_info('{}')", - self.name - ))? - .query_map([], |row| row.get::<_, String>(0))? - .collect::>>()?; - Ok(columns - .into_iter() - .filter(|column| column != crate::change_feed::REVISION_COLUMN) - .collect()) + crate::change_feed::stored_columns(conn, self.name) } pub fn payload(&self, conn: &Connection, row: &str) -> Result { diff --git a/crates/ai-hist/src/lib.rs b/crates/ai-hist/src/lib.rs index ca129df7..37c15f18 100644 --- a/crates/ai-hist/src/lib.rs +++ b/crates/ai-hist/src/lib.rs @@ -71,7 +71,7 @@ pub use store::*; pub(crate) use store::*; pub use change_feed::{ - Change, ChangeKind, ChangeOp, ChangeQuery, Changes, EvidenceRow, Watermark, + Change, ChangeKind, ChangeOp, ChangeQuery, Changes, EvidenceRow, StoredRow, Watermark, DEFAULT_CHANGE_BATCH, MAX_CHANGE_BATCH, }; pub use discover::{declared_evidence_kinds, missing_evidence_kinds, ShallowSession}; diff --git a/crates/ai-hist/src/source_evidence.rs b/crates/ai-hist/src/source_evidence.rs index d2a74c14..4caef46c 100644 --- a/crates/ai-hist/src/source_evidence.rs +++ b/crates/ai-hist/src/source_evidence.rs @@ -842,6 +842,8 @@ mod tests { ("file_edits", "revision"), ("session_relationships", "revision"), ("session_markers", "revision"), + ("history", "revision"), + ("session_commit_links", "revision"), ]; let conn = rusqlite::Connection::open_in_memory().unwrap(); crate::init_db(&conn).unwrap(); diff --git a/crates/ai-hist/src/store.rs b/crates/ai-hist/src/store.rs index 6b8601f7..fbb38470 100644 --- a/crates/ai-hist/src/store.rs +++ b/crates/ai-hist/src/store.rs @@ -334,7 +334,10 @@ CREATE TRIGGER IF NOT EXISTS history_ai AFTER INSERT ON history BEGIN INSERT INTO history_fts(rowid, prompt, project) VALUES (new.id, new.prompt, new.project); END; -CREATE TRIGGER IF NOT EXISTS history_au AFTER UPDATE ON history BEGIN +-- `UPDATE OF`, for the same reason as `session_events_au` below: the +-- change feed re-stamps `revision` after each insert, and that touches +-- nothing the index holds. +CREATE TRIGGER IF NOT EXISTS history_au AFTER UPDATE OF id, prompt, project ON history BEGIN INSERT INTO history_fts(history_fts, rowid, prompt, project) VALUES('delete', old.id, old.prompt, old.project); INSERT INTO history_fts(rowid, prompt, project) @@ -786,6 +789,7 @@ const REQUIRED_SCHEMA_MIGRATIONS: &[&str] = &[ // TRIGGER IF NOT EXISTS` keeps an existing database's unconditional body, // so the marker is what makes the rebuild happen exactly once. "session_events_fts_update_of_v1", + "history_fts_update_of_v1", ]; #[cfg(feature = "export")] const REQUIRED_EXPORT_MIGRATIONS: &[&str] = &["delivery_v1"]; @@ -1097,6 +1101,9 @@ fn init_db_locked(conn: &Connection) -> Result<()> { if !migration_applied(conn, "session_events_fts_update_of_v1")? { conn.execute_batch("DROP TRIGGER IF EXISTS session_events_au;")?; } + if !migration_applied(conn, "history_fts_update_of_v1")? { + conn.execute_batch("DROP TRIGGER IF EXISTS history_au;")?; + } conn.execute_batch(SCHEMA)?; // Before the trigger below, whose body deletes from these tables. conn.execute_batch(SESSION_RELATIONSHIPS_DDL)?; @@ -1257,7 +1264,8 @@ END; conn.execute_batch( "INSERT OR IGNORE INTO schema_migrations (name) \ VALUES ('session_delete_continuity_reopen_v1'), \ - ('session_events_fts_update_of_v1');", + ('session_events_fts_update_of_v1'), \ + ('history_fts_update_of_v1');", )?; migrate_session_relationships_v2(conn)?; migrate_tool_result_fidelity_v1(conn)?; diff --git a/crates/ai-hist/tests/change_feed.rs b/crates/ai-hist/tests/change_feed.rs index d678490a..826e9805 100644 --- a/crates/ai-hist/tests/change_feed.rs +++ b/crates/ai-hist/tests/change_feed.rs @@ -10,7 +10,9 @@ use ai_hist::{ Change, ChangeKind, ChangeOp, ChangeQuery, EvidenceRow, SessionQuery, SessionRef, SessionStore, Source, StoreOptions, Watermark, }; +use rusqlite::types::ValueRef; use rusqlite::{Connection, OpenFlags}; +use serde_json::Value; use std::collections::{BTreeMap, BTreeSet}; use std::fs; use std::path::{Path, PathBuf}; @@ -106,6 +108,12 @@ impl Home { fn raw(&self) -> Connection { Connection::open_with_flags(self.db(), OpenFlags::SQLITE_OPEN_READ_ONLY).unwrap() } + + /// A raw write, standing in for a migration or a writer this crate does + /// not own. + fn raw_writer(&self) -> Connection { + Connection::open(self.db()).unwrap() + } } fn drain(store: &SessionStore, from: Watermark) -> Vec { @@ -182,7 +190,7 @@ fn an_in_progress_message_reaches_the_feed_only_once_complete() { ); assert!(tick .iter() - .all(|change| change.source == Source::Claude && change.session_id == SESSION)); + .all(|change| change.source == Some(Source::Claude) && change.session_id == SESSION)); // Nothing changed: a tick is empty, and the head stays put. let head = store.head_revision().unwrap(); @@ -249,10 +257,14 @@ fn filetime_now_plus(seconds: u64) -> std::time::SystemTime { std::time::SystemTime::now() + std::time::Duration::from_secs(seconds) } -type Key = (ChangeKind, String, String, String); +/// A record as the feed and the tables both name it: its kind and +/// [`Change::key`], serialized. +type Key = (ChangeKind, String); +/// A stored row: every column but `revision`, in table order. +type Row = Vec<(String, Value)>; /// Apply a feed in order: an upsert replaces, a delete removes. -fn replay(changes: &[Change]) -> BTreeMap { +fn replay(changes: &[Change]) -> BTreeMap { let mut state = BTreeMap::new(); let mut last = 0; for change in changes { @@ -261,15 +273,17 @@ fn replay(changes: &[Change]) -> BTreeMap { "the feed is strictly ordered by revision: {change:?}" ); last = change.revision; - let key = ( - change.kind, - change.source.as_str().to_string(), - change.session_id.clone(), - change.record_key.clone(), - ); + let key = (change.kind, serde_json::to_string(&change.key).unwrap()); match &change.op { - ChangeOp::Upsert(row) => { - state.insert(key, row.clone()); + ChangeOp::Upsert(_) => { + let columns = change.columns.as_ref().expect("an upsert carries its row"); + state.insert( + key, + columns + .iter() + .map(|(name, value)| (name.to_string(), value.clone())) + .collect(), + ); } ChangeOp::Delete => { state.remove(&key); @@ -280,69 +294,141 @@ fn replay(changes: &[Change]) -> BTreeMap { state } -/// The tables as they stand, keyed the way the feed keys them. -fn direct(conn: &Connection) -> BTreeMap> { - let mut rows = BTreeMap::new(); - let reads: [(ChangeKind, &str); 6] = [ - ( - ChangeKind::Session, - "SELECT source, session_id, session_id, cwd FROM sessions", - ), - ( - ChangeKind::SessionEvent, - "SELECT source, session_id, event_uid, text FROM session_events", - ), - ( - ChangeKind::ToolCall, - "SELECT source, session_id, tool_use_id, name FROM tool_calls", - ), - ( - ChangeKind::FileEdit, - "SELECT source, session_id, tool_use_id, file_path FROM file_edits", - ), - ( - ChangeKind::SessionMarker, - "SELECT source, session_id, marker_uid, kind FROM session_markers", - ), - ( - ChangeKind::Relationship, - "SELECT source, parent_session_id, relationship_uid, relationship \ - FROM session_relationships", - ), - ]; - for (kind, sql) in reads { - let mut statement = conn.prepare(sql).unwrap(); - let read = statement - .query_map([], |row| { - Ok(( - ( - kind, - row.get::<_, String>(0)?, - row.get::<_, String>(1)?, - row.get::<_, String>(2)?, - ), - row.get::<_, Option>(3)?, - )) - }) - .unwrap(); - for entry in read { - let (key, value) = entry.unwrap(); - rows.insert(key, value); - } +/// Each kind's table, its source (a column, or the literal a table without +/// one shares) and the columns of its identity after the source. +const TABLES: &[(ChangeKind, &str, Option<&str>, &[&str])] = &[ + (ChangeKind::Session, "sessions", None, &["session_id"]), + ( + ChangeKind::SessionEvent, + "session_events", + None, + &["session_id", "event_uid"], + ), + ( + ChangeKind::ToolCall, + "tool_calls", + None, + &["session_id", "tool_use_id"], + ), + ( + ChangeKind::FileEdit, + "file_edits", + None, + &["session_id", "tool_use_id"], + ), + ( + ChangeKind::SessionMarker, + "session_markers", + None, + &["session_id", "marker_uid"], + ), + ( + ChangeKind::Relationship, + "session_relationships", + None, + &["parent_session_id", "relationship_uid"], + ), + ( + ChangeKind::History, + "history", + None, + &["timestamp_ms", "prompt"], + ), + ( + ChangeKind::Presence, + "session_presences", + None, + &["session_id", "location"], + ), + ( + ChangeKind::CommitLink, + "session_commit_links", + None, + &["session_id", "commit_sha", "match_method"], + ), + ( + ChangeKind::Trajectory, + "trajectories", + Some("trajectory"), + &["id"], + ), + ( + ChangeKind::SourceObservation, + "session_observations", + None, + &[ + "session_id", + "location", + "connector_id", + "connector_instance", + ], + ), + ( + ChangeKind::ObservationEvidence, + "observation_evidence", + None, + &[ + "session_id", + "location", + "connector_id", + "connector_instance", + "evidence_uid", + ], + ), +]; + +fn sqlite_value(value: ValueRef<'_>) -> Value { + match value { + ValueRef::Null => Value::Null, + ValueRef::Integer(integer) => Value::from(integer), + ValueRef::Real(real) => Value::from(real), + ValueRef::Text(text) => Value::from(String::from_utf8(text.to_vec()).unwrap()), + ValueRef::Blob(bytes) => Value::from(bytes.to_vec()), } - rows } -fn value_of(row: &EvidenceRow) -> Option { - match row { - EvidenceRow::Session(session) => session.cwd.clone(), - EvidenceRow::SessionEvent(event) => event.text.clone(), - EvidenceRow::ToolCall(call) => Some(call.name.clone()), - EvidenceRow::FileEdit(edit) => Some(edit.file_path.clone()), - EvidenceRow::SessionMarker(marker) => Some(marker.kind.clone()), - EvidenceRow::Relationship(relationship) => Some(relationship.relationship.clone()), - other => panic!("a row this consumer does not know: {other:?}"), +/// Every row of one table as stored, every column but `revision`. +fn stored_rows(conn: &Connection, table: &str) -> Vec { + let mut statement = conn.prepare(&format!("SELECT * FROM {table}")).unwrap(); + let names: Vec = statement + .column_names() + .into_iter() + .map(str::to_string) + .collect(); + let rows = statement + .query_map([], |row| { + let mut stored = Row::new(); + for (index, name) in names.iter().enumerate() { + if name != "revision" { + stored.push((name.clone(), sqlite_value(row.get_ref(index)?))); + } + } + Ok(stored) + }) + .unwrap(); + rows.collect::>().unwrap() +} + +/// Every fed table as it stands, keyed the way the feed keys it. +fn direct(conn: &Connection) -> BTreeMap { + let mut state = BTreeMap::new(); + for (kind, table, literal_source, identity) in TABLES { + for row in stored_rows(conn, table) { + let value = |column: &str| { + row.iter() + .find(|(name, _)| name == column) + .map(|(_, value)| value.clone()) + .unwrap() + }; + let mut key = vec![Value::from(kind.as_str())]; + if literal_source.is_none() { + key.push(value("source")); + } + key.extend(identity.iter().map(|column| value(column))); + state.insert((*kind, serde_json::to_string(&key).unwrap()), row); + } } + state } fn assert_replay_matches_tables(home: &Home, store: &SessionStore, step: &str) { @@ -356,9 +442,8 @@ fn assert_replay_matches_tables(home: &Home, store: &SessionStore, step: &str) { ); for (key, row) in &replayed { assert_eq!( - &value_of(row), - &tables[key], - "{step}: the replayed row for {key:?} is the row the table holds" + row, &tables[key], + "{step}: the replayed row for {key:?} is the row the table holds, column for column" ); } assert!( @@ -458,7 +543,7 @@ fn replaying_the_feed_reconstructs_the_tables_after_every_sync() { assert_replay_matches_tables(&home, &store, "a second provider"); let codex: Vec = drain(&store, Watermark::START) .into_iter() - .filter(|change| change.source == Source::Codex) + .filter(|change| change.source == Some(Source::Codex)) .collect(); assert!(!codex.is_empty()); kinds_seen.extend(codex.iter().map(|change| change.kind)); @@ -474,10 +559,21 @@ fn replaying_the_feed_reconstructs_the_tables_after_every_sync() { assert!(drain(&store, head).is_empty()); assert_replay_matches_tables(&home, &store, "an idle sync"); - // The corpus exercised every kind the feed reports, so the replay check - // above covered every table rather than only the easy ones. - let expected: BTreeSet = ChangeKind::ALL.iter().copied().collect(); - assert_eq!(kinds_seen, expected, "every kind was fed by the corpus"); + // The corpus exercised every kind it writes, so the replay check above + // covered every table the syncs filled rather than only the easy ones. + let filled: BTreeSet = direct(&home.raw()).keys().map(|(kind, _)| *kind).collect(); + assert_eq!(kinds_seen, filled, "every kind the corpus wrote was fed"); + for kind in [ + ChangeKind::Session, + ChangeKind::SessionEvent, + ChangeKind::ToolCall, + ChangeKind::FileEdit, + ChangeKind::SessionMarker, + ChangeKind::Relationship, + ChangeKind::Presence, + ] { + assert!(filled.contains(&kind), "the corpus writes {kind:?}"); + } // A second consumer, starting now, replays everything independently of // the first one's commit. @@ -493,3 +589,183 @@ fn replaying_the_feed_reconstructs_the_tables_after_every_sync() { direct(&home.raw()).len() ); } + +fn only(store: &SessionStore, from: Watermark, kind: ChangeKind) -> Vec { + store + .changes_since(from, ChangeQuery::default().kinds([kind])) + .unwrap() + .map(|change| change.unwrap()) + .collect() +} + +/// A prompt from a provider's prompt log is a `history` change, keyed the +/// way the prompt log is unique, with its typed entry and its stored row. +#[test] +fn a_history_prompt_reaches_the_feed_with_its_stored_row() { + let home = Home::new(); + fs::create_dir_all(home.path().join(".claude")).unwrap(); + fs::write( + home.path().join(".claude/history.jsonl"), + "{\"display\":\"ship the feed\",\"timestamp\":1756634400000,\ + \"project\":\"/tmp/project\",\"sessionId\":\"history-session\"}\n", + ) + .unwrap(); + let store = home.store(); + store.sync(Default::default()).unwrap(); + + let history = only(&store, Watermark::START, ChangeKind::History); + assert_eq!(history.len(), 1, "{history:?}"); + let change = &history[0]; + assert_eq!(change.source, Some(Source::Claude)); + assert_eq!(change.source_name, "claude"); + assert_eq!(change.session_id, "history-session"); + assert_eq!( + change.key, + vec![ + Value::from("history"), + Value::from("claude"), + Value::from(1_756_634_400_000i64), + Value::from("ship the feed"), + ] + ); + match &change.op { + ChangeOp::Upsert(EvidenceRow::History(entry)) => { + assert_eq!(entry.prompt, "ship the feed"); + assert_eq!(entry.timestamp_ms, 1_756_634_400_000); + } + other => panic!("a typed history row: {other:?}"), + } + let columns = change.columns.as_ref().unwrap(); + let stored = stored_rows(&home.raw(), "history"); + let row: Row = columns + .iter() + .map(|(name, value)| (name.to_string(), value.clone())) + .collect(); + assert_eq!(vec![row], stored, "the stored row, column for column"); + assert_eq!(columns.get("project"), Some(&Value::from("/tmp/project"))); + assert!(columns.get("revision").is_none()); +} + +/// A column a table gains after the feed was built is carried as soon as it +/// exists, with every other column exactly as stored: JSON text stays text, +/// integers stay integers, NULL stays null. +#[test] +fn a_column_a_table_gains_is_carried_verbatim() { + let home = Home::new(); + home.stage_claude("multi-block-turn.jsonl"); + let store = home.store(); + store.sync(Default::default()).unwrap(); + + let head = store.head_revision().unwrap(); + let writer = home.raw_writer(); + writer + .execute_batch( + "ALTER TABLE tool_calls ADD COLUMN review_note TEXT; \ + UPDATE tool_calls SET review_note = 'looked fine' \ + WHERE rowid = (SELECT MIN(rowid) FROM tool_calls);", + ) + .unwrap(); + let changes = only(&store, head, ChangeKind::ToolCall); + assert_eq!(changes.len(), 1, "{changes:?}"); + let columns = changes[0].columns.as_ref().unwrap(); + assert_eq!( + columns.get("review_note"), + Some(&Value::from("looked fine")) + ); + let names: Vec<&str> = columns.iter().map(|(name, _)| name).collect(); + assert_eq!(names.last(), Some(&"review_note"), "in table order"); + let row: Row = columns + .iter() + .map(|(name, value)| (name.to_string(), value.clone())) + .collect(); + let stored = stored_rows(&home.raw(), "tool_calls"); + assert!(stored.contains(&row), "the stored row, column for column"); + assert!( + columns.get("args_json").is_some_and(Value::is_string), + "JSON text stays text" + ); + assert!(columns.get("id").is_some_and(Value::is_i64)); + + // The session row carries every catalog column, `parser_version` + // included, and its JSON columns unparsed. + let sessions = only(&store, Watermark::START, ChangeKind::Session); + let catalog = sessions[0].columns.as_ref().unwrap(); + assert!(catalog.get("parser_version").is_some_and(Value::is_i64)); + assert!(catalog + .get("models_json") + .is_some_and(|models| models.is_string() || models.is_null())); + assert!(catalog.get("locations").is_none(), "no derived column"); +} + +/// A delete carries the key the record's upserts carried, including for a +/// kind whose identity spans several columns of different types. +#[test] +fn a_tombstone_carries_the_record_key() { + let home = Home::new(); + home.stage_claude("multi-block-turn.jsonl"); + fs::write( + home.path().join(".claude/history.jsonl"), + "{\"display\":\"to be removed\",\"timestamp\":42,\"sessionId\":\"gone\"}\n", + ) + .unwrap(); + let store = home.store(); + store.sync(Default::default()).unwrap(); + let upserts: Vec = drain(&store, Watermark::START); + let call = upserts + .iter() + .find(|change| change.kind == ChangeKind::ToolCall) + .unwrap(); + let prompt = upserts + .iter() + .find(|change| change.kind == ChangeKind::History) + .unwrap(); + + let head = store.head_revision().unwrap(); + let writer = home.raw_writer(); + writer + .execute( + "DELETE FROM tool_calls WHERE tool_use_id = ?", + [&call.record_key], + ) + .unwrap(); + writer + .execute("DELETE FROM history WHERE timestamp_ms = 42", []) + .unwrap(); + let deletes = drain(&store, head); + assert_eq!(deletes.len(), 2, "{deletes:?}"); + for (delete, upsert) in deletes.iter().zip([call, prompt]) { + assert_eq!(delete.op, ChangeOp::Delete); + assert!(delete.columns.is_none()); + assert_eq!(delete.kind, upsert.kind); + assert_eq!(delete.key, upsert.key, "the same identity"); + assert_eq!(delete.record_key, upsert.record_key); + assert_eq!(delete.source_name, upsert.source_name); + } + assert_eq!(deletes[1].key[2], Value::from(42), "an integer stays one"); +} + +/// A row from a source this build does not know is carried, not a failed +/// drain: `source` is `None` and `source_name` names it. +#[test] +fn a_row_from_an_unknown_source_is_carried() { + let home = Home::new(); + let store = home.store(); + home.raw_writer() + .execute( + "INSERT INTO history (source, session_id, prompt, timestamp_ms) \ + VALUES ('some-new-agent', 's1', 'hi', 7)", + [], + ) + .unwrap(); + let changes = drain(&store, Watermark::START); + assert_eq!(changes.len(), 1, "{changes:?}"); + assert_eq!(changes[0].source, None); + assert_eq!(changes[0].source_name, "some-new-agent"); + assert_eq!(changes[0].key[1], Value::from("some-new-agent")); + home.raw_writer() + .execute("DELETE FROM history", []) + .unwrap(); + let changes = drain(&store, Watermark::START); + assert_eq!(changes[0].op, ChangeOp::Delete); + assert_eq!(changes[0].source_name, "some-new-agent"); +} diff --git a/docs/sourcing-sdk.md b/docs/sourcing-sdk.md index 488adc31..2cf55179 100644 --- a/docs/sourcing-sdk.md +++ b/docs/sourcing-sdk.md @@ -343,16 +343,52 @@ Two facts about identity a consumer must not paper over: ### `changes_since` The revision-stamped change feed: every row of `sessions`, `session_events`, -`tool_calls`, `file_edits`, `session_markers` and `session_relationships` -carries a `revision` drawn from the database-wide `observation_clock` and -stamped by a trigger on every insert and update, so no write site can forget -one; a deleted row leaves a tombstone at its own revision, which a later insert -of the same key clears. `changes_since(from, ChangeQuery)` drains -`Change { kind, source, session_id, record_key, revision, op }` in -`(revision, kind, record_key)` order, bounded to the head at open, in +`tool_calls`, `file_edits`, `session_markers`, `session_relationships`, +`history`, `session_presences`, `session_commit_links`, `trajectories`, +`session_observations` and `observation_evidence` carries a `revision` drawn +from the database-wide `observation_clock` and stamped by a trigger on every +insert and update, so no write site can forget one; a deleted row leaves a +tombstone at its own revision, which a later insert of the same key clears. +`changes_since(from, ChangeQuery)` drains +`Change { kind, source, source_name, session_id, record_key, key, revision, op, columns }` +in `(revision, kind, record_key)` order, bounded to the head at open, in `batch`-sized indexed reads of at most `MAX_CHANGE_BATCH`; `op` is -`Upsert(EvidenceRow)` — the typed row, so no second read is needed — or -`Delete`. A re-seen `record_key` is a replace, never a duplicate. +`Upsert(EvidenceRow)` — the typed row, so no second read is needed, or +`EvidenceRow::Untyped` for a kind with none — or `Delete`. A re-seen key is a +replace, never a duplicate. + +| `ChangeKind` | Table | `key` after the kind | Typed row | +| --------------------- | ----------------------- | --------------------------------------------------------------------- | --------------------- | +| `Session` | `sessions` | `source, session_id` | `ShallowSession` | +| `SessionEvent` | `session_events` | `source, session_id, event_uid` | `SessionEvent` | +| `ToolCall` | `tool_calls` | `source, session_id, tool_use_id` | `SessionToolCall` | +| `FileEdit` | `file_edits` | `source, session_id, tool_use_id` | `SessionFileEdit` | +| `SessionMarker` | `session_markers` | `source, session_id, marker_uid` | `SessionMarker` | +| `Relationship` | `session_relationships` | `source, parent_session_id, relationship_uid` | `SessionRelationship` | +| `History` | `history` | `source, timestamp_ms, prompt` | `HistoryEntry` | +| `Presence` | `session_presences` | `source, session_id, location` | — | +| `CommitLink` | `session_commit_links` | `source, session_id, commit_sha, match_method` | — | +| `Trajectory` | `trajectories` | `id` | — | +| `SourceObservation` | `session_observations` | `source, session_id, location, connector_id, connector_instance` | — | +| `ObservationEvidence` | `observation_evidence` | `source, session_id, location, connector_id, connector_instance, evidence_uid` | — | + +`columns` is the row as stored, on every upsert: a `StoredRow` of every column +but `revision`, in table order, each value as SQLite holds it — JSON text stays +text, integers stay integers, NULL stays `null` — read from the live table, so +a column a migration adds is carried without a code change. It serializes as a +JSON object in column order, the same object the export journal records for the +row. `key` is the record's identity as that journal keys it: the kind's wire +name, then the stored values of the table's uniqueness columns, identical on an +upsert and on the delete that retracts it — `["history", "claude", +1756634400000, "ship the feed"]`, `["trajectory", "traj-1"]`. `record_key` is +the part of `key` inside a session: the one identity column's text, or a JSON +array of several (`[1756634400000,"ship the feed"]` for a prompt). +`session_id` is the parent session for a relationship, the id for a trajectory +and empty for a prompt that names no session. `source` is `None` for a source +this build does not know — a row written by a newer release — and +`source_name` is the stored name either way, so such a row is carried rather +than failing the drain. The feed applies no consent or exclusion rule: an +embedder that uploads applies its own selection. `from` is `Watermark::START` to replay everything, an explicit watermark to resume from one a consumer stored itself, or `Watermark::CONSUMER` with @@ -526,7 +562,8 @@ it to fold requests by model. `ShallowSession`, `SessionRelationship`, `SessionLocation` and `SessionScope` are re-exported on the default features because the change feed's `EvidenceRow` carries them: a `Change` hands back the typed row it is about, -so a consumer needs no second read. `session()` returns the facade's own +so a consumer needs no second read. `StoredRow` carries every kind's row as +stored. `session()` returns the facade's own structs (`Prompt`, `Message` and its `Block`s, `ToolCall`, `ToolResult`, `FileEdit`, `Marker`, `Relationship`) — those are the read-side shapes, with JSON columns parsed. `CommitLink` is carried by neither. @@ -607,7 +644,7 @@ embedder reads before bumping. | Feature | Default | What it adds | For | | --- | --- | --- | --- | -| *(none)* | ✓ | `SessionStore` and its ten operations, the change feed (`Change`, `ChangeQuery`, `Watermark`, `EvidenceRow`), `Source` and `SourceCapabilities`, `Error`, the evidence structs above, `NormalizedUsage` and the usage normalizers, `project_identity`, `declared_evidence_kinds` | Embedders | +| *(none)* | ✓ | `SessionStore` and its ten operations, the change feed (`Change`, `ChangeQuery`, `Watermark`, `EvidenceRow`, `StoredRow`), `Source` and `SourceCapabilities`, `Error`, the evidence structs above, `NormalizedUsage` and the usage normalizers, `project_identity`, `declared_evidence_kinds` | Embedders | | `fs-events` | — | The `notify` backend behind `watch`; without it `watch` polls at `poll_interval_ms`. `WatchOptions::use_fs_events` selects it when it is compiled in | The CLI, and an embedder that wants event-driven ticks | | `delivery` | — | Durable delivery of captured evidence to a destination | The CLI, napi, the relayhistory plugin | | `opencode-backup` | — | Snapshot a live OpenCode SQLite store through `rusqlite`'s backup API before reading it | The CLI, napi | From b6e6a75ac09902c24c068292f5970d5d98928b6a Mon Sep 17 00:00:00 2001 From: Will Washburn Date: Thu, 24 Sep 2026 16:30:52 -0700 Subject: [PATCH 2/2] feat(sdk): a prompt's session is not part of its feed identity A history row is unique on source, timestamp and prompt, so a prompt that gains or changes its session is an upsert of the same record. Its tombstones store an empty session and the identity check skips the session column, so no transient delete sits between the two writes. Document that a cursor stored as `*` now spans all twelve kinds. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 10 ++++++-- crates/ai-hist/src/change_feed.rs | 41 ++++++++++++++++++++++++++----- docs/sourcing-sdk.md | 4 ++- 3 files changed, 46 insertions(+), 9 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index df0b1722..97cd1ad5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -189,7 +189,12 @@ Notable changes to the native `ai-hist` CLI are documented here. `Trajectory`, `SourceObservation` and `ObservationEvidence` (`ALL` lists twelve kinds); their rows are stamped and tombstoned by the same triggers, and a database the six-kind feed reached stamps them once on open, above its - head, so a cursor bound to every kind resumes into all of them. + head, so a cursor bound to every kind resumes into all of them. A named + cursor stored as `*` now spans all twelve kinds. A consumer that passed the + six original kinds as an explicit list was stored as `*` too; that list no + longer equals `ChangeKind::ALL`, so its resume fails with + `ConsumerKindsMismatch` — drain it under a new consumer name, or resync from + `Watermark::START`. `EvidenceRow::History(HistoryEntry)` is the typed prompt row and `EvidenceRow::Untyped` marks a kind with none. `Change` gains `columns: Option` — every column but `revision`, in table order, values as @@ -197,7 +202,8 @@ Notable changes to the native `ai-hist` CLI are documented here. without a code change; the export journal's payload is built from the same column list — and `key: Vec`, the record's identity as the journal keys it (`["history", source, timestamp_ms, prompt]`), on - upserts and tombstones alike. `Change::source` is `Option` and + upserts and tombstones alike; a prompt's session is not part of its key, so + a prompt gaining one is an upsert, never a delete. `Change::source` is `Option` and `Change::source_name` holds the stored name, so a row from a source this build does not know is carried instead of failing the drain. The `history` FTS update trigger fires only on the columns it indexes. diff --git a/crates/ai-hist/src/change_feed.rs b/crates/ai-hist/src/change_feed.rs index 6c3cc37c..7ff4e3ab 100644 --- a/crates/ai-hist/src/change_feed.rs +++ b/crates/ai-hist/src/change_feed.rs @@ -409,7 +409,7 @@ struct FedTable { /// A column, or a quoted literal for a table that stores no source. source: &'static str, session: &'static str, - /// Whether `session` may be NULL; a tombstone stores it as `''`. + /// Whether `session` may be NULL; a change reports NULL as `''`. optional_session: bool, key: &'static [&'static str], record: &'static [&'static str], @@ -429,7 +429,7 @@ impl FedTable { } } - /// The session as SQL over `row`, as a tombstone stores it. + /// The session as SQL over `row`, as an upsert reports it. fn session_sql(&self, row: &str) -> String { if self.optional_session { format!("COALESCE({row}.{}, '')", self.session) @@ -438,6 +438,25 @@ impl FedTable { } } + /// Whether the session is part of the record's identity. A prompt's is + /// not: `history` is unique on source, time and text, and a prompt that + /// gains or changes its session is the same record. + fn session_keyed(&self) -> bool { + self.key.contains(&self.session) + } + + /// The session as a tombstone stores it: the row's, when the session is + /// part of the identity, else `''`, so a tombstone names exactly the + /// record's key and a later insert of that key clears it whatever + /// session it carries. + fn tombstone_session_sql(&self, row: &str) -> String { + if self.session_keyed() { + format!("{row}.{}", self.session) + } else { + "''".to_string() + } + } + /// The `record_key` as SQL over `row`. fn record_key_sql(&self, row: &str) -> String { match self.record { @@ -459,7 +478,9 @@ impl FedTable { if !self.source.starts_with('\'') { changed.push(format!("OLD.{0} IS NOT NEW.{0}", self.source)); } - changed.push(format!("OLD.{0} IS NOT NEW.{0}", self.session)); + if self.session_keyed() { + changed.push(format!("OLD.{0} IS NOT NEW.{0}", self.session)); + } changed.push(format!( "{} IS NOT {}", self.record_key_sql("OLD"), @@ -666,7 +687,9 @@ pub struct Change { #[serde(default)] pub source_name: String, /// For a relationship, the parent session; for a trajectory, its id; for - /// a prompt that names no session, empty. + /// a prompt that names no session, empty. A prompt's session is not part + /// of its identity, so a prompt's delete carries it empty too, and a + /// prompt gaining a session is an upsert, never a delete. pub session_id: String, /// The record's identity within its source, session and kind: the one /// identity column's text (`event_uid`, `tool_use_id`, `marker_uid`, @@ -1603,8 +1626,8 @@ pub(crate) fn init_schema(conn: &Connection) -> Result<()> { let kind = kind.as_str(); let new_source = table.source_sql("NEW"); let old_source = table.source_sql("OLD"); - let new_session = table.session_sql("NEW"); - let old_session = table.session_sql("OLD"); + let new_session = table.tombstone_session_sql("NEW"); + let old_session = table.tombstone_session_sql("OLD"); let new_record = table.record_key_sql("NEW"); let old_record = table.record_key_sql("OLD"); let moved = table.identity_changed_sql(); @@ -3403,6 +3426,12 @@ UPDATE observation_evidence SET payload_json = '{"a":2}'; ) .unwrap(); assert_parity(&store, &conn, "updates"); + assert!( + !all(&store) + .iter() + .any(|change| change.kind == ChangeKind::History && change.op == ChangeOp::Delete), + "a prompt gaining a session is an upsert of the same record, never a delete" + ); conn.execute_batch( "DELETE FROM history WHERE prompt = 'hello'; diff --git a/docs/sourcing-sdk.md b/docs/sourcing-sdk.md index 2cf55179..3a58c9db 100644 --- a/docs/sourcing-sdk.md +++ b/docs/sourcing-sdk.md @@ -384,7 +384,9 @@ upsert and on the delete that retracts it — `["history", "claude", the part of `key` inside a session: the one identity column's text, or a JSON array of several (`[1756634400000,"ship the feed"]` for a prompt). `session_id` is the parent session for a relationship, the id for a trajectory -and empty for a prompt that names no session. `source` is `None` for a source +and empty for a prompt that names no session; a prompt's session is not part of +its key, so its delete carries it empty and a prompt gaining a session is an +upsert, never a delete. `source` is `None` for a source this build does not know — a row written by a newer release — and `source_name` is the stored name either way, so such a row is carried rather than failing the drain. The feed applies no consent or exclusion rule: an