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