Skip to main content

turnframe_store/
commit.rs

1//! The atomicity model: one [`CommitBundle`], one all-or-nothing write
2//! (spec §16.3, §23 step N, ADR-006 point 8).
3//!
4//! After the executor has committed the domain state, the runtime still has the
5//! journal outcomes, the committed events, the card resolutions, the cards the
6//! new state requires and the ones the new revision invalidates, the outbox
7//! rows, the replay record and the phase marker to write. A partial write across
8//! those is an invariant violation, not a degraded success.
9//!
10//! So [`CommitStore::commit`] takes the whole bundle and applies it all or
11//! nothing, and "nothing" covers the records a bundle *changes*, not only the
12//! ones it creates — which is what
13//! [`crate::conformance::check_commit_bundle_restores_modified_records`] holds
14//! an implementation to.
15//!
16//! Why the executor's own commit is deliberately outside that transaction, and
17//! what makes the seam safe anyway, is in
18//! [`docs/persistence.md`](https://github.com/turnframe-rs/turnframe/blob/main/docs/persistence.md).
19
20use async_trait::async_trait;
21use chrono::{DateTime, Utc};
22use serde::{Deserialize, Serialize};
23use turnframe_core::case::CaseKey;
24use turnframe_core::event::OutboxEntry;
25use turnframe_core::ids::{AccountId, CaseRevision, CommandId, EventId, InteractionId, TurnId};
26use turnframe_core::interaction::Interaction;
27use turnframe_core::replay::{ReplayRecord, TurnPhase};
28
29use crate::error::{StoreError, bundle_account_mismatch, invalid_record};
30use crate::events::EventBatch;
31use crate::interaction::{InvalidationReason, ResolutionOutcome};
32use crate::journal::JournalOutcome;
33
34/// Records the outcome of one journaled command.
35#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
36pub struct JournalCompletion {
37    /// The command.
38    pub command_id: CommandId,
39    /// Its outcome.
40    pub outcome: JournalOutcome,
41}
42
43/// Settles one `Resolving` interaction.
44#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
45pub struct InteractionFinish {
46    /// The interaction.
47    pub interaction_id: InteractionId,
48    /// How it is settled.
49    pub outcome: ResolutionOutcome,
50}
51
52/// Inserts one new interaction.
53#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
54pub struct InteractionInsert {
55    /// The interaction, in status `Active`.
56    pub interaction: Interaction,
57    /// `true` to invalidate an `Active` blocking occupant of the same case
58    /// (see `InteractionWriter::insert_replacing_blocking`); `false` to fail with
59    /// `Conflict` when the slot is taken.
60    pub replace_blocking: bool,
61}
62
63/// Invalidates the revision-bound interactions of one case.
64#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65pub struct CaseInvalidation {
66    /// The case.
67    pub case_key: CaseKey,
68    /// The revision the case moved to.
69    pub new_revision: CaseRevision,
70    /// The reason recorded on every invalidated interaction.
71    pub reason: InvalidationReason,
72}
73
74/// Writes the phase marker of a turn.
75#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
76pub struct TurnPhaseUpdate {
77    /// The turn.
78    pub turn_id: TurnId,
79    /// The new phase.
80    pub phase: TurnPhase,
81}
82
83/// Everything one commit writes, applied all or nothing.
84///
85/// Compared with [`PartialEq`] only, because the replay record it may carry
86/// holds a sampling temperature and therefore a float.
87#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
88#[non_exhaustive]
89pub struct CommitBundle {
90    /// Journal outcomes to record.
91    #[serde(default)]
92    pub journal_completions: Vec<JournalCompletion>,
93    /// Event batches to append.
94    #[serde(default)]
95    pub events: Vec<EventBatch>,
96    /// Interactions to settle.
97    #[serde(default)]
98    pub interaction_finishes: Vec<InteractionFinish>,
99    /// Cases whose revision-bound interactions are invalidated.
100    #[serde(default)]
101    pub interaction_invalidations: Vec<CaseInvalidation>,
102    /// Interactions to insert.
103    #[serde(default)]
104    pub interaction_inserts: Vec<InteractionInsert>,
105    /// Outbox rows to enqueue.
106    #[serde(default)]
107    pub outbox_entries: Vec<OutboxEntry>,
108    /// Replay record to upsert.
109    #[serde(default, skip_serializing_if = "Option::is_none")]
110    pub replay_record: Option<ReplayRecord>,
111    /// Turn phase to write.
112    #[serde(default, skip_serializing_if = "Option::is_none")]
113    pub turn_phase: Option<TurnPhaseUpdate>,
114}
115
116impl CommitBundle {
117    /// An empty bundle.
118    #[must_use]
119    pub fn new() -> Self {
120        Self::default()
121    }
122
123    /// Adds a journal completion.
124    #[must_use]
125    pub fn with_journal_completion(
126        mut self,
127        command_id: CommandId,
128        outcome: JournalOutcome,
129    ) -> Self {
130        self.journal_completions.push(JournalCompletion {
131            command_id,
132            outcome,
133        });
134        self
135    }
136
137    /// Adds an event batch.
138    #[must_use]
139    pub fn with_events(mut self, batch: EventBatch) -> Self {
140        self.events.push(batch);
141        self
142    }
143
144    /// Adds an interaction finish.
145    #[must_use]
146    pub fn with_interaction_finish(
147        mut self,
148        interaction_id: InteractionId,
149        outcome: ResolutionOutcome,
150    ) -> Self {
151        self.interaction_finishes.push(InteractionFinish {
152            interaction_id,
153            outcome,
154        });
155        self
156    }
157
158    /// Adds a case invalidation.
159    #[must_use]
160    pub fn with_invalidation(
161        mut self,
162        case_key: CaseKey,
163        new_revision: CaseRevision,
164        reason: InvalidationReason,
165    ) -> Self {
166        self.interaction_invalidations.push(CaseInvalidation {
167            case_key,
168            new_revision,
169            reason,
170        });
171        self
172    }
173
174    /// Adds an interaction insert.
175    #[must_use]
176    pub fn with_interaction_insert(
177        mut self,
178        interaction: Interaction,
179        replace_blocking: bool,
180    ) -> Self {
181        self.interaction_inserts.push(InteractionInsert {
182            interaction,
183            replace_blocking,
184        });
185        self
186    }
187
188    /// Adds an outbox entry.
189    #[must_use]
190    pub fn with_outbox_entry(mut self, entry: OutboxEntry) -> Self {
191        self.outbox_entries.push(entry);
192        self
193    }
194
195    /// Sets the replay record.
196    #[must_use]
197    pub fn with_replay_record(mut self, record: ReplayRecord) -> Self {
198        self.replay_record = Some(record);
199        self
200    }
201
202    /// Sets the turn phase update.
203    #[must_use]
204    pub fn with_turn_phase(mut self, turn_id: TurnId, phase: TurnPhase) -> Self {
205        self.turn_phase = Some(TurnPhaseUpdate { turn_id, phase });
206        self
207    }
208
209    /// Returns `true` when the bundle writes nothing.
210    #[must_use]
211    pub fn is_empty(&self) -> bool {
212        self.journal_completions.is_empty()
213            && self.events.is_empty()
214            && self.interaction_finishes.is_empty()
215            && self.interaction_invalidations.is_empty()
216            && self.interaction_inserts.is_empty()
217            && self.outbox_entries.is_empty()
218            && self.replay_record.is_none()
219            && self.turn_phase.is_none()
220    }
221
222    /// Checks the bundle before anything is written: every account-bearing item
223    /// must belong to `account`, and no event batch may be empty.
224    ///
225    /// Implementations call this first; adopters writing their own store should
226    /// too, so a foreign item is refused identically everywhere.
227    ///
228    /// # Errors
229    /// * `Other(BUNDLE_ACCOUNT_MISMATCH)` for a foreign item.
230    /// * `Other(INVALID_RECORD)` for an empty event batch.
231    pub fn validate(&self, account: &AccountId) -> Result<(), StoreError> {
232        for batch in &self.events {
233            if &batch.account_id != account {
234                return Err(bundle_account_mismatch());
235            }
236            if batch.is_empty() {
237                return Err(invalid_record());
238            }
239        }
240        if self
241            .interaction_inserts
242            .iter()
243            .any(|insert| &insert.interaction.account_id != account)
244        {
245            return Err(bundle_account_mismatch());
246        }
247        if self
248            .replay_record
249            .as_ref()
250            .is_some_and(|record| &record.account_id != account)
251        {
252            return Err(bundle_account_mismatch());
253        }
254        Ok(())
255    }
256}
257
258/// What a successful commit reports back.
259#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
260pub struct CommitReceipt {
261    /// Every event appended, in append order.
262    pub event_ids: Vec<EventId>,
263    /// Interactions inserted, in bundle order.
264    pub inserted_interactions: Vec<InteractionId>,
265    /// Interactions invalidated by invalidations and replacing inserts, in
266    /// application order.
267    pub invalidated_interactions: Vec<InteractionId>,
268    /// When the commit was written, stamped by the store's clock.
269    pub committed_at: DateTime<Utc>,
270}
271
272/// The all-or-nothing write of a [`CommitBundle`].
273///
274/// See the module documentation for the atomicity model. The conformance suite
275/// proves the contract with [`crate::conformance::check_commit_bundle_applies_all`],
276/// [`crate::conformance::check_commit_bundle_atomic_on_invalid_item`] and
277/// [`crate::conformance::check_commit_bundle_rejects_foreign_account_items`].
278///
279/// This is the one persistence trait with no `…Reader` / `…Writer` split,
280/// because it is entirely the write half: a commit *is* the write. It is
281/// therefore absent from
282/// [`ReadOnlyStores`](crate::stores::ReadOnlyStores) altogether, which is the
283/// point — a path that only holds the read half cannot name it.
284#[async_trait]
285pub trait CommitStore: Send + Sync {
286    /// Applies the bundle for `account` atomically.
287    ///
288    /// # Errors
289    /// Any error an item would raise through its own store trait, plus the
290    /// refusals of [`CommitBundle::validate`]. On error nothing of the bundle
291    /// is visible.
292    async fn commit(
293        &self,
294        account: &AccountId,
295        bundle: CommitBundle,
296    ) -> Result<CommitReceipt, StoreError>;
297}
298
299#[cfg(test)]
300mod tests {
301    use super::*;
302    use crate::error::{codes, has_code};
303    use turnframe_core::ids::ConversationId;
304
305    #[test]
306    fn empty_bundle_is_empty_and_valid() {
307        let bundle = CommitBundle::new();
308        assert!(bundle.is_empty());
309        assert!(bundle.validate(&AccountId::from("a")).is_ok());
310        let with_phase = bundle.with_turn_phase(TurnId::nil(), TurnPhase::Committed);
311        assert!(!with_phase.is_empty());
312    }
313
314    #[test]
315    fn validate_refuses_foreign_and_empty_items() {
316        let account = AccountId::from("a");
317        let foreign = CommitBundle::new().with_events(EventBatch::new(
318            AccountId::from("b"),
319            CaseKey::new("w", "c"),
320            CommandId::nil(),
321            CaseRevision(1),
322            vec![],
323        ));
324        assert!(has_code(
325            &foreign.validate(&account).unwrap_err(),
326            codes::BUNDLE_ACCOUNT_MISMATCH
327        ));
328        let empty = CommitBundle::new().with_events(EventBatch::new(
329            account.clone(),
330            CaseKey::new("w", "c"),
331            CommandId::nil(),
332            CaseRevision(1),
333            vec![],
334        ));
335        assert!(has_code(
336            &empty.validate(&account).unwrap_err(),
337            codes::INVALID_RECORD
338        ));
339        let replay = CommitBundle::new().with_replay_record(ReplayRecord::received(
340            TurnId::nil(),
341            ConversationId::nil(),
342            AccountId::from("b"),
343            DateTime::<Utc>::UNIX_EPOCH,
344        ));
345        assert!(has_code(
346            &replay.validate(&account).unwrap_err(),
347            codes::BUNDLE_ACCOUNT_MISMATCH
348        ));
349    }
350}