Skip to main content

coven_database/store/store_session/
store_acknowledgements.rs

1use crate::*;
2use coven_protocol::store_commit::{
3    StoreAck, StoreAckRef, StoreBatchCommitRef, StoreDeviceRegistrationRef,
4};
5use rusqlite::OptionalExtension;
6
7use super::*;
8use crate::store_ack_records::{load_expected_outbound_store_ack_on, load_outbound_store_ack_on};
9
10impl StoreSession<'_> {
11    fn latest_local_store_ack(&self) -> Result<Option<PublishedStoreAck>, DbError> {
12        load_published_store_ack_on(self.conn)
13    }
14
15    fn activated_store_ack(
16        &self,
17        registration: &StoreDeviceRegistrationRef,
18    ) -> Result<Option<ActivatedStoreAck>, DbError> {
19        self.conn
20            .query_row(
21                "SELECT ack_ref, activating_commit FROM activated_store_acks WHERE device_id = ?1",
22                [registration.device_id.to_string()],
23                |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
24            )
25            .optional()
26            .map_err(DbError::from)?
27            .map(|(raw, activating_commit)| {
28                let reference: StoreAckRef = serde_json::from_str(&raw).map_err(|error| {
29                    DbError::context("activated Store acknowledgement ref", error)
30                })?;
31                if &reference.registration != registration {
32                    return Err(DbError::Message(
33                        "activated Store acknowledgement names another registration".to_string(),
34                    ));
35                }
36                Ok(ActivatedStoreAck {
37                    reference,
38                    activating_commit: serde_json::from_str(&activating_commit).map_err(
39                        |error| {
40                            DbError::context(
41                                "activated Store acknowledgement activating commit",
42                                error,
43                            )
44                        },
45                    )?,
46                })
47            })
48            .transpose()
49    }
50
51    fn stage_store_ack(
52        &mut self,
53        ack: StoreAck,
54        prepared: PreparedExactObject,
55    ) -> Result<StoreAckRef, DbError> {
56        let authority = self.local_store_authority()?;
57        let bytes = ack.to_bytes();
58        let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
59        let (reference, verified) =
60            verify_next_local_store_ack_on(&tx, &authority, &bytes, &prepared)?;
61        if verified != ack {
62            return Err(DbError::Message(
63                "staged Store acknowledgement changed during exact verification".to_string(),
64            ));
65        }
66        let ack_ref = serde_json::to_string(&reference).map_err(|error| {
67            DbError::context("serialize exact Store acknowledgement ref", error)
68        })?;
69        let prepared = serde_json::to_string(&prepared)
70            .map_err(|error| DbError::context("serialize prepared Store acknowledgement", error))?;
71        let activation = serde_json::to_string(&OutboundStoreAckActivation::AwaitingCandidate)
72            .map_err(|error| {
73                DbError::context("serialize Store acknowledgement activation state", error)
74            })?;
75        tx.execute(
76            "INSERT INTO outbound_store_acks
77             (singleton, ack_ref, ack_bytes, prepared_object, activation)
78             VALUES (1, ?1, ?2, ?3, ?4)",
79            rusqlite::params![ack_ref, bytes, prepared, activation],
80        )
81        .map_err(DbError::from)?;
82        tx.commit().map_err(DbError::from)?;
83        Ok(reference)
84    }
85
86    fn adopt_outbound_store_ack_slot_winner(
87        &mut self,
88        expected: &StoreAckRef,
89        winner_bytes: Vec<u8>,
90        winner_prepared: PreparedExactObject,
91    ) -> Result<(), DbError> {
92        let authority = self.local_store_authority()?;
93        let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
94        let outbound = load_expected_outbound_store_ack_on(
95            &tx,
96            &authority,
97            expected,
98            "acknowledgement slot winner names another queued object",
99        )?;
100        let OutboundStoreAckActivation::Prepared(candidate) = &outbound.activation else {
101            return Err(DbError::Message(
102                "acknowledgement slot collision has no prepared activation candidate".to_string(),
103            ));
104        };
105        if candidate.commit.acknowledgement() != Some(expected) {
106            return Err(DbError::Message(
107                "prepared activation candidate names another acknowledgement".to_string(),
108            ));
109        }
110        if winner_prepared.reference().slot() != expected.object.slot()
111            || winner_prepared.reference() == &expected.object
112        {
113            return Err(DbError::Message(
114                "acknowledgement slot winner is not a distinct object at the occupied slot"
115                    .to_string(),
116            ));
117        }
118        let (winner_reference, _) =
119            verify_next_local_store_ack_on(&tx, &authority, &winner_bytes, &winner_prepared)?;
120        let expected_records = candidate
121            .acknowledgement_remote_objects(&outbound.ack)
122            .map_err(DbError::from)?;
123        for expected_record in &expected_records {
124            let object_id = expected_record.object_id();
125            let stored = load_remote_object_on(&tx, object_id)?;
126            if stored != **expected_record {
127                return Err(DbError::Message(
128                    "losing acknowledgement candidate is no longer wholly unuploaded".to_string(),
129                ));
130            }
131        }
132        for expected_record in expected_records {
133            if !crate::remote_object_records::delete_remote_object_on(
134                &tx,
135                expected_record.object_id(),
136            )? {
137                return Err(DbError::Message(
138                    "losing acknowledgement candidate object disappeared".to_string(),
139                ));
140            }
141        }
142        let activation = serde_json::to_string(&OutboundStoreAckActivation::AwaitingCandidate)
143            .map_err(|error| {
144                DbError::context("serialize adopted Store acknowledgement activation", error)
145            })?;
146        let winner_ref = serde_json::to_string(&winner_reference).map_err(|error| {
147            DbError::context("serialize adopted Store acknowledgement ref", error)
148        })?;
149        let winner_prepared = serde_json::to_string(&winner_prepared).map_err(|error| {
150            DbError::context("serialize adopted prepared Store acknowledgement", error)
151        })?;
152        let updated = tx
153            .execute(
154                "UPDATE outbound_store_acks
155                 SET ack_ref = ?2, ack_bytes = ?3, prepared_object = ?4, activation = ?5
156                 WHERE singleton = 1 AND ack_ref = ?1",
157                rusqlite::params![
158                    serde_json::to_string(expected).map_err(|error| {
159                        DbError::context("serialize losing Store acknowledgement ref", error)
160                    })?,
161                    winner_ref,
162                    winner_bytes,
163                    winner_prepared,
164                    activation,
165                ],
166            )
167            .map_err(DbError::from)?;
168        if updated != 1 {
169            return Err(DbError::Message(
170                "outbound Store acknowledgement changed during winner adoption".to_string(),
171            ));
172        }
173        tx.commit().map_err(DbError::from)
174    }
175
176    fn oldest_outbound_store_ack(&mut self) -> Result<Option<OutboundStoreAck>, DbError> {
177        let authority = self.local_store_authority()?;
178        load_outbound_store_ack_on(self.conn, &authority)
179    }
180
181    fn complete_outbound_store_ack(
182        &mut self,
183        accepted: &StoreAckRef,
184        activating_commit: &StoreBatchCommitRef,
185    ) -> Result<(), DbError> {
186        let authority = self.local_store_authority()?;
187        let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
188        let outbound = load_expected_outbound_store_ack_on(
189            &tx,
190            &authority,
191            accepted,
192            "accepted Store acknowledgement differs from the prepared exact object",
193        )?;
194        finish_outbound_store_ack_on(
195            &tx,
196            accepted,
197            &outbound.ack.value.successor.next_slot,
198            &coven_protocol::store_commit::StandingStoreAck {
199                assertion: outbound.ack.value.assertion(),
200                activating_commit: Some(activating_commit.clone()),
201            },
202        )?;
203        for circle in &outbound.circle_acknowledgements {
204            let circle_id = circle.reference.circle_id.to_string();
205            let removed = tx
206                .execute(
207                    "DELETE FROM outbound_circle_acks WHERE circle_id = ?1",
208                    [&circle_id],
209                )
210                .map_err(DbError::from)?;
211            if removed != 1 {
212                return Err(DbError::Message(
213                    "outbound Circle acknowledgement disappeared during completion".to_string(),
214                ));
215            }
216            let successor_slot = serde_json::to_string(&circle.ack.value.successor.next_slot)
217                .map_err(|error| {
218                    DbError::context("serialize Circle acknowledgement successor slot", error)
219                })?;
220            let store_cut = serde_json::to_string(&circle.ack.value.store_cut)
221                .map_err(|error| DbError::context("serialize Circle acknowledgement cut", error))?;
222            let control_coord =
223                serde_json::to_string(&circle.ack.value.control).map_err(|error| {
224                    DbError::context("serialize Circle acknowledgement control", error)
225                })?;
226            tx.execute(
227                "INSERT INTO published_circle_acks
228                   (circle_id, ack_ref, successor_slot, store_cut, control_coord)
229                 VALUES (?1, ?2, ?3, ?4, ?5)
230                 ON CONFLICT(circle_id) DO UPDATE SET
231                   ack_ref = excluded.ack_ref, successor_slot = excluded.successor_slot,
232                   store_cut = excluded.store_cut, control_coord = excluded.control_coord",
233                rusqlite::params![
234                    circle_id,
235                    serde_json::to_string(&circle.reference).map_err(|error| {
236                        DbError::context("serialize published Circle acknowledgement ref", error)
237                    })?,
238                    successor_slot,
239                    store_cut,
240                    control_coord,
241                ],
242            )
243            .map_err(DbError::from)?;
244        }
245        tx.commit().map_err(DbError::from)
246    }
247}
248
249impl StoreDatabase {
250    pub async fn latest_local_store_ack(&self) -> Result<Option<PublishedStoreAck>, DbError> {
251        self.call_store(|session| session.latest_local_store_ack())
252            .await
253    }
254
255    pub async fn activated_store_ack(
256        &self,
257        registration: &StoreDeviceRegistrationRef,
258    ) -> Result<Option<ActivatedStoreAck>, DbError> {
259        let registration = registration.clone();
260        self.call_store(move |session| session.activated_store_ack(&registration))
261            .await
262    }
263
264    pub async fn stage_store_ack(
265        &self,
266        ack: StoreAck,
267        prepared: PreparedExactObject,
268    ) -> Result<StoreAckRef, DbError> {
269        self.call_store(move |session| session.stage_store_ack(ack, prepared))
270            .await
271    }
272
273    pub async fn adopt_outbound_store_ack_slot_winner(
274        &self,
275        expected: StoreAckRef,
276        winner_bytes: Vec<u8>,
277        winner_prepared: PreparedExactObject,
278    ) -> Result<(), DbError> {
279        self.call_store(move |session| {
280            session.adopt_outbound_store_ack_slot_winner(&expected, winner_bytes, winner_prepared)
281        })
282        .await
283    }
284
285    pub async fn oldest_outbound_store_ack(&self) -> Result<Option<OutboundStoreAck>, DbError> {
286        self.call_store(|session| session.oldest_outbound_store_ack())
287            .await
288    }
289
290    pub async fn complete_outbound_store_ack(
291        &self,
292        accepted: StoreAckRef,
293        activating_commit: StoreBatchCommitRef,
294    ) -> Result<(), DbError> {
295        self.call_store(move |session| {
296            session.complete_outbound_store_ack(&accepted, &activating_commit)
297        })
298        .await
299    }
300}