1use 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#[cfg(test)]
40fn mark_unknown_append_usage(error: &StorageError) {
41 khive_storage::usage::account_event_write(Err(error));
42}
43
44pub 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 pub fn new_scoped(
57 pool: Arc<ConnectionPool>,
58 is_file_backed: bool,
59 namespace: impl Into<String>,
60 ) -> Self {
61 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 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 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 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 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
175fn 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
245fn 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
341fn 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
405pub 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
416fn 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
458fn 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
503fn 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
539pub 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
639pub 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
718pub 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
771fn 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
800fn 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
828fn 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
921fn 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
929fn 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
961fn 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 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
1100fn 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#[async_trait]
1297impl EventStore for SqlEventStore {
1298 async fn append_event(&self, event: Event) -> Result<(), StorageError> {
1299 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 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 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
1503const 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
1516pub(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;