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#[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 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#[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#[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 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 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 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 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 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 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 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; 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 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 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 #[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 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}