Skip to main content

coven_replication/sync/store/authorization/
registration.rs

1//! Durable append-only Store device registration and recovery.
2
3use 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                // Registration, initial ACK, recovery node, authority entry/head,
339                // commit, publication, result, and the interrupted create.
340                9,
341                "retry after boundary {failed_call} creates only the absent suffix",
342            );
343            let completed = database
344                .latest_local_store_device_registration()
345                .await
346                .expect("read completed recovery journal")
347                .expect("completed recovery journal exists");
348            assert_eq!(
349                completed.prepared.reference(),
350                interrupted.prepared.reference(),
351            );
352            assert_eq!(
353                completed.prepared.stored_bytes(),
354                interrupted.prepared.stored_bytes(),
355            );
356            if failed_call >= 3 {
357                assert_eq!(
358                    completed.initial_ack.prepared.reference(),
359                    interrupted.initial_ack.prepared.reference(),
360                );
361                assert_eq!(
362                    completed.initial_ack.prepared.stored_bytes(),
363                    interrupted.initial_ack.prepared.stored_bytes(),
364                );
365            }
366            if let Some(interrupted_node) = interrupted_node {
367                let registration = StoreDeviceRegistration::parse_at(
368                    &completed.registration_bytes,
369                    &store.root(),
370                    completed.device_id,
371                )
372                .expect("parse completed recovery registration");
373                let StoreDeviceRegistrationOrigin::Recovery { recovery_slot, .. } =
374                    registration.origin.clone()
375                else {
376                    panic!("completed registration is not a Recovery registration");
377                };
378                let context = coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
379                    store.root().store_root_hash,
380                    ProtocolObjectDomain::OwnerRecoveryNode,
381                );
382                let completed_node = cloud_storage
383                    .read_prepared_protocol_slot(
384                        &context,
385                        &recovery_slot,
386                        &owner_recovery_semantic_prefix(
387                            &coven_keys::keys::public_key_hex(&signer),
388                            authority.owner_grant.clone(),
389                            1,
390                        ),
391                    )
392                    .await
393                    .expect("read completed recovery node")
394                    .1;
395                assert_eq!(completed_node.reference(), interrupted_node.reference());
396                assert_eq!(
397                    completed_node.stored_bytes(),
398                    interrupted_node.stored_bytes(),
399                );
400            }
401        }
402    }
403
404    #[tokio::test]
405    async fn owner_recovery_retry_reuses_its_staged_activation_after_history_advances() {
406        let founder = UserKeypair::generate();
407        let founder_store_dir = crate::sync::test_helpers::test_store_dir();
408        let founder_db = crate::sync::test_helpers::open_test_db(founder_store_dir.clone());
409        let home = crate::sync::test_helpers::test_cloud_home();
410        let (store, cloud_storage) = TestStore::create_with_connection(
411            &founder_db,
412            founder_store_dir.clone(),
413            "staged-recovery-retry",
414            founder.clone(),
415            home.clone(),
416        )
417        .await
418        .expect("create recovery retry Store");
419        let peer = UserKeypair::generate();
420        let peer_store_dir = crate::sync::test_helpers::test_store_dir();
421        let peer_db = crate::sync::test_helpers::open_test_db(peer_store_dir.clone());
422        let peer_device = store
423            .admit_and_activate_peer(
424                &founder_db,
425                founder_store_dir.clone(),
426                &peer_db,
427                peer_store_dir,
428                &peer,
429            )
430            .await
431            .expect("activate peer writer");
432        let founder_device = store
433            .bind_device(&founder_db, founder_store_dir.clone(), &founder)
434            .await
435            .expect("bind recovery Store");
436        let authority = store.founder_recovery_authority().await;
437        let database = coven_database::StoreDatabase::new(&founder_db);
438        let mut recovery = founder_device
439            .owner_recovery_for_test()
440            .await
441            .expect("authorize Owner recovery Store");
442
443        home.fail_exact_create_before_call(4);
444        recovery
445            .recover_owner_device(&authority, None)
446            .await
447            .expect_err("activation commit publication is interrupted");
448        let staged_before = database
449            .owner_recovery_publication()
450            .await
451            .expect("read staged recovery publication")
452            .expect("recovery activation is staged before publication");
453
454        peer_device
455            .publish_fixture_position("staged-recovery")
456            .await;
457        let _recovered = recovery
458            .recover_owner_device(&authority, None)
459            .await
460            .expect("retry publishes the exact staged activation");
461        let commit_value = staged_before.commit.value.value();
462        let commit_prefix = coven_protocol::store_commit::commit_semantic_prefix(
463            commit_value.candidate_family(),
464            &staged_before
465                .commit
466                .value
467                .reference()
468                .coord
469                .stream_id
470                .to_string(),
471            commit_value.seq(),
472            commit_value.commit_hash(),
473        );
474        let published_commit = cloud_storage
475            .read_prepared_protocol_slot(
476                &coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
477                    store.root().store_root_hash,
478                    ProtocolObjectDomain::StoreCommit,
479                ),
480                staged_before.commit.prepared.reference().slot(),
481                &commit_prefix,
482            )
483            .await
484            .expect("read published recovery commit")
485            .1;
486        assert_eq!(
487            published_commit.reference(),
488            staged_before.commit.prepared.reference(),
489        );
490        assert_eq!(
491            published_commit.stored_bytes(),
492            staged_before.commit.prepared.stored_bytes(),
493        );
494        let accepted_publication = database
495            .store_publication_entries()
496            .await
497            .expect("read accepted recovery publication")
498            .into_iter()
499            .find(|entry| {
500                entry.value.payload
501                    == coven_protocol::store_commit::StorePublicationPayload::Commit(
502                        staged_before.commit.value.reference().clone(),
503                    )
504            })
505            .expect("the unchanged recovery commit has an accepted publication");
506        assert_ne!(accepted_publication.value, staged_before.publication.entry);
507        let publication_prefix =
508            coven_protocol::store_commit::store_publication_entry_semantic_prefix(
509                &accepted_publication.value,
510            );
511        let published_publication = cloud_storage
512            .read_prepared_protocol_slot(
513                &coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
514                    store.root().store_root_hash,
515                    ProtocolObjectDomain::StorePublicationEntry,
516                ),
517                accepted_publication.prepared.reference().slot(),
518                &publication_prefix,
519            )
520            .await
521            .expect("read published recovery publication entry")
522            .1;
523        assert_eq!(
524            published_publication.reference(),
525            accepted_publication.prepared.reference(),
526        );
527        assert_eq!(
528            published_publication.stored_bytes(),
529            accepted_publication.value.to_bytes(),
530        );
531    }
532
533    #[tokio::test]
534    async fn owner_recovery_activation_covers_history_published_after_its_initial_ack() {
535        let founder = UserKeypair::generate();
536        let founder_store_dir = crate::sync::test_helpers::test_store_dir();
537        let founder_db = crate::sync::test_helpers::open_test_db(founder_store_dir.clone());
538        let home = crate::sync::test_helpers::test_cloud_home();
539        let (store, _cloud_storage) = TestStore::create_with_connection(
540            &founder_db,
541            founder_store_dir.clone(),
542            "recovery-predecessor-barrier",
543            founder.clone(),
544            home.clone(),
545        )
546        .await
547        .expect("create recovery barrier Store");
548        let peer = UserKeypair::generate();
549        let peer_store_dir = crate::sync::test_helpers::test_store_dir();
550        let peer_db = crate::sync::test_helpers::open_test_db(peer_store_dir.clone());
551        let peer_device = store
552            .admit_and_activate_peer(
553                &founder_db,
554                founder_store_dir.clone(),
555                &peer_db,
556                peer_store_dir,
557                &peer,
558            )
559            .await
560            .expect("activate peer writer");
561        let founder_device = store
562            .bind_device(&founder_db, founder_store_dir.clone(), &founder)
563            .await
564            .expect("bind recovery Store");
565        let authority = store.founder_recovery_authority().await;
566        let mut recovery = founder_device
567            .owner_recovery_for_test()
568            .await
569            .expect("authorize Owner recovery Store");
570        let (node_published, release_recovery) = home.pause_after_exact_create_call(3);
571
572        let recover = recovery.recover_owner_device(&authority, None);
573        let publish = async {
574            node_published.notified().await;
575            peer_device
576                .publish_fixture_position("recovery-predecessor")
577                .await;
578            let reference = peer_device
579                .latest_local_store_position()
580                .await
581                .expect("read peer position")
582                .expect("peer position exists");
583            release_recovery.notify_one();
584            reference
585        };
586        let (recovered, peer_reference) = tokio::join!(recover, publish);
587        let recovered_registration = recovered.expect("recover over the fixed predecessor cut");
588
589        let loaded = store
590            .bind_device(&founder_db, founder_store_dir, &founder)
591            .await
592            .expect("reload recovered Store");
593        let database = coven_database::StoreDatabase::new(&founder_db);
594        let mut activation = None;
595        for reference in database
596            .materialized_frontier()
597            .await
598            .expect("read recovery frontier")
599            .into_values()
600        {
601            let commit = loaded
602                .load_commit_for_test(&reference)
603                .await
604                .expect("load recovery frontier commit");
605            if commit.value().author_registration == recovered_registration {
606                activation = Some(commit);
607                break;
608            }
609        }
610        let activation = activation.expect("recovery activation is materialized");
611        assert_eq!(
612            activation
613                .value()
614                .order
615                .dependencies
616                .get(&peer_reference.coord.stream_id),
617            Some(&peer_reference),
618            "the recovery activation orders itself after the peer history visible before staging",
619        );
620    }
621
622    #[tokio::test]
623    async fn owner_recovery_does_not_stage_activation_across_held_history() {
624        let founder = UserKeypair::generate();
625        let founder_store_dir = crate::sync::test_helpers::test_store_dir();
626        let founder_db = crate::sync::test_helpers::open_test_db(founder_store_dir.clone());
627        let home = crate::sync::test_helpers::test_cloud_home();
628        let (store, _cloud_storage) = TestStore::create_with_connection(
629            &founder_db,
630            founder_store_dir.clone(),
631            "held-recovery-predecessor",
632            founder.clone(),
633            home.clone(),
634        )
635        .await
636        .expect("create held recovery Store");
637        let peer = UserKeypair::generate();
638        let peer_store_dir = crate::sync::test_helpers::test_store_dir();
639        let peer_db = crate::sync::test_helpers::open_test_db(peer_store_dir.clone());
640        let peer_device = store
641            .admit_and_activate_peer(
642                &founder_db,
643                founder_store_dir.clone(),
644                &peer_db,
645                peer_store_dir,
646                &peer,
647            )
648            .await
649            .expect("activate peer writer");
650        peer_device.publish_fixture_position("held-recovery").await;
651        let peer_reference = peer_device
652            .latest_local_store_position()
653            .await
654            .expect("read held peer position")
655            .expect("held peer position exists");
656        let peer_commit = peer_device
657            .load_commit_for_test(&peer_reference)
658            .await
659            .expect("load held peer commit");
660        let package = peer_commit
661            .value()
662            .store_package()
663            .expect("peer commit carries its Store package");
664        home.remove_exact_object(package.object.slot());
665
666        let founder_device = store
667            .bind_device(&founder_db, founder_store_dir, &founder)
668            .await
669            .expect("bind recovery Store");
670        let authority = store.founder_recovery_authority().await;
671        let database = coven_database::StoreDatabase::new(&founder_db);
672        let error = founder_device
673            .owner_recovery_for_test()
674            .await
675            .expect("authorize Owner recovery Store")
676            .recover_owner_device(&authority, None)
677            .await
678            .expect_err("held predecessor history blocks activation staging");
679
680        assert!(
681            error.to_string().contains("is held at"),
682            "unexpected held recovery error: {error}",
683        );
684        assert!(
685            database
686                .owner_recovery_publication()
687                .await
688                .expect("read recovery publication state")
689                .is_none(),
690            "no activation is staged without its complete predecessor history",
691        );
692    }
693
694    #[tokio::test]
695    async fn published_owner_recovery_survives_snapshot_retirement() {
696        let founder = UserKeypair::generate();
697        let founder_store_dir = crate::sync::test_helpers::test_store_dir();
698        let founder_db = crate::sync::test_helpers::open_test_db(founder_store_dir.clone());
699        let home = crate::sync::test_helpers::test_cloud_home();
700        let (store, _cloud_storage) = TestStore::create_with_connection(
701            &founder_db,
702            founder_store_dir.clone(),
703            "pending-recovery-retirement",
704            founder.clone(),
705            home.clone(),
706        )
707        .await
708        .expect("create recovery retirement Store");
709        let peer_store_dir = crate::sync::test_helpers::test_store_dir();
710        let peer_db = crate::sync::test_helpers::open_test_db(peer_store_dir.clone());
711        let recovering_device = store
712            .activate_joined_device(
713                &founder_db,
714                founder_store_dir.clone(),
715                &peer_db,
716                peer_store_dir,
717                &founder,
718                "2026-07-16T00:00:00Z",
719            )
720            .await
721            .expect("activate another Owner device");
722        let administrator = store
723            .bind_device(&founder_db, founder_store_dir, &founder)
724            .await
725            .expect("bind original provider administrator");
726        recovering_device
727            .publish_fixture_position("retirement-snapshot-input")
728            .await;
729        let original = recovering_device
730            .publish_snapshot_generation_for_test()
731            .await
732            .expect("publish the recovering device's snapshot");
733        let (_, pulled) = administrator
734            .pull_store()
735            .await
736            .expect("administrator adopts the snapshot");
737        assert!(pulled.held_positions.is_empty());
738        assert!(home.contains_exact_object(&original.meta.image.object));
739        assert!(home.contains_exact_object(&original.reference.object));
740
741        let authority = store.founder_recovery_authority().await;
742        let encryption = coven_keys::encryption::EncryptionService::from_key([42; 32]);
743        let mut recovery = recovering_device
744            .owner_recovery_for_test()
745            .await
746            .expect("authorize Owner recovery Store");
747        let (node_published, release_recovery) = home.pause_after_exact_create_call(3);
748        let recover = recovery.recover_owner_device(&authority, Some(&encryption));
749        tokio::pin!(recover);
750        tokio::time::timeout(std::time::Duration::from_secs(10), async {
751            tokio::select! {
752                _ = node_published.notified() => {},
753                result = &mut recover => panic!("recovery returned before its node upload barrier: {result:?}"),
754            }
755        })
756        .await
757        .expect("pause recovery after its node upload");
758        let pending = coven_database::StoreDatabase::new(&peer_db)
759            .latest_local_store_device_registration()
760            .await
761            .expect("read pending recovery registration")
762            .expect("recovery readiness is durable before its node upload");
763        assert!(!pending.is_activated());
764        let registration = StoreDeviceRegistration::parse_at(
765            &pending.registration_bytes,
766            &store.root(),
767            pending.device_id,
768        )
769        .expect("verify the pending recovery registration");
770        assert!(matches!(
771            registration.origin,
772            StoreDeviceRegistrationOrigin::Recovery { .. }
773        ));
774
775        administrator
776            .publish_fixture_position("successor-snapshot-input")
777            .await;
778        let successor = administrator
779            .publish_snapshot_generation_for_test()
780            .await
781            .expect("publish successor coverage while recovery is pending");
782        assert_ne!(successor.reference, original.reference);
783        administrator
784            .reclaim_packages()
785            .await
786            .expect("retire the original snapshot while recovery is pending");
787        let image_retired = !home.contains_exact_object(&original.meta.image.object);
788        let metadata_retired = !home.contains_exact_object(&original.reference.object);
789        release_recovery.notify_one();
790        assert!(
791            image_retired,
792            "the original snapshot image must be physically retired"
793        );
794        assert!(
795            metadata_retired,
796            "the original snapshot metadata must be physically retired"
797        );
798        tokio::time::timeout(std::time::Duration::from_secs(10), &mut recover)
799            .await
800            .expect("resume recovery after snapshot retirement")
801            .expect("recovery completes after its published snapshot is retired");
802    }
803}