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}