Skip to main content

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}