turnframe_store_postgres/commit.rs
1//! The all-or-nothing write of a [`CommitBundle`] (spec §16.3, §23 step N).
2//!
3//! Everything one commit produces — the journal outcome of each command, the
4//! events, the resolution of the card that authorized them, the cards the new
5//! revision invalidates, the cards the new state requires, the outbox rows, the
6//! replay record and the phase marker — travels in one bundle and lands in one
7//! PostgreSQL transaction, in the normative order of
8//! [`turnframe_store::commit`].
9//!
10//! Every item goes through exactly the same function its own trait uses, so it
11//! obeys exactly the same rules: the same compare-and-swap, the same uniqueness,
12//! the same tenant scoping. The only difference is that the transaction is not
13//! committed until the last item has succeeded. An item that fails returns its
14//! error, the transaction is dropped, PostgreSQL rolls it back, and nothing the
15//! bundle carried was ever visible to another connection.
16//!
17//! What is deliberately *outside* this transaction is the domain executor's own
18//! state commit, which may live in another database; Turnframe does not attempt
19//! a distributed transaction. When the domain tables do live here, enlist them
20//! with [`PgStores::commit_in`] and the two become one transaction.
21
22use async_trait::async_trait;
23use chrono::{DateTime, Utc};
24use sqlx::PgConnection;
25use turnframe_core::ids::AccountId;
26use turnframe_store::commit::{CommitBundle, CommitReceipt, CommitStore};
27use turnframe_store::error::StoreError;
28
29use crate::codec::now;
30use crate::store::{PgStores, commit as commit_transaction};
31use crate::{conversations, events, interactions, journal, outbox, replay};
32
33/// Applies every item of `bundle` on `conn`, in the contract's order.
34///
35/// The caller owns the transaction: nothing here commits, and returning an error
36/// leaves it to be rolled back.
37pub(crate) async fn apply_bundle(
38 conn: &mut PgConnection,
39 account: &AccountId,
40 bundle: CommitBundle,
41 at: DateTime<Utc>,
42) -> Result<CommitReceipt, StoreError> {
43 // The account and empty-batch checks every implementation owes, before a
44 // single row is touched.
45 bundle.validate(account)?;
46 let mut event_ids = Vec::new();
47 let mut inserted_interactions = Vec::new();
48 let mut invalidated_interactions = Vec::new();
49
50 for completion in bundle.journal_completions {
51 journal::complete(
52 &mut *conn,
53 account,
54 &completion.command_id,
55 completion.outcome,
56 at,
57 )
58 .await?;
59 }
60
61 for batch in bundle.events {
62 event_ids.extend(events::append(&mut *conn, batch).await?);
63 }
64
65 for finish in bundle.interaction_finishes {
66 interactions::finish_resolution(
67 &mut *conn,
68 account,
69 &finish.interaction_id,
70 finish.outcome,
71 )
72 .await?;
73 }
74
75 // Invalidation comes before the inserts: a card the new revision retires
76 // must leave the blocking slot before the card that replaces it takes it.
77 for invalidation in bundle.interaction_invalidations {
78 invalidated_interactions.extend(
79 interactions::invalidate_for_case(
80 &mut *conn,
81 account,
82 &invalidation.case_key,
83 invalidation.new_revision,
84 invalidation.reason,
85 at,
86 )
87 .await?,
88 );
89 }
90
91 for insert in bundle.interaction_inserts {
92 let id = insert.interaction.id;
93 invalidated_interactions.extend(
94 interactions::insert_interaction(
95 &mut *conn,
96 insert.interaction,
97 insert.replace_blocking,
98 at,
99 )
100 .await?,
101 );
102 inserted_interactions.push(id);
103 }
104
105 for entry in bundle.outbox_entries {
106 outbox::enqueue(&mut *conn, entry).await?;
107 }
108
109 if let Some(record) = bundle.replay_record {
110 replay::put(&mut *conn, record).await?;
111 }
112
113 if let Some(update) = bundle.turn_phase {
114 conversations::set_turn_phase(&mut *conn, account, &update.turn_id, update.phase, at)
115 .await?;
116 }
117
118 Ok(CommitReceipt {
119 event_ids,
120 inserted_interactions,
121 invalidated_interactions,
122 committed_at: at,
123 })
124}
125
126impl PgStores {
127 /// Applies a bundle inside a transaction the caller owns.
128 ///
129 /// This is the seam an adopter needs when the domain tables live in the same
130 /// database as the stores: begin one transaction, run the workflow
131 /// executor's own writes on it, hand it to this method, and commit once. The
132 /// journal admission still happens before execution, so recovery works the
133 /// same way if the process dies — the transaction simply removes the window
134 /// in which the domain state and the bookkeeping could disagree.
135 ///
136 /// Nothing is committed here. Commit the transaction to make the bundle
137 /// visible; drop it to discard the bundle along with your own writes.
138 ///
139 /// ```rust,no_run
140 /// use turnframe_core::ids::AccountId;
141 /// use turnframe_store::commit::CommitBundle;
142 /// use turnframe_store_postgres::PgStores;
143 ///
144 /// # async fn example(store: &PgStores, bundle: CommitBundle) -> Result<(), Box<dyn std::error::Error>> {
145 /// let mut transaction = store.pool().begin().await?;
146 /// // ... the executor's own writes go here, on the same transaction ...
147 /// let receipt = store
148 /// .commit_in(&mut transaction, &AccountId::from("aurora"), bundle)
149 /// .await?;
150 /// transaction.commit().await?;
151 /// # let _ = receipt;
152 /// # Ok(())
153 /// # }
154 /// ```
155 ///
156 /// # Errors
157 /// Any error an item would raise through its own store trait, plus the
158 /// refusals of [`CommitBundle::validate`].
159 pub async fn commit_in(
160 &self,
161 conn: &mut PgConnection,
162 account: &AccountId,
163 bundle: CommitBundle,
164 ) -> Result<CommitReceipt, StoreError> {
165 apply_bundle(conn, account, bundle, now()).await
166 }
167}
168
169#[async_trait]
170impl CommitStore for PgStores {
171 async fn commit(
172 &self,
173 account: &AccountId,
174 bundle: CommitBundle,
175 ) -> Result<CommitReceipt, StoreError> {
176 // Refuse a foreign item before a connection is even taken.
177 bundle.validate(account)?;
178 let mut transaction = self.transaction().await?;
179 let receipt = apply_bundle(&mut transaction, account, bundle, now()).await?;
180 commit_transaction(transaction).await?;
181 Ok(receipt)
182 }
183}