Skip to main content

coven_database/store/store_session/candidate_lifecycle/
cleanup.rs

1use super::*;
2use crate::query_mapped_rows;
3use crate::store::StoreSession;
4
5impl StoreSession<'_> {
6    fn merge_candidate_cleanup_pending(&mut self, write_id: &WriteId) -> Result<bool, DbError> {
7        let records = crate::store::store_session::StoreRecords::new(self.conn, self.store_dir);
8        let verified_authority = &mut *self.verified_store_authority;
9        let conn = self.conn;
10        let (raw_status, raw_prepared): (String, Option<String>) = conn
11            .query_row(
12                "SELECT status, prepared FROM store_writes WHERE write_id = ?1",
13                [write_id.as_str()],
14                |row| Ok((row.get(0)?, row.get(1)?)),
15            )
16            .map_err(DbError::from)?;
17        let status: WriteStatus = serde_json::from_str(&raw_status)
18            .map_err(|error| DbError::context("Merge cleanup status", error))?;
19        if let WriteStatus::Resolved(WriteResolution::Retracted { witness }) = status {
20            witness.validate().map_err(DbError::from)?;
21            let candidate = witness.original_position().commit();
22            let StoreCommitCoord {
23                stream_id,
24                sequence,
25            } = &candidate.coord;
26            let exists: bool = conn
27                .query_row(
28                    "SELECT EXISTS(
29                         SELECT 1 FROM merge_retraction_cleanups
30                         WHERE device_id = ?1 AND seq = ?2 AND commit_ref = ?3
31                     )",
32                    rusqlite::params![
33                        stream_id.to_string(),
34                        Database::sequence_to_sqlite(&stream_id.to_string(), *sequence)?,
35                        serde_json::to_string(candidate).map_err(|error| {
36                            DbError::context("serialize Merge retraction cleanup ref", error)
37                        })?,
38                    ],
39                    |row| row.get(0),
40                )
41                .map_err(DbError::from)?;
42            if exists {
43                crate::StoreDatabase::load_merge_retraction_cleanup_on(
44                    records,
45                    verified_authority,
46                    candidate,
47                )?;
48            }
49            return Ok(exists);
50        }
51        let Some(raw_prepared) = raw_prepared else {
52            return Ok(false);
53        };
54        let prepared: PreparedStoreWriteState = serde_json::from_str(&raw_prepared)
55            .map_err(|error| DbError::context("prepared Merge cleanup", error))?;
56        let candidate = parse_prepared_merge_candidate_on(records, verified_authority, &prepared)?;
57        let cleanup_pending = |candidate: &PreparedMergeCandidate| -> Result<bool, DbError> {
58            let remote =
59                load_remote_object_on(conn, remote_object_id(&candidate.reference.object))?;
60            Ok(matches!(
61                remote,
62                RemoteObjectRecord::CandidateCommit(
63                    coven_protocol::remote_object::CandidateCommitRecord {
64                        state:
65                            coven_protocol::remote_object::CandidateCommitState::CleanupPending {
66                                proof: coven_protocol::remote_object::CandidateNonactivationProof::MergeWinner { .. }
67                                    | coven_protocol::remote_object::CandidateNonactivationProof::AuthorExclusion { .. }
68                            },
69                        ..
70                    }
71                )
72            ))
73        };
74        match &prepared {
75            PreparedStoreWriteState::Publication { .. } => cleanup_pending(&candidate),
76            PreparedStoreWriteState::MergeAbandonment {
77                outcome,
78                authority_commit,
79                authority_head,
80                ..
81            } => {
82                let authority = parse_prepared_merge_candidate_parts_on(
83                    records,
84                    verified_authority,
85                    authority_commit.semantic_bytes(),
86                    authority_commit.prepared().reference(),
87                    authority_head.semantic_bytes(),
88                    authority_head.prepared().reference(),
89                )?;
90                match outcome {
91                    MergeAbandonmentOutcome::Prepared => Ok(false),
92                    MergeAbandonmentOutcome::Accepted { .. } => cleanup_pending(&candidate),
93                    MergeAbandonmentOutcome::AuthorExcluded => {
94                        Ok(cleanup_pending(&candidate)? || cleanup_pending(&authority)?)
95                    }
96                    MergeAbandonmentOutcome::Lost { winner_commit, .. } => Ok((winner_commit
97                        != &candidate.reference
98                        && cleanup_pending(&candidate)?)
99                        || cleanup_pending(&authority)?),
100                }
101            }
102        }
103    }
104
105    fn merge_candidate_cleanup_targets(
106        &mut self,
107        write_id: &WriteId,
108    ) -> Result<Vec<CandidateCleanupObject>, DbError> {
109        let records = crate::store::store_session::StoreRecords::new(self.conn, self.store_dir);
110        let verified_authority = &mut *self.verified_store_authority;
111        let conn = self.conn;
112        let (raw_status, raw_prepared): (String, Option<String>) = conn
113            .query_row(
114                "SELECT status, prepared FROM store_writes WHERE write_id = ?1",
115                [write_id.as_str()],
116                |row| Ok((row.get(0)?, row.get(1)?)),
117            )
118            .map_err(DbError::from)?;
119        let status: WriteStatus = serde_json::from_str(&raw_status)
120            .map_err(|error| DbError::context("Merge cleanup status", error))?;
121        if let WriteStatus::Resolved(WriteResolution::Retracted { witness }) = &status {
122            witness.validate().map_err(DbError::from)?;
123            let candidate = crate::StoreDatabase::load_merge_retraction_cleanup_on(
124                records,
125                verified_authority,
126                witness.original_position().commit(),
127            )?;
128            if candidate.commit.write_id != *write_id {
129                return Err(DbError::Message(
130                    "Merge retraction cleanup names another write".to_string(),
131                ));
132            }
133            return merge_candidate_cleanup_targets_on(conn, write_id, &candidate, false, &[]);
134        }
135        if !matches!(status, WriteStatus::Blocked(_)) {
136            return Err(DbError::Message(format!(
137                "Merge cleanup write {write_id} is not blocked"
138            )));
139        }
140        let raw_prepared = raw_prepared.ok_or_else(|| {
141            DbError::Message("blocked Merge cleanup has no prepared candidate".to_string())
142        })?;
143        let prepared: PreparedStoreWriteState = serde_json::from_str(&raw_prepared)
144            .map_err(|error| DbError::context("prepared Merge cleanup", error))?;
145        let candidate = parse_prepared_merge_candidate_on(records, verified_authority, &prepared)?;
146        match &prepared {
147            PreparedStoreWriteState::Publication { .. } => {
148                merge_candidate_cleanup_targets_on(conn, write_id, &candidate, true, &[])
149            }
150            PreparedStoreWriteState::MergeAbandonment {
151                outcome,
152                authority_commit,
153                authority_head,
154                ..
155            } => {
156                let authority = parse_prepared_merge_candidate_parts_on(
157                    records,
158                    verified_authority,
159                    authority_commit.semantic_bytes(),
160                    authority_commit.prepared().reference(),
161                    authority_head.semantic_bytes(),
162                    authority_head.prepared().reference(),
163                )?;
164                let mut targets = Vec::new();
165                match outcome {
166                    MergeAbandonmentOutcome::Prepared => {
167                        return Err(DbError::Message(
168                            "Merge abandonment has no accepted winner".to_string(),
169                        ));
170                    }
171                    MergeAbandonmentOutcome::Accepted { .. } => {
172                        targets.extend(merge_candidate_cleanup_targets_on(
173                            conn,
174                            write_id,
175                            &candidate,
176                            true,
177                            &[],
178                        )?);
179                    }
180                    MergeAbandonmentOutcome::AuthorExcluded => {
181                        targets.extend(merge_candidate_cleanup_targets_on(
182                            conn,
183                            write_id,
184                            &candidate,
185                            true,
186                            &[],
187                        )?);
188                        targets.extend(merge_candidate_cleanup_targets_on(
189                            conn,
190                            write_id,
191                            &authority,
192                            false,
193                            &[],
194                        )?);
195                    }
196                    MergeAbandonmentOutcome::Lost { winner_commit, .. } => {
197                        if winner_commit != &candidate.reference {
198                            targets.extend(merge_candidate_cleanup_targets_on(
199                                conn,
200                                write_id,
201                                &candidate,
202                                true,
203                                &[],
204                            )?);
205                        }
206                        targets.extend(merge_candidate_cleanup_targets_on(
207                            conn,
208                            write_id,
209                            &authority,
210                            false,
211                            &[],
212                        )?);
213                    }
214                }
215                Ok(targets)
216            }
217        }
218    }
219
220    fn finish_retracted_merge_candidate_cleanup(
221        &mut self,
222        write_id: &WriteId,
223    ) -> Result<(), DbError> {
224        let verified_authority = &mut *self.verified_store_authority;
225        let conn = self.conn;
226        let tx = conn.unchecked_transaction().map_err(DbError::from)?;
227        let raw_status: String = tx
228            .query_row(
229                "SELECT status FROM store_writes WHERE write_id = ?1",
230                [write_id.as_str()],
231                |row| row.get(0),
232            )
233            .map_err(DbError::from)?;
234        let status: WriteStatus = serde_json::from_str(&raw_status)
235            .map_err(|error| DbError::context("Merge retraction cleanup status", error))?;
236        let WriteStatus::Resolved(WriteResolution::Retracted { witness }) = status else {
237            return Ok(());
238        };
239        witness.validate().map_err(DbError::from)?;
240        let candidate_ref = witness.original_position().commit().clone();
241        let candidate = crate::store::store_session::StoreTransaction::new(&tx, self.store_dir)
242            .load_merge_retraction_cleanup(verified_authority, &candidate_ref)?;
243        if candidate.commit.write_id != *write_id {
244            return Err(DbError::Message(
245                "Merge retraction cleanup names another write".to_string(),
246            ));
247        }
248        finish_merge_retraction_cleanup_on(&tx, &candidate)?;
249        tx.commit().map_err(DbError::from)
250    }
251
252    fn pending_merge_retraction_cleanups(&mut self) -> Result<Vec<StoreBatchCommitRef>, DbError> {
253        let records = crate::store::store_session::StoreRecords::new(self.conn, self.store_dir);
254        let verified_authority = &mut *self.verified_store_authority;
255        let conn = self.conn;
256        let rows = query_mapped_rows(
257            conn,
258            "SELECT device_id, seq, commit_ref
259             FROM merge_retraction_cleanups
260             ORDER BY device_id, seq",
261            [],
262            |row| {
263                Ok((
264                    row.get::<_, String>(0)?,
265                    row.get::<_, i64>(1)?,
266                    row.get::<_, String>(2)?,
267                ))
268            },
269        )?;
270        rows.into_iter()
271            .map(|(stream_id, sequence, encoded_ref)| {
272                let sequence = Database::sequence_from_sqlite(&stream_id, sequence)?;
273                let candidate = crate::store::materialized_commit_index::parse_stored_commit_ref(
274                    &stream_id,
275                    sequence,
276                    &encoded_ref,
277                )?;
278                crate::StoreDatabase::load_merge_retraction_cleanup_on(
279                    records,
280                    verified_authority,
281                    &candidate,
282                )?;
283                Ok(candidate)
284            })
285            .collect()
286    }
287
288    fn merge_retraction_cleanup_verification(
289        &mut self,
290        root: &coven_protocol::store_commit::StoreRootRef,
291        candidate: &StoreBatchCommitRef,
292    ) -> Result<TerminalCandidateCleanupVerification, DbError> {
293        let records = crate::store::store_session::StoreRecords::new(self.conn, self.store_dir);
294        let verified_authority = &mut *self.verified_store_authority;
295        let conn = self.conn;
296        let prepared = crate::StoreDatabase::load_merge_retraction_cleanup_on(
297            records,
298            verified_authority,
299            candidate,
300        )?;
301        let remote = load_remote_object_on(conn, remote_object_id(&candidate.object))?;
302        let proof = remote
303            .candidate_nonactivation_proof(candidate)
304            .map_err(DbError::from)?
305            .ok_or_else(|| {
306                DbError::Message("Merge retraction cleanup has no terminal proof".to_string())
307            })?;
308        let authority = match proof {
309            coven_protocol::remote_object::CandidateNonactivationProof::AuthorExclusion {
310                exclusion,
311                ..
312            } => TerminalCandidateAuthority::AuthorExclusion(
313                load_author_exclusion_activation_locator_on(
314                    records,
315                    verified_authority,
316                    root,
317                    exclusion,
318                )?,
319            ),
320            coven_protocol::remote_object::CandidateNonactivationProof::MergeMembershipGrantRevocation {
321                grant_id,
322                membership,
323                activation_commit,
324                activation_head,
325            } => TerminalCandidateAuthority::MembershipGrantRevocation {
326                grant_id: grant_id.clone(),
327                membership: membership.clone(),
328                activation_commit: activation_commit.clone(),
329                activation_head: activation_head.clone(),
330            },
331            coven_protocol::remote_object::CandidateNonactivationProof::MergeDependencyRetraction { .. } => {
332                let durable = coven_protocol::remote_object::CandidateNonactivation::from_durable_parts(
333                    candidate,
334                    &prepared.commit,
335                    proof.clone(),
336                )
337                .map_err(DbError::from)?;
338                validate_terminal_nonactivation_authority_on(
339                    records,
340                    verified_authority,
341                    root,
342                    &durable,
343                )?;
344                TerminalCandidateAuthority::DependencyRetraction(
345                    coven_protocol::remote_object::VerifiedDependencyRetractionAuthority::after_live_authority_check(durable)
346                        .map_err(DbError::from)?,
347                )
348            }
349            coven_protocol::remote_object::CandidateNonactivationProof::MergeWinner { .. } => {
350                return Err(DbError::Message(
351                    "Merge retraction cleanup has nonterminal proof".to_string(),
352                ));
353            }
354        };
355        Ok(TerminalCandidateCleanupVerification {
356            authority,
357            candidate: blocked_merge_candidate_from_prepared(prepared),
358        })
359    }
360
361    fn merge_retraction_cleanup_targets(
362        &mut self,
363        candidate: &StoreBatchCommitRef,
364    ) -> Result<Vec<CandidateCleanupObject>, DbError> {
365        let records = crate::store::store_session::StoreRecords::new(self.conn, self.store_dir);
366        let verified_authority = &mut *self.verified_store_authority;
367        let conn = self.conn;
368        let prepared = crate::StoreDatabase::load_merge_retraction_cleanup_on(
369            records,
370            verified_authority,
371            candidate,
372        )?;
373        merge_candidate_cleanup_targets_on(conn, &prepared.commit.write_id, &prepared, false, &[])
374    }
375
376    fn confirm_merge_retraction_cleanup_nonactivation(
377        &mut self,
378        root: &coven_protocol::store_commit::StoreRootRef,
379        candidate: &StoreBatchCommitRef,
380        durable: &coven_protocol::remote_object::CandidateNonactivation,
381        head_nonactivation: &coven_protocol::remote_object::VerifiedCandidateHeadNonactivation,
382    ) -> Result<(), DbError> {
383        let records = crate::store::store_session::StoreRecords::new(self.conn, self.store_dir);
384        let verified_authority = &mut *self.verified_store_authority;
385        let conn = self.conn;
386        let prepared = crate::StoreDatabase::load_merge_retraction_cleanup_on(
387            records,
388            verified_authority,
389            candidate,
390        )?;
391        if durable.reference().map_err(DbError::from)? != *candidate
392            || head_nonactivation.head().object() != &prepared.head_object
393        {
394            return Err(DbError::Message(
395                "verified Merge retraction cleanup names another candidate".to_string(),
396            ));
397        }
398        if let coven_protocol::remote_object::CandidateNonactivationProof::AuthorExclusion {
399            exclusion,
400            accepted_cut,
401            activation_head,
402        } = durable.proof()
403        {
404            let locator = load_author_exclusion_activation_locator_on(
405                records,
406                verified_authority,
407                root,
408                exclusion,
409            )?;
410            if locator.accepted_cut() != accepted_cut
411                || locator.activation_head() != activation_head
412            {
413                return Err(DbError::Message(
414                    "verified Merge retraction differs from durable exclusion authority"
415                        .to_string(),
416                ));
417            }
418        }
419        let remote = load_remote_object_on(conn, remote_object_id(&candidate.object))?;
420        if remote
421            .candidate_nonactivation_proof(candidate)
422            .map_err(DbError::from)?
423            != Some(durable.proof())
424        {
425            return Err(DbError::Message(
426                "verified Merge retraction differs from candidate ownership".to_string(),
427            ));
428        }
429        if !matches!(
430            load_merge_candidate_head_cleanup_on(conn, &prepared.head_object, candidate)?,
431            MergeCandidateHeadCleanup::ProtocolInert
432        ) {
433            return Err(DbError::Message(
434                "retracted Merge activation head is not retained as inert authority".to_string(),
435            ));
436        }
437        Ok(())
438    }
439
440    fn finish_merge_retraction_cleanup(
441        &mut self,
442        candidate: &StoreBatchCommitRef,
443    ) -> Result<(), DbError> {
444        let verified_authority = &mut *self.verified_store_authority;
445        let conn = self.conn;
446        let tx = conn.unchecked_transaction().map_err(DbError::from)?;
447        let prepared = crate::store::store_session::StoreTransaction::new(&tx, self.store_dir)
448            .load_merge_retraction_cleanup(verified_authority, candidate)?;
449        finish_merge_retraction_cleanup_on(&tx, &prepared)?;
450        tx.commit().map_err(DbError::from)
451    }
452}
453
454impl StoreDatabase {
455    pub async fn merge_candidate_cleanup_pending(
456        &self,
457        write_id: &WriteId,
458    ) -> Result<bool, DbError> {
459        let write_id = write_id.clone();
460        self.call_store(move |session| session.merge_candidate_cleanup_pending(&write_id))
461            .await
462    }
463
464    pub async fn merge_candidate_cleanup_targets(
465        &self,
466        write_id: WriteId,
467    ) -> Result<Vec<CandidateCleanupObject>, DbError> {
468        self.call_store(move |session| session.merge_candidate_cleanup_targets(&write_id))
469            .await
470    }
471
472    pub async fn finish_retracted_merge_candidate_cleanup(
473        &self,
474        write_id: WriteId,
475    ) -> Result<(), DbError> {
476        self.call_store(move |session| session.finish_retracted_merge_candidate_cleanup(&write_id))
477            .await
478    }
479
480    pub async fn pending_merge_retraction_cleanups(
481        &self,
482    ) -> Result<Vec<StoreBatchCommitRef>, DbError> {
483        self.call_store(|session| session.pending_merge_retraction_cleanups())
484            .await
485    }
486
487    pub async fn merge_retraction_cleanup_verification(
488        &self,
489        root: coven_protocol::store_commit::StoreRootRef,
490        candidate: StoreBatchCommitRef,
491    ) -> Result<TerminalCandidateCleanupVerification, DbError> {
492        self.call_store(move |session| {
493            session.merge_retraction_cleanup_verification(&root, &candidate)
494        })
495        .await
496    }
497
498    pub async fn merge_retraction_cleanup_targets(
499        &self,
500        candidate: StoreBatchCommitRef,
501    ) -> Result<Vec<CandidateCleanupObject>, DbError> {
502        self.call_store(move |session| session.merge_retraction_cleanup_targets(&candidate))
503            .await
504    }
505
506    pub async fn confirm_merge_retraction_cleanup_nonactivation(
507        &self,
508        root: coven_protocol::store_commit::StoreRootRef,
509        candidate: StoreBatchCommitRef,
510        verified: coven_protocol::remote_object::VerifiedCandidateNonactivation,
511    ) -> Result<(), DbError> {
512        let (durable, head_nonactivation) = verified
513            .into_terminal_head_nonactivation()
514            .map_err(DbError::from)?;
515        self.call_store(move |session| {
516            session.confirm_merge_retraction_cleanup_nonactivation(
517                &root,
518                &candidate,
519                &durable,
520                &head_nonactivation,
521            )
522        })
523        .await
524    }
525
526    pub async fn finish_merge_retraction_cleanup(
527        &self,
528        candidate: StoreBatchCommitRef,
529    ) -> Result<(), DbError> {
530        self.call_store(move |session| session.finish_merge_retraction_cleanup(&candidate))
531            .await
532    }
533}