kmp_adapter_embedded/adapter/
replay.rs1use kmp_domain::{ContextUpdatedEvent, PortError, ProjectionMutation};
2
3use super::projection_write::apply_mutations_in_transaction;
4use super::store::{
5 ANCHORS, DETAILS, EmbeddedKernelStore, NODES, RELATIONS, RELATIONS_BY_TARGET, commit_error,
6 table_error,
7};
8
9#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub struct ProjectionRebuildReport {
12 pub events_replayed: u64,
13 pub mutations_applied: u64,
14}
15
16impl EmbeddedKernelStore {
17 pub async fn rebuild_projections<F>(
26 &self,
27 derive: F,
28 ) -> Result<ProjectionRebuildReport, PortError>
29 where
30 F: Fn(&ContextUpdatedEvent) -> Result<Vec<ProjectionMutation>, PortError> + Send + 'static,
31 {
32 let events = self.run(EmbeddedKernelStore::read_event_log).await?;
33
34 let mut mutations = Vec::new();
35 for event in &events {
36 mutations.extend(derive(event)?);
37 }
38 let events_replayed = events.len() as u64;
39
40 self.run(move |store| {
41 let tx = store.begin_write()?;
42 tx.delete_table(NODES).map_err(table_error)?;
43 tx.delete_table(RELATIONS).map_err(table_error)?;
44 tx.delete_table(RELATIONS_BY_TARGET).map_err(table_error)?;
45 tx.delete_table(DETAILS).map_err(table_error)?;
46 tx.delete_table(ANCHORS).map_err(table_error)?;
47
48 let mutations_applied = apply_mutations_in_transaction(&tx, mutations)?;
49 tx.commit().map_err(commit_error)?;
50
51 Ok(ProjectionRebuildReport {
52 events_replayed,
53 mutations_applied,
54 })
55 })
56 .await
57 }
58}