1use 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
25pub 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#[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 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#[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#[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 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 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 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 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 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 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 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; 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 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 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 #[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 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}