1use super::{
2 candidate_records::begin_merge_candidate_nonactivation_on,
3 publication_state::{MergeAbandonmentOutcome, PreparedStoreWriteState},
4 MergeMaterializationTransaction, StoreDatabase, StoreSession, StoreTransactionOutcome,
5 VerifiedStoreTransaction,
6};
7use crate::{
8 candidate_graph_exact_objects, load_prepared_audience_objects_on, load_remote_object_on,
9 update_remote_object_on, CloudOutboxRecords, CompletePreparedStoreWriteOutcome, Database,
10 DbError, PreparedAudienceBlob, RetainedPackageApplication, LOCAL_DEVICE_ID_STATE_KEY,
11};
12use coven_protocol::remote_object::{remote_object_id, CandidateNonactivation};
13use coven_protocol::store_commit::{
14 StoreBatchCommit, StoreBatchCommitRef, StoreDeviceHead, StoreDeviceHeadRef,
15 VerifiedStoreBatchCommit,
16};
17use coven_protocol::write::{PublishedPosition, WriteId, WriteResolution, WriteStatus};
18
19impl VerifiedStoreTransaction<'_, '_, '_> {
20 fn complete_prepared_store_write(
21 &mut self,
22 root: coven_protocol::store_commit::StoreRootRef,
23 accepted: StoreBatchCommitRef,
24 nonactivations: std::collections::BTreeMap<StoreBatchCommitRef, CandidateNonactivation>,
25 routing_key: Option<coven_protocol::circle::RowRoutingKey>,
26 ) -> Result<
27 (
28 CompletePreparedStoreWriteOutcome,
29 Option<(WriteId, WriteStatus)>,
30 ),
31 DbError,
32 > {
33 let state = &mut *self.authority;
34 let gates = self.gates;
35 let synced_tables = self.synced_tables;
36 let store_transaction = self.store;
37 let tx = store_transaction.transaction;
38 let local_device_id = crate::required_protocol_state_on(tx, LOCAL_DEVICE_ID_STATE_KEY)?;
39 let prepared_count: i64 = tx
40 .query_row(
41 "SELECT COUNT(*) FROM store_writes WHERE prepared IS NOT NULL",
42 [],
43 |row| row.get(0),
44 )
45 .map_err(DbError::from)?;
46 if prepared_count != 1 {
47 return Err(DbError::Message(format!(
48 "Store publication expected one prepared write, found {prepared_count}"
49 )));
50 }
51 let (stored_write_id, raw_status, raw_prepared): (String, String, String) = tx
52 .query_row(
53 "SELECT write_id, status, prepared FROM store_writes
54 WHERE prepared IS NOT NULL ORDER BY ordinal LIMIT 1",
55 [],
56 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
57 )
58 .map_err(DbError::from)?;
59 let current_status: WriteStatus = serde_json::from_str(&raw_status)
60 .map_err(|error| DbError::context("prepared Store write status", error))?;
61 let prepared: PreparedStoreWriteState = serde_json::from_str(&raw_prepared)
62 .map_err(|error| DbError::context("prepared Store write", error))?;
63 let exclusion_candidate = state.prepared_merge_candidate_on(
64 crate::store::store_session::StoreRecords::new(
65 self.store.transaction,
66 self.store.store_dir,
67 ),
68 &prepared,
69 )?;
70 if store_transaction
71 .author_exclusion_activation_for_candidate(
72 state,
73 &root,
74 &exclusion_candidate.reference,
75 &exclusion_candidate.commit.author_registration,
76 )?
77 .is_some()
78 {
79 let device_id = exclusion_candidate.commit.author_registration.device_id;
80 let write_id = WriteId::from_generated(stored_write_id.clone());
81 if let WriteStatus::Resolved(WriteResolution::Retracted { witness }) = ¤t_status {
82 witness.validate().map_err(DbError::from)?;
83 if witness.original_position().commit() != &exclusion_candidate.reference {
84 return Err(DbError::Message(
85 "terminal write retraction names another prepared candidate".to_string(),
86 ));
87 }
88 tx.execute(
89 "DELETE FROM store_write_blob_leases WHERE write_id = ?1",
90 [write_id.as_str()],
91 )
92 .map_err(DbError::from)?;
93 tx.execute(
94 "DELETE FROM store_write_packages WHERE write_id = ?1",
95 [write_id.as_str()],
96 )
97 .map_err(DbError::from)?;
98 tx.execute(
99 "DELETE FROM store_write_blobs WHERE write_id = ?1",
100 [write_id.as_str()],
101 )
102 .map_err(DbError::from)?;
103 let updated = tx
104 .execute(
105 "UPDATE store_writes SET prepared = NULL
106 WHERE write_id = ?1 AND status = ?2 AND prepared = ?3",
107 rusqlite::params![write_id.as_str(), &raw_status, &raw_prepared],
108 )
109 .map_err(DbError::from)?;
110 if updated != 1 {
111 return Err(DbError::Message(
112 "terminally retracted Store write changed during completion".to_string(),
113 ));
114 }
115 return Ok((
116 CompletePreparedStoreWriteOutcome::AuthorExcluded { device_id },
117 None,
118 ));
119 }
120 let status =
121 WriteStatus::Blocked(coven_protocol::write::WriteBlock::InvalidProtocolState {
122 reason: format!(
123 "Store author {device_id} was excluded before candidate activation"
124 ),
125 });
126 Database::set_write_status_on(tx, &write_id, &status)?;
127 return Ok((
128 CompletePreparedStoreWriteOutcome::AuthorExcluded { device_id },
129 Some((write_id, status)),
130 ));
131 }
132 if let PreparedStoreWriteState::MergeAbandonment {
133 candidate_commit,
134 candidate_head,
135 authority_commit,
136 authority_head,
137 authority_history_evidence,
138 ..
139 } = &prepared
140 {
141 let root = state.root().clone();
142 let candidate = state.prepared_merge_candidate_parts_on(
143 crate::store::store_session::StoreRecords::new(
144 self.store.transaction,
145 self.store.store_dir,
146 ),
147 candidate_commit.semantic_bytes(),
148 candidate_commit.prepared().reference(),
149 candidate_head.semantic_bytes(),
150 candidate_head.prepared().reference(),
151 )?;
152 let authority = state.prepared_merge_candidate_parts_on(
153 crate::store::store_session::StoreRecords::new(
154 self.store.transaction,
155 self.store.store_dir,
156 ),
157 authority_commit.semantic_bytes(),
158 authority_commit.prepared().reference(),
159 authority_head.semantic_bytes(),
160 authority_head.prepared().reference(),
161 )?;
162 if authority.commit.write_id.as_str() != stored_write_id
163 || accepted != authority.reference
164 || !matches!(
165 &authority.commit.body,
166 coven_protocol::store_commit::StoreCommitBody::AbandonCandidates { .. }
167 )
168 {
169 return Err(DbError::Message(
170 "accepted Merge abandonment differs from its durable authority".to_string(),
171 ));
172 }
173 let registration = super::verified_store_authority::VerifiedRegistrationLookup::activated_registration_on(
174 state,
175 crate::store::store_session::StoreRecords::new(self.store.transaction, self.store.store_dir),
176 &root,
177 &authority.commit.author_registration,
178 )?;
179 StoreDeviceHead::parse_at(
180 &authority.head.to_bytes(),
181 root.store_root_hash,
182 ®istration,
183 &accepted,
184 )
185 .map_err(|error| DbError::context("verify accepted Merge abandonment head", error))?;
186 for object in [
187 authority_commit.prepared().reference(),
188 authority_head.prepared().reference(),
189 ] {
190 let object_id = remote_object_id(object);
191 let remote = load_remote_object_on(tx, object_id)?
192 .into_activated(&accepted)
193 .map_err(|error| {
194 DbError::context(
195 format!("activate Merge abandonment object {object_id}"),
196 error,
197 )
198 })?;
199 update_remote_object_on(tx, object_id, &remote)?;
200 }
201 let nonactivation = nonactivations.get(&candidate.reference).ok_or_else(|| {
202 DbError::Message(
203 "accepted Merge abandonment has no verified candidate nonactivation"
204 .to_string(),
205 )
206 })?;
207 begin_merge_candidate_nonactivation_on(
208 tx,
209 &WriteId::from_generated(stored_write_id.clone()),
210 &candidate,
211 nonactivation,
212 true,
213 &[],
214 )?;
215 let retained = MergeMaterializationTransaction::from_store(self.store)
216 .record_materialized_merge_commit(
217 state,
218 &root,
219 &authority.commit,
220 &[],
221 &authority.head,
222 &authority.head_object,
223 authority_history_evidence,
224 &[],
225 None,
226 )?;
227 state.insert_verified(retained)?;
228 let mut completed_preparation = prepared.clone();
229 let PreparedStoreWriteState::MergeAbandonment { outcome, .. } =
230 &mut completed_preparation
231 else {
232 unreachable!("matched Merge abandonment")
233 };
234 *outcome = MergeAbandonmentOutcome::Accepted {
235 authority: accepted.clone(),
236 };
237 let completed_preparation = serde_json::to_string(&completed_preparation)
238 .map_err(|error| DbError::context("serialize accepted Merge abandonment", error))?;
239 let updated = tx
240 .execute(
241 "UPDATE store_writes SET prepared = ?2
242 WHERE write_id = ?1 AND prepared = ?3",
243 rusqlite::params![
244 stored_write_id.as_str(),
245 completed_preparation,
246 raw_prepared
247 ],
248 )
249 .map_err(DbError::from)?;
250 if updated != 1 {
251 return Err(DbError::Message(
252 "Merge abandonment changed during activation".to_string(),
253 ));
254 }
255 let blocked =
256 WriteStatus::Blocked(coven_protocol::write::WriteBlock::InvalidProtocolState {
257 reason: format!(
258 "candidate abandonment {} is accepted; exact cleanup is pending",
259 authority.head.head_hash()
260 ),
261 });
262 let write_id = authority.commit.write_id.clone();
263 Database::set_write_status_on(tx, &write_id, &blocked)?;
264 return Ok((
265 CompletePreparedStoreWriteOutcome::Published,
266 Some((write_id, blocked)),
267 ));
268 }
269 let PreparedStoreWriteState::Publication {
270 commit,
271 head,
272 history_evidence,
273 local_cleanup,
274 ..
275 } = prepared
276 else {
277 return Err(DbError::Message(
278 "Merge abandonment reached ordinary publication completion".to_string(),
279 ));
280 };
281 let root = state.root().clone();
282 let unverified: StoreBatchCommit = serde_json::from_slice(commit.semantic_bytes())
283 .map_err(|error| DbError::context("prepared Store commit", error))?;
284 let registration =
285 super::verified_store_authority::VerifiedRegistrationLookup::activated_registration_on(
286 state,
287 crate::store::store_session::StoreRecords::new(
288 self.store.transaction,
289 self.store.store_dir,
290 ),
291 &root,
292 &unverified.author_registration,
293 )?;
294 let expected_stream =
295 coven_protocol::store_commit::StreamActivation::device_authorized_stream_id(
296 root.store_root_hash,
297 &unverified.author_registration,
298 coven_protocol::store_commit::StreamAnchorDomain::StoreAnnouncements,
299 );
300 if accepted.coord.stream_id != expected_stream
301 || accepted.object != *commit.prepared().reference()
302 {
303 return Err(DbError::Message(
304 "accepted Merge head differs from the exact prepared commit".to_string(),
305 ));
306 }
307 let commit_value = VerifiedStoreBatchCommit::parse(
308 commit.semantic_bytes(),
309 root.store_root_hash,
310 &accepted,
311 ®istration,
312 )
313 .map_err(|error| DbError::context("outbound commit", error))?;
314 let head_value = StoreDeviceHead::parse_at(
315 head.semantic_bytes(),
316 root.store_root_hash,
317 ®istration,
318 &accepted,
319 )
320 .map_err(|error| DbError::context("outbound Store head", error))?;
321 if commit_value.write_id.as_str() != stored_write_id {
322 return Err(DbError::Message(
323 "prepared write id differs from signed commit".to_string(),
324 ));
325 }
326 let write_id = commit_value.write_id.clone();
327 let head_object_id = remote_object_id(head.prepared().reference());
328 let commit = commit_value.value();
329 let commit_ref = commit_value.reference();
330 let remaining_spools: i64 = tx
331 .query_row(
332 "SELECT COUNT(*) FROM store_write_blobs
333 WHERE write_id = ?1 AND spool_path IS NOT NULL",
334 [write_id.as_str()],
335 |row| row.get(0),
336 )
337 .map_err(DbError::from)?;
338 if remaining_spools != 0 {
339 return Err(DbError::Message(format!(
340 "prepared write {write_id} retains {remaining_spools} uploaded blob spool(s)"
341 )));
342 }
343 let audiences = load_prepared_audience_objects_on(tx, self.store.store_dir, &write_id)?;
344 let retained_packages = audiences
345 .packages
346 .iter()
347 .map(|package| package.package().clone())
348 .collect::<Vec<_>>();
349 for package in &audiences.packages {
350 package
351 .package()
352 .validate_blob_uploader(&commit.author_registration)
353 .map_err(DbError::from)?;
354 }
355 let mut object_ids = std::collections::BTreeSet::new();
356 object_ids.insert(remote_object_id(&commit_ref.object));
357 object_ids.extend(
358 candidate_graph_exact_objects(commit)?
359 .iter()
360 .map(remote_object_id),
361 );
362 object_ids.extend(
363 audiences
364 .blobs
365 .iter()
366 .map(PreparedAudienceBlob::remote_object_id),
367 );
368 object_ids.insert(head_object_id);
369 for object_id in object_ids {
370 let remote = load_remote_object_on(tx, object_id)?
371 .into_activated(commit_ref)
372 .map_err(|error| {
373 DbError::context(format!("activate remote object {object_id}"), error)
374 })?;
375 let state = serde_json::to_string(&remote)
376 .map_err(|error| DbError::context("serialize activated remote object", error))?;
377 let updated = tx
378 .execute(
379 "UPDATE remote_objects SET state = ?2 WHERE object_id = ?1",
380 (object_id.to_string(), state),
381 )
382 .map_err(DbError::from)?;
383 if updated != 1 {
384 return Err(DbError::Message(format!(
385 "remote object {object_id} disappeared during activation"
386 )));
387 }
388 }
389 let merge_transaction = MergeMaterializationTransaction::from_store(self.store);
390 let retained = merge_transaction.record_materialized_merge_commit(
391 state,
392 &root,
393 &commit_value,
394 &[],
395 &head_value,
396 head.prepared().reference(),
397 &history_evidence,
398 &retained_packages,
399 (!retained_packages.is_empty()).then_some(RetainedPackageApplication::LocallyAuthored),
400 )?;
401 state.insert_verified(retained)?;
402 let replayed = state.replay_projection_watching_on(
403 store_transaction,
404 self.blob_decls,
405 gates,
406 synced_tables,
407 routing_key.as_ref(),
408 &std::collections::BTreeSet::new(),
409 crate::ReplayJournal::Owed,
410 coven_protocol::membership::LocalStoreMembership::Current,
411 commit_ref,
412 )?;
413 match replayed.watched_outcome() {
414 Some(super::WatchedReplayOutcome::Applied { .. }) => {}
415 Some(super::WatchedReplayOutcome::Held(reason)) => {
416 return Err(DbError::Message(format!(
417 "accepted local Store publication held during replay: {reason:?}"
418 )))
419 }
420 None => {
421 return Err(DbError::Message(
422 "accepted local Store publication was absent from replay".to_string(),
423 ))
424 }
425 }
426 replayed.install_on(self, &root)?;
427 let cloud_outbox = CloudOutboxRecords::new(tx);
428 let mut consumed_uploads = 0;
429 for package in &audiences.packages {
430 for binding in package.package().blob_bindings() {
431 if cloud_outbox.consume_created_upload_handoff(package.package(), binding)? {
432 consumed_uploads += 1;
433 }
434 }
435 }
436 match Database::make_remote_publication_root_on(tx, &write_id)? {
437 Some((root_table, root_id)) => {
438 if consumed_uploads == 0 {
439 return Err(DbError::Message(format!(
440 "make_remote publication {write_id} for {root_table:?}/{root_id:?} contains no Created upload handoff"
441 )));
442 }
443 let remaining: i64 = tx
444 .query_row(
445 "SELECT COUNT(*) FROM cloud_outbox
446 WHERE operation = 'upload' AND root_table = ?1 AND root_id = ?2",
447 (&root_table, &root_id),
448 |row| row.get(0),
449 )
450 .map_err(DbError::from)?;
451 if remaining != 0 {
452 return Err(DbError::Message(format!(
453 "make_remote publication {write_id} left {remaining} upload handoff(s) for {root_table:?}/{root_id:?}"
454 )));
455 }
456 Database::complete_make_remote_publication_on(tx, &write_id)?;
457 }
458 None if consumed_uploads != 0 => {
459 return Err(DbError::Message(format!(
460 "Store write {write_id} consumed Created upload handoffs without a make_remote publication intent"
461 )));
462 }
463 None => {}
464 }
465 for drop in local_cleanup.drops {
466 tx.execute(
467 "INSERT INTO published_blob_drop_intents
468 (seq, namespace, blob_id, size, plaintext_hash, locator_hash, disposition)
469 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
470 ON CONFLICT(seq, namespace, blob_id, locator_hash) DO NOTHING",
471 rusqlite::params![
472 Database::sequence_to_sqlite(
473 &commit_ref.coord.stream_id.to_string(),
474 commit_ref.coord.sequence(),
475 )?,
476 drop.namespace,
477 drop.id,
478 i64::try_from(drop.size).map_err(|_| DbError::Message(
479 "outbound local cleanup size exceeds SQLite integer".to_string()
480 ))?,
481 drop.plaintext_hash.to_string(),
482 drop.locator_hash.to_string(),
483 drop.disposition.as_db(),
484 ],
485 )
486 .map_err(DbError::from)?;
487 }
488 tx.execute(
489 "DELETE FROM store_write_packages WHERE write_id = ?1",
490 [write_id.as_str()],
491 )
492 .map_err(DbError::from)?;
493 tx.execute(
494 "DELETE FROM store_write_blobs WHERE write_id = ?1",
495 [write_id.as_str()],
496 )
497 .map_err(DbError::from)?;
498 retain_local_replay_blob_leases(tx, self.store.store_dir, &write_id)?;
499 let cleared = tx
500 .execute(
501 "UPDATE store_writes SET prepared = NULL
502 WHERE write_id = ?1 AND prepared IS NOT NULL",
503 [stored_write_id.as_str()],
504 )
505 .map_err(DbError::from)?;
506 if cleared != 1 {
507 return Err(DbError::Message(
508 "prepared Store write disappeared".to_string(),
509 ));
510 }
511 let status = WriteStatus::Published(Box::new(PublishedPosition {
512 device_id: local_device_id,
513 commit: accepted.clone(),
514 }));
515 Database::set_write_status_on(tx, &write_id, &status)?;
516 Ok((
517 CompletePreparedStoreWriteOutcome::Published,
518 Some((write_id, status)),
519 ))
520 }
521}
522
523fn retain_local_replay_blob_leases(
524 tx: &rusqlite::Transaction<'_>,
525 store_dir: &coven_foundation::store_dir::StoreDir,
526 write_id: &WriteId,
527) -> Result<(), DbError> {
528 let records = super::StoreRecords::new(tx, store_dir);
529 let partitions = records.store_write_partitions(write_id.as_str())?;
530 let local_rows = partitions
531 .local
532 .iter()
533 .map(|partition| crate::walk_changeset(&partition.changeset))
534 .collect::<Result<Vec<_>, _>>()?
535 .into_iter()
536 .flatten()
537 .filter(|change| {
538 !crate::is_routing_table(&change.table)
539 && !matches!(change.op, coven_foundation::changeset::ChangeOp::Delete)
540 })
541 .filter_map(|change| {
542 let row_id = change.pk()?.to_string();
543 Some((change.table, row_id))
544 })
545 .collect::<std::collections::BTreeSet<_>>();
546 let raw_facts: String = tx
547 .query_row(
548 "SELECT blob_facts FROM store_writes WHERE write_id = ?1",
549 [write_id.as_str()],
550 |row| row.get(0),
551 )
552 .map_err(DbError::from)?;
553 let facts: crate::StoreWriteBlobFacts = serde_json::from_str(&raw_facts)
554 .map_err(|error| DbError::context("published Store write blob facts", error))?;
555 let retained = facts
556 .blobs
557 .into_iter()
558 .filter(|fact| {
559 fact.blob.provenance == coven_protocol::blob::Provenance::HostProvided
560 && local_rows.contains(&(fact.table.clone(), fact.row_id.clone()))
561 })
562 .map(|fact| (fact.blob.namespace, fact.blob.id))
563 .collect::<std::collections::BTreeSet<_>>();
564 let leases = crate::query_mapped_rows(
565 tx,
566 "SELECT namespace, blob_id FROM store_write_blob_leases
567 WHERE write_id = ?1 ORDER BY namespace, blob_id",
568 [write_id.as_str()],
569 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
570 )?;
571 for (namespace, blob_id) in leases {
572 if retained.contains(&(namespace.clone(), blob_id.clone())) {
573 continue;
574 }
575 tx.execute(
576 "DELETE FROM store_write_blob_leases
577 WHERE write_id = ?1 AND namespace = ?2 AND blob_id = ?3",
578 (write_id.as_str(), namespace, blob_id),
579 )
580 .map_err(DbError::from)?;
581 }
582 Ok(())
583}
584
585impl StoreSession<'_> {
586 fn complete_prepared_store_write(
587 &mut self,
588 root: coven_protocol::store_commit::StoreRootRef,
589 accepted: StoreBatchCommitRef,
590 nonactivations: std::collections::BTreeMap<StoreBatchCommitRef, CandidateNonactivation>,
591 routing_key: Option<coven_protocol::circle::RowRoutingKey>,
592 ) -> Result<
593 (
594 CompletePreparedStoreWriteOutcome,
595 Option<(WriteId, WriteStatus)>,
596 ),
597 DbError,
598 > {
599 self.verified_store_transaction(move |transaction| {
600 let result = transaction.complete_prepared_store_write(
601 root,
602 accepted,
603 nonactivations,
604 routing_key,
605 )?;
606 Ok(StoreTransactionOutcome::Commit(result))
607 })
608 }
609
610 fn mark_merge_candidate_conflict(
611 &mut self,
612 write_id: WriteId,
613 winner_commit: StoreBatchCommitRef,
614 winner_head: StoreDeviceHeadRef,
615 nonactivations: std::collections::BTreeMap<StoreBatchCommitRef, CandidateNonactivation>,
616 ) -> Result<WriteStatus, DbError> {
617 let verified_authority = &mut *self.verified_store_authority;
618 let conn = self.conn;
619 let tx = conn.unchecked_transaction().map_err(DbError::from)?;
620 let (raw_status, raw_prepared): (String, String) = tx
621 .query_row(
622 "SELECT status, prepared FROM store_writes WHERE write_id = ?1",
623 [write_id.as_str()],
624 |row| Ok((row.get(0)?, row.get(1)?)),
625 )
626 .map_err(DbError::from)?;
627 let status: WriteStatus = serde_json::from_str(&raw_status)
628 .map_err(|error| DbError::context("Merge candidate status", error))?;
629 if !matches!(status, WriteStatus::Publishing) {
630 return Err(DbError::Message(format!(
631 "Merge candidate {write_id} is not publishing"
632 )));
633 }
634 let prepared: PreparedStoreWriteState = serde_json::from_str(&raw_prepared)
635 .map_err(|error| DbError::context("prepared Merge candidate", error))?;
636 let store_transaction =
637 crate::store::store_session::StoreTransaction::new(&tx, self.store_dir);
638 let prepared_candidate =
639 store_transaction.prepared_merge_candidate(verified_authority, &prepared)?;
640 let publication =
641 store_transaction.prepared_merge_publication(verified_authority, &prepared)?;
642 if winner_head.object.slot() != publication.head_object.slot()
643 || winner_head.object == publication.head_object
644 {
645 return Err(DbError::Message(
646 "Merge winner does not replace the prepared exact head slot".to_string(),
647 ));
648 }
649 if prepared_candidate.commit.write_id != write_id {
650 return Err(DbError::Message(
651 "prepared Merge graph differs from its write identity".to_string(),
652 ));
653 }
654 if matches!(&prepared, PreparedStoreWriteState::MergeAbandonment { .. }) {
655 let publication_nonactivation =
656 nonactivations.get(&publication.reference).ok_or_else(|| {
657 DbError::Message(
658 "Merge abandonment authority has no verified nonactivation".to_string(),
659 )
660 })?;
661 begin_merge_candidate_nonactivation_on(
662 &tx,
663 &write_id,
664 &publication,
665 publication_nonactivation,
666 false,
667 &[],
668 )?;
669 if winner_commit != prepared_candidate.reference {
670 let candidate_nonactivation = nonactivations
671 .get(&prepared_candidate.reference)
672 .ok_or_else(|| {
673 DbError::Message(
674 "Merge abandonment candidate has no verified nonactivation".to_string(),
675 )
676 })?;
677 begin_merge_candidate_nonactivation_on(
678 &tx,
679 &write_id,
680 &prepared_candidate,
681 candidate_nonactivation,
682 true,
683 &[],
684 )?;
685 }
686 let mut lost_preparation = prepared.clone();
687 let PreparedStoreWriteState::MergeAbandonment { outcome, .. } = &mut lost_preparation
688 else {
689 unreachable!("matched Merge abandonment")
690 };
691 *outcome = MergeAbandonmentOutcome::Lost {
692 winner_commit: winner_commit.clone(),
693 winner_head: winner_head.clone(),
694 };
695 let lost_preparation = serde_json::to_string(&lost_preparation)
696 .map_err(|error| DbError::context("serialize lost Merge abandonment", error))?;
697 let updated = tx
698 .execute(
699 "UPDATE store_writes SET prepared = ?2
700 WHERE write_id = ?1 AND prepared = ?3",
701 rusqlite::params![write_id.as_str(), lost_preparation, raw_prepared],
702 )
703 .map_err(DbError::from)?;
704 if updated != 1 {
705 return Err(DbError::Message(
706 "Merge abandonment changed while recording its winner".to_string(),
707 ));
708 }
709 } else {
710 let candidate_nonactivation = nonactivations
711 .get(&prepared_candidate.reference)
712 .ok_or_else(|| {
713 DbError::Message("Merge candidate has no verified nonactivation".to_string())
714 })?;
715 begin_merge_candidate_nonactivation_on(
716 &tx,
717 &write_id,
718 &prepared_candidate,
719 candidate_nonactivation,
720 true,
721 &[],
722 )?;
723 }
724 let blocked =
725 WriteStatus::Blocked(coven_protocol::write::WriteBlock::InvalidProtocolState {
726 reason: format!(
727 "Merge successor slot is occupied by signed head {}",
728 winner_head.head_hash
729 ),
730 });
731 Database::set_write_status_on(&tx, &write_id, &blocked)?;
732 tx.commit().map_err(DbError::from)?;
733 Ok(blocked)
734 }
735}
736
737impl StoreDatabase {
738 pub async fn complete_prepared_store_write(
739 &self,
740 root: coven_protocol::store_commit::StoreRootRef,
741 accepted: StoreBatchCommitRef,
742 nonactivations: Vec<coven_protocol::remote_object::VerifiedCandidateNonactivation>,
743 routing_key: Option<coven_protocol::circle::RowRoutingKey>,
744 ) -> Result<CompletePreparedStoreWriteOutcome, DbError> {
745 let nonactivations = nonactivations
746 .into_iter()
747 .map(|verified| {
748 verified
749 .candidate_reference()
750 .map(|reference| (reference, verified.into_durable()))
751 .map_err(DbError::from)
752 })
753 .collect::<Result<std::collections::BTreeMap<_, _>, _>>()?;
754 let (outcome, notification) = self
755 .call_store(move |session| {
756 session.complete_prepared_store_write(root, accepted, nonactivations, routing_key)
757 })
758 .await?;
759 if let Some((write_id, status)) = notification {
760 self.notify_write_status(write_id, status);
761 }
762 Ok(outcome)
763 }
764
765 pub async fn mark_merge_candidate_conflict(
766 &self,
767 write_id: WriteId,
768 nonactivations: Vec<coven_protocol::remote_object::VerifiedCandidateNonactivation>,
769 ) -> Result<(), DbError> {
770 let first = nonactivations.first().ok_or_else(|| {
771 DbError::Message("Merge candidate conflict has no verified candidates".to_string())
772 })?;
773 let winner_commit = first
774 .merge_winner_commit()
775 .cloned()
776 .map_err(DbError::from)?;
777 let winner_head = match first.proof() {
778 coven_protocol::remote_object::CandidateNonactivationProof::MergeWinner {
779 winner_head,
780 } => winner_head.clone(),
781 coven_protocol::remote_object::CandidateNonactivationProof::AuthorExclusion { .. } => {
782 return Err(DbError::Message(
783 "Merge slot conflict cannot carry author-exclusion evidence".to_string(),
784 ));
785 }
786 coven_protocol::remote_object::CandidateNonactivationProof::MergeMembershipGrantRevocation { .. } => {
787 return Err(DbError::Message(
788 "Merge slot conflict cannot carry membership-grant revocation evidence"
789 .to_string(),
790 ));
791 }
792 coven_protocol::remote_object::CandidateNonactivationProof::MergeDependencyRetraction { .. } => {
793 return Err(DbError::Message(
794 "Merge slot conflict cannot carry dependent-retraction evidence".to_string(),
795 ));
796 }
797 };
798 let winner_proof = first.proof().clone();
799 let nonactivations = nonactivations
800 .into_iter()
801 .map(|verified| {
802 if verified.merge_winner_commit().map_err(DbError::from)? != &winner_commit
803 || verified.proof() != &winner_proof
804 {
805 return Err(DbError::Message(
806 "Merge candidate conflict observations name different winners".to_string(),
807 ));
808 }
809 verified
810 .candidate_reference()
811 .map(|reference| (reference, verified.into_durable()))
812 .map_err(DbError::from)
813 })
814 .collect::<Result<std::collections::BTreeMap<_, _>, _>>()?;
815 let notified_write_id = write_id.clone();
816 let blocked = self
817 .call_store(move |session| {
818 session.mark_merge_candidate_conflict(
819 write_id,
820 winner_commit,
821 winner_head,
822 nonactivations,
823 )
824 })
825 .await?;
826 self.notify_write_status(notified_write_id, blocked);
827 Ok(())
828 }
829}