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    pub(super) fn activated_store_ack(
16        &mut self,
17        registration: &StoreDeviceRegistrationRef,
18    ) -> Result<Option<ActivatedStoreAck>, DbError> {
19        let activated = self
20            .conn
21            .query_row(
22                "SELECT ack_ref, activating_commit FROM activated_store_acks WHERE device_id = ?1",
23                [registration.device_id.to_string()],
24                |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
25            )
26            .optional()
27            .map_err(DbError::from)?
28            .map(|(raw, activating_commit)| {
29                let reference: StoreAckRef = serde_json::from_str(&raw).map_err(|error| {
30                    DbError::context("activated Store acknowledgement ref", error)
31                })?;
32                if &reference.registration != registration {
33                    return Err(DbError::Message(
34                        "activated Store acknowledgement names another registration".to_string(),
35                    ));
36                }
37                Ok(ActivatedStoreAck {
38                    reference,
39                    activating_commit: serde_json::from_str(&activating_commit).map_err(
40                        |error| {
41                            DbError::context(
42                                "activated Store acknowledgement activating commit",
43                                error,
44                            )
45                        },
46                    )?,
47                })
48            })
49            .transpose()?;
50        let baseline = self
51            .verified_store_authority
52            .retained_replay_baseline_on(StoreRecords::new(self.conn, self.store_dir))?;
53        let RetainedReplayAuthority::InstalledSnapshot(snapshot) = &baseline.authority else {
54            return Ok(activated);
55        };
56        let Some(chain) = snapshot
57            .metadata
58            .history_summary
59            .acknowledgements
60            .get(&registration.device_id)
61        else {
62            return Ok(activated);
63        };
64        let (reference, _) = chain.latest().ok_or_else(|| {
65            DbError::Message("installed snapshot acknowledgement chain is empty".into())
66        })?;
67        if &reference.registration != registration {
68            return Err(DbError::Message(
69                "snapshot acknowledgement names another registration".into(),
70            ));
71        }
72        if let Some(current) = &activated {
73            if current.reference.sequence > reference.sequence {
74                return Ok(activated);
75            }
76            if chain
77                .chain
78                .get(&current.reference.sequence)
79                .map(|(reference, _)| reference)
80                != Some(&current.reference)
81                || (current.reference.sequence == reference.sequence
82                    && current.activating_commit != chain.activating_commit)
83            {
84                return Err(DbError::Message(
85                    "activated acknowledgement conflicts with the installed snapshot".into(),
86                ));
87            }
88        }
89        Ok(Some(ActivatedStoreAck {
90            reference: reference.clone(),
91            activating_commit: chain.activating_commit.clone(),
92        }))
93    }
94
95    fn stage_store_ack(
96        &mut self,
97        ack: StoreAck,
98        prepared: PreparedExactObject,
99    ) -> Result<StoreAckRef, DbError> {
100        let authority = self.local_store_authority()?;
101        let bytes = ack.to_bytes();
102        let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
103        let reference = verify_next_local_store_ack_on(&tx, &authority, &bytes, &prepared)?;
104        let ack_ref = serde_json::to_string(&reference).map_err(|error| {
105            DbError::context("serialize exact Store acknowledgement ref", error)
106        })?;
107        let prepared = serde_json::to_string(&prepared)
108            .map_err(|error| DbError::context("serialize prepared Store acknowledgement", error))?;
109        let activation = serde_json::to_string(&OutboundStoreAckActivation::AwaitingCandidate)
110            .map_err(|error| {
111                DbError::context("serialize Store acknowledgement activation state", error)
112            })?;
113        tx.execute(
114            "INSERT INTO outbound_store_acks
115             (singleton, ack_ref, ack_bytes, prepared_object, activation)
116             VALUES (1, ?1, ?2, ?3, ?4)",
117            rusqlite::params![ack_ref, bytes, prepared, activation],
118        )
119        .map_err(DbError::from)?;
120        tx.commit().map_err(DbError::from)?;
121        Ok(reference)
122    }
123
124    fn adopt_outbound_store_ack_slot_winner(
125        &mut self,
126        expected: &StoreAckRef,
127        winner_bytes: Vec<u8>,
128        winner_prepared: PreparedExactObject,
129    ) -> Result<(), DbError> {
130        let authority = self.local_store_authority()?;
131        let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
132        let outbound = load_expected_outbound_store_ack_on(
133            &tx,
134            &authority,
135            expected,
136            "acknowledgement slot winner names another queued object",
137        )?;
138        let OutboundStoreAckActivation::Prepared(candidate) = &outbound.activation else {
139            return Err(DbError::Message(
140                "acknowledgement slot collision has no prepared activation candidate".to_string(),
141            ));
142        };
143        let active_publication = ActiveStorePublication::for_commit(
144            crate::ActiveStorePublicationOwner::StoreAcknowledgement,
145            candidate,
146        )?;
147        if candidate.commit.acknowledgement() != Some(expected) {
148            return Err(DbError::Message(
149                "prepared activation candidate names another acknowledgement".to_string(),
150            ));
151        }
152        if winner_prepared.reference().slot() != expected.object.slot()
153            || winner_prepared.reference() == &expected.object
154        {
155            return Err(DbError::Message(
156                "acknowledgement slot winner is not a distinct object at the occupied slot"
157                    .to_string(),
158            ));
159        }
160        let winner_reference =
161            verify_next_local_store_ack_on(&tx, &authority, &winner_bytes, &winner_prepared)?;
162        let mut expected_records = candidate
163            .acknowledgement_remote_objects(&outbound.ack)
164            .map_err(DbError::from)?;
165        for reference in candidate.commit.circle_acknowledgements() {
166            let circle = outbound
167                .circle_acknowledgements
168                .iter()
169                .find(|circle| &circle.reference == reference)
170                .ok_or_else(|| {
171                    DbError::Message(
172                        "losing acknowledgement candidate lost its queued Circle statement".into(),
173                    )
174                })?;
175            expected_records.extend(candidate.circle_acknowledgement_remote_objects(&circle.ack)?);
176        }
177        let expected_records = expected_records
178            .into_iter()
179            .map(|record| (record.object_id(), record))
180            .collect::<std::collections::BTreeMap<_, _>>();
181        for expected_record in expected_records.values() {
182            let object_id = expected_record.object_id();
183            let stored = load_remote_object_on(&tx, object_id)?;
184            if stored != **expected_record {
185                return Err(DbError::Message(
186                    "losing acknowledgement candidate is no longer wholly unuploaded".to_string(),
187                ));
188            }
189        }
190        for expected_record in expected_records.into_values() {
191            if !crate::remote_object_records::delete_remote_object_on(
192                &tx,
193                expected_record.object_id(),
194            )? {
195                return Err(DbError::Message(
196                    "losing acknowledgement candidate object disappeared".to_string(),
197                ));
198            }
199        }
200        let activation =
201            serde_json::to_string(&OutboundStoreAckActivation::Created).map_err(|error| {
202                DbError::context("serialize adopted Store acknowledgement activation", error)
203            })?;
204        let winner_ref = serde_json::to_string(&winner_reference).map_err(|error| {
205            DbError::context("serialize adopted Store acknowledgement ref", error)
206        })?;
207        let winner_prepared = serde_json::to_string(&winner_prepared).map_err(|error| {
208            DbError::context("serialize adopted prepared Store acknowledgement", error)
209        })?;
210        let updated = tx
211            .execute(
212                "UPDATE outbound_store_acks
213                 SET ack_ref = ?2, ack_bytes = ?3, prepared_object = ?4, activation = ?5
214                 WHERE singleton = 1 AND ack_ref = ?1",
215                rusqlite::params![
216                    serde_json::to_string(expected).map_err(|error| {
217                        DbError::context("serialize losing Store acknowledgement ref", error)
218                    })?,
219                    winner_ref,
220                    winner_bytes,
221                    winner_prepared,
222                    activation,
223                ],
224            )
225            .map_err(DbError::from)?;
226        if updated != 1 {
227            return Err(DbError::Message(
228                "outbound Store acknowledgement changed during winner adoption".to_string(),
229            ));
230        }
231        super::active_store_publication::clear_active_store_publication_on(
232            &tx,
233            &active_publication,
234        )?;
235        tx.commit().map_err(DbError::from)
236    }
237
238    fn oldest_outbound_store_ack(&mut self) -> Result<Option<OutboundStoreAck>, DbError> {
239        let authority = self.local_store_authority()?;
240        load_outbound_store_ack_on(self.conn, &authority)
241    }
242
243    fn complete_outbound_store_ack(
244        &mut self,
245        accepted: &StoreAckRef,
246        activating_commit: &StoreBatchCommitRef,
247    ) -> Result<(), DbError> {
248        let authority = self.local_store_authority()?;
249        let activated = self.activated_store_ack(&accepted.registration)?;
250        if !activated.is_some_and(|activated| {
251            &activated.reference == accepted && &activated.activating_commit == activating_commit
252        }) {
253            return Err(DbError::Message(
254                "Store acknowledgement completion requires its exact installed activation"
255                    .to_string(),
256            ));
257        }
258        let tx = self.conn.unchecked_transaction().map_err(DbError::from)?;
259        let outbound = load_expected_outbound_store_ack_on(
260            &tx,
261            &authority,
262            accepted,
263            "accepted Store acknowledgement differs from the prepared exact object",
264        )?;
265        let candidate = match &outbound.activation {
266            OutboundStoreAckActivation::Prepared(candidate) => Some(candidate),
267            OutboundStoreAckActivation::Created => None,
268            OutboundStoreAckActivation::AwaitingCandidate => {
269                return Err(DbError::Message(
270                    "accepted Store acknowledgement has no prepared activation".to_string(),
271                ));
272            }
273        };
274        if candidate.is_some_and(|candidate| &candidate.reference != activating_commit) {
275            return Err(DbError::Message(
276                "accepted Store acknowledgement differs from its prepared activating commit"
277                    .to_string(),
278            ));
279        }
280        finish_outbound_store_ack_on(
281            &tx,
282            accepted,
283            &outbound.ack.value.successor.next_slot,
284            &coven_protocol::store_commit::StandingStoreAck {
285                assertion: outbound.ack.value.assertion(),
286                activating_commit: Some(activating_commit.clone()),
287            },
288        )?;
289        for circle in outbound.circle_acknowledgements.iter().filter(|circle| {
290            candidate.is_some_and(|candidate| {
291                candidate
292                    .commit
293                    .circle_acknowledgements()
294                    .contains(&circle.reference)
295            })
296        }) {
297            let circle_id = circle.reference.circle_id.to_string();
298            let removed = tx
299                .execute(
300                    "DELETE FROM outbound_circle_acks WHERE circle_id = ?1",
301                    [&circle_id],
302                )
303                .map_err(DbError::from)?;
304            if removed != 1 {
305                return Err(DbError::Message(
306                    "outbound Circle acknowledgement disappeared during completion".to_string(),
307                ));
308            }
309            let successor_slot = serde_json::to_string(&circle.ack.value.successor.next_slot)
310                .map_err(|error| {
311                    DbError::context("serialize Circle acknowledgement successor slot", error)
312                })?;
313            let store_cut = serde_json::to_string(&circle.ack.value.store_cut)
314                .map_err(|error| DbError::context("serialize Circle acknowledgement cut", error))?;
315            let control_coord =
316                serde_json::to_string(&circle.ack.value.control).map_err(|error| {
317                    DbError::context("serialize Circle acknowledgement control", error)
318                })?;
319            tx.execute(
320                "INSERT INTO published_circle_acks
321                   (circle_id, ack_ref, successor_slot, store_cut, control_coord)
322                 VALUES (?1, ?2, ?3, ?4, ?5)
323                 ON CONFLICT(circle_id) DO UPDATE SET
324                   ack_ref = excluded.ack_ref, successor_slot = excluded.successor_slot,
325                   store_cut = excluded.store_cut, control_coord = excluded.control_coord",
326                rusqlite::params![
327                    circle_id,
328                    serde_json::to_string(&circle.reference).map_err(|error| {
329                        DbError::context("serialize published Circle acknowledgement ref", error)
330                    })?,
331                    successor_slot,
332                    store_cut,
333                    control_coord,
334                ],
335            )
336            .map_err(DbError::from)?;
337        }
338        // A peer may have won the original publication position. The exact
339        // commit stays fixed while its publication entry is replaced.
340        if candidate.is_some() {
341            super::active_store_publication::clear_active_store_commit_for_owner_on(
342                &tx,
343                &crate::ActiveStorePublicationOwner::StoreAcknowledgement,
344                activating_commit,
345            )?;
346        }
347        tx.commit().map_err(DbError::from)
348    }
349}
350
351impl StoreDatabase {
352    pub async fn latest_local_store_ack(&self) -> Result<Option<PublishedStoreAck>, DbError> {
353        self.call_store(|session| session.latest_local_store_ack())
354            .await
355    }
356
357    pub async fn activated_store_ack(
358        &self,
359        registration: &StoreDeviceRegistrationRef,
360    ) -> Result<Option<ActivatedStoreAck>, DbError> {
361        let registration = registration.clone();
362        self.call_store(move |session| session.activated_store_ack(&registration))
363            .await
364    }
365
366    pub async fn stage_store_ack(
367        &self,
368        ack: StoreAck,
369        prepared: PreparedExactObject,
370    ) -> Result<StoreAckRef, DbError> {
371        self.call_store(move |session| session.stage_store_ack(ack, prepared))
372            .await
373    }
374
375    pub async fn adopt_outbound_store_ack_slot_winner(
376        &self,
377        expected: StoreAckRef,
378        winner_bytes: Vec<u8>,
379        winner_prepared: PreparedExactObject,
380    ) -> Result<(), DbError> {
381        self.call_store(move |session| {
382            session.adopt_outbound_store_ack_slot_winner(&expected, winner_bytes, winner_prepared)
383        })
384        .await
385    }
386
387    pub async fn oldest_outbound_store_ack(&self) -> Result<Option<OutboundStoreAck>, DbError> {
388        self.call_store(|session| session.oldest_outbound_store_ack())
389            .await
390    }
391
392    pub async fn complete_outbound_store_ack(
393        &self,
394        accepted: StoreAckRef,
395        activating_commit: StoreBatchCommitRef,
396    ) -> Result<(), DbError> {
397        self.call_store(move |session| {
398            session.complete_outbound_store_ack(&accepted, &activating_commit)
399        })
400        .await
401    }
402}