Skip to main content

khive_db/stores/
event.rs

1//! SQL-backed `EventStore` implementation.
2//!
3//! FILE SIZE JUSTIFICATION: Event store covers append, query-by-filter,
4//! observation recording, and paginated listing with shared row-mapping and
5//! timestamp serialization helpers. The event schema has complex JSON data
6//! columns (observations, referent kinds, outcomes) whose parsing is shared
7//! across all read paths, making a split impractical without duplicating the
8//! deserialization logic.
9
10use std::sync::Arc;
11
12use async_trait::async_trait;
13use uuid::Uuid;
14
15use khive_storage::error::StorageError;
16#[cfg(test)]
17use khive_storage::error::WriterTaskRequestState;
18use khive_storage::event::{
19    Event, EventAppendDisposition, EventFilter, EventObservation, IdempotentEventBatchResult,
20    ObservationRole, ReferentKind,
21};
22use khive_storage::types::{BatchWriteSummary, Page, PageRequest, SqlStatement, SqlValue};
23use khive_storage::EventStore;
24use khive_storage::SqlWriter;
25use khive_storage::StorageCapability;
26use khive_types::{EventKind, EventOutcome, SubstrateKind};
27
28use crate::pool::ConnectionPool;
29use crate::writer_task::WriterTaskHandle;
30
31#[path = "event_cursor.rs"]
32mod cursor;
33
34fn map_err(e: rusqlite::Error, op: &'static str) -> StorageError {
35    StorageError::driver(StorageCapability::Events, op, e)
36}
37
38// Preserve the existing error-classification fixtures at their original seam.
39#[cfg(test)]
40fn mark_unknown_append_usage(error: &StorageError) {
41    khive_storage::usage::account_event_write(Err(error));
42}
43
44/// An EventStore backed by SQLite tables.
45pub struct SqlEventStore {
46    pool: Arc<ConnectionPool>,
47    is_file_backed: bool,
48    namespace: String,
49    writer_task: Option<WriterTaskHandle>,
50    #[cfg(test)]
51    lose_next_fallback_reply: std::sync::atomic::AtomicBool,
52}
53
54impl SqlEventStore {
55    /// Create a new store scoped to one namespace.
56    pub fn new_scoped(
57        pool: Arc<ConnectionPool>,
58        is_file_backed: bool,
59        namespace: impl Into<String>,
60    ) -> Self {
61        // Enabled by default for file-backed pools; explicit off/degraded
62        // construction remains synchronous (ADR-067 Component A, mirrors
63        // entity.rs policy): a missing writer task is cached without failing
64        // construction. Every write re-resolves it and applies
65        // strict/compatibility policy then.
66        let writer_task = pool.writer_task_handle().ok().flatten();
67        Self {
68            pool,
69            is_file_backed,
70            namespace: namespace.into(),
71            writer_task,
72            #[cfg(test)]
73            lose_next_fallback_reply: std::sync::atomic::AtomicBool::new(false),
74        }
75    }
76
77    #[cfg(test)]
78    fn after_fallback_write_for_test<T>(
79        &self,
80        result: Result<T, StorageError>,
81    ) -> Result<T, StorageError> {
82        // Lose only the reply: the real fallback transaction already committed.
83        if result.is_ok()
84            && self
85                .lose_next_fallback_reply
86                .swap(false, std::sync::atomic::Ordering::SeqCst)
87        {
88            return Err(StorageError::writer_task_terminated(
89                WriterTaskRequestState::SideEffectsUnknown,
90            ));
91        }
92        result
93    }
94
95    fn current_writer_task(
96        &self,
97        operation: &'static str,
98    ) -> Result<Option<WriterTaskHandle>, StorageError> {
99        self.pool
100            .writer_task_for_write(self.writer_task.as_ref(), operation)
101    }
102
103    /// Route a single-row write through the pool-wide `WriterTask` when
104    /// the write queue is enabled and a handle is available. Strict mode
105    /// refuses a missing handle; compatibility mode falls back to the legacy
106    /// standalone-connection / pool-mutex path (ADR-067 Component A, Fork C
107    /// slice 2).
108    ///
109    /// `append_event`/`append_events` perform the same write-time lookup
110    /// first; a non-strict `None` then falls through this helper, which
111    /// records the actual compatibility fallback. Strict mode returns before
112    /// the direct-writer seam.
113    /// `f` must be DML-only on the flag-on path (no bare `BEGIN IMMEDIATE`)
114    /// since it runs inside the WriterTask's own transaction.
115    async fn with_writer<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
116    where
117        F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
118        R: Send + 'static,
119    {
120        if let Some(writer_task) = self.current_writer_task(op)? {
121            return writer_task
122                .send_bounded(move |conn| f(conn).map_err(|e| map_err(e, op)))
123                .await;
124        }
125
126        self.pool
127            .record_direct_route(crate::timeout_sink::Site::DirectRouteEventGeneralWrite);
128        // Atomic SQL units release an in-memory connection between statements;
129        // share their unit budget so this transaction cannot overlap one.
130        let unit_slot = if self.is_file_backed {
131            None
132        } else {
133            Some(crate::sql_bridge::acquire_in_memory_write_unit(&self.pool, op).await?)
134        };
135        let pool = Arc::clone(&self.pool);
136        let db = crate::timeout_sink::db_label(&pool);
137        let is_file_backed = self.is_file_backed;
138        tokio::task::spawn_blocking(move || {
139            // Cancellation of the caller must not release the slot before this job ends.
140            let _unit_slot = unit_slot;
141            let result =
142                pool.execute_direct_transaction(StorageCapability::Events, op, move |conn| {
143                    f(conn).map_err(|error| map_err(error, op))
144                });
145            if is_file_backed {
146                if let Err(error) = &result {
147                    crate::timeout_sink::maybe_emit_busy_storage_error(
148                        &db,
149                        crate::timeout_sink::Site::StandaloneEvent,
150                        error,
151                    );
152                }
153            }
154            result
155        })
156        .await
157        .map_err(|e| StorageError::driver(StorageCapability::Events, op, e))?
158    }
159
160    async fn with_reader<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
161    where
162        F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
163        R: Send + 'static,
164    {
165        super::run_pooled_store_read(
166            Arc::clone(&self.pool),
167            StorageCapability::Events,
168            op,
169            move |conn| f(conn).map_err(|error| map_err(error, op)),
170        )
171        .await
172    }
173}
174
175// =============================================================================
176// Helpers: parse SubstrateKind / EventOutcome / EventKind from DB strings
177// =============================================================================
178
179fn substrate_from_str(s: &str) -> Result<SubstrateKind, rusqlite::Error> {
180    s.parse::<SubstrateKind>().map_err(|_| {
181        rusqlite::Error::FromSqlConversionFailure(
182            0,
183            rusqlite::types::Type::Text,
184            format!("unknown SubstrateKind: {s}").into(),
185        )
186    })
187}
188
189fn outcome_from_str(s: &str) -> Result<EventOutcome, rusqlite::Error> {
190    match s {
191        "success" => Ok(EventOutcome::Success),
192        "denied" => Ok(EventOutcome::Denied),
193        "error" => Ok(EventOutcome::Error),
194        other => Err(rusqlite::Error::FromSqlConversionFailure(
195            0,
196            rusqlite::types::Type::Text,
197            format!("unknown EventOutcome: {other}").into(),
198        )),
199    }
200}
201
202fn kind_from_str(s: &str) -> Result<EventKind, rusqlite::Error> {
203    s.parse::<EventKind>().map_err(|_| {
204        rusqlite::Error::FromSqlConversionFailure(
205            0,
206            rusqlite::types::Type::Text,
207            format!("unknown EventKind: {s}").into(),
208        )
209    })
210}
211
212fn referent_kind_from_str(s: &str) -> Result<ReferentKind, rusqlite::Error> {
213    match s {
214        "entity" => Ok(ReferentKind::Entity),
215        "note" => Ok(ReferentKind::Note),
216        "edge" => Ok(ReferentKind::Edge),
217        other => Err(rusqlite::Error::FromSqlConversionFailure(
218            0,
219            rusqlite::types::Type::Text,
220            format!("unknown ReferentKind: {other}").into(),
221        )),
222    }
223}
224
225fn observation_role_from_str(s: &str) -> Result<ObservationRole, rusqlite::Error> {
226    match s {
227        "candidate" => Ok(ObservationRole::Candidate),
228        "selected" => Ok(ObservationRole::Selected),
229        "target" => Ok(ObservationRole::Target),
230        "signal" => Ok(ObservationRole::Signal),
231        other => Err(rusqlite::Error::FromSqlConversionFailure(
232            0,
233            rusqlite::types::Type::Text,
234            format!("unknown ObservationRole: {other}").into(),
235        )),
236    }
237}
238
239fn parse_uuid(s: &str) -> Result<Uuid, rusqlite::Error> {
240    Uuid::parse_str(s).map_err(|e| {
241        rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(e))
242    })
243}
244
245// Column order: id(0), namespace(1), verb(2), substrate(3), actor(4),
246//               kind(5), outcome(6), payload(7), payload_schema_version(8),
247//               profile_state_version(9), duration_us(10), target_id(11),
248//               session_id(12), aggregate_kind(13), aggregate_id(14), created_at(15),
249//               op_index(16), ref_resolution(17)
250fn read_event(row: &rusqlite::Row<'_>) -> Result<Event, rusqlite::Error> {
251    let id_str: String = row.get(0)?;
252    let namespace: String = row.get(1)?;
253    let verb: String = row.get(2)?;
254    let substrate_str: String = row.get(3)?;
255    let actor: String = row.get(4)?;
256    let kind_str: String = row.get(5)?;
257    let outcome_str: String = row.get(6)?;
258    let payload_str: String = row.get(7)?;
259    let payload_schema_version: i64 = row.get(8)?;
260    let profile_state_version: Option<i64> = row.get(9)?;
261    let duration_us: i64 = row.get(10)?;
262    let target_str: Option<String> = row.get(11)?;
263    let session_str: Option<String> = row.get(12)?;
264    let aggregate_kind: Option<String> = row.get(13)?;
265    let aggregate_str: Option<String> = row.get(14)?;
266    let created_at: i64 = row.get(15)?;
267    let op_index = row.get::<_, Option<u32>>(16)?;
268    let ref_resolution = row
269        .get::<_, Option<String>>(17)?
270        .map(|value| {
271            value
272                .parse::<khive_types::RefResolution>()
273                .map_err(|error| {
274                    rusqlite::Error::FromSqlConversionFailure(
275                        17,
276                        rusqlite::types::Type::Text,
277                        Box::new(error),
278                    )
279                })
280        })
281        .transpose()?;
282    if op_index.is_some() != ref_resolution.is_some() {
283        return Err(rusqlite::Error::FromSqlConversionFailure(
284            16,
285            rusqlite::types::Type::Integer,
286            "event operation attribution must be present or absent together".into(),
287        ));
288    }
289
290    let id = parse_uuid(&id_str)?;
291    let substrate = substrate_from_str(&substrate_str)?;
292    let kind = kind_from_str(&kind_str)?;
293    let outcome = outcome_from_str(&outcome_str)?;
294    let payload: serde_json::Value = serde_json::from_str(&payload_str).map_err(|e| {
295        rusqlite::Error::FromSqlConversionFailure(7, rusqlite::types::Type::Text, Box::new(e))
296    })?;
297    let target_id = target_str.as_deref().map(parse_uuid).transpose()?;
298    let session_id = session_str.as_deref().map(parse_uuid).transpose()?;
299    let aggregate_id = aggregate_str.as_deref().map(parse_uuid).transpose()?;
300    let payload_schema_version_u32: u32 = payload_schema_version.try_into().map_err(|_| {
301        rusqlite::Error::FromSqlConversionFailure(
302            8,
303            rusqlite::types::Type::Integer,
304            format!("payload_schema_version {payload_schema_version} out of u32 range").into(),
305        )
306    })?;
307    let profile_state_version_u64: Option<u64> = profile_state_version
308        .map(|v| {
309            u64::try_from(v).map_err(|_| {
310                rusqlite::Error::FromSqlConversionFailure(
311                    9,
312                    rusqlite::types::Type::Integer,
313                    format!("profile_state_version {v} out of u64 range").into(),
314                )
315            })
316        })
317        .transpose()?;
318
319    Ok(Event {
320        id,
321        namespace,
322        verb,
323        substrate,
324        actor,
325        kind,
326        outcome,
327        payload,
328        payload_schema_version: payload_schema_version_u32,
329        profile_state_version: profile_state_version_u64,
330        duration_us,
331        target_id,
332        session_id,
333        aggregate_kind,
334        aggregate_id,
335        created_at,
336        op_index,
337        ref_resolution,
338    })
339}
340
341// =============================================================================
342// Helpers: observation projection write path
343// =============================================================================
344
345fn insert_event_with_observations(
346    conn: &rusqlite::Connection,
347    event: &Event,
348) -> Result<(), rusqlite::Error> {
349    validate_operation_pair(event)?;
350    let id_str = event.id.to_string();
351    let substrate_str = event.substrate.name().to_string();
352    let kind_str = event.kind.name().to_string();
353    let outcome_str = event.outcome.name().to_string();
354    let payload_str = event.payload.to_string();
355    let target_str = event.target_id.map(|u| u.to_string());
356    let session_str = event.session_id.map(|u| u.to_string());
357    let aggregate_str = event.aggregate_id.map(|u| u.to_string());
358    let profile_state_version = profile_state_version_to_sql(event)?;
359
360    conn.execute(
361        "INSERT INTO events \
362         (id, namespace, verb, substrate, actor, kind, outcome, payload, payload_schema_version, \
363          profile_state_version, duration_us, target_id, session_id, aggregate_kind, aggregate_id, created_at, op_index, ref_resolution) \
364         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18)",
365        rusqlite::params![
366            id_str,
367            &event.namespace,
368            &event.verb,
369            substrate_str,
370            &event.actor,
371            kind_str,
372            outcome_str,
373            payload_str,
374            event.payload_schema_version as i64,
375            profile_state_version,
376            event.duration_us,
377            target_str,
378            session_str,
379            &event.aggregate_kind,
380            aggregate_str,
381            event.created_at,
382            event.op_index,
383            event.ref_resolution.map(|value| value.name()),
384        ],
385    )?;
386
387    for observation in decode_event_observations(event)? {
388        conn.execute(
389            "INSERT INTO event_observations \
390             (event_id, entity_id, referent_kind, role, position) \
391             VALUES (?1, ?2, ?3, ?4, ?5)",
392            rusqlite::params![
393                observation.event_id.to_string(),
394                observation.entity_id.to_string(),
395                observation.referent_kind.name(),
396                observation.role.name(),
397                observation.position as i64,
398            ],
399        )?;
400    }
401
402    Ok(())
403}
404
405/// Append an event on the caller's existing SQLite transaction, including the
406/// same observation projection as `SqlEventStore::append_event`. This does
407/// not open or commit a transaction: the caller must roll back its domain
408/// changes if this insert fails.
409pub fn append_event_in_transaction(
410    conn: &rusqlite::Connection,
411    event: &Event,
412) -> Result<(), rusqlite::Error> {
413    insert_event_with_observations(conn, event)
414}
415
416/// DML-only batch append loop shared by both the legacy (flag-off) and
417/// WriterTask-routed (flag-on) `append_events` paths (ADR-067 Component A).
418///
419/// Issues no `BEGIN` / `COMMIT` / `ROLLBACK` itself — the caller owns the
420/// enclosing transaction. All-or-nothing: the first failed insert returns
421/// `Err` immediately (matching the pre-existing `append_events` contract) —
422/// the caller's transaction wrapper issues the ROLLBACK.
423fn batch_append_events_dml(
424    conn: &rusqlite::Connection,
425    events: &[Event],
426    attempted: u64,
427) -> Result<BatchWriteSummary, rusqlite::Error> {
428    let mut affected = 0u64;
429    for event in events {
430        insert_event_with_observations(conn, event)?;
431        affected += 1;
432    }
433    Ok(BatchWriteSummary {
434        attempted,
435        affected,
436        ..BatchWriteSummary::default()
437    })
438}
439
440fn fetch_event_by_id(
441    conn: &rusqlite::Connection,
442    id: Uuid,
443) -> Result<Option<Event>, rusqlite::Error> {
444    let id_str = id.to_string();
445    let mut stmt = conn.prepare(
446        "SELECT id, namespace, verb, substrate, actor, kind, outcome, payload, \
447                payload_schema_version, profile_state_version, duration_us, target_id, \
448                session_id, aggregate_kind, aggregate_id, created_at, op_index, ref_resolution \
449         FROM events WHERE id = ?1",
450    )?;
451    let mut rows = stmt.query(rusqlite::params![id_str])?;
452    match rows.next()? {
453        Some(row) => Ok(Some(read_event(row)?)),
454        None => Ok(None),
455    }
456}
457
458/// Fetch `event_id`'s persisted observation projection in the same order it
459/// was written: `event_observations`'s primary key is `(event_id, role,
460/// position)`, and every `decode_*_observations` builder emits one role
461/// group at a time in a fixed sequence that sorts identically by `role`
462/// text, so `ORDER BY role, position` reproduces the original emission
463/// order without a separate sequence column.
464fn fetch_event_observations(
465    conn: &rusqlite::Connection,
466    event_id: Uuid,
467) -> Result<Vec<EventObservation>, rusqlite::Error> {
468    let id_str = event_id.to_string();
469    let mut stmt = conn.prepare(
470        "SELECT event_id, entity_id, referent_kind, role, position \
471         FROM event_observations WHERE event_id = ?1 ORDER BY role, position",
472    )?;
473    let rows = stmt.query_map(rusqlite::params![id_str], |row| {
474        let event_id: String = row.get(0)?;
475        let entity_id: String = row.get(1)?;
476        let referent_kind: String = row.get(2)?;
477        let role: String = row.get(3)?;
478        let position: i64 = row.get(4)?;
479        Ok((event_id, entity_id, referent_kind, role, position))
480    })?;
481
482    let mut out = Vec::new();
483    for row in rows {
484        let (event_id, entity_id, referent_kind, role, position) = row?;
485        let position_u32: u32 = position.try_into().map_err(|_| {
486            rusqlite::Error::FromSqlConversionFailure(
487                4,
488                rusqlite::types::Type::Integer,
489                format!("position {position} out of u32 range").into(),
490            )
491        })?;
492        out.push(EventObservation {
493            event_id: parse_uuid(&event_id)?,
494            entity_id: parse_uuid(&entity_id)?,
495            referent_kind: referent_kind_from_str(&referent_kind)?,
496            role: observation_role_from_str(&role)?,
497            position: position_u32,
498        });
499    }
500    Ok(out)
501}
502
503/// DML-only idempotent batch loop shared by both the writer-task-routed and
504/// standalone-transaction `append_events_idempotent` paths, mirroring
505/// [`batch_append_events_dml`]'s split. A row whose id is unseen is
506/// inserted; a row whose id already exists is compared column-for-column
507/// plus its ordered observation projection against the submitted event —
508/// exact equality reports `AlreadyPresentIdentical` and skips the write,
509/// any mismatch reports `IdentityConflict` and skips the write for that row
510/// alone. A genuine store failure on a fresh insert still aborts the whole
511/// batch via `Err`, matching `append_events`'s all-or-nothing contract; only
512/// identity comparison is per-row.
513fn idempotent_batch_dml(
514    conn: &rusqlite::Connection,
515    events: &[Event],
516) -> Result<IdempotentEventBatchResult, rusqlite::Error> {
517    let mut rows = Vec::with_capacity(events.len());
518    for event in events {
519        validate_operation_pair(event)?;
520        match fetch_event_by_id(conn, event.id)? {
521            None => {
522                insert_event_with_observations(conn, event)?;
523                rows.push(EventAppendDisposition::Inserted);
524            }
525            Some(existing) => {
526                let existing_observations = fetch_event_observations(conn, event.id)?;
527                let submitted_observations = decode_event_observations(event)?;
528                if existing == *event && existing_observations == submitted_observations {
529                    rows.push(EventAppendDisposition::AlreadyPresentIdentical);
530                } else {
531                    rows.push(EventAppendDisposition::IdentityConflict);
532                }
533            }
534        }
535    }
536    Ok(IdempotentEventBatchResult { rows })
537}
538
539/// Pure statement builder (ADR-099 B3 r6 structural cut): the exact
540/// `events` row insert plus every derived `event_observations` insert for
541/// `event`, as plain [`SqlStatement`]s with no I/O. This is the ONE place
542/// that knows the `events`/`event_observations` insert shape at the
543/// statement-text level; both [`append_event_on_writer`] below (the async
544/// `SqlWriter` execution path) and `khive-runtime`'s ADR-099 `--atomic`
545/// prepare path (which needs plain-data statements it can fold into a
546/// synchronous commit-phase plan, not an async writer call) build on this
547/// function rather than each hand-writing the INSERT text — the divergence
548/// that produced the drift this cut fixes.
549pub fn event_insert_statements(event: &Event) -> Result<Vec<SqlStatement>, rusqlite::Error> {
550    validate_operation_pair(event)?;
551    let id_str = event.id.to_string();
552    let substrate_str = event.substrate.name().to_string();
553    let kind_str = event.kind.name().to_string();
554    let outcome_str = event.outcome.name().to_string();
555    let payload_str = event.payload.to_string();
556    let target_str = event.target_id.map(|u| u.to_string());
557    let session_str = event.session_id.map(|u| u.to_string());
558    let aggregate_str = event.aggregate_id.map(|u| u.to_string());
559    let profile_state_version = profile_state_version_to_sql(event)?;
560
561    let mut statements = vec![SqlStatement {
562        sql: "INSERT INTO events \
563              (id, namespace, verb, substrate, actor, kind, outcome, payload, payload_schema_version, \
564               profile_state_version, duration_us, target_id, session_id, aggregate_kind, aggregate_id, created_at, op_index, ref_resolution) \
565              VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18)"
566            .into(),
567        params: vec![
568            SqlValue::Text(id_str),
569            SqlValue::Text(event.namespace.clone()),
570            SqlValue::Text(event.verb.clone()),
571            SqlValue::Text(substrate_str),
572            SqlValue::Text(event.actor.clone()),
573            SqlValue::Text(kind_str),
574            SqlValue::Text(outcome_str),
575            SqlValue::Text(payload_str),
576            SqlValue::Integer(event.payload_schema_version as i64),
577            profile_state_version
578                .map(SqlValue::Integer)
579                .unwrap_or(SqlValue::Null),
580            SqlValue::Integer(event.duration_us),
581            target_str.map(SqlValue::Text).unwrap_or(SqlValue::Null),
582            session_str.map(SqlValue::Text).unwrap_or(SqlValue::Null),
583            event
584                .aggregate_kind
585                .clone()
586                .map(SqlValue::Text)
587                .unwrap_or(SqlValue::Null),
588            aggregate_str.map(SqlValue::Text).unwrap_or(SqlValue::Null),
589            SqlValue::Integer(event.created_at),
590            event.op_index.map(|value| SqlValue::Integer(i64::from(value))).unwrap_or(SqlValue::Null),
591            event.ref_resolution.map(|value| SqlValue::Text(value.name().into())).unwrap_or(SqlValue::Null),
592        ],
593        label: Some("event_insert_on_writer".into()),
594    }];
595
596    for observation in decode_event_observations(event)? {
597        statements.push(SqlStatement {
598            sql: "INSERT INTO event_observations \
599                  (event_id, entity_id, referent_kind, role, position) \
600                  VALUES (?1, ?2, ?3, ?4, ?5)"
601                .into(),
602            params: vec![
603                SqlValue::Text(observation.event_id.to_string()),
604                SqlValue::Text(observation.entity_id.to_string()),
605                SqlValue::Text(observation.referent_kind.name().to_string()),
606                SqlValue::Text(observation.role.name().to_string()),
607                SqlValue::Integer(observation.position as i64),
608            ],
609            label: Some("event_observation_insert_on_writer".into()),
610        });
611    }
612
613    Ok(statements)
614}
615
616fn validate_operation_pair(event: &Event) -> Result<(), rusqlite::Error> {
617    if event.op_index.is_some() != event.ref_resolution.is_some() {
618        return Err(rusqlite::Error::ToSqlConversionFailure(
619            "event operation attribution must be present or absent together".into(),
620        ));
621    }
622    profile_state_version_to_sql(event)?;
623    Ok(())
624}
625
626fn profile_state_version_to_sql(event: &Event) -> Result<Option<i64>, rusqlite::Error> {
627    event
628        .profile_state_version
629        .map(|version| {
630            i64::try_from(version).map_err(|_| {
631                rusqlite::Error::ToSqlConversionFailure(
632                    format!("profile_state_version {version} exceeds i64::MAX").into(),
633                )
634            })
635        })
636        .transpose()
637}
638
639/// Build commit-time warning inserts for lineage-sensitive incident edges
640/// removed by a hard-delete cascade (ADR-002).
641///
642/// Each statement handles one protected relation and reads `graph_edges`
643/// when it executes, inside the caller's delete transaction. This is
644/// intentionally not a prepare-time edge query: a guarded edge write that
645/// commits immediately before the delete must be represented in the warning
646/// payload before the subsequent cascade statement removes it. Relations
647/// with no matching incident edge insert no event.
648pub fn hard_delete_lineage_warning_statements(
649    namespace: &str,
650    actor: &str,
651    target_id: Uuid,
652    substrate: SubstrateKind,
653) -> Vec<SqlStatement> {
654    const WARNINGS: [(&str, &str); 5] = [
655        ("derived_from", "provenance_loss"),
656        ("supersedes", "replacement_lineage_loss"),
657        ("precedes", "temporal_sequence_loss"),
658        ("supports", "evidential_link_loss"),
659        ("refutes", "evidential_link_loss"),
660    ];
661
662    let target_id = target_id.to_string();
663    let created_at = chrono::Utc::now().timestamp_micros();
664    let operation = khive_storage::operation_context::current_operation_attribution();
665    WARNINGS
666        .into_iter()
667        .map(|(relation, warning)| SqlStatement {
668            sql: "INSERT INTO events \
669                  (id, namespace, verb, substrate, actor, kind, outcome, payload, \
670                   payload_schema_version, profile_state_version, duration_us, target_id, \
671                   session_id, aggregate_kind, aggregate_id, created_at, op_index, ref_resolution) \
672                  SELECT ?1, ?2, 'delete', ?3, ?4, ?5, ?6, \
673                         json_object( \
674                           'severity', 'warning', \
675                           'warning', ?7, \
676                           'deleted_id', ?8, \
677                           'relation', ?9, \
678                           'edge_count', count(*), \
679                           'edges', json_group_array(json_object( \
680                             'id', incident.id, \
681                             'namespace', incident.namespace, \
682                             'source_id', incident.source_id, \
683                             'target_id', incident.target_id, \
684                             'deleted_at', incident.deleted_at)) \
685                         ), \
686                         1, NULL, 0, ?8, NULL, NULL, NULL, ?10, ?11, ?12 \
687                  FROM ( \
688                    SELECT id, namespace, source_id, target_id, deleted_at \
689                    FROM graph_edges \
690                    WHERE (source_id = ?8 OR target_id = ?8) AND relation = ?9 \
691                    ORDER BY namespace, id \
692                  ) AS incident \
693                  HAVING count(*) > 0"
694                .into(),
695            params: vec![
696                SqlValue::Text(Uuid::new_v4().to_string()),
697                SqlValue::Text(namespace.to_string()),
698                SqlValue::Text(substrate.name().to_string()),
699                SqlValue::Text(actor.to_string()),
700                SqlValue::Text(EventKind::Audit.name().to_string()),
701                SqlValue::Text(EventOutcome::Success.name().to_string()),
702                SqlValue::Text(warning.to_string()),
703                SqlValue::Text(target_id.clone()),
704                SqlValue::Text(relation.to_string()),
705                SqlValue::Integer(created_at),
706                operation
707                    .map(|value| SqlValue::Integer(i64::from(value.op_index)))
708                    .unwrap_or(SqlValue::Null),
709                operation
710                    .map(|value| SqlValue::Text(value.ref_resolution.name().into()))
711                    .unwrap_or(SqlValue::Null),
712            ],
713            label: Some(format!("hard-delete-{relation}-warning")),
714        })
715        .collect()
716}
717
718/// Insert `event` (and any derived `event_observations` rows) through the
719/// `khive-storage` `SqlWriter` seam, on a transaction the CALLER already
720/// opened and controls the commit/rollback boundary for. Issues no
721/// `BEGIN`/`COMMIT`/`ROLLBACK` of its own.
722///
723/// This exists alongside `insert_event_with_observations` (the raw
724/// `rusqlite::Connection` path `SqlEventStore::append_event` uses, whose
725/// caller now owns an admitted transaction for the ordinary case) rather
726/// than replacing it: the two run on
727/// different connection abstractions — a standalone `rusqlite::Connection`
728/// vs. a `Box<dyn SqlWriter>` — so they cannot share one function body. Both
729/// build on [`event_insert_statements`] for the actual insert shape.
730/// Callers that need the event append to be part of a larger atomic unit —
731/// e.g. ADR-081's brain fold gate (`khive-pack-brain/src/fold_gate.rs`),
732/// which holds its own `BEGIN IMMEDIATE` transaction on a `SqlWriter` for
733/// the dedup claim + mass fold and needs the feedback event to land in that
734/// same transaction — call this instead of duplicating the insert shape
735/// into their own crate.
736pub async fn append_event_on_writer(
737    writer: &mut dyn SqlWriter,
738    event: &Event,
739) -> Result<(), StorageError> {
740    let statements =
741        event_insert_statements(event).map_err(|e| map_err(e, "decode_event_observations"))?;
742    for statement in statements {
743        writer.execute(statement).await?;
744    }
745    Ok(())
746}
747
748fn decode_event_observations(event: &Event) -> Result<Vec<EventObservation>, rusqlite::Error> {
749    match event.kind {
750        EventKind::RerankExecuted => decode_rerank_observations(event),
751        EventKind::RecallExecuted => decode_recall_observations(event),
752        EventKind::SearchExecuted => decode_search_observations(event),
753        EventKind::LinkCreated => decode_link_observations(event),
754        EventKind::EdgeUpdated | EventKind::EdgeDeleted => decode_edge_target_observation(event),
755        EventKind::EntityCreated
756        | EventKind::EntityUpdated
757        | EventKind::EntityDeleted
758        | EventKind::NoteCreated
759        | EventKind::NoteUpdated
760        | EventKind::NoteDeleted
761        | EventKind::TaskTransitioned => decode_target_observation(event),
762        EventKind::FeedbackExplicit => decode_signal_observation(event),
763        _ => Ok(Vec::new()),
764    }
765}
766
767fn payload_uuid_array(event: &Event, field: &'static str) -> Result<Vec<Uuid>, rusqlite::Error> {
768    Ok(payload_uuid_array_opt(event, field)?.unwrap_or_default())
769}
770
771/// Like [`payload_uuid_array`], but distinguishes an absent `field` (`Ok(None)`)
772/// from a present-but-malformed one (`Err`). Callers that fall back across a
773/// chain of alternative field names (e.g. `decode_rerank_observations`'s
774/// `final_scores`/`reranked`) need that distinction: collapsing
775/// "missing" into `Ok(vec![])` makes every later field in the chain
776/// unreachable, since `Result::or_else` only fires on `Err`.
777fn payload_uuid_array_opt(
778    event: &Event,
779    field: &'static str,
780) -> Result<Option<Vec<Uuid>>, rusqlite::Error> {
781    let Some(values) = event.payload.get(field) else {
782        return Ok(None);
783    };
784    let Some(array) = values.as_array() else {
785        return Err(invalid_payload(event.kind, field, "expected array"));
786    };
787
788    array
789        .iter()
790        .map(|value| {
791            value
792                .as_str()
793                .ok_or_else(|| invalid_payload(event.kind, field, "expected UUID string"))
794                .and_then(|s| Uuid::parse_str(s).map_err(|e| invalid_payload(event.kind, field, e)))
795        })
796        .collect::<Result<Vec<_>, _>>()
797        .map(Some)
798}
799
800/// Decode `RerankExecutedPayload::final_scores` (`Vec<(Id128, f32)>` per
801/// `khive_types::event::RerankExecutedPayload`) into the UUID leading each
802/// tuple. Deserializes through that exact typed tuple shape, so a tuple that
803/// is not precisely two elements, or whose second element is not a finite
804/// score, is rejected rather than silently accepted. Absent field is
805/// `Ok(None)`; present-but-wrong-shape is `Err`, matching
806/// [`payload_uuid_array_opt`]'s missing-vs-malformed contract.
807fn payload_final_scores_uuid_array_opt(
808    event: &Event,
809    field: &'static str,
810) -> Result<Option<Vec<Uuid>>, rusqlite::Error> {
811    let Some(values) = event.payload.get(field) else {
812        return Ok(None);
813    };
814    let tuples: Vec<(khive_types::Id128, f32)> = serde_json::from_value(values.clone())
815        .map_err(|e| invalid_payload(event.kind, field, e))?;
816    tuples
817        .into_iter()
818        .map(|(id, score)| {
819            if !score.is_finite() {
820                return Err(invalid_payload(event.kind, field, "score is not finite"));
821            }
822            Ok(Uuid::from_bytes(*id.as_bytes()))
823        })
824        .collect::<Result<Vec<_>, _>>()
825        .map(Some)
826}
827
828/// Decode `RerankExecutedPayload::reranked` (`Vec<(Id128, Vec<(String, f32)>)>`
829/// per `khive_types::event::RerankExecutedPayload`) into the UUID leading each
830/// tuple. Same exact-shape decoding contract as
831/// [`payload_final_scores_uuid_array_opt`], applied to `reranked`'s
832/// per-reranker sub-score shape instead of a single scalar score.
833fn payload_reranked_uuid_array_opt(
834    event: &Event,
835    field: &'static str,
836) -> Result<Option<Vec<Uuid>>, rusqlite::Error> {
837    let Some(values) = event.payload.get(field) else {
838        return Ok(None);
839    };
840    let tuples: Vec<(khive_types::Id128, Vec<(String, f32)>)> =
841        serde_json::from_value(values.clone())
842            .map_err(|e| invalid_payload(event.kind, field, e))?;
843    tuples
844        .into_iter()
845        .map(|(id, scores)| {
846            if !scores.iter().all(|(_, s)| s.is_finite()) {
847                return Err(invalid_payload(event.kind, field, "score is not finite"));
848            }
849            Ok(Uuid::from_bytes(*id.as_bytes()))
850        })
851        .collect::<Result<Vec<_>, _>>()
852        .map(Some)
853}
854
855fn payload_uuid(event: &Event, field: &'static str) -> Result<Option<Uuid>, rusqlite::Error> {
856    let Some(value) = event.payload.get(field) else {
857        return Ok(None);
858    };
859    let Some(s) = value.as_str() else {
860        return Err(invalid_payload(event.kind, field, "expected UUID string"));
861    };
862    Uuid::parse_str(s)
863        .map(Some)
864        .map_err(|e| invalid_payload(event.kind, field, e))
865}
866
867fn decode_candidate_observations(
868    event: &Event,
869    referent_kind: ReferentKind,
870) -> Result<Vec<EventObservation>, rusqlite::Error> {
871    let mut rows = Vec::new();
872
873    for (position, entity_id) in payload_uuid_array(event, "candidates")?
874        .into_iter()
875        .enumerate()
876    {
877        let position_u32 = u32::try_from(position).map_err(|_| {
878            invalid_payload(
879                event.kind,
880                "candidates[position]",
881                "position out of u32 range",
882            )
883        })?;
884        rows.push(EventObservation {
885            event_id: event.id,
886            entity_id,
887            referent_kind,
888            role: ObservationRole::Candidate,
889            position: position_u32,
890        });
891    }
892
893    Ok(rows)
894}
895
896fn push_selected_observations(
897    event: &Event,
898    selected: Vec<Uuid>,
899    referent_kind: ReferentKind,
900    rows: &mut Vec<EventObservation>,
901) -> Result<(), rusqlite::Error> {
902    for (position, entity_id) in selected.into_iter().enumerate() {
903        let position_u32 = u32::try_from(position).map_err(|_| {
904            invalid_payload(
905                event.kind,
906                "selected[position]",
907                "position out of u32 range",
908            )
909        })?;
910        rows.push(EventObservation {
911            event_id: event.id,
912            entity_id,
913            referent_kind,
914            role: ObservationRole::Selected,
915            position: position_u32,
916        });
917    }
918    Ok(())
919}
920
921/// `RecallExecuted` payloads carry flat candidate and selected note UUID lists.
922fn decode_recall_observations(event: &Event) -> Result<Vec<EventObservation>, rusqlite::Error> {
923    let mut rows = decode_candidate_observations(event, ReferentKind::Note)?;
924    let selected = payload_uuid_array_opt(event, "selected")?.unwrap_or_default();
925    push_selected_observations(event, selected, ReferentKind::Note, &mut rows)?;
926    Ok(rows)
927}
928
929/// `SearchExecuted.result_kind` identifies which substrate owns every UUID in
930/// the candidate and selected lists. A missing key is the historical note
931/// shape under ADR-041 A3; present unknown or non-string values are invalid.
932fn decode_search_observations(event: &Event) -> Result<Vec<EventObservation>, rusqlite::Error> {
933    if event.payload.as_object().is_none() {
934        return Err(invalid_payload(event.kind, "payload", "expected object"));
935    }
936    let referent_kind = match event.payload.get("result_kind") {
937        Some(serde_json::Value::String(kind)) if kind == "entity" => ReferentKind::Entity,
938        Some(serde_json::Value::String(kind)) if kind == "note" => ReferentKind::Note,
939        Some(serde_json::Value::String(_)) => {
940            return Err(invalid_payload(
941                event.kind,
942                "result_kind",
943                "expected \"entity\" or \"note\"",
944            ));
945        }
946        Some(_) => {
947            return Err(invalid_payload(
948                event.kind,
949                "result_kind",
950                "expected string \"entity\" or \"note\"",
951            ));
952        }
953        None => ReferentKind::Note,
954    };
955    let mut rows = decode_candidate_observations(event, referent_kind)?;
956    let selected = payload_uuid_array_opt(event, "selected")?.unwrap_or_default();
957    push_selected_observations(event, selected, referent_kind, &mut rows)?;
958    Ok(rows)
959}
960
961/// `RerankExecutedPayload` (`khive_types::event::RerankExecutedPayload`) has no
962/// `selected` field at all — a stray `selected` key in the raw payload is not
963/// part of its typed contract and must never be consulted. Per ADR-042 §5,
964/// `final_scores` is the ordered rerank output (positions match output
965/// order); `reranked` is per-reranker audit/debug data with no ordering
966/// guarantee, so it is only a legacy fallback for events emitted before
967/// `final_scores` existed. `reranked` is tried only when `final_scores` is
968/// absent (`None`) — a present-but-malformed `final_scores` errors
969/// immediately instead of masking the problem by falling through.
970fn decode_rerank_observations(event: &Event) -> Result<Vec<EventObservation>, rusqlite::Error> {
971    let mut rows = decode_candidate_observations(event, ReferentKind::Note)?;
972    let selected = payload_final_scores_uuid_array_opt(event, "final_scores")?
973        .or(payload_reranked_uuid_array_opt(event, "reranked")?)
974        .unwrap_or_default();
975    push_selected_observations(event, selected, ReferentKind::Note, &mut rows)?;
976
977    Ok(rows)
978}
979
980fn decode_link_observations(event: &Event) -> Result<Vec<EventObservation>, rusqlite::Error> {
981    let mut rows = Vec::new();
982    let source_kind = payload_link_referent_kind(event, "source_kind")?;
983    let target_kind = payload_link_referent_kind(event, "target_kind")?;
984    if let Some(source) = payload_uuid(event, "source_id")? {
985        if let Some(referent_kind) = source_kind {
986            rows.push(EventObservation {
987                event_id: event.id,
988                entity_id: source,
989                referent_kind,
990                role: ObservationRole::Target,
991                position: 0,
992            });
993        }
994    }
995    if let Some(target) = payload_uuid(event, "target_id")? {
996        if let Some(referent_kind) = target_kind {
997            rows.push(EventObservation {
998                event_id: event.id,
999                entity_id: target,
1000                referent_kind,
1001                role: ObservationRole::Target,
1002                position: 1,
1003            });
1004        }
1005    }
1006    if let Some(edge_id) = event.target_id.or(payload_uuid(event, "id")?) {
1007        rows.push(EventObservation {
1008            event_id: event.id,
1009            entity_id: edge_id,
1010            referent_kind: ReferentKind::Edge,
1011            role: ObservationRole::Target,
1012            position: 2,
1013        });
1014    }
1015    Ok(rows)
1016}
1017
1018fn payload_link_referent_kind(
1019    event: &Event,
1020    field: &'static str,
1021) -> Result<Option<ReferentKind>, rusqlite::Error> {
1022    match event.payload.get(field) {
1023        None => Ok(Some(ReferentKind::Entity)),
1024        Some(value) => match value.as_str() {
1025            Some("entity") => Ok(Some(ReferentKind::Entity)),
1026            Some("note") => Ok(Some(ReferentKind::Note)),
1027            Some("edge") => Ok(Some(ReferentKind::Edge)),
1028            Some("event") => Ok(None),
1029            _ => Err(invalid_payload(
1030                event.kind,
1031                field,
1032                "expected \"entity\", \"note\", \"edge\", or \"event\"",
1033            )),
1034        },
1035    }
1036}
1037
1038fn decode_edge_target_observation(event: &Event) -> Result<Vec<EventObservation>, rusqlite::Error> {
1039    let Some(edge_id) = event.target_id.or(payload_uuid(event, "id")?) else {
1040        return Ok(Vec::new());
1041    };
1042    Ok(vec![EventObservation {
1043        event_id: event.id,
1044        entity_id: edge_id,
1045        referent_kind: ReferentKind::Edge,
1046        role: ObservationRole::Target,
1047        position: 0,
1048    }])
1049}
1050
1051fn decode_target_observation(event: &Event) -> Result<Vec<EventObservation>, rusqlite::Error> {
1052    let Some(entity_id) = event.target_id.or(payload_uuid(event, "target_id")?) else {
1053        return Ok(Vec::new());
1054    };
1055    Ok(vec![EventObservation {
1056        event_id: event.id,
1057        entity_id,
1058        referent_kind: if event.substrate == SubstrateKind::Note {
1059            ReferentKind::Note
1060        } else {
1061            ReferentKind::Entity
1062        },
1063        role: ObservationRole::Target,
1064        position: 0,
1065    }])
1066}
1067
1068fn decode_signal_observation(event: &Event) -> Result<Vec<EventObservation>, rusqlite::Error> {
1069    let Some(entity_id) = event.target_id else {
1070        return Ok(Vec::new());
1071    };
1072    Ok(vec![EventObservation {
1073        event_id: event.id,
1074        entity_id,
1075        // ADR-041 permits both entity and note signal targets. `brain.feedback`
1076        // threads the resolved target's substrate onto the event (#831);
1077        // pre-fix events all carry the old
1078        // `SubstrateKind::Event` placeholder and fall back to Entity, matching
1079        // the decoder's prior hard-coded behavior for that historical data.
1080        referent_kind: if event.substrate == SubstrateKind::Note {
1081            ReferentKind::Note
1082        } else {
1083            ReferentKind::Entity
1084        },
1085        role: ObservationRole::Signal,
1086        position: 0,
1087    }])
1088}
1089
1090fn invalid_payload(
1091    kind: EventKind,
1092    field: &'static str,
1093    reason: impl std::fmt::Display,
1094) -> rusqlite::Error {
1095    rusqlite::Error::ToSqlConversionFailure(
1096        format!("invalid payload for {}.{field}: {reason}", kind.name()).into(),
1097    )
1098}
1099
1100// =============================================================================
1101// Helpers: filter SQL builder
1102// =============================================================================
1103
1104fn build_event_filter_sql(
1105    conn: &rusqlite::Connection,
1106    default_namespace: &str,
1107    filter: &EventFilter,
1108) -> Result<(String, Vec<Box<dyn rusqlite::types::ToSql>>), rusqlite::Error> {
1109    reject_missing_event_filter_schema(conn, filter)?;
1110
1111    let mut conditions: Vec<String> = Vec::new();
1112    let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
1113
1114    params.push(Box::new(default_namespace.to_string()));
1115    conditions.push(format!("namespace = ?{}", params.len()));
1116
1117    push_in_clause(
1118        &mut conditions,
1119        &mut params,
1120        "id",
1121        filter.ids.iter().map(Uuid::to_string),
1122    );
1123    push_in_clause(
1124        &mut conditions,
1125        &mut params,
1126        "kind",
1127        filter.kinds.iter().map(|kind| kind.name().to_string()),
1128    );
1129    push_in_clause(
1130        &mut conditions,
1131        &mut params,
1132        "verb",
1133        filter.verbs.iter().cloned(),
1134    );
1135    push_in_clause(
1136        &mut conditions,
1137        &mut params,
1138        "substrate",
1139        filter.substrates.iter().map(|s| s.name().to_string()),
1140    );
1141    push_in_clause(
1142        &mut conditions,
1143        &mut params,
1144        "actor",
1145        filter.actors.iter().cloned(),
1146    );
1147
1148    if let Some(outcome) = filter.outcome {
1149        params.push(Box::new(outcome.name().to_string()));
1150        conditions.push(format!("outcome = ?{}", params.len()));
1151    }
1152    super::append_json_equalities(
1153        &mut conditions,
1154        &mut params,
1155        "payload",
1156        &filter.payload_equalities,
1157    );
1158
1159    if let Some(after) = filter.after {
1160        params.push(Box::new(after));
1161        conditions.push(format!("created_at > ?{}", params.len()));
1162    }
1163
1164    if let Some(before) = filter.before {
1165        params.push(Box::new(before));
1166        conditions.push(format!("created_at < ?{}", params.len()));
1167    }
1168
1169    if let Some(target_id) = filter.target_id {
1170        params.push(Box::new(target_id.to_string()));
1171        conditions.push(format!("target_id = ?{}", params.len()));
1172    }
1173
1174    if let Some(session_id) = filter.session_id {
1175        params.push(Box::new(session_id.to_string()));
1176        conditions.push(format!("session_id = ?{}", params.len()));
1177    }
1178
1179    push_observation_exists(&mut conditions, &mut params, None, &filter.observed);
1180    push_observation_exists(
1181        &mut conditions,
1182        &mut params,
1183        Some("selected"),
1184        &filter.selected,
1185    );
1186
1187    if let Some(proposal_id) = filter.payload_proposal_id {
1188        params.push(Box::new(proposal_id.to_string()));
1189        conditions.push(format!(
1190            "json_extract(payload, '$.proposal_id') = ?{}",
1191            params.len()
1192        ));
1193    }
1194
1195    let clause = format!(" WHERE {}", conditions.join(" AND "));
1196    Ok((clause, params))
1197}
1198
1199fn push_in_clause<I>(
1200    conditions: &mut Vec<String>,
1201    params: &mut Vec<Box<dyn rusqlite::types::ToSql>>,
1202    column: &'static str,
1203    values: I,
1204) where
1205    I: IntoIterator<Item = String>,
1206{
1207    let placeholders: Vec<String> = values
1208        .into_iter()
1209        .map(|value| {
1210            params.push(Box::new(value));
1211            format!("?{}", params.len())
1212        })
1213        .collect();
1214    if !placeholders.is_empty() {
1215        conditions.push(format!("{column} IN ({})", placeholders.join(",")));
1216    }
1217}
1218
1219fn push_observation_exists(
1220    conditions: &mut Vec<String>,
1221    params: &mut Vec<Box<dyn rusqlite::types::ToSql>>,
1222    role: Option<&'static str>,
1223    entity_ids: &[Uuid],
1224) {
1225    if entity_ids.is_empty() {
1226        return;
1227    }
1228    let placeholders: Vec<String> = entity_ids
1229        .iter()
1230        .map(|id| {
1231            params.push(Box::new(id.to_string()));
1232            format!("?{}", params.len())
1233        })
1234        .collect();
1235    let role_clause = role
1236        .map(|role| format!(" AND o.role = '{role}'"))
1237        .unwrap_or_default();
1238    conditions.push(format!(
1239        "EXISTS (SELECT 1 FROM event_observations o \
1240         WHERE o.event_id = events.id{role_clause} AND o.entity_id IN ({}))",
1241        placeholders.join(",")
1242    ));
1243}
1244
1245fn reject_missing_event_filter_schema(
1246    conn: &rusqlite::Connection,
1247    filter: &EventFilter,
1248) -> Result<(), rusqlite::Error> {
1249    if filter.target_id.is_some() && !has_column(conn, "events", "target_id")? {
1250        return Err(schema_absent("events.target_id"));
1251    }
1252    if filter.session_id.is_some() && !has_column(conn, "events", "session_id")? {
1253        return Err(schema_absent("events.session_id"));
1254    }
1255    if (!filter.observed.is_empty() || !filter.selected.is_empty())
1256        && !has_table(conn, "event_observations")?
1257    {
1258        return Err(schema_absent("event_observations"));
1259    }
1260    if filter.payload_proposal_id.is_some() && !has_column(conn, "events", "payload")? {
1261        return Err(schema_absent("events.payload"));
1262    }
1263    Ok(())
1264}
1265
1266fn has_table(conn: &rusqlite::Connection, table: &'static str) -> Result<bool, rusqlite::Error> {
1267    conn.query_row(
1268        "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = 'table' AND name = ?1",
1269        [table],
1270        |row| row.get(0),
1271    )
1272}
1273
1274fn has_column(
1275    conn: &rusqlite::Connection,
1276    table: &'static str,
1277    column: &'static str,
1278) -> Result<bool, rusqlite::Error> {
1279    conn.query_row(
1280        "SELECT COUNT(*) > 0 FROM pragma_table_info(?1) WHERE name = ?2",
1281        rusqlite::params![table, column],
1282        |row| row.get(0),
1283    )
1284}
1285
1286fn schema_absent(name: &'static str) -> rusqlite::Error {
1287    rusqlite::Error::ToSqlConversionFailure(
1288        format!("event filter requires missing schema element {name}; run migrations").into(),
1289    )
1290}
1291
1292// =============================================================================
1293// EventStore implementation
1294// =============================================================================
1295
1296#[async_trait]
1297impl EventStore for SqlEventStore {
1298    async fn append_event(&self, event: Event) -> Result<(), StorageError> {
1299        // ADR-067 Component A: when the write queue is enabled, route
1300        // through the pool-wide WriterTask. DML-only closure — no BEGIN
1301        // IMMEDIATE/COMMIT/ROLLBACK here, since the WriterTask's run loop
1302        // owns the transaction.
1303        //
1304        // `event_rows` is counted here at the store seam (on success, both
1305        // paths) so every request-owned append — proposal lifecycle,
1306        // curation, mutation events — is covered without per-call-site
1307        // instrumentation. The enclosing per-dispatch audit row is appended
1308        // only after the usage snapshot is frozen, so it never counts itself.
1309        if let Some(writer_task) = self.current_writer_task("append_event")? {
1310            let result = writer_task
1311                .send_bounded(move |conn| {
1312                    insert_event_with_observations(conn, &event)
1313                        .map_err(|e| map_err(e, "append_event"))
1314                })
1315                .await;
1316            khive_storage::usage::account_event_write(result.as_ref().map(|()| 1));
1317            return result;
1318        }
1319
1320        let result = self
1321            .with_writer("append_event", move |conn| {
1322                insert_event_with_observations(conn, &event)
1323            })
1324            .await;
1325        #[cfg(test)]
1326        let result = self.after_fallback_write_for_test(result);
1327        khive_storage::usage::account_event_write(result.as_ref().map(|()| 1));
1328        result
1329    }
1330
1331    async fn append_events(&self, events: Vec<Event>) -> Result<BatchWriteSummary, StorageError> {
1332        let attempted = events.len() as u64;
1333
1334        // ADR-067 Component A: when the write queue is enabled, route
1335        // through the pool-wide WriterTask. DML-only closure preserving the
1336        // all-or-nothing semantics (first failed insert aborts the whole
1337        // batch) — the WriterTask's run loop owns the enclosing transaction
1338        // and issues the ROLLBACK on `Err`.
1339        if let Some(writer_task) = self.current_writer_task("append_events")? {
1340            let result = writer_task
1341                .send_bounded(move |conn| {
1342                    batch_append_events_dml(conn, &events, attempted)
1343                        .map_err(|e| map_err(e, "append_events"))
1344                })
1345                .await;
1346            khive_storage::usage::account_event_write(
1347                result.as_ref().map(|summary| summary.affected),
1348            );
1349            return result;
1350        }
1351
1352        let result = self
1353            .with_writer("append_events", move |conn| {
1354                batch_append_events_dml(conn, &events, attempted)
1355            })
1356            .await;
1357        #[cfg(test)]
1358        let result = self.after_fallback_write_for_test(result);
1359        khive_storage::usage::account_event_write(result.as_ref().map(|summary| summary.affected));
1360        result
1361    }
1362
1363    async fn get_event(&self, id: Uuid) -> Result<Option<Event>, StorageError> {
1364        let namespace = self.namespace.clone();
1365        let id_str = id.to_string();
1366
1367        self.with_reader("get_event", move |conn| {
1368            let mut stmt = conn.prepare(
1369                "SELECT id, namespace, verb, substrate, actor, kind, outcome, payload, \
1370                        payload_schema_version, profile_state_version, duration_us, target_id, \
1371                        session_id, aggregate_kind, aggregate_id, created_at, op_index, ref_resolution \
1372                 FROM events WHERE namespace = ?1 AND id = ?2",
1373            )?;
1374            let mut rows = stmt.query(rusqlite::params![namespace, id_str])?;
1375            match rows.next()? {
1376                Some(row) => Ok(Some(read_event(row)?)),
1377                None => Ok(None),
1378            }
1379        })
1380        .await
1381    }
1382
1383    async fn query_events(
1384        &self,
1385        filter: EventFilter,
1386        page: PageRequest,
1387    ) -> Result<Page<Event>, StorageError> {
1388        super::validate_json_equality_paths(
1389            &filter.payload_equalities,
1390            StorageCapability::Events,
1391            "query_events",
1392        )?;
1393        let namespace = self.namespace.clone();
1394        let limit_i64 = i64::from(page.limit);
1395        let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
1396            capability: StorageCapability::Events,
1397            operation: "query_events".into(),
1398            message: format!(
1399                "PageRequest: offset must be <= i64::MAX, got {}",
1400                page.offset
1401            ),
1402        })?;
1403
1404        self.with_reader("query_events", move |conn| {
1405            // No `COUNT(*)` here, and `total` is therefore `None`. The count
1406            // carried no `LIMIT`, so it scanned the whole filtered set on every
1407            // paged read while the data query below fetches only
1408            // `offset + limit` rows — and on the merged event plane a single
1409            // read paid it twice, once per underlying store. `Page.total` is
1410            // `Option<u64>` precisely so a store may decline to compute it, and
1411            // the merged fold in `khive-runtime::events_split` propagates `None`
1412            // rather than inventing a number. `count_events` below remains for
1413            // callers that genuinely want a cardinality.
1414            let (where_clause, filter_params) = build_event_filter_sql(conn, &namespace, &filter)?;
1415
1416            let mut all_params: Vec<Box<dyn rusqlite::types::ToSql>> = filter_params;
1417            all_params.push(Box::new(limit_i64));
1418            all_params.push(Box::new(offset_i64));
1419
1420            let limit_idx = all_params.len() - 1;
1421            let offset_idx = all_params.len();
1422
1423            let data_sql = format!(
1424                "SELECT id, namespace, verb, substrate, actor, kind, outcome, payload, \
1425                        payload_schema_version, profile_state_version, duration_us, target_id, \
1426                        session_id, aggregate_kind, aggregate_id, created_at, op_index, ref_resolution \
1427                 FROM events{} ORDER BY created_at DESC, id DESC LIMIT ?{} OFFSET ?{}",
1428                where_clause, limit_idx, offset_idx,
1429            );
1430
1431            let mut stmt = conn.prepare(&data_sql)?;
1432            let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1433                all_params.iter().map(|p| p.as_ref()).collect();
1434            let rows = stmt.query_map(param_refs.as_slice(), read_event)?;
1435
1436            let mut items = Vec::new();
1437            for row in rows {
1438                items.push(row?);
1439            }
1440
1441            Ok(Page { items, total: None })
1442        })
1443        .await
1444    }
1445
1446    async fn query_event_page(
1447        &self,
1448        query: khive_storage::event::EventPageQuery,
1449    ) -> Result<khive_storage::event::EventPageWindow, StorageError> {
1450        cursor::query_event_page(self, query).await
1451    }
1452
1453    async fn count_events(&self, filter: EventFilter) -> Result<u64, StorageError> {
1454        super::validate_json_equality_paths(
1455            &filter.payload_equalities,
1456            StorageCapability::Events,
1457            "count_events",
1458        )?;
1459        let namespace = self.namespace.clone();
1460
1461        self.with_reader("count_events", move |conn| {
1462            let (where_clause, params) = build_event_filter_sql(conn, &namespace, &filter)?;
1463            let sql = format!("SELECT COUNT(*) FROM events{}", where_clause);
1464            let mut stmt = conn.prepare(&sql)?;
1465            let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1466                params.iter().map(|p| p.as_ref()).collect();
1467            let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
1468            Ok(count as u64)
1469        })
1470        .await
1471    }
1472
1473    fn preflight_event(&self, event: &Event) -> Result<(), StorageError> {
1474        event_insert_statements(event)
1475            .map(|_| ())
1476            .map_err(|e| map_err(e, "preflight_event"))
1477    }
1478
1479    async fn append_events_idempotent(
1480        &self,
1481        events: Vec<Event>,
1482    ) -> Result<IdempotentEventBatchResult, StorageError> {
1483        if let Some(writer_task) = self.current_writer_task("append_events_idempotent")? {
1484            return writer_task
1485                .send_bounded(move |conn| {
1486                    idempotent_batch_dml(conn, &events)
1487                        .map_err(|e| map_err(e, "append_events_idempotent"))
1488                })
1489                .await;
1490        }
1491
1492        self.with_writer("append_events_idempotent", move |conn| {
1493            idempotent_batch_dml(conn, &events)
1494        })
1495        .await
1496    }
1497
1498    fn supports_idempotent_audit_batch(&self) -> bool {
1499        true
1500    }
1501}
1502
1503// =============================================================================
1504// DDL
1505// =============================================================================
1506
1507const EVENTS_DDL: &str = include_str!("../../sql/events-ddl.sql");
1508const OPERATION_ATTRIBUTION_COLUMNS: &str =
1509    include_str!("../../sql/036-events-operation-attribution.sql");
1510
1511pub(crate) fn ensure_events_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
1512    conn.execute_batch(EVENTS_DDL)?;
1513    ensure_operation_attribution_columns(conn)
1514}
1515
1516/// Add `op_index` and `ref_resolution` to an `events` table created before
1517/// they existed. `CREATE TABLE IF NOT EXISTS` leaves such a table unchanged,
1518/// and a standalone events database has no migration chain to run V36, so
1519/// every insert naming the two columns would fail. V36 calls this too, which
1520/// keeps it valid on a table the store DDL has already upgraded. Both columns
1521/// are added in one transaction (the caller's, if one is open).
1522pub(crate) fn ensure_operation_attribution_columns(
1523    conn: &rusqlite::Connection,
1524) -> Result<(), rusqlite::Error> {
1525    match (
1526        has_column(conn, "events", "op_index")?,
1527        has_column(conn, "events", "ref_resolution")?,
1528    ) {
1529        (true, true) => Ok(()),
1530        (false, false) if conn.is_autocommit() => {
1531            let tx = conn.unchecked_transaction()?;
1532            tx.execute_batch(OPERATION_ATTRIBUTION_COLUMNS)?;
1533            tx.commit()
1534        }
1535        (false, false) => conn.execute_batch(OPERATION_ATTRIBUTION_COLUMNS),
1536        (op_index, ref_resolution) => Err(rusqlite::Error::ToSqlConversionFailure(
1537            format!(
1538                "events table has only one operation attribution column \
1539                 (op_index={op_index}, ref_resolution={ref_resolution}); refusing to guess"
1540            )
1541            .into(),
1542        )),
1543    }
1544}
1545
1546#[cfg(test)]
1547#[path = "event_tests.rs"]
1548mod tests;
1549
1550#[cfg(test)]
1551#[path = "event_busy_tests.rs"]
1552mod direct_busy_tests;
1553
1554#[cfg(test)]
1555#[path = "event_fallback_usage_tests.rs"]
1556mod fallback_usage_tests;