Skip to main content

arc_es_sqlite/
lib.rs

1//! # Arc ES SQLite
2//!
3//! SQLite implementation of the [`EventStore`] trait from `arc-core`.
4//!
5//! Persists [`AuditMetadata`] inline alongside each event (see HIPAA-1 in
6//! `docs/ark/refactor-plan.md`). `append` calls
7//! [`validate_audit_batch`](arc_core::event_store::validate_audit_batch)
8//! before any write — defense-in-depth against an upstream that forgot to
9//! stamp.
10
11use arc_core::audit::AuditMetadata;
12use arc_core::event::Event;
13use arc_core::event_store::{
14    validate_audit_batch, EventStore, EventStoreError, EventStoreResult, VersionCheck,
15};
16use arc_core::integrity::{EventSignature, HmacSha256Chain, IntegrityChain, IntegrityError};
17use arc_core::snapshot::Snapshot;
18use async_trait::async_trait;
19use diesel::prelude::*;
20use diesel::r2d2::{self, ConnectionManager};
21use diesel::sqlite::SqliteConnection;
22use std::sync::Arc;
23use uuid::Uuid;
24
25// Re-export for convenience
26pub use arc_core::{Deserialize, Serialize};
27
28pub mod session;
29pub use session::SqliteSessionStore;
30
31pub mod read_model_store;
32pub use read_model_store::SqliteReadModelStore;
33
34#[cfg(test)]
35mod test_support;
36
37/// Database row used for inserting events.
38#[derive(Debug, Insertable, Clone)]
39#[diesel(table_name = events)]
40struct NewEventRecord {
41    pub event_id: String,
42    pub aggregate_type: String,
43    pub aggregate_id: String,
44    pub sequence: i64,
45    pub event_type: String,
46    pub payload: String,
47    pub timestamp: i64,
48    pub actor_id: String,
49    pub actor_session_id: Option<String>,
50    pub source_ip: Option<String>,
51    pub user_agent: Option<String>,
52    pub timestamp_utc_us: i64,
53    pub causation_id: Option<String>,
54    pub correlation_id: String,
55    pub integrity_signature: Option<String>,
56    pub integrity_key_id: Option<String>,
57}
58
59#[derive(Debug, Queryable, Clone)]
60struct EventRecord {
61    #[allow(dead_code)]
62    pub id: Option<i32>,
63    pub event_id: String,
64    pub aggregate_type: String,
65    pub aggregate_id: String,
66    pub sequence: i64,
67    pub event_type: String,
68    pub payload: String,
69    pub timestamp: i64,
70    pub actor_id: String,
71    pub actor_session_id: Option<String>,
72    pub source_ip: Option<String>,
73    pub user_agent: Option<String>,
74    pub timestamp_utc_us: i64,
75    pub causation_id: Option<String>,
76    pub correlation_id: String,
77    pub integrity_signature: Option<String>,
78    pub integrity_key_id: Option<String>,
79}
80
81impl NewEventRecord {
82    fn from_event(
83        event: &Event,
84        integrity_signature: Option<String>,
85        integrity_key_id: Option<String>,
86    ) -> Result<Self, EventStoreError> {
87        // sequence and timestamp are i64 end-to-end now — no truncation.
88        let timestamp_seconds: i64 = (event.timestamp / 1000) as i64;
89        Ok(NewEventRecord {
90            event_id: event.event_id.to_string(),
91            aggregate_type: event.aggregate_type.clone(),
92            aggregate_id: event.aggregate_id.clone(),
93            sequence: event.sequence,
94            event_type: event.event_type.clone(),
95            payload: serde_json::to_string(&event.payload)
96                .map_err(|e| EventStoreError::serialization(e.to_string()))?,
97            timestamp: timestamp_seconds,
98            actor_id: event.audit.actor_id.clone(),
99            actor_session_id: event.audit.actor_session_id.clone(),
100            source_ip: event.audit.source_ip.clone(),
101            user_agent: event.audit.user_agent.clone(),
102            timestamp_utc_us: event.audit.timestamp_utc_us,
103            causation_id: event.audit.causation_id.map(|u| u.to_string()),
104            correlation_id: event.audit.correlation_id.to_string(),
105            integrity_signature,
106            integrity_key_id,
107        })
108    }
109}
110
111impl EventRecord {
112    fn to_event(&self) -> Result<Event, EventStoreError> {
113        let event_id = Uuid::parse_str(&self.event_id)
114            .map_err(|e| EventStoreError::serialization(format!("Invalid UUID: {}", e)))?;
115
116        let payload: serde_json::Value = serde_json::from_str(&self.payload)
117            .map_err(|e| EventStoreError::serialization(e.to_string()))?;
118
119        let causation_id = match self.causation_id.as_deref() {
120            Some(s) => Some(Uuid::parse_str(s).map_err(|e| {
121                EventStoreError::serialization(format!("Invalid causation UUID: {}", e))
122            })?),
123            None => None,
124        };
125
126        let correlation_id = Uuid::parse_str(&self.correlation_id).map_err(|e| {
127            EventStoreError::serialization(format!("Invalid correlation UUID: {}", e))
128        })?;
129
130        let audit = AuditMetadata {
131            actor_id: self.actor_id.clone(),
132            actor_session_id: self.actor_session_id.clone(),
133            source_ip: self.source_ip.clone(),
134            user_agent: self.user_agent.clone(),
135            timestamp_utc_us: self.timestamp_utc_us,
136            causation_id,
137            correlation_id,
138        };
139
140        Ok(Event {
141            event_id,
142            aggregate_type: self.aggregate_type.clone(),
143            aggregate_id: self.aggregate_id.clone(),
144            sequence: self.sequence,
145            event_type: self.event_type.clone(),
146            payload,
147            audit,
148            timestamp: (self.timestamp as u64) * 1000,
149        })
150    }
151}
152
153/// Database row used for upserting snapshots. `state` holds the aggregate's
154/// serialized JSON; `created_at` is milliseconds since epoch (same unit as the
155/// core `Snapshot`), stored as i64 to match the events table's `timestamp`.
156#[derive(Debug, Insertable, Clone)]
157#[diesel(table_name = snapshots)]
158struct NewSnapshotRecord {
159    pub aggregate_id: String,
160    pub aggregate_type: String,
161    pub version: i64,
162    pub state: String,
163    pub created_at: i64,
164}
165
166#[derive(Debug, Queryable, Clone)]
167struct SnapshotRecord {
168    pub aggregate_id: String,
169    pub aggregate_type: String,
170    pub version: i64,
171    pub state: String,
172    pub created_at: i64,
173}
174
175impl NewSnapshotRecord {
176    fn from_snapshot(snapshot: &Snapshot) -> Result<Self, EventStoreError> {
177        Ok(NewSnapshotRecord {
178            aggregate_id: snapshot.aggregate_id.clone(),
179            aggregate_type: snapshot.aggregate_type.clone(),
180            version: snapshot.version,
181            state: serde_json::to_string(&snapshot.state)
182                .map_err(|e| EventStoreError::serialization(e.to_string()))?,
183            created_at: snapshot.created_at as i64,
184        })
185    }
186}
187
188impl SnapshotRecord {
189    fn to_snapshot(&self) -> Result<Snapshot, EventStoreError> {
190        let state: serde_json::Value = serde_json::from_str(&self.state)
191            .map_err(|e| EventStoreError::serialization(e.to_string()))?;
192        Ok(Snapshot {
193            aggregate_id: self.aggregate_id.clone(),
194            aggregate_type: self.aggregate_type.clone(),
195            version: self.version,
196            state,
197            created_at: self.created_at as u64,
198        })
199    }
200}
201
202mod schema {
203    diesel::table! {
204        events (id) {
205            id -> Nullable<Integer>,
206            event_id -> Text,
207            aggregate_type -> Text,
208            aggregate_id -> Text,
209            sequence -> BigInt,
210            event_type -> Text,
211            payload -> Text,
212            timestamp -> BigInt,
213            actor_id -> Text,
214            actor_session_id -> Nullable<Text>,
215            source_ip -> Nullable<Text>,
216            user_agent -> Nullable<Text>,
217            timestamp_utc_us -> BigInt,
218            causation_id -> Nullable<Text>,
219            correlation_id -> Text,
220            integrity_signature -> Nullable<Text>,
221            integrity_key_id -> Nullable<Text>,
222        }
223    }
224
225    diesel::table! {
226        snapshots (aggregate_type, aggregate_id) {
227            aggregate_id -> Text,
228            aggregate_type -> Text,
229            version -> BigInt,
230            state -> Text,
231            created_at -> BigInt,
232        }
233    }
234}
235
236use schema::{events, snapshots};
237
238type Pool = r2d2::Pool<ConnectionManager<SqliteConnection>>;
239
240/// SQLite implementation of EventStore.
241#[derive(Clone)]
242pub struct SqliteEventStore {
243    pool: Arc<Pool>,
244    integrity: Option<Arc<IntegrityConfig>>,
245}
246
247struct IntegrityConfig {
248    chain: Arc<dyn IntegrityChain>,
249    key_id: String,
250}
251
252impl SqliteEventStore {
253    pub async fn new(database_url: &str) -> EventStoreResult<Self> {
254        let manager = ConnectionManager::<SqliteConnection>::new(database_url);
255        let pool = Pool::builder()
256            .max_size(10)
257            .build(manager)
258            .map_err(|e| EventStoreError::database(format!("Failed to create pool: {}", e)))?;
259
260        Ok(SqliteEventStore {
261            pool: Arc::new(pool),
262            integrity: None,
263        })
264    }
265
266    pub async fn new_with_integrity_key(
267        database_url: &str,
268        key: impl Into<Vec<u8>>,
269        key_id: impl Into<String>,
270    ) -> EventStoreResult<Self> {
271        let mut store = Self::new(database_url).await?;
272        store.integrity = Some(Arc::new(IntegrityConfig {
273            chain: Arc::new(HmacSha256Chain::new(key).map_err(EventStoreError::from)?),
274            key_id: key_id.into(),
275        }));
276        Ok(store)
277    }
278
279    pub fn with_pool(pool: Pool) -> Self {
280        SqliteEventStore {
281            pool: Arc::new(pool),
282            integrity: None,
283        }
284    }
285
286    pub fn with_pool_and_integrity_key(
287        pool: Pool,
288        key: impl Into<Vec<u8>>,
289        key_id: impl Into<String>,
290    ) -> EventStoreResult<Self> {
291        Ok(SqliteEventStore {
292            pool: Arc::new(pool),
293            integrity: Some(Arc::new(IntegrityConfig {
294                chain: Arc::new(HmacSha256Chain::new(key).map_err(EventStoreError::from)?),
295                key_id: key_id.into(),
296            })),
297        })
298    }
299}
300
301fn required_signature(
302    record: &EventRecord,
303    aggregate_id: &str,
304    sequence: i64,
305) -> EventStoreResult<EventSignature> {
306    let _key_id = record.integrity_key_id.as_ref().ok_or_else(|| {
307        EventStoreError::from(IntegrityError::BrokenAt {
308            aggregate_id: aggregate_id.to_string(),
309            sequence,
310        })
311    })?;
312
313    record
314        .integrity_signature
315        .as_ref()
316        .map(|s| EventSignature(s.clone()))
317        .ok_or_else(|| {
318            EventStoreError::from(IntegrityError::BrokenAt {
319                aggregate_id: aggregate_id.to_string(),
320                sequence,
321            })
322        })
323}
324
325fn verify_integrity_records(
326    integrity: &IntegrityConfig,
327    records: &[EventRecord],
328    previous_signature: EventSignature,
329) -> EventStoreResult<Vec<Event>> {
330    let mut previous = previous_signature;
331    let mut events = Vec::with_capacity(records.len());
332
333    for record in records {
334        let event = record.to_event()?;
335        let expected = integrity.chain.sign_event(&previous, &event)?;
336        let claimed = required_signature(record, &event.aggregate_id, event.sequence)?;
337
338        if expected != claimed {
339            return Err(EventStoreError::from(IntegrityError::BrokenAt {
340                aggregate_id: event.aggregate_id,
341                sequence: event.sequence,
342            }));
343        }
344
345        previous = claimed;
346        events.push(event);
347    }
348
349    Ok(events)
350}
351
352fn previous_signature_for_aggregate(
353    conn: &mut SqliteConnection,
354    aggregate_type: Option<&str>,
355    aggregate_id: &str,
356    before_sequence: i64,
357) -> EventStoreResult<EventSignature> {
358    if before_sequence <= 1 {
359        return Ok(EventSignature::genesis());
360    }
361
362    let mut query = events::table
363        .filter(events::aggregate_id.eq(aggregate_id))
364        .filter(events::sequence.lt(before_sequence))
365        .into_boxed();
366    if let Some(aggregate_type) = aggregate_type {
367        query = query.filter(events::aggregate_type.eq(aggregate_type));
368    }
369    let record = query
370        .order(events::sequence.desc())
371        .first::<EventRecord>(conn)
372        .optional()
373        .map_err(|e| EventStoreError::database(e.to_string()))?;
374
375    match record {
376        Some(record) => required_signature(&record, aggregate_id, record.sequence),
377        None => Ok(EventSignature::genesis()),
378    }
379}
380
381fn verify_stream_integrity_records(
382    conn: &mut SqliteConnection,
383    integrity: &IntegrityConfig,
384    records: &[EventRecord],
385) -> EventStoreResult<Vec<Event>> {
386    use std::collections::HashMap;
387
388    let mut previous_by_aggregate: HashMap<(String, String), EventSignature> = HashMap::new();
389    let mut events = Vec::with_capacity(records.len());
390
391    for record in records {
392        let event = record.to_event()?;
393        let stream = (event.aggregate_type.clone(), event.aggregate_id.clone());
394        let previous = match previous_by_aggregate.get(&stream) {
395            Some(sig) => sig.clone(),
396            None => previous_signature_for_aggregate(
397                conn,
398                Some(&event.aggregate_type),
399                &event.aggregate_id,
400                event.sequence,
401            )?,
402        };
403
404        let expected = integrity.chain.sign_event(&previous, &event)?;
405        let claimed = required_signature(record, &event.aggregate_id, event.sequence)?;
406
407        if expected != claimed {
408            return Err(EventStoreError::from(IntegrityError::BrokenAt {
409                aggregate_id: event.aggregate_id,
410                sequence: event.sequence,
411            }));
412        }
413
414        previous_by_aggregate.insert(stream, claimed);
415        events.push(event);
416    }
417
418    Ok(events)
419}
420
421#[async_trait]
422impl EventStore for SqliteEventStore {
423    async fn append(
424        &self,
425        aggregate_id: &str,
426        version_check: VersionCheck,
427        new_events: Vec<Event>,
428    ) -> EventStoreResult<()> {
429        let aggregate_type = new_events
430            .first()
431            .map(|event| event.aggregate_type.clone())
432            .unwrap_or_default();
433        self.append_to(&aggregate_type, aggregate_id, version_check, new_events)
434            .await
435    }
436
437    async fn append_to(
438        &self,
439        aggregate_type: &str,
440        aggregate_id: &str,
441        version_check: VersionCheck,
442        new_events: Vec<Event>,
443    ) -> EventStoreResult<()> {
444        if new_events.is_empty() {
445            return Ok(());
446        }
447
448        // Defense-in-depth: reject any event with invalid audit before touching the DB.
449        validate_audit_batch(aggregate_id, &new_events)?;
450
451        let aggregate_type = aggregate_type.to_string();
452        let aggregate_id = aggregate_id.to_string();
453        let pool = self.pool.clone();
454        let integrity = self.integrity.clone();
455
456        tokio::task::spawn_blocking(move || -> EventStoreResult<()> {
457            use diesel::connection::AnsiTransactionManager;
458            use diesel::connection::TransactionManager;
459
460            let mut conn = pool.get().map_err(|e| {
461                EventStoreError::database(format!("Failed to get connection: {}", e))
462            })?;
463
464            AnsiTransactionManager::begin_transaction(&mut *conn)
465                .map_err(|e| EventStoreError::database(e.to_string()))?;
466
467            let result = (|| -> EventStoreResult<()> {
468                let current_version = events::table
469                    .filter(events::aggregate_type.eq(&aggregate_type))
470                    .filter(events::aggregate_id.eq(&aggregate_id))
471                    .select(diesel::dsl::max(events::sequence))
472                    .first::<Option<i64>>(&mut *conn)
473                    .map_err(|e| EventStoreError::database(e.to_string()))?
474                    .unwrap_or(0);
475
476                if let Some(expected) = version_check.version() {
477                    if current_version != expected {
478                        return Err(EventStoreError::ConcurrencyConflict {
479                            aggregate_id: aggregate_id.clone(),
480                            expected,
481                            actual: current_version,
482                        });
483                    }
484                }
485
486                for (expected_sequence, event) in (current_version + 1..).zip(new_events.iter()) {
487                    if event.sequence != expected_sequence {
488                        return Err(EventStoreError::InvalidSequence {
489                            aggregate_id: aggregate_id.clone(),
490                            expected: expected_sequence,
491                            actual: event.sequence,
492                        });
493                    }
494                }
495
496                let mut previous_signature = if integrity.is_some() {
497                    previous_signature_for_aggregate(
498                        &mut conn,
499                        Some(&aggregate_type),
500                        &aggregate_id,
501                        current_version + 1,
502                    )?
503                } else {
504                    EventSignature::genesis()
505                };
506
507                for event in &new_events {
508                    let mut record = NewEventRecord::from_event(event, None, None)?;
509
510                    if let Some(integrity) = integrity.as_ref() {
511                        let mut persisted_event = event.clone();
512                        persisted_event.timestamp = (record.timestamp as u64) * 1000;
513                        let signature = integrity
514                            .chain
515                            .sign_event(&previous_signature, &persisted_event)
516                            .map_err(EventStoreError::from)?;
517                        previous_signature = signature.clone();
518                        record.integrity_signature = Some(signature.0);
519                        record.integrity_key_id = Some(integrity.key_id.clone());
520                    }
521
522                    diesel::insert_into(events::table)
523                        .values(&record)
524                        .execute(&mut *conn)
525                        .map_err(|e| EventStoreError::database(e.to_string()))?;
526                }
527
528                Ok(())
529            })();
530
531            match result {
532                Ok(_) => {
533                    AnsiTransactionManager::commit_transaction(&mut *conn)
534                        .map_err(|e| EventStoreError::database(e.to_string()))?;
535                    Ok(())
536                }
537                Err(e) => {
538                    let _ = AnsiTransactionManager::rollback_transaction(&mut *conn);
539                    Err(e)
540                }
541            }
542        })
543        .await
544        .map_err(|e| EventStoreError::other(format!("Task join error: {}", e)))?
545    }
546
547    async fn load(&self, aggregate_id: &str) -> EventStoreResult<Vec<Event>> {
548        self.load_from(aggregate_id, 1).await
549    }
550
551    async fn load_stream(
552        &self,
553        aggregate_type: &str,
554        aggregate_id: &str,
555    ) -> EventStoreResult<Vec<Event>> {
556        self.load_stream_from(aggregate_type, aggregate_id, 1).await
557    }
558
559    async fn load_from(
560        &self,
561        aggregate_id: &str,
562        from_sequence: i64,
563    ) -> EventStoreResult<Vec<Event>> {
564        let aggregate_id = aggregate_id.to_string();
565        let pool = self.pool.clone();
566        let integrity = self.integrity.clone();
567
568        tokio::task::spawn_blocking(move || {
569            let mut conn = pool.get().map_err(|e| {
570                EventStoreError::database(format!("Failed to get connection: {}", e))
571            })?;
572
573            let records: Vec<EventRecord> = events::table
574                .filter(events::aggregate_id.eq(&aggregate_id))
575                .filter(events::sequence.ge(from_sequence))
576                .order(events::sequence.asc())
577                .load(&mut conn)
578                .map_err(|e| EventStoreError::database(e.to_string()))?;
579
580            match integrity.as_ref() {
581                Some(integrity) => {
582                    let previous = previous_signature_for_aggregate(
583                        &mut conn,
584                        None,
585                        &aggregate_id,
586                        from_sequence,
587                    )?;
588                    verify_integrity_records(integrity, &records, previous)
589                }
590                None => records.iter().map(|r| r.to_event()).collect(),
591            }
592        })
593        .await
594        .map_err(|e| EventStoreError::other(format!("Task join error: {}", e)))?
595    }
596
597    async fn load_stream_from(
598        &self,
599        aggregate_type: &str,
600        aggregate_id: &str,
601        from_sequence: i64,
602    ) -> EventStoreResult<Vec<Event>> {
603        let aggregate_type = aggregate_type.to_string();
604        let aggregate_id = aggregate_id.to_string();
605        let pool = self.pool.clone();
606        let integrity = self.integrity.clone();
607
608        tokio::task::spawn_blocking(move || {
609            let mut conn = pool.get().map_err(|e| {
610                EventStoreError::database(format!("Failed to get connection: {}", e))
611            })?;
612
613            let records: Vec<EventRecord> = events::table
614                .filter(events::aggregate_type.eq(&aggregate_type))
615                .filter(events::aggregate_id.eq(&aggregate_id))
616                .filter(events::sequence.ge(from_sequence))
617                .order(events::sequence.asc())
618                .load(&mut conn)
619                .map_err(|e| EventStoreError::database(e.to_string()))?;
620
621            match integrity.as_ref() {
622                Some(integrity) => {
623                    let previous = previous_signature_for_aggregate(
624                        &mut conn,
625                        Some(&aggregate_type),
626                        &aggregate_id,
627                        from_sequence,
628                    )?;
629                    verify_integrity_records(integrity, &records, previous)
630                }
631                None => records.iter().map(|record| record.to_event()).collect(),
632            }
633        })
634        .await
635        .map_err(|e| EventStoreError::other(format!("Task join error: {}", e)))?
636    }
637
638    async fn stream_all(&self, from_position: i64) -> EventStoreResult<Vec<Event>> {
639        let pool = self.pool.clone();
640        let integrity = self.integrity.clone();
641
642        tokio::task::spawn_blocking(move || {
643            let mut conn = pool.get().map_err(|e| {
644                EventStoreError::database(format!("Failed to get connection: {}", e))
645            })?;
646
647            let records: Vec<EventRecord> = events::table
648                .filter(events::id.ge(from_position as i32))
649                .order(events::id.asc())
650                .load(&mut conn)
651                .map_err(|e| EventStoreError::database(e.to_string()))?;
652
653            match integrity.as_ref() {
654                Some(integrity) => verify_stream_integrity_records(&mut conn, integrity, &records),
655                None => records.iter().map(|r| r.to_event()).collect(),
656            }
657        })
658        .await
659        .map_err(|e| EventStoreError::other(format!("Task join error: {}", e)))?
660    }
661
662    async fn get_version(&self, aggregate_id: &str) -> EventStoreResult<i64> {
663        let aggregate_id = aggregate_id.to_string();
664        let pool = self.pool.clone();
665
666        tokio::task::spawn_blocking(move || {
667            let mut conn = pool.get().map_err(|e| {
668                EventStoreError::database(format!("Failed to get connection: {}", e))
669            })?;
670
671            let version = events::table
672                .filter(events::aggregate_id.eq(&aggregate_id))
673                .select(diesel::dsl::max(events::sequence))
674                .first::<Option<i64>>(&mut conn)
675                .map_err(|e| EventStoreError::database(e.to_string()))?
676                .unwrap_or(0);
677
678            Ok(version)
679        })
680        .await
681        .map_err(|e| EventStoreError::other(format!("Task join error: {}", e)))?
682    }
683
684    async fn get_stream_version(
685        &self,
686        aggregate_type: &str,
687        aggregate_id: &str,
688    ) -> EventStoreResult<i64> {
689        let aggregate_type = aggregate_type.to_string();
690        let aggregate_id = aggregate_id.to_string();
691        let pool = self.pool.clone();
692
693        tokio::task::spawn_blocking(move || {
694            let mut conn = pool.get().map_err(|e| {
695                EventStoreError::database(format!("Failed to get connection: {}", e))
696            })?;
697            events::table
698                .filter(events::aggregate_type.eq(&aggregate_type))
699                .filter(events::aggregate_id.eq(&aggregate_id))
700                .select(diesel::dsl::max(events::sequence))
701                .first::<Option<i64>>(&mut conn)
702                .map(|version| version.unwrap_or(0))
703                .map_err(|e| EventStoreError::database(e.to_string()))
704        })
705        .await
706        .map_err(|e| EventStoreError::other(format!("Task join error: {}", e)))?
707    }
708
709    async fn save_snapshot(&self, snapshot: &Snapshot) -> EventStoreResult<()> {
710        let record = NewSnapshotRecord::from_snapshot(snapshot)?;
711        let pool = self.pool.clone();
712
713        tokio::task::spawn_blocking(move || -> EventStoreResult<()> {
714            let mut conn = pool.get().map_err(|e| {
715                EventStoreError::database(format!("Failed to get connection: {}", e))
716            })?;
717
718            // One snapshot per aggregate: replace the stored row in place rather
719            // than accumulating stale versions.
720            diesel::insert_into(snapshots::table)
721                .values(&record)
722                .on_conflict((snapshots::aggregate_type, snapshots::aggregate_id))
723                .do_update()
724                .set((
725                    snapshots::version.eq(record.version),
726                    snapshots::state.eq(&record.state),
727                    snapshots::created_at.eq(record.created_at),
728                ))
729                .execute(&mut *conn)
730                .map_err(|e| EventStoreError::database(e.to_string()))?;
731
732            Ok(())
733        })
734        .await
735        .map_err(|e| EventStoreError::other(format!("Task join error: {}", e)))?
736    }
737
738    async fn load_snapshot(&self, aggregate_id: &str) -> EventStoreResult<Option<Snapshot>> {
739        let aggregate_id = aggregate_id.to_string();
740        let pool = self.pool.clone();
741
742        tokio::task::spawn_blocking(move || {
743            let mut conn = pool.get().map_err(|e| {
744                EventStoreError::database(format!("Failed to get connection: {}", e))
745            })?;
746
747            let record: Option<SnapshotRecord> = snapshots::table
748                .filter(snapshots::aggregate_id.eq(&aggregate_id))
749                .first::<SnapshotRecord>(&mut conn)
750                .optional()
751                .map_err(|e| EventStoreError::database(e.to_string()))?;
752
753            record.map(|r| r.to_snapshot()).transpose()
754        })
755        .await
756        .map_err(|e| EventStoreError::other(format!("Task join error: {}", e)))?
757    }
758
759    async fn load_snapshot_for(
760        &self,
761        aggregate_type: &str,
762        aggregate_id: &str,
763    ) -> EventStoreResult<Option<Snapshot>> {
764        let aggregate_type = aggregate_type.to_string();
765        let aggregate_id = aggregate_id.to_string();
766        let pool = self.pool.clone();
767
768        tokio::task::spawn_blocking(move || {
769            let mut conn = pool.get().map_err(|e| {
770                EventStoreError::database(format!("Failed to get connection: {}", e))
771            })?;
772
773            let record: Option<SnapshotRecord> = snapshots::table
774                .filter(snapshots::aggregate_type.eq(&aggregate_type))
775                .filter(snapshots::aggregate_id.eq(&aggregate_id))
776                .first::<SnapshotRecord>(&mut conn)
777                .optional()
778                .map_err(|e| EventStoreError::database(e.to_string()))?;
779
780            record.map(|record| record.to_snapshot()).transpose()
781        })
782        .await
783        .map_err(|e| EventStoreError::other(format!("Task join error: {}", e)))?
784    }
785}
786
787#[cfg(test)]
788mod tests {
789    use super::*;
790    use arc_core::audit::AuditMetadata;
791    use serde_json::json;
792
793    async fn setup_test_store() -> SqliteEventStore {
794        let manager = ConnectionManager::<SqliteConnection>::new(":memory:");
795        let pool = Pool::builder()
796            .max_size(1)
797            .build(manager)
798            .expect("Failed to create pool");
799
800        let mut conn = pool.get().expect("Failed to get connection");
801        crate::test_support::migrate(&mut conn);
802        drop(conn);
803
804        SqliteEventStore::with_pool(pool)
805    }
806
807    async fn setup_integrity_test_store() -> SqliteEventStore {
808        let manager = ConnectionManager::<SqliteConnection>::new(":memory:");
809        let pool = Pool::builder()
810            .max_size(1)
811            .build(manager)
812            .expect("Failed to create pool");
813
814        let mut conn = pool.get().expect("Failed to get connection");
815        crate::test_support::migrate(&mut conn);
816        drop(conn);
817
818        SqliteEventStore::with_pool_and_integrity_key(pool, integrity_key(), "test-key")
819            .expect("integrity store")
820    }
821
822    fn integrity_key() -> Vec<u8> {
823        b"012345678901234567890123456789AB".to_vec()
824    }
825
826    /// Helper: build an event with stamped audit.
827    fn stamped_event(
828        agg_type: &str,
829        agg_id: &str,
830        sequence: i64,
831        event_type: &str,
832        payload: serde_json::Value,
833    ) -> Event {
834        Event::new(agg_type, agg_id, sequence, event_type, payload)
835            .with_audit(AuditMetadata::test_default())
836    }
837
838    #[derive(QueryableByName, Debug)]
839    struct SignatureRow {
840        #[diesel(sql_type = diesel::sql_types::Nullable<diesel::sql_types::Text>)]
841        integrity_signature: Option<String>,
842        #[diesel(sql_type = diesel::sql_types::Nullable<diesel::sql_types::Text>)]
843        integrity_key_id: Option<String>,
844    }
845
846    #[tokio::test]
847    async fn test_append_and_load_single_event() {
848        let store = setup_test_store().await;
849        let event = stamped_event(
850            "User",
851            "user-123",
852            1,
853            "UserCreated",
854            json!({ "name": "Alice" }),
855        );
856
857        store
858            .append("user-123", VersionCheck::New, vec![event.clone()])
859            .await
860            .unwrap();
861        let loaded = store.load("user-123").await.unwrap();
862
863        assert_eq!(loaded.len(), 1);
864        assert_eq!(loaded[0].aggregate_id, "user-123");
865        assert_eq!(loaded[0].event_type, "UserCreated");
866        assert_eq!(loaded[0].sequence, 1);
867        assert_eq!(loaded[0].audit.actor_id, "test");
868    }
869
870    #[tokio::test]
871    async fn same_instance_id_is_isolated_by_aggregate_type() {
872        let store = setup_test_store().await;
873        store
874            .append_to(
875                "Product",
876                "shared-id",
877                VersionCheck::New,
878                vec![stamped_event(
879                    "Product",
880                    "shared-id",
881                    1,
882                    "ProductCreated",
883                    json!({}),
884                )],
885            )
886            .await
887            .unwrap();
888        store
889            .append_to(
890                "Order",
891                "shared-id",
892                VersionCheck::New,
893                vec![stamped_event(
894                    "Order",
895                    "shared-id",
896                    1,
897                    "OrderPlaced",
898                    json!({}),
899                )],
900            )
901            .await
902            .unwrap();
903
904        assert_eq!(
905            store.load_stream("Product", "shared-id").await.unwrap()[0].event_type,
906            "ProductCreated"
907        );
908        assert_eq!(
909            store.load_stream("Order", "shared-id").await.unwrap()[0].event_type,
910            "OrderPlaced"
911        );
912    }
913
914    #[tokio::test]
915    async fn test_append_multiple_events() {
916        let store = setup_test_store().await;
917        let events = vec![
918            stamped_event("User", "user-456", 1, "UserCreated", json!({})),
919            stamped_event("User", "user-456", 2, "ProfileUpdated", json!({})),
920            stamped_event("User", "user-456", 3, "EmailChanged", json!({})),
921        ];
922
923        store
924            .append("user-456", VersionCheck::New, events)
925            .await
926            .unwrap();
927        let loaded = store.load("user-456").await.unwrap();
928
929        assert_eq!(loaded.len(), 3);
930        assert_eq!(loaded[0].sequence, 1);
931        assert_eq!(loaded[2].sequence, 3);
932    }
933
934    #[tokio::test]
935    async fn test_integrity_append_persists_signatures() {
936        let store = setup_integrity_test_store().await;
937        store
938            .append(
939                "signed-1",
940                VersionCheck::New,
941                vec![
942                    stamped_event("User", "signed-1", 1, "UserCreated", json!({})),
943                    stamped_event("User", "signed-1", 2, "ProfileUpdated", json!({})),
944                ],
945            )
946            .await
947            .unwrap();
948
949        let pool = store.pool.clone();
950        let rows = tokio::task::spawn_blocking(move || -> EventStoreResult<Vec<SignatureRow>> {
951            let mut conn = pool
952                .get()
953                .map_err(|e| EventStoreError::database(e.to_string()))?;
954            diesel::sql_query(
955                "SELECT integrity_signature, integrity_key_id
956                 FROM events WHERE aggregate_id = 'signed-1' ORDER BY sequence",
957            )
958            .load(&mut *conn)
959            .map_err(|e| EventStoreError::database(e.to_string()))
960        })
961        .await
962        .unwrap()
963        .unwrap();
964
965        assert_eq!(rows.len(), 2);
966        for row in rows {
967            assert_eq!(row.integrity_signature.as_deref().map(str::len), Some(64));
968            assert_eq!(row.integrity_key_id.as_deref(), Some("test-key"));
969        }
970    }
971
972    #[tokio::test]
973    async fn test_integrity_load_rejects_tampered_payload() {
974        let store = setup_integrity_test_store().await;
975        store
976            .append(
977                "tamper-1",
978                VersionCheck::New,
979                vec![stamped_event(
980                    "User",
981                    "tamper-1",
982                    1,
983                    "UserCreated",
984                    json!({"ok": true}),
985                )],
986            )
987            .await
988            .unwrap();
989
990        let pool = store.pool.clone();
991        tokio::task::spawn_blocking(move || -> EventStoreResult<()> {
992            let mut conn = pool
993                .get()
994                .map_err(|e| EventStoreError::database(e.to_string()))?;
995            diesel::sql_query(
996                "UPDATE events SET payload = '{\"ok\": false}' WHERE aggregate_id = 'tamper-1'",
997            )
998            .execute(&mut *conn)
999            .map_err(|e| EventStoreError::database(e.to_string()))?;
1000            Ok(())
1001        })
1002        .await
1003        .unwrap()
1004        .unwrap();
1005
1006        let err = store.load("tamper-1").await.unwrap_err();
1007        assert!(
1008            matches!(err, EventStoreError::Integrity { .. }),
1009            "expected integrity error, got {err:?}"
1010        );
1011    }
1012
1013    #[tokio::test]
1014    async fn test_integrity_load_rejects_missing_signature() {
1015        let store = setup_integrity_test_store().await;
1016        store
1017            .append(
1018                "missing-sig",
1019                VersionCheck::New,
1020                vec![stamped_event(
1021                    "User",
1022                    "missing-sig",
1023                    1,
1024                    "UserCreated",
1025                    json!({}),
1026                )],
1027            )
1028            .await
1029            .unwrap();
1030
1031        let pool = store.pool.clone();
1032        tokio::task::spawn_blocking(move || -> EventStoreResult<()> {
1033            let mut conn = pool
1034                .get()
1035                .map_err(|e| EventStoreError::database(e.to_string()))?;
1036            diesel::sql_query(
1037                "UPDATE events SET integrity_signature = NULL WHERE aggregate_id = 'missing-sig'",
1038            )
1039            .execute(&mut *conn)
1040            .map_err(|e| EventStoreError::database(e.to_string()))?;
1041            Ok(())
1042        })
1043        .await
1044        .unwrap()
1045        .unwrap();
1046
1047        let err = store.load("missing-sig").await.unwrap_err();
1048        assert!(
1049            matches!(err, EventStoreError::Integrity { .. }),
1050            "expected integrity error, got {err:?}"
1051        );
1052    }
1053
1054    #[tokio::test]
1055    async fn test_integrity_load_from_uses_previous_signature() {
1056        let store = setup_integrity_test_store().await;
1057        store
1058            .append(
1059                "load-from-signed",
1060                VersionCheck::New,
1061                vec![
1062                    stamped_event("User", "load-from-signed", 1, "UserCreated", json!({})),
1063                    stamped_event("User", "load-from-signed", 2, "ProfileUpdated", json!({})),
1064                ],
1065            )
1066            .await
1067            .unwrap();
1068
1069        let loaded = store.load_from("load-from-signed", 2).await.unwrap();
1070        assert_eq!(loaded.len(), 1);
1071        assert_eq!(loaded[0].sequence, 2);
1072    }
1073
1074    #[tokio::test]
1075    async fn test_integrity_stream_all_verifies_per_aggregate() {
1076        let store = setup_integrity_test_store().await;
1077        store
1078            .append(
1079                "signed-a",
1080                VersionCheck::New,
1081                vec![stamped_event(
1082                    "User",
1083                    "signed-a",
1084                    1,
1085                    "UserCreated",
1086                    json!({}),
1087                )],
1088            )
1089            .await
1090            .unwrap();
1091        store
1092            .append(
1093                "signed-b",
1094                VersionCheck::New,
1095                vec![stamped_event(
1096                    "User",
1097                    "signed-b",
1098                    1,
1099                    "UserCreated",
1100                    json!({}),
1101                )],
1102            )
1103            .await
1104            .unwrap();
1105        store
1106            .append(
1107                "signed-a",
1108                VersionCheck::Expected(1),
1109                vec![stamped_event(
1110                    "User",
1111                    "signed-a",
1112                    2,
1113                    "ProfileUpdated",
1114                    json!({}),
1115                )],
1116            )
1117            .await
1118            .unwrap();
1119
1120        let loaded = store.stream_all(0).await.unwrap();
1121        assert_eq!(loaded.len(), 3);
1122    }
1123
1124    #[tokio::test]
1125    async fn test_optimistic_concurrency_control() {
1126        let store = setup_test_store().await;
1127        store
1128            .append(
1129                "user-789",
1130                VersionCheck::New,
1131                vec![stamped_event(
1132                    "User",
1133                    "user-789",
1134                    1,
1135                    "UserCreated",
1136                    json!({}),
1137                )],
1138            )
1139            .await
1140            .unwrap();
1141        store
1142            .append(
1143                "user-789",
1144                VersionCheck::Expected(1),
1145                vec![stamped_event(
1146                    "User",
1147                    "user-789",
1148                    2,
1149                    "ProfileUpdated",
1150                    json!({}),
1151                )],
1152            )
1153            .await
1154            .unwrap();
1155        let result = store
1156            .append(
1157                "user-789",
1158                VersionCheck::Expected(1),
1159                vec![stamped_event(
1160                    "User",
1161                    "user-789",
1162                    3,
1163                    "EmailChanged",
1164                    json!({}),
1165                )],
1166            )
1167            .await;
1168        assert!(matches!(
1169            result,
1170            Err(EventStoreError::ConcurrencyConflict {
1171                expected: 1,
1172                actual: 2,
1173                ..
1174            })
1175        ));
1176    }
1177
1178    #[tokio::test]
1179    async fn test_invalid_sequence() {
1180        let store = setup_test_store().await;
1181        let result = store
1182            .append(
1183                "user-999",
1184                VersionCheck::New,
1185                vec![stamped_event(
1186                    "User",
1187                    "user-999",
1188                    5,
1189                    "UserCreated",
1190                    json!({}),
1191                )],
1192            )
1193            .await;
1194        assert!(matches!(
1195            result,
1196            Err(EventStoreError::InvalidSequence {
1197                expected: 1,
1198                actual: 5,
1199                ..
1200            })
1201        ));
1202    }
1203
1204    #[tokio::test]
1205    async fn test_load_from_sequence() {
1206        let store = setup_test_store().await;
1207        let events = vec![
1208            stamped_event("Order", "order-1", 1, "OrderCreated", json!({})),
1209            stamped_event("Order", "order-1", 2, "ItemAdded", json!({})),
1210            stamped_event("Order", "order-1", 3, "ItemAdded", json!({})),
1211            stamped_event("Order", "order-1", 4, "OrderShipped", json!({})),
1212        ];
1213        store
1214            .append("order-1", VersionCheck::New, events)
1215            .await
1216            .unwrap();
1217        let loaded = store.load_from("order-1", 3).await.unwrap();
1218        assert_eq!(loaded.len(), 2);
1219        assert_eq!(loaded[0].sequence, 3);
1220    }
1221
1222    #[tokio::test]
1223    async fn test_get_version() {
1224        let store = setup_test_store().await;
1225        assert_eq!(store.get_version("nope").await.unwrap(), 0);
1226        let events = vec![
1227            stamped_event("User", "u1", 1, "UserCreated", json!({})),
1228            stamped_event("User", "u1", 2, "ProfileUpdated", json!({})),
1229            stamped_event("User", "u1", 3, "EmailChanged", json!({})),
1230        ];
1231        store.append("u1", VersionCheck::New, events).await.unwrap();
1232        assert_eq!(store.get_version("u1").await.unwrap(), 3);
1233    }
1234
1235    #[tokio::test]
1236    async fn test_stream_all() {
1237        let store = setup_test_store().await;
1238        store
1239            .append(
1240                "user-1",
1241                VersionCheck::New,
1242                vec![
1243                    stamped_event("User", "user-1", 1, "UserCreated", json!({})),
1244                    stamped_event("User", "user-1", 2, "ProfileUpdated", json!({})),
1245                ],
1246            )
1247            .await
1248            .unwrap();
1249        store
1250            .append(
1251                "order-1",
1252                VersionCheck::New,
1253                vec![
1254                    stamped_event("Order", "order-1", 1, "OrderCreated", json!({})),
1255                    stamped_event("Order", "order-1", 2, "OrderShipped", json!({})),
1256                ],
1257            )
1258            .await
1259            .unwrap();
1260        assert_eq!(store.stream_all(0).await.unwrap().len(), 4);
1261    }
1262
1263    #[tokio::test]
1264    async fn test_empty_aggregate() {
1265        let store = setup_test_store().await;
1266        assert_eq!(store.load("nothing").await.unwrap().len(), 0);
1267    }
1268
1269    #[tokio::test]
1270    async fn test_audit_roundtrip_preserves_all_fields() {
1271        let store = setup_test_store().await;
1272        let mut audit = AuditMetadata::test_default();
1273        audit.actor_id = "user-uuid-42".to_string();
1274        audit.actor_session_id = Some("sess-XYZ".to_string());
1275        audit.source_ip = Some("10.0.0.42".to_string());
1276        audit.user_agent = Some("Mozilla/5.0 (test)".to_string());
1277        audit.causation_id = Some(Uuid::new_v4());
1278        let expected_corr = audit.correlation_id;
1279        let expected_caus = audit.causation_id;
1280
1281        let event =
1282            Event::new("User", "u-audit", 1, "UserCreated", json!({})).with_audit(audit.clone());
1283
1284        store
1285            .append("u-audit", VersionCheck::New, vec![event])
1286            .await
1287            .unwrap();
1288        let loaded = store.load("u-audit").await.unwrap();
1289
1290        assert_eq!(loaded[0].audit.actor_id, "user-uuid-42");
1291        assert_eq!(
1292            loaded[0].audit.actor_session_id.as_deref(),
1293            Some("sess-XYZ")
1294        );
1295        assert_eq!(loaded[0].audit.source_ip.as_deref(), Some("10.0.0.42"));
1296        assert_eq!(
1297            loaded[0].audit.user_agent.as_deref(),
1298            Some("Mozilla/5.0 (test)")
1299        );
1300        assert_eq!(loaded[0].audit.correlation_id, expected_corr);
1301        assert_eq!(loaded[0].audit.causation_id, expected_caus);
1302        assert!(loaded[0].audit.timestamp_utc_us > 0);
1303    }
1304
1305    #[tokio::test]
1306    async fn test_append_rejects_pending_audit() {
1307        let store = setup_test_store().await;
1308        // Built without with_audit — audit stays pending.
1309        let event = Event::new("User", "u-bad", 1, "UserCreated", json!({}));
1310        let err = store
1311            .append("u-bad", VersionCheck::New, vec![event])
1312            .await
1313            .unwrap_err();
1314        assert!(matches!(err, EventStoreError::InvalidAudit { .. }));
1315
1316        // No row should have been written.
1317        assert_eq!(store.load("u-bad").await.unwrap().len(), 0);
1318    }
1319
1320    #[tokio::test]
1321    async fn test_actor_id_index_used() {
1322        let store = setup_test_store().await;
1323        let mut a = AuditMetadata::test_default();
1324        a.actor_id = "alice-uuid".to_string();
1325        let event = Event::new("User", "u1", 1, "UserCreated", json!({})).with_audit(a);
1326        store
1327            .append("u1", VersionCheck::New, vec![event])
1328            .await
1329            .unwrap();
1330
1331        // EXPLAIN QUERY PLAN must show an index search on actor_id
1332        let pool = store.pool.clone();
1333        let plan = tokio::task::spawn_blocking(move || -> EventStoreResult<Vec<String>> {
1334            let mut conn = pool
1335                .get()
1336                .map_err(|e| EventStoreError::database(e.to_string()))?;
1337            let plan: Vec<(i32, i32, i32, String)> = diesel::sql_query(
1338                "EXPLAIN QUERY PLAN SELECT * FROM events WHERE actor_id = 'alice-uuid'",
1339            )
1340            .load::<ExplainRow>(&mut *conn)
1341            .map_err(|e| EventStoreError::database(e.to_string()))?
1342            .into_iter()
1343            .map(|r| (r.id, r.parent, r.notused, r.detail))
1344            .collect();
1345            Ok(plan.into_iter().map(|(_, _, _, d)| d).collect())
1346        })
1347        .await
1348        .unwrap()
1349        .unwrap();
1350
1351        assert!(
1352            plan.iter().any(|d| d.contains("idx_events_actor_id")),
1353            "actor_id query did not use index; plan: {:?}",
1354            plan
1355        );
1356    }
1357
1358    #[derive(QueryableByName, Debug)]
1359    struct ExplainRow {
1360        #[diesel(sql_type = diesel::sql_types::Integer)]
1361        id: i32,
1362        #[diesel(sql_type = diesel::sql_types::Integer)]
1363        parent: i32,
1364        #[diesel(sql_type = diesel::sql_types::Integer)]
1365        notused: i32,
1366        #[diesel(sql_type = diesel::sql_types::Text)]
1367        detail: String,
1368    }
1369
1370    #[tokio::test]
1371    async fn test_concurrent_appends() {
1372        let store = setup_test_store().await;
1373        store
1374            .append(
1375                "uc",
1376                VersionCheck::New,
1377                vec![stamped_event("User", "uc", 1, "UserCreated", json!({}))],
1378            )
1379            .await
1380            .unwrap();
1381
1382        let s1 = store.clone();
1383        let s2 = store.clone();
1384        let h1 = tokio::spawn(async move {
1385            s1.append(
1386                "uc",
1387                VersionCheck::Expected(1),
1388                vec![stamped_event("User", "uc", 2, "U1", json!({}))],
1389            )
1390            .await
1391        });
1392        let h2 = tokio::spawn(async move {
1393            s2.append(
1394                "uc",
1395                VersionCheck::Expected(1),
1396                vec![stamped_event("User", "uc", 2, "U2", json!({}))],
1397            )
1398            .await
1399        });
1400        let r1 = h1.await.unwrap();
1401        let r2 = h2.await.unwrap();
1402        assert!(r1.is_ok() != r2.is_ok());
1403    }
1404
1405    #[tokio::test]
1406    async fn test_sequence_above_i32_max_roundtrips_without_truncation() {
1407        // Pre-fix bug: sequence was cast to i32 on insert and read back as i32.
1408        // A value above i32::MAX (2_147_483_647) would silently overflow.
1409        // After widening, this round-trips intact.
1410        let store = setup_test_store().await;
1411        let pool = store.pool.clone();
1412
1413        let huge_seq: i64 = (i32::MAX as i64) + 1234;
1414        let huge_ts: i64 = 9_999_999_999; // year 2286 — would never fit in i32
1415
1416        let inserted = tokio::task::spawn_blocking(move || -> EventStoreResult<usize> {
1417            let mut conn = pool
1418                .get()
1419                .map_err(|e| EventStoreError::database(e.to_string()))?;
1420            diesel::sql_query(format!(
1421                "INSERT INTO events (event_id, aggregate_type, aggregate_id, sequence,
1422                    event_type, payload, timestamp,
1423                    actor_id, timestamp_utc_us, correlation_id)
1424                 VALUES ('aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa', 'User', 'u-big', {seq},
1425                    'Event', '{{}}', {ts},
1426                    'tester', {seq}, '00000000-0000-0000-0000-000000000001')",
1427                seq = huge_seq,
1428                ts = huge_ts,
1429            ))
1430            .execute(&mut *conn)
1431            .map_err(|e| EventStoreError::database(e.to_string()))
1432        })
1433        .await
1434        .unwrap()
1435        .unwrap();
1436        assert_eq!(inserted, 1);
1437
1438        let loaded = store.load("u-big").await.expect("load");
1439        assert_eq!(loaded.len(), 1);
1440        assert_eq!(loaded[0].sequence, huge_seq, "sequence must not truncate");
1441        // timestamp stored as seconds; round-trip back to milliseconds in Event::timestamp
1442        assert_eq!(loaded[0].timestamp, (huge_ts as u64) * 1000);
1443
1444        let v = store.get_version("u-big").await.expect("version");
1445        assert_eq!(v, huge_seq, "get_version must not truncate either");
1446    }
1447
1448    #[tokio::test]
1449    async fn test_legacy_backfilled_row_roundtrips() {
1450        // Simulates a row written before HIPAA-1, then backfilled by the migration:
1451        // actor_id='legacy-pre-hipaa', timestamp_utc_us derived from seconds*1_000_000,
1452        // correlation_id = nil UUID. The row must load without panicking.
1453        let store = setup_test_store().await;
1454        let pool = store.pool.clone();
1455
1456        let inserted_count = tokio::task::spawn_blocking(move || -> EventStoreResult<usize> {
1457            let mut conn = pool
1458                .get()
1459                .map_err(|e| EventStoreError::database(e.to_string()))?;
1460            diesel::sql_query(
1461                "INSERT INTO events (event_id, aggregate_type, aggregate_id, sequence,
1462                    event_type, payload, timestamp,
1463                    actor_id, timestamp_utc_us, correlation_id)
1464                 VALUES ('11111111-1111-1111-1111-111111111111', 'User', 'u-legacy', 1,
1465                    'UserCreated', '{}', 1700000000,
1466                    'legacy-pre-hipaa', 1700000000000000,
1467                    '00000000-0000-0000-0000-000000000000')",
1468            )
1469            .execute(&mut *conn)
1470            .map_err(|e| EventStoreError::database(e.to_string()))
1471        })
1472        .await
1473        .unwrap()
1474        .unwrap();
1475        assert_eq!(inserted_count, 1);
1476
1477        let loaded = store.load("u-legacy").await.expect("legacy row must load");
1478        assert_eq!(loaded.len(), 1);
1479        assert_eq!(loaded[0].audit.actor_id, "legacy-pre-hipaa");
1480        assert_eq!(loaded[0].audit.timestamp_utc_us, 1_700_000_000_000_000);
1481        assert_eq!(loaded[0].audit.correlation_id, Uuid::nil());
1482        assert!(loaded[0].audit.causation_id.is_none());
1483    }
1484
1485    #[tokio::test]
1486    async fn test_load_rejects_malformed_correlation_uuid() {
1487        let store = setup_test_store().await;
1488        let pool = store.pool.clone();
1489
1490        tokio::task::spawn_blocking(move || -> EventStoreResult<()> {
1491            let mut conn = pool
1492                .get()
1493                .map_err(|e| EventStoreError::database(e.to_string()))?;
1494            diesel::sql_query(
1495                "INSERT INTO events (event_id, aggregate_type, aggregate_id, sequence,
1496                    event_type, payload, timestamp,
1497                    actor_id, timestamp_utc_us, correlation_id)
1498                 VALUES ('22222222-2222-2222-2222-222222222222', 'User', 'u-bad', 1,
1499                    'X', '{}', 1700000000,
1500                    'tester', 1700000000000000,
1501                    'not-a-uuid')",
1502            )
1503            .execute(&mut *conn)
1504            .map_err(|e| EventStoreError::database(e.to_string()))?;
1505            Ok(())
1506        })
1507        .await
1508        .unwrap()
1509        .unwrap();
1510
1511        let err = store.load("u-bad").await.unwrap_err();
1512        assert!(
1513            matches!(err, EventStoreError::SerializationError { ref message } if message.contains("Invalid correlation UUID")),
1514            "expected SerializationError on malformed correlation_id, got {:?}",
1515            err
1516        );
1517    }
1518
1519    #[tokio::test]
1520    async fn test_caused_by_chain_roundtrips_through_sqlite() {
1521        use arc_core::aggregate::{Aggregate, Command};
1522        use arc_core::command_bus::{CommandBus, CommandContext};
1523        use arc_core::event::Event as CoreEvent;
1524        use arc_core::event_bus::InProcessEventBus;
1525
1526        // Trivial aggregate so we can exercise CommandBus + SqliteEventStore together.
1527        #[derive(Default)]
1528        struct Counter {
1529            v: i64,
1530        }
1531        struct Cmd {
1532            id: String,
1533        }
1534        impl Command for Cmd {
1535            fn aggregate_id(&self) -> &str {
1536                &self.id
1537            }
1538        }
1539        #[derive(Debug, thiserror::Error)]
1540        #[error("never")]
1541        struct Never;
1542        #[async_trait]
1543        impl Aggregate for Counter {
1544            type Command = Cmd;
1545            type Event = ();
1546            type Error = Never;
1547            fn aggregate_type() -> &'static str {
1548                "Counter"
1549            }
1550            fn version(&self) -> i64 {
1551                self.v
1552            }
1553            async fn handle(&self, c: Self::Command) -> Result<Vec<CoreEvent>, Self::Error> {
1554                Ok(vec![CoreEvent::new(
1555                    "Counter",
1556                    &c.id,
1557                    self.v + 1,
1558                    "Incremented",
1559                    serde_json::json!({}),
1560                )])
1561            }
1562            fn apply(&mut self, e: &CoreEvent) {
1563                self.v = e.sequence;
1564            }
1565        }
1566
1567        let store = setup_test_store().await;
1568        let bus =
1569            CommandBus::<Counter>::new(Box::new(store.clone()), Box::new(InProcessEventBus::new()));
1570
1571        let first_ctx = CommandContext::for_actor("alice");
1572        let trigger_corr = first_ctx.correlation_id;
1573        let triggers = bus
1574            .dispatch(Cmd { id: "c1".into() }, first_ctx)
1575            .await
1576            .unwrap();
1577
1578        let follow_ctx = CommandContext::caused_by("worker", &triggers[0]);
1579        let follow_corr = follow_ctx.correlation_id;
1580        let _ = bus
1581            .dispatch(Cmd { id: "c2".into() }, follow_ctx)
1582            .await
1583            .unwrap();
1584
1585        // Reload from SQLite — chain must survive the trip
1586        let loaded_c2 = store.load("c2").await.unwrap();
1587        assert_eq!(loaded_c2.len(), 1);
1588        assert_eq!(loaded_c2[0].audit.correlation_id, trigger_corr);
1589        assert_eq!(loaded_c2[0].audit.correlation_id, follow_corr);
1590        assert_eq!(loaded_c2[0].audit.causation_id, Some(triggers[0].event_id));
1591    }
1592
1593    #[tokio::test]
1594    async fn test_event_ordering_within_aggregate() {
1595        let store = setup_test_store().await;
1596        store
1597            .append(
1598                "uo",
1599                VersionCheck::New,
1600                vec![
1601                    stamped_event("User", "uo", 1, "UserCreated", json!({})),
1602                    stamped_event("User", "uo", 2, "EmailChanged", json!({})),
1603                ],
1604            )
1605            .await
1606            .unwrap();
1607        store
1608            .append(
1609                "uo",
1610                VersionCheck::Expected(2),
1611                vec![
1612                    stamped_event("User", "uo", 3, "ProfileUpdated", json!({})),
1613                    stamped_event("User", "uo", 4, "PasswordChanged", json!({})),
1614                ],
1615            )
1616            .await
1617            .unwrap();
1618        let loaded = store.load("uo").await.unwrap();
1619        for (i, e) in loaded.iter().enumerate() {
1620            assert_eq!(e.sequence, (i + 1) as i64);
1621        }
1622    }
1623
1624    #[tokio::test]
1625    async fn test_save_then_load_snapshot() {
1626        let store = setup_test_store().await;
1627        let snap = Snapshot::new("agg-1", "User", 5, json!({ "name": "Alice" }));
1628        store.save_snapshot(&snap).await.unwrap();
1629        let loaded = store.load_snapshot("agg-1").await.unwrap();
1630        assert_eq!(loaded, Some(snap));
1631    }
1632
1633    #[tokio::test]
1634    async fn test_load_snapshot_unknown_returns_none() {
1635        let store = setup_test_store().await;
1636        assert_eq!(store.load_snapshot("missing").await.unwrap(), None);
1637    }
1638
1639    #[tokio::test]
1640    async fn test_save_snapshot_upserts_newer_version() {
1641        let store = setup_test_store().await;
1642        store
1643            .save_snapshot(&Snapshot::new("agg-2", "User", 3, json!({ "v": 3 })))
1644            .await
1645            .unwrap();
1646        let newer = Snapshot::new("agg-2", "User", 9, json!({ "v": 9 }));
1647        store.save_snapshot(&newer).await.unwrap();
1648        let loaded = store.load_snapshot("agg-2").await.unwrap().unwrap();
1649        assert_eq!(loaded.version, 9);
1650        assert_eq!(loaded.state["v"], 9);
1651    }
1652}