1use coven_protocol::objects::StoreObjectError;
4
5#[cfg(test)]
6use super::RegistrationOutbox;
7#[cfg(test)]
8use coven_keys::keys::UserKeypair;
9#[cfg(test)]
10use coven_protocol::objects::ProtocolObjectDomain;
11#[cfg(test)]
12use coven_protocol::store_commit::{
13 owner_recovery_semantic_prefix, StoreCommitCoord, StoreDeviceRegistration,
14 StoreDeviceRegistrationOrigin, StoreDeviceRegistrationRef,
15};
16
17#[derive(Debug, thiserror::Error)]
18pub enum StoreRegistrationError {
19 #[error("Store device registration database state: {0}")]
20 Database(#[from] coven_database::DbError),
21 #[error("Store device registration protocol: {0}")]
22 Protocol(#[from] coven_protocol::store_commit::StoreProtocolError),
23 #[error("Store device registration JSON: {0}")]
24 Json(#[from] serde_json::Error),
25 #[error("Store device registration author stream: {0}")]
26 AuthorStreamId(#[from] coven_protocol::causal_grants::AuthorStreamIdParseError),
27 #[error("Store device registration history: {0}")]
28 History(#[source] Box<crate::sync::store::StorePullError>),
29 #[error("Store device registration snapshot stream: {0}")]
30 SnapshotStream(#[source] Box<crate::sync::store::SnapshotError>),
31 #[error("{0}")]
32 Object(#[from] StoreObjectError),
33 #[error("exact Store root authority is absent")]
34 ExactRootAuthorityMissing,
35 #[error("Store device registration bytes are invalid: {0}")]
36 Invalid(String),
37 #[error("this Store installation requires an activated Join or Recovery registration")]
38 ActivationRequired,
39 #[error("registration publish count has no representable successor")]
40 PublishCountExhausted,
41 #[error("Store device registration activation: {0}")]
42 Outbound(#[from] crate::sync::store::StoreError),
43}
44
45impl From<crate::sync::store::StorePullError> for StoreRegistrationError {
46 fn from(error: crate::sync::store::StorePullError) -> Self {
47 Self::History(Box::new(error))
48 }
49}
50
51#[cfg(test)]
52mod tests {
53 use super::*;
54 use crate::sync::test_helpers::TestStore;
55 use coven_database::Database;
56 use coven_protocol::store_commit::StoreBatchCommitRef;
57 use coven_storage::CloudSyncObjectStorage;
58
59 async fn initialized() -> (
60 std::sync::Arc<TestStore>,
61 Database,
62 coven_foundation::store_dir::StoreDir,
63 UserKeypair,
64 ) {
65 let signer = UserKeypair::generate();
66 let db_store_dir = crate::sync::test_helpers::test_store_dir();
67 let db = crate::sync::test_helpers::open_test_db(db_store_dir.clone());
68 let store = TestStore::create(
69 &db,
70 db_store_dir.clone(),
71 "registration-store-test",
72 signer.clone(),
73 crate::sync::test_helpers::test_cloud_home(),
74 )
75 .await
76 .expect("create exact registration test Store");
77 (store, db, db_store_dir, signer)
78 }
79
80 async fn recovered_author() -> (
81 std::sync::Arc<TestStore>,
82 Database,
83 coven_foundation::store_dir::StoreDir,
84 StoreDeviceRegistrationRef,
85 StoreBatchCommitRef,
86 ) {
87 let (store, db, db_store_dir, signer) = initialized().await;
88 let loaded = store
89 .bind_device(&db, db_store_dir.clone(), &signer)
90 .await
91 .expect("load recovery Store");
92 let authority = store.founder_recovery_authority().await;
93 let database = coven_database::StoreDatabase::new(&db);
94 let registration = loaded
95 .owner_recovery_for_test()
96 .await
97 .expect("authorize Owner recovery Store")
98 .recover_owner_device(&authority, None)
99 .await
100 .expect("recover Owner device");
101 let loaded = store
102 .bind_device(&db, db_store_dir.clone(), &signer)
103 .await
104 .expect("reload recovered Store");
105 for reference in database
106 .materialized_frontier()
107 .await
108 .expect("load materialized Store frontier")
109 .into_values()
110 {
111 let commit = loaded
112 .load_commit_for_test(&reference)
113 .await
114 .expect("load materialized recovery commit");
115 if commit.value().author_registration == registration {
116 return (store, db, db_store_dir, registration, reference);
117 }
118 }
119 panic!("recovery commit is materialized")
120 }
121
122 #[tokio::test]
123 async fn store_root_state_failures_keep_registration_error_variants() {
124 let db_store_dir = crate::sync::test_helpers::test_store_dir();
125 let db = crate::sync::test_helpers::open_test_db(db_store_dir.clone());
126 let database = coven_database::StoreDatabase::new(&db);
127 let initialized_db_store_dir = crate::sync::test_helpers::test_store_dir();
128 let initialized_db =
129 crate::sync::test_helpers::open_test_db(initialized_db_store_dir.clone());
130 let (_, cloud_storage) = TestStore::create_with_connection(
131 &initialized_db,
132 initialized_db_store_dir.clone(),
133 "registration-missing-root-storage",
134 UserKeypair::generate(),
135 crate::sync::test_helpers::test_cloud_home(),
136 )
137 .await
138 .expect("create registration failure test Store");
139
140 let result = super::RegistrationOutbox::new(database, &*cloud_storage)
141 .drain()
142 .await;
143 assert!(
144 matches!(
145 result,
146 Err(StoreRegistrationError::ExactRootAuthorityMissing)
147 ),
148 "unexpected registration outbox result: {result:?}",
149 );
150 }
151
152 #[tokio::test]
153 async fn exact_founder_registration_is_already_activated() {
154 let (store, db, db_store_dir, signer) = initialized().await;
155 let database = coven_database::StoreDatabase::new(&db);
156 let loaded = store
157 .bind_device(&db, db_store_dir.clone(), &signer)
158 .await
159 .expect("load founder Store");
160 loaded
161 .authorize_writer()
162 .await
163 .expect("founder registration remains active");
164 let activated = database
165 .activated_store_device_registrations()
166 .await
167 .unwrap();
168 assert_eq!(activated.len(), 1);
169 assert_eq!(activated[0].store_root, store.root());
170 }
171
172 #[tokio::test]
173 async fn owner_recovery_publishes_and_activates_replacement_device() {
174 let (store, db, db_store_dir, signer) = initialized().await;
175 let loaded = store
176 .bind_device(&db, db_store_dir.clone(), &signer)
177 .await
178 .expect("load recovery Store");
179 let authority = store.founder_recovery_authority().await;
180 let database = coven_database::StoreDatabase::new(&db);
181 let registration = loaded
182 .owner_recovery_for_test()
183 .await
184 .expect("authorize Owner recovery Store")
185 .recover_owner_device(&authority, None)
186 .await
187 .expect("recover Owner device");
188
189 let durable = database
190 .latest_local_store_device_registration()
191 .await
192 .expect("load replacement registration")
193 .expect("replacement registration exists");
194 assert_eq!(durable.device_id, registration.device_id);
195 assert!(durable.is_activated());
196 let loaded = store
197 .bind_device(&db, db_store_dir.clone(), &signer)
198 .await
199 .expect("load recovered Owner Store");
200 loaded
201 .authorize_writer()
202 .await
203 .expect("replacement registration is usable");
204 }
205
206 #[tokio::test]
207 async fn recovery_materialization_reopens_its_retained_introduced_author() {
208 let (_store, db, _db_store_dir, registration, reference) = recovered_author().await;
209 let registration = registration.clone();
210 db.corrupt_store_device_registration_bytes_for_test(registration)
211 .await
212 .expect("corrupt activated recovery registration fixture");
213
214 let frontier = coven_database::StoreDatabase::new(&db)
215 .materialized_frontier()
216 .await
217 .expect("retained recovery author does not depend on mutable registration rows");
218 let StoreCommitCoord { stream_id, .. } = &reference.coord;
219 assert_eq!(frontier.get(&stream_id.to_string()), Some(&reference));
220 }
221
222 #[tokio::test]
223 async fn recovery_materialization_rejects_tampered_retained_registration_bytes() {
224 let (store, db, _db_store_dir, _registration, reference) = recovered_author().await;
225 db.tamper_retained_recovery_registration_for_test(
226 &reference,
227 coven_database::RetainedRegistrationTamper::CanonicalRegistration,
228 )
229 .await;
230
231 let root = store.root().clone();
232 db.validate_retained_merge_replay_for_test(root)
233 .await
234 .expect_err(
235 "tampered retained recovery registration bytes must fail durable history verification",
236 );
237 }
238
239 #[tokio::test]
240 async fn recovery_materialization_rejects_tampered_retained_registration_authority() {
241 let (store, db, _db_store_dir, _registration, reference) = recovered_author().await;
242 db.tamper_retained_recovery_registration_for_test(
243 &reference,
244 coven_database::RetainedRegistrationTamper::ActivationAuthority,
245 )
246 .await;
247
248 let root = store.root().clone();
249 db
250 .validate_retained_merge_replay_for_test(root)
251 .await
252 .expect_err(
253 "tampered retained recovery registration authority must fail durable history verification",
254 );
255 }
256
257 #[tokio::test]
258 async fn owner_recovery_retry_reuses_each_published_readiness_prefix() {
259 for failed_call in [2, 3, 4] {
260 let signer = UserKeypair::generate();
261 let db_store_dir = crate::sync::test_helpers::test_store_dir();
262 let db = crate::sync::test_helpers::open_test_db(db_store_dir.clone());
263 let home = crate::sync::test_helpers::test_cloud_home();
264 let (store, cloud_storage) = TestStore::create_with_connection(
265 &db,
266 db_store_dir.clone(),
267 &format!("recovery-prefix-{failed_call}"),
268 signer.clone(),
269 home.clone(),
270 )
271 .await
272 .expect("create recovery prefix Store");
273 let loaded = store
274 .bind_device(&db, db_store_dir.clone(), &signer)
275 .await
276 .expect("load recovery Store");
277 let authority = store.founder_recovery_authority().await;
278 let database = coven_database::StoreDatabase::new(&db);
279 let mut recovery = loaded
280 .owner_recovery_for_test()
281 .await
282 .expect("authorize Owner recovery Store");
283 home.fail_exact_create_before_call(failed_call);
284 assert!(
285 recovery
286 .recover_owner_device(&authority, None)
287 .await
288 .is_err(),
289 "failure before exact create {failed_call} interrupts recovery",
290 );
291
292 let interrupted = database
293 .latest_local_store_device_registration()
294 .await
295 .expect("read interrupted recovery journal")
296 .expect("interrupted recovery journal exists");
297 let interrupted_node = if failed_call == 4 {
298 let registration = StoreDeviceRegistration::parse_at(
299 &interrupted.registration_bytes,
300 &store.root(),
301 interrupted.device_id,
302 )
303 .expect("parse interrupted recovery registration");
304 let StoreDeviceRegistrationOrigin::Recovery { recovery_slot, .. } =
305 registration.origin.clone()
306 else {
307 panic!("interrupted registration is not a Recovery registration");
308 };
309 let context = coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
310 store.root().store_root_hash,
311 ProtocolObjectDomain::OwnerRecoveryNode,
312 );
313 Some(
314 cloud_storage
315 .read_prepared_protocol_slot(
316 &context,
317 &recovery_slot,
318 &owner_recovery_semantic_prefix(
319 &coven_keys::keys::public_key_hex(&signer),
320 authority.owner_grant.clone(),
321 1,
322 ),
323 )
324 .await
325 .expect("read published recovery node")
326 .1,
327 )
328 } else {
329 None
330 };
331
332 recovery
333 .recover_owner_device(&authority, None)
334 .await
335 .expect("retry completes absent recovery suffix");
336 assert_eq!(
337 home.exact_create_count(),
338 6,
339 "retry after boundary {failed_call} creates only the absent suffix",
340 );
341 let completed = database
342 .latest_local_store_device_registration()
343 .await
344 .expect("read completed recovery journal")
345 .expect("completed recovery journal exists");
346 assert_eq!(
347 completed.prepared.reference(),
348 interrupted.prepared.reference(),
349 );
350 assert_eq!(
351 completed.prepared.stored_bytes(),
352 interrupted.prepared.stored_bytes(),
353 );
354 if failed_call >= 3 {
355 assert_eq!(
356 completed.initial_ack.prepared.reference(),
357 interrupted.initial_ack.prepared.reference(),
358 );
359 assert_eq!(
360 completed.initial_ack.prepared.stored_bytes(),
361 interrupted.initial_ack.prepared.stored_bytes(),
362 );
363 }
364 if let Some(interrupted_node) = interrupted_node {
365 let registration = StoreDeviceRegistration::parse_at(
366 &completed.registration_bytes,
367 &store.root(),
368 completed.device_id,
369 )
370 .expect("parse completed recovery registration");
371 let StoreDeviceRegistrationOrigin::Recovery { recovery_slot, .. } =
372 registration.origin.clone()
373 else {
374 panic!("completed registration is not a Recovery registration");
375 };
376 let context = coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
377 store.root().store_root_hash,
378 ProtocolObjectDomain::OwnerRecoveryNode,
379 );
380 let completed_node = cloud_storage
381 .read_prepared_protocol_slot(
382 &context,
383 &recovery_slot,
384 &owner_recovery_semantic_prefix(
385 &coven_keys::keys::public_key_hex(&signer),
386 authority.owner_grant.clone(),
387 1,
388 ),
389 )
390 .await
391 .expect("read completed recovery node")
392 .1;
393 assert_eq!(completed_node.reference(), interrupted_node.reference());
394 assert_eq!(
395 completed_node.stored_bytes(),
396 interrupted_node.stored_bytes(),
397 );
398 }
399 }
400 }
401
402 #[tokio::test]
403 async fn owner_recovery_retry_reuses_its_staged_activation_after_history_advances() {
404 let founder = UserKeypair::generate();
405 let founder_store_dir = crate::sync::test_helpers::test_store_dir();
406 let founder_db = crate::sync::test_helpers::open_test_db(founder_store_dir.clone());
407 let home = crate::sync::test_helpers::test_cloud_home();
408 let (store, cloud_storage) = TestStore::create_with_connection(
409 &founder_db,
410 founder_store_dir.clone(),
411 "staged-recovery-retry",
412 founder.clone(),
413 home.clone(),
414 )
415 .await
416 .expect("create recovery retry Store");
417 let peer = UserKeypair::generate();
418 let peer_store_dir = crate::sync::test_helpers::test_store_dir();
419 let peer_db = crate::sync::test_helpers::open_test_db(peer_store_dir.clone());
420 let peer_device = store
421 .admit_and_activate_peer(
422 &founder_db,
423 founder_store_dir.clone(),
424 &peer_db,
425 peer_store_dir,
426 &peer,
427 )
428 .await
429 .expect("activate peer writer");
430 let founder_device = store
431 .bind_device(&founder_db, founder_store_dir.clone(), &founder)
432 .await
433 .expect("bind recovery Store");
434 let authority = store.founder_recovery_authority().await;
435 let database = coven_database::StoreDatabase::new(&founder_db);
436 let mut recovery = founder_device
437 .owner_recovery_for_test()
438 .await
439 .expect("authorize Owner recovery Store");
440
441 home.fail_exact_create_before_call(4);
442 recovery
443 .recover_owner_device(&authority, None)
444 .await
445 .expect_err("activation commit publication is interrupted");
446 let staged_before = database
447 .owner_recovery_publication()
448 .await
449 .expect("read staged recovery publication")
450 .expect("recovery activation is staged before publication");
451
452 peer_device
453 .publish_fixture_position("staged-recovery")
454 .await;
455 let recovered = recovery
456 .recover_owner_device(&authority, None)
457 .await
458 .expect("retry publishes the exact staged activation");
459 let commit_value = staged_before.commit.value.value();
460 let commit_prefix = coven_protocol::store_commit::commit_semantic_prefix(
461 commit_value.candidate_family(),
462 &staged_before
463 .commit
464 .value
465 .reference()
466 .coord
467 .stream_id
468 .to_string(),
469 commit_value.seq(),
470 commit_value.commit_hash(),
471 );
472 let published_commit = cloud_storage
473 .read_prepared_protocol_slot(
474 &coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
475 store.root().store_root_hash,
476 ProtocolObjectDomain::StoreCommit,
477 ),
478 staged_before.commit.prepared.reference().slot(),
479 &commit_prefix,
480 )
481 .await
482 .expect("read published recovery commit")
483 .1;
484 assert_eq!(
485 published_commit.reference(),
486 staged_before.commit.prepared.reference(),
487 );
488 assert_eq!(
489 published_commit.stored_bytes(),
490 staged_before.commit.prepared.stored_bytes(),
491 );
492 let head_prefix = coven_protocol::store_commit::head_slot_prefix(
493 &recovered.device_id.to_string(),
494 staged_before.head.value.slot_sequence(),
495 );
496 let published_head = cloud_storage
497 .read_prepared_protocol_slot(
498 &coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
499 store.root().store_root_hash,
500 ProtocolObjectDomain::StoreHead,
501 ),
502 staged_before.head.prepared.reference().slot(),
503 &head_prefix,
504 )
505 .await
506 .expect("read published recovery head")
507 .1;
508 assert_eq!(
509 published_head.reference(),
510 staged_before.head.prepared.reference(),
511 );
512 assert_eq!(
513 published_head.stored_bytes(),
514 staged_before.head.prepared.stored_bytes(),
515 );
516 }
517
518 #[tokio::test]
519 async fn owner_recovery_activation_covers_history_published_after_its_initial_ack() {
520 let founder = UserKeypair::generate();
521 let founder_store_dir = crate::sync::test_helpers::test_store_dir();
522 let founder_db = crate::sync::test_helpers::open_test_db(founder_store_dir.clone());
523 let home = crate::sync::test_helpers::test_cloud_home();
524 let (store, _cloud_storage) = TestStore::create_with_connection(
525 &founder_db,
526 founder_store_dir.clone(),
527 "recovery-predecessor-barrier",
528 founder.clone(),
529 home.clone(),
530 )
531 .await
532 .expect("create recovery barrier Store");
533 let peer = UserKeypair::generate();
534 let peer_store_dir = crate::sync::test_helpers::test_store_dir();
535 let peer_db = crate::sync::test_helpers::open_test_db(peer_store_dir.clone());
536 let peer_device = store
537 .admit_and_activate_peer(
538 &founder_db,
539 founder_store_dir.clone(),
540 &peer_db,
541 peer_store_dir,
542 &peer,
543 )
544 .await
545 .expect("activate peer writer");
546 let founder_device = store
547 .bind_device(&founder_db, founder_store_dir.clone(), &founder)
548 .await
549 .expect("bind recovery Store");
550 let authority = store.founder_recovery_authority().await;
551 let mut recovery = founder_device
552 .owner_recovery_for_test()
553 .await
554 .expect("authorize Owner recovery Store");
555 let (node_published, release_recovery) = home.pause_after_exact_create_call(3);
556
557 let recover = recovery.recover_owner_device(&authority, None);
558 let publish = async {
559 node_published.notified().await;
560 peer_device
561 .publish_fixture_position("recovery-predecessor")
562 .await;
563 let reference = peer_device
564 .latest_local_store_position()
565 .await
566 .expect("read peer position")
567 .expect("peer position exists");
568 release_recovery.notify_one();
569 reference
570 };
571 let (recovered, peer_reference) = tokio::join!(recover, publish);
572 let recovered_registration = recovered.expect("recover over the fixed predecessor cut");
573
574 let loaded = store
575 .bind_device(&founder_db, founder_store_dir, &founder)
576 .await
577 .expect("reload recovered Store");
578 let database = coven_database::StoreDatabase::new(&founder_db);
579 let mut activation = None;
580 for reference in database
581 .materialized_frontier()
582 .await
583 .expect("read recovery frontier")
584 .into_values()
585 {
586 let commit = loaded
587 .load_commit_for_test(&reference)
588 .await
589 .expect("load recovery frontier commit");
590 if commit.value().author_registration == recovered_registration {
591 activation = Some(commit);
592 break;
593 }
594 }
595 let activation = activation.expect("recovery activation is materialized");
596 assert_eq!(
597 activation
598 .value()
599 .order
600 .dependencies
601 .get(&peer_reference.coord.stream_id),
602 Some(&peer_reference),
603 "the recovery activation orders itself after the peer history visible before staging",
604 );
605 }
606
607 #[tokio::test]
608 async fn owner_recovery_does_not_stage_activation_across_held_history() {
609 let founder = UserKeypair::generate();
610 let founder_store_dir = crate::sync::test_helpers::test_store_dir();
611 let founder_db = crate::sync::test_helpers::open_test_db(founder_store_dir.clone());
612 let home = crate::sync::test_helpers::test_cloud_home();
613 let (store, _cloud_storage) = TestStore::create_with_connection(
614 &founder_db,
615 founder_store_dir.clone(),
616 "held-recovery-predecessor",
617 founder.clone(),
618 home.clone(),
619 )
620 .await
621 .expect("create held recovery Store");
622 let peer = UserKeypair::generate();
623 let peer_store_dir = crate::sync::test_helpers::test_store_dir();
624 let peer_db = crate::sync::test_helpers::open_test_db(peer_store_dir.clone());
625 let peer_device = store
626 .admit_and_activate_peer(
627 &founder_db,
628 founder_store_dir.clone(),
629 &peer_db,
630 peer_store_dir,
631 &peer,
632 )
633 .await
634 .expect("activate peer writer");
635 peer_device.publish_fixture_position("held-recovery").await;
636 let peer_reference = peer_device
637 .latest_local_store_position()
638 .await
639 .expect("read held peer position")
640 .expect("held peer position exists");
641 let peer_commit = peer_device
642 .load_commit_for_test(&peer_reference)
643 .await
644 .expect("load held peer commit");
645 let package = peer_commit
646 .value()
647 .store_package()
648 .expect("peer commit carries its Store package");
649 home.remove_exact_object(package.object.slot());
650
651 let founder_device = store
652 .bind_device(&founder_db, founder_store_dir, &founder)
653 .await
654 .expect("bind recovery Store");
655 let authority = store.founder_recovery_authority().await;
656 let database = coven_database::StoreDatabase::new(&founder_db);
657 let error = founder_device
658 .owner_recovery_for_test()
659 .await
660 .expect("authorize Owner recovery Store")
661 .recover_owner_device(&authority, None)
662 .await
663 .expect_err("held predecessor history blocks activation staging");
664
665 assert!(
666 error.to_string().contains("is held at"),
667 "unexpected held recovery error: {error}",
668 );
669 assert!(
670 database
671 .owner_recovery_publication()
672 .await
673 .expect("read recovery publication state")
674 .is_none(),
675 "no activation is staged without its complete predecessor history",
676 );
677 }
678
679 #[tokio::test]
680 async fn published_owner_recovery_blocks_snapshot_retirement_until_activation() {
681 let founder = UserKeypair::generate();
682 let founder_store_dir = crate::sync::test_helpers::test_store_dir();
683 let founder_db = crate::sync::test_helpers::open_test_db(founder_store_dir.clone());
684 let home = crate::sync::test_helpers::test_cloud_home();
685 let (store, _cloud_storage) = TestStore::create_with_connection(
686 &founder_db,
687 founder_store_dir.clone(),
688 "pending-recovery-retirement",
689 founder.clone(),
690 home.clone(),
691 )
692 .await
693 .expect("create recovery retirement Store");
694 let peer = UserKeypair::generate();
695 let peer_store_dir = crate::sync::test_helpers::test_store_dir();
696 let peer_db = crate::sync::test_helpers::open_test_db(peer_store_dir.clone());
697 let peer_device = store
698 .admit_and_activate_peer(
699 &founder_db,
700 founder_store_dir.clone(),
701 &peer_db,
702 peer_store_dir,
703 &peer,
704 )
705 .await
706 .expect("activate peer writer");
707 let founder_device = store
708 .bind_device(&founder_db, founder_store_dir.clone(), &founder)
709 .await
710 .expect("bind founder writer");
711 peer_device
712 .publish_fixture_position("retirement-snapshot-input")
713 .await;
714 let (_, founder_pull) = founder_device
715 .pull_store()
716 .await
717 .expect("pull snapshot input into founder");
718 assert!(founder_pull.held_positions.is_empty());
719
720 let founder_database = coven_database::StoreDatabase::new(&founder_db);
721 let coverage = coven_protocol::store_commit::CommitFrontier::from_refs(
722 founder_database
723 .materialized_frontier()
724 .await
725 .expect("read snapshot coverage"),
726 )
727 .expect("shape snapshot coverage");
728 let image_dir = tempfile::tempdir().expect("create snapshot image directory");
729 let encryption = coven_keys::encryption::EncryptionService::from_key([42; 32]);
730 let image = founder_database
731 .capture_snapshot_image_for_test(
732 store.root(),
733 image_dir.path().to_path_buf(),
734 Some(encryption.clone()),
735 )
736 .await
737 .expect("capture snapshot image");
738 founder_device
739 .publish_snapshot(image, coverage.clone())
740 .await
741 .expect("publish retirement snapshot");
742 founder_device
743 .publish_acknowledgement_without_advancing(coverage.clone())
744 .await
745 .expect("publish founder crossing acknowledgement");
746 peer_device
747 .publish_acknowledgement_without_advancing(coverage.clone())
748 .await
749 .expect("publish peer crossing acknowledgement");
750 founder_device
751 .publish_fixture_position("founder-acknowledgement-activation")
752 .await;
753 peer_device
754 .publish_fixture_position("peer-acknowledgement-activation")
755 .await;
756 let (_, peer_pull) = peer_device
757 .pull_store()
758 .await
759 .expect("materialize the acknowledgement closure");
760 assert!(peer_pull.held_positions.is_empty());
761
762 let authority = store.founder_recovery_authority().await;
763 let mut recovery = founder_device
764 .owner_recovery_for_test()
765 .await
766 .expect("authorize Owner recovery Store");
767 let (node_published, release_recovery) = home.pause_after_exact_create_call(3);
768 let recover = recovery.recover_owner_device(&authority, Some(&encryption));
769 let retire = async {
770 node_published.notified().await;
771 let outcome = peer_device
772 .stand_on_acknowledged_snapshot()
773 .await
774 .expect("evaluate retirement while recovery is pending");
775 release_recovery.notify_one();
776 outcome
777 };
778 let (recovered, retirement) = tokio::join!(recover, retire);
779 recovered.expect("recovery completes after retirement declines");
780 assert!(
781 matches!(
782 retirement,
783 crate::sync::store::ReplayBaselineAdvance::Declined(
784 crate::sync::store::ReplayBaselineDecline::PendingOwnerRecovery { .. }
785 )
786 ),
787 "published recovery registration must block retirement: {retirement:?}",
788 );
789 }
790}