Skip to main content

kmp_adapter_embedded/adapter/
replay.rs

1use 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/// Outcome of a projection rebuild from the append-only event log.
10#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub struct ProjectionRebuildReport {
12    pub events_replayed: u64,
13    pub mutations_applied: u64,
14}
15
16impl EmbeddedKernelStore {
17    /// Drops every projection table and rebuilds them by replaying the event
18    /// log in sequence order — the recovery and migration story in one.
19    ///
20    /// The mutation derivation is injected so this adapter stays free of
21    /// application-layer dependencies; the composition root passes
22    /// `kmp_application::projection_mutations_for_context_event`.
23    /// The whole rebuild is one transaction: a crash mid-rebuild leaves the
24    /// previous projections intact.
25    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}