Skip to main content

coven_domain/joining/
transport.rs

1//! The joining device's side of the storage-mediated device-join transport.
2//!
3//! The admitting side's driver lives beside the Store it advances; this is its
4//! counterpart on the device being admitted, where the join steps hang off
5//! [`DeviceJoinClient`] rather than an open Store.
6//!
7//! Each step is the same call a host driving the join by hand would make. The
8//! transport only replaces handing the artifacts across: publish what the step
9//! produced, wait for what the next step needs.
10
11use base64::engine::general_purpose::URL_SAFE_NO_PAD;
12use base64::Engine;
13use tokio::sync::watch;
14
15use crate::joining::client::{
16    enrollment_oauth_tokens, BootstrapError, DeviceJoinClient, EnrollmentProviderAccess,
17};
18use coven_foundation::config::Config;
19use coven_replication::sync::store::{
20    DeviceJoinAbandonment, DeviceJoinAction, DeviceJoinActivation, DeviceJoinOfferBundle,
21    DeviceJoinRole, DeviceJoinStatus, DeviceJoinStep, DeviceJoinTransport,
22    DeviceJoinTransportTiming, DeviceProviderAdmissionApproval, ProviderReadyDeviceBootstrap,
23};
24
25use coven_replication::sync::MemberAdmission;
26
27/// Complete the joining side after scanning the existing device's one pairing
28/// code. The local session returns the invitation sealed to this attempt; the
29/// existing Store transport then performs registration and bootstrap. The
30/// caller's Coven migration policy controls every writer open during bootstrap.
31#[allow(clippy::too_many_arguments)]
32pub async fn join_with_device_pairing(
33    pairing: &crate::joining::PreparedDevicePairing,
34    layout: coven_foundation::store_dir::StoreLayout,
35    synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
36    migrations: Vec<coven_database::Migration>,
37    coven_migration_policy: coven_database::CovenMigrationPolicy,
38    exact_upload_verification: coven_foundation::config::ExactUploadVerification,
39    transfer_limits: coven_protocol::blob::TransferLimits,
40    key_custody: coven_keys::custody::KeyCustody,
41    identity_custody: coven_keys::identity_custody::IdentityCustody,
42    oauth_clients: coven_storage::oauth::OAuthClients,
43    oauth_tokens: Option<coven_storage::oauth::OAuthTokens>,
44    cloudkit_ops: Option<std::sync::Arc<dyn coven_storage::cloud::cloudkit::CloudKitOps>>,
45    clock: coven_foundation::clock::ClockRef,
46    on_progress: coven_replication::sync::JoiningDeviceJoinProgressObserver,
47    cancel: &watch::Receiver<bool>,
48) -> Result<DeviceJoinTransportOutcome, BootstrapError> {
49    let timing = DeviceJoinTransportTiming::interactive();
50    let continuation = durable_invitation(
51        pairing,
52        &layout,
53        timing,
54        clock.clone(),
55        &on_progress,
56        cancel,
57    )
58    .await?;
59    let (pairing, invitation, provider_access) = match continuation {
60        EnrollmentContinuation::ProviderAccessPending {
61            pairing,
62            invitation,
63        } => (
64            pairing,
65            invitation,
66            EnrollmentProviderAccess::Supplied(oauth_tokens),
67        ),
68        EnrollmentContinuation::LibraryInstallationPending {
69            pairing,
70            invitation,
71        } => (pairing, invitation, EnrollmentProviderAccess::Stored),
72    };
73    let invite = DeviceJoinInvite::from_bytes(&invitation)?;
74    let client = invitation_client(
75        &invite,
76        pairing.request().public_key(),
77        layout.clone(),
78        synced_tables,
79        migrations,
80        coven_migration_policy,
81        exact_upload_verification,
82        transfer_limits,
83        key_custody,
84        identity_custody,
85        oauth_clients,
86        provider_access,
87        cloudkit_ops,
88        clock,
89    )?;
90    let pairing = pairing.record_library_installation_pending(&layout)?;
91    let outcome = client
92        .join_via_transport(&invite.bundle, timing, on_progress, cancel)
93        .await?;
94    pairing.finish(&layout)?;
95    Ok(outcome)
96}
97
98enum EnrollmentContinuation {
99    ProviderAccessPending {
100        pairing: crate::joining::PreparedDevicePairing,
101        invitation: Vec<u8>,
102    },
103    LibraryInstallationPending {
104        pairing: crate::joining::PreparedDevicePairing,
105        invitation: Vec<u8>,
106    },
107}
108
109async fn durable_invitation(
110    pairing: &crate::joining::PreparedDevicePairing,
111    layout: &coven_foundation::store_dir::StoreLayout,
112    timing: DeviceJoinTransportTiming,
113    clock: coven_foundation::clock::ClockRef,
114    on_progress: &coven_replication::sync::JoiningDeviceJoinProgressObserver,
115    cancel: &watch::Receiver<bool>,
116) -> Result<EnrollmentContinuation, BootstrapError> {
117    let pairing = match pairing.phase() {
118        crate::joining::DevicePairingPhase::AwaitingInvitation => {
119            on_progress(coven_replication::sync::JoiningDeviceJoinProgress::WaitingForApproval);
120            let invitation = crate::joining::receive_device_invitation(
121                pairing.offer(),
122                pairing.sealed_request(),
123                timing,
124                clock,
125                cancel,
126            )
127            .await
128            .map_err(BootstrapError::Pairing)?;
129            pairing.record_invitation_received(layout, &invitation)?
130        }
131        crate::joining::DevicePairingPhase::ProviderAccessPending
132        | crate::joining::DevicePairingPhase::LibraryInstallationPending => pairing.clone(),
133    };
134    let invitation = pairing
135        .pending_invitation()
136        .expect("a durable invitation phase carries its invitation")
137        .to_vec();
138    Ok(match pairing.phase() {
139        crate::joining::DevicePairingPhase::ProviderAccessPending => {
140            EnrollmentContinuation::ProviderAccessPending {
141                pairing,
142                invitation,
143            }
144        }
145        crate::joining::DevicePairingPhase::LibraryInstallationPending => {
146            EnrollmentContinuation::LibraryInstallationPending {
147                pairing,
148                invitation,
149            }
150        }
151        crate::joining::DevicePairingPhase::AwaitingInvitation => {
152            unreachable!("recording the invitation advances the durable phase")
153        }
154    })
155}
156
157/// Everything a joining device needs, sealed to the pending identity named by
158/// its signed pairing request. The transport bundle is public kickoff data; the member
159/// admission inside `sealed_invitation` carries provider credentials and can
160/// be opened only by that joining device.
161#[derive(Clone, Debug)]
162pub struct DeviceJoinInvite {
163    sealed_invitation: Vec<u8>,
164    pub bundle: DeviceJoinOfferBundle,
165}
166
167#[derive(serde::Serialize, serde::Deserialize)]
168#[serde(deny_unknown_fields)]
169struct DeviceJoinInviteWire {
170    version: u32,
171    sealed_invitation: String,
172    bundle: DeviceJoinOfferBundle,
173}
174
175impl DeviceJoinInvite {
176    pub fn new(
177        admission: MemberAdmission,
178        bundle: DeviceJoinOfferBundle,
179    ) -> Result<Self, DeviceInviteError> {
180        validate_admission(&admission)?;
181        require_admission_matches_bundle(&admission, &bundle)?;
182        let recipient =
183            coven_keys::keys::ed25519_hex_to_x25519_public_key(&bundle.offer.member_pubkey)?;
184        let plaintext =
185            serde_json::to_vec(&admission).expect("member admission serialization cannot fail");
186        Ok(Self {
187            sealed_invitation: coven_keys::keys::seal_box_encrypt(&plaintext, &recipient),
188            bundle,
189        })
190    }
191
192    pub(crate) fn open_admission(
193        &self,
194        member_pubkey: &str,
195    ) -> Result<MemberAdmission, DeviceInviteError> {
196        if member_pubkey != self.bundle.offer.member_pubkey {
197            return Err(DeviceInviteError::RecipientMismatch);
198        }
199        let recipient = coven_keys::keys::peek_pending_identity(member_pubkey)?;
200        let plaintext = coven_keys::keys::seal_box_decrypt(
201            &self.sealed_invitation,
202            &recipient.to_x25519_secret_key(),
203        )?;
204        let admission: MemberAdmission =
205            serde_json::from_slice(&plaintext).map_err(DeviceInviteError::AdmissionJson)?;
206        validate_admission(&admission)?;
207        require_admission_matches_bundle(&admission, &self.bundle)?;
208        Ok(admission)
209    }
210
211    pub fn to_bytes(&self) -> Vec<u8> {
212        serde_json::to_vec(&DeviceJoinInviteWire {
213            version: coven_protocol::store_commit::STORE_PROTOCOL_VERSION,
214            sealed_invitation: URL_SAFE_NO_PAD.encode(&self.sealed_invitation),
215            bundle: self.bundle.clone(),
216        })
217        .expect("device join invite serialization cannot fail")
218    }
219
220    pub fn from_bytes(bytes: &[u8]) -> Result<Self, BootstrapError> {
221        let wire: DeviceJoinInviteWire =
222            serde_json::from_slice(bytes).map_err(DeviceInviteError::WireJson)?;
223        if wire.version != coven_protocol::store_commit::STORE_PROTOCOL_VERSION {
224            return Err(BootstrapError::UnsupportedDeviceInviteVersion(wire.version));
225        }
226        Ok(Self {
227            sealed_invitation: URL_SAFE_NO_PAD
228                .decode(wire.sealed_invitation)
229                .map_err(DeviceInviteError::Ciphertext)?,
230            bundle: wire.bundle,
231        })
232    }
233}
234
235#[derive(Debug, thiserror::Error)]
236pub enum DeviceInviteError {
237    #[error("device invitation wire is not valid JSON: {0}")]
238    WireJson(#[source] serde_json::Error),
239    #[error("device invitation ciphertext is not valid base64: {0}")]
240    Ciphertext(#[source] base64::DecodeError),
241    #[error("device invitation admission payload is not valid JSON: {0}")]
242    AdmissionJson(#[source] serde_json::Error),
243    #[error("device invitation key: {0}")]
244    Key(#[from] coven_keys::keys::KeyError),
245    #[error("device invitation is for a different pairing identity")]
246    RecipientMismatch,
247    #[error("device invitation does not match its signed join offer")]
248    OfferMismatch,
249    #[error("device invitation admission payload is invalid: {0}")]
250    Admission(#[from] AdmissionPayloadError),
251}
252
253#[derive(Debug, thiserror::Error)]
254pub enum AdmissionPayloadError {
255    #[error("store id: {0}")]
256    StoreId(#[from] coven_foundation::store_dir::PathTokenError),
257    #[error("owner public key: {0}")]
258    OwnerPublicKey(#[source] coven_foundation::code_envelope::FixedHexError),
259    #[error("wrapped-key material: {0}")]
260    WrappedKeyMaterial(#[source] coven_foundation::code_envelope::FixedHexError),
261    #[error("wrapped-key identity: {0}")]
262    WrappedKeyIdentity(#[source] coven_protocol::objects::StorageError),
263    #[error("membership floor is empty")]
264    EmptyMembershipFloor,
265    #[error("membership floor: {0}")]
266    MembershipFloor(#[source] coven_protocol::membership::MembershipFloorError),
267}
268
269fn validate_admission(admission: &MemberAdmission) -> Result<(), AdmissionPayloadError> {
270    coven_foundation::store_dir::validate_path_token(&admission.store_id)?;
271    coven_foundation::code_envelope::decode_fixed_hex(
272        "owner public key",
273        &admission.owner_pubkey,
274        32,
275    )
276    .map_err(AdmissionPayloadError::OwnerPublicKey)?;
277    for (subject, value) in [
278        (
279            "wrapped-key author public key",
280            &admission.wrapped_key.owner_pubkey,
281        ),
282        (
283            "wrapped-key recipient public key",
284            &admission.wrapped_key.recipient_pubkey,
285        ),
286    ] {
287        coven_foundation::code_envelope::decode_fixed_hex(subject, value, 32)
288            .map_err(AdmissionPayloadError::WrappedKeyMaterial)?;
289    }
290    admission
291        .wrapped_key
292        .validate_identity()
293        .map_err(AdmissionPayloadError::WrappedKeyIdentity)?;
294    if admission.membership_floor.0.is_empty() {
295        return Err(AdmissionPayloadError::EmptyMembershipFloor);
296    }
297    admission
298        .membership_floor
299        .validate()
300        .map_err(AdmissionPayloadError::MembershipFloor)
301}
302
303fn require_admission_matches_bundle(
304    admission: &MemberAdmission,
305    bundle: &DeviceJoinOfferBundle,
306) -> Result<(), DeviceInviteError> {
307    if admission.wrapped_key.recipient_pubkey != bundle.offer.member_pubkey
308        || admission.store_root != bundle.offer.store_root
309    {
310        return Err(DeviceInviteError::OfferMismatch);
311    }
312    Ok(())
313}
314
315/// Build the joining device's client from a scanned payload.
316///
317/// The arguments after the payload are this device's own — its pending public
318/// key, and the same custody, schema, and provider wiring
319/// [`DeviceJoinClient::new`] takes, since that is what this constructs.
320#[allow(clippy::too_many_arguments)]
321fn invitation_client(
322    invite: &DeviceJoinInvite,
323    member_pubkey: &str,
324    layout: coven_foundation::store_dir::StoreLayout,
325    synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
326    migrations: Vec<coven_database::Migration>,
327    coven_migration_policy: coven_database::CovenMigrationPolicy,
328    exact_upload_verification: coven_foundation::config::ExactUploadVerification,
329    transfer_limits: coven_protocol::blob::TransferLimits,
330    key_custody: coven_keys::custody::KeyCustody,
331    identity_custody: coven_keys::identity_custody::IdentityCustody,
332    oauth_clients: coven_storage::oauth::OAuthClients,
333    provider_access: EnrollmentProviderAccess,
334    cloudkit_ops: Option<std::sync::Arc<dyn coven_storage::cloud::cloudkit::CloudKitOps>>,
335    clock: coven_foundation::clock::ClockRef,
336) -> Result<DeviceJoinClient, BootstrapError> {
337    let admission = invite.open_admission(member_pubkey)?;
338    let store_keys = coven_keys::keys::StoreKeys::bind(admission.store_id.clone());
339    let oauth_tokens = enrollment_oauth_tokens(&admission.join_info, &store_keys, provider_access)?;
340    DeviceJoinClient::new(
341        admission,
342        member_pubkey.to_string(),
343        layout,
344        synced_tables,
345        migrations,
346        coven_migration_policy,
347        exact_upload_verification,
348        transfer_limits,
349        key_custody,
350        identity_custody,
351        oauth_clients,
352        oauth_tokens,
353        cloudkit_ops,
354        clock,
355    )
356}
357
358/// Test-only: the joining device's client over an injected cloud home, the way
359/// the host's own test entry point injects one for the admitting side. The
360/// provider knobs a real device reads from its invitation are fixed here,
361/// since the home is supplied outright — including the exact-slot capability,
362/// which the injected home has by construction.
363#[cfg(any(test, feature = "test-utils"))]
364fn invitation_test_client(
365    invite: &DeviceJoinInvite,
366    member_pubkey: &str,
367    layout: coven_foundation::store_dir::StoreLayout,
368    synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
369    migrations: Vec<coven_database::Migration>,
370    coven_migration_policy: coven_database::CovenMigrationPolicy,
371    clock: coven_foundation::clock::ClockRef,
372    home: std::sync::Arc<dyn coven_storage::cloud::ExactCloudHome>,
373) -> Result<DeviceJoinClient, BootstrapError> {
374    Ok(invitation_client(
375        invite,
376        member_pubkey,
377        layout,
378        synced_tables,
379        migrations,
380        coven_migration_policy,
381        coven_foundation::config::ExactUploadVerification::MetadataHash,
382        coven_protocol::blob::TransferLimits::one_at_a_time(),
383        coven_keys::custody::KeyCustody::Keyring,
384        coven_keys::identity_custody::IdentityCustody::Keyring,
385        coven_storage::oauth::OAuthClients::empty(),
386        EnrollmentProviderAccess::InjectedHome,
387        None,
388        clock,
389    )?
390    .with_test_bootstrap_home(home))
391}
392
393/// Test-only counterpart of [`join_with_device_pairing`].
394#[cfg(any(test, feature = "test-utils"))]
395#[allow(clippy::too_many_arguments)]
396pub async fn join_with_device_pairing_over_test_home(
397    pairing: &crate::joining::PreparedDevicePairing,
398    layout: coven_foundation::store_dir::StoreLayout,
399    synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
400    migrations: Vec<coven_database::Migration>,
401    coven_migration_policy: coven_database::CovenMigrationPolicy,
402    clock: coven_foundation::clock::ClockRef,
403    home: std::sync::Arc<dyn coven_storage::cloud::ExactCloudHome>,
404    timing: DeviceJoinTransportTiming,
405    on_progress: coven_replication::sync::JoiningDeviceJoinProgressObserver,
406    cancel: &watch::Receiver<bool>,
407) -> Result<DeviceJoinTransportOutcome, BootstrapError> {
408    let continuation = durable_invitation(
409        pairing,
410        &layout,
411        timing,
412        clock.clone(),
413        &on_progress,
414        cancel,
415    )
416    .await?;
417    let (pairing, invitation) = match continuation {
418        EnrollmentContinuation::ProviderAccessPending {
419            pairing,
420            invitation,
421        }
422        | EnrollmentContinuation::LibraryInstallationPending {
423            pairing,
424            invitation,
425        } => (pairing, invitation),
426    };
427    let invite = DeviceJoinInvite::from_bytes(&invitation)?;
428    let client = invitation_test_client(
429        &invite,
430        pairing.request().public_key(),
431        layout.clone(),
432        synced_tables,
433        migrations,
434        coven_migration_policy,
435        clock,
436        home,
437    )?;
438    let pairing = pairing.record_library_installation_pending(&layout)?;
439    let outcome = client
440        .join_via_transport(&invite.bundle, timing, on_progress, cancel)
441        .await?;
442    pairing.finish(&layout)?;
443    Ok(outcome)
444}
445
446/// How a join driven through the transport ended for the joining device.
447#[derive(Clone, Debug)]
448pub enum DeviceJoinTransportOutcome {
449    /// The device is a member: its store is saved and its config returned.
450    Joined(Config),
451    /// The owner gave up on the attempt before it completed.
452    Abandoned(DeviceJoinAbandonment),
453}
454
455async fn publish_once(
456    transport: &DeviceJoinTransport<'_>,
457    published: &mut Vec<DeviceJoinAction>,
458    action: DeviceJoinAction,
459) -> Result<(), coven_replication::sync::DeviceJoinTransportError> {
460    if published.contains(&action) {
461        return Ok(());
462    }
463    transport.publish(&action).await?;
464    published.push(action);
465    Ok(())
466}
467
468impl DeviceJoinClient {
469    /// Join through the transport: one call from the scanned offer bundle to a
470    /// saved member [`Config`], or to the owner's abandonment of the attempt.
471    ///
472    /// Every step resumes from the joiner journal, so calling this again after
473    /// a crash picks up where the last durable step left off — a republished
474    /// artifact that is already at its slot is accepted as the same transfer,
475    /// and an awaited artifact is simply read again.
476    ///
477    /// The admitting side can give up until it approves the join, and every
478    /// wait below watches for that. After it approves there is nothing to watch
479    /// for: the approval is what grants this device storage access, and taking
480    /// it back is member removal and a key rotation, not a message.
481    pub(crate) async fn join_via_transport(
482        &self,
483        bundle: &DeviceJoinOfferBundle,
484        timing: DeviceJoinTransportTiming,
485        on_progress: coven_replication::sync::JoiningDeviceJoinProgressObserver,
486        cancel: &watch::Receiver<bool>,
487    ) -> Result<DeviceJoinTransportOutcome, BootstrapError> {
488        let storage = self.transport_storage().await?;
489        let transport = DeviceJoinTransport::open(&storage, bundle, DeviceJoinRole::Joiner)?;
490        self.drive_join_via_transport(&transport, bundle, timing, &on_progress, cancel)
491            .await
492    }
493
494    async fn drive_join_via_transport(
495        &self,
496        transport: &DeviceJoinTransport<'_>,
497        bundle: &DeviceJoinOfferBundle,
498        timing: DeviceJoinTransportTiming,
499        on_progress: &coven_replication::sync::JoiningDeviceJoinProgressObserver,
500        cancel: &watch::Receiver<bool>,
501    ) -> Result<DeviceJoinTransportOutcome, BootstrapError> {
502        let attempt_id = bundle.offer.attempt_id;
503        let mut published = Vec::new();
504
505        // A finished join leaves no journal row, so the library it produced is
506        // what says it finished. Without this the loop below would read "no
507        // record of this attempt" as "not started" and begin it again.
508        if let Some(config) = self.completed_library()? {
509            transport.delete_attempt_slots().await?;
510            return Ok(DeviceJoinTransportOutcome::Joined(config));
511        }
512
513        // Each pass takes the joiner journal's durable state and performs the
514        // one step that follows it — never an earlier step, which the journal
515        // refuses once it is past. A step that produced an artifact but died
516        // before publishing it republishes here; a step whose artifact is
517        // already at its slot publishes the same transfer again for nothing.
518        loop {
519            match self.device_join_status(attempt_id)? {
520                None | Some(DeviceJoinStatus::AwaitingAccessRequest { .. }) => {
521                    on_progress(
522                        coven_replication::sync::JoiningDeviceJoinProgress::RequestingProviderAccess,
523                    );
524                    let request = self
525                        .prepare_provider_access_request(bundle.offer.clone())
526                        .await?;
527                    publish_once(
528                        transport,
529                        &mut published,
530                        DeviceJoinAction::TransferProviderAccessRequest(request),
531                    )
532                    .await?;
533                }
534                Some(DeviceJoinStatus::AwaitingProviderAdmission { request }) => {
535                    let same_principal =
536                        request.offer.provider_admin.provider == request.peer_provider;
537                    publish_once(
538                        transport,
539                        &mut published,
540                        DeviceJoinAction::TransferProviderAccessRequest(request),
541                    )
542                    .await?;
543                    if same_principal {
544                        on_progress(
545                            coven_replication::sync::JoiningDeviceJoinProgress::WaitingForLibrary,
546                        );
547                        let join = match transport
548                            .await_step::<coven_replication::sync::SamePrincipalDeviceJoin>(timing)
549                            .await?
550                        {
551                            DeviceJoinStep::Continue(join) => join,
552                            DeviceJoinStep::Abandoned(abandonment) => {
553                                return self.accept_abandonment(transport, abandonment).await;
554                            }
555                        };
556                        self.record_same_principal_registration_request(
557                            join.bootstrap.bootstrap.request.approval().clone(),
558                        )?;
559                        return self
560                            .finish_same_principal(transport, join, on_progress, cancel)
561                            .await;
562                    } else {
563                        on_progress(
564                            coven_replication::sync::JoiningDeviceJoinProgress::WaitingForProviderAccess,
565                        );
566                        // The owner may give up on the attempt while this device
567                        // waits, so the wait watches the abandonment slot alongside
568                        // the approval rather than sitting out its deadline.
569                        let approval = match transport
570                            .await_step::<DeviceProviderAdmissionApproval>(timing)
571                            .await?
572                        {
573                            DeviceJoinStep::Continue(approval) => approval,
574                            DeviceJoinStep::Abandoned(abandonment) => {
575                                return self.accept_abandonment(transport, abandonment).await;
576                            }
577                        };
578                        on_progress(
579                            coven_replication::sync::JoiningDeviceJoinProgress::RegisteringDevice,
580                        );
581                        let registration_request =
582                            self.prepare_registration_request(approval).await?;
583                        publish_once(
584                            transport,
585                            &mut published,
586                            DeviceJoinAction::TransferRegistrationRequest(registration_request),
587                        )
588                        .await?;
589                    }
590                }
591                Some(DeviceJoinStatus::AwaitingRegistrationRequest { approval }) => {
592                    on_progress(
593                        coven_replication::sync::JoiningDeviceJoinProgress::RegisteringDevice,
594                    );
595                    let registration_request = self.prepare_registration_request(approval).await?;
596                    publish_once(
597                        transport,
598                        &mut published,
599                        DeviceJoinAction::TransferRegistrationRequest(registration_request),
600                    )
601                    .await?;
602                }
603                Some(DeviceJoinStatus::AwaitingBootstrap { request }) => {
604                    let same_principal = matches!(
605                        &request,
606                        coven_replication::sync::DeviceRegistrationRequest::SamePrincipal { .. }
607                    );
608                    publish_once(
609                        transport,
610                        &mut published,
611                        DeviceJoinAction::TransferRegistrationRequest(request),
612                    )
613                    .await?;
614                    on_progress(
615                        coven_replication::sync::JoiningDeviceJoinProgress::WaitingForLibrary,
616                    );
617                    if same_principal {
618                        let join = match transport
619                            .await_step::<coven_replication::sync::SamePrincipalDeviceJoin>(timing)
620                            .await?
621                        {
622                            DeviceJoinStep::Continue(join) => join,
623                            DeviceJoinStep::Abandoned(abandonment) => {
624                                return self.accept_abandonment(transport, abandonment).await;
625                            }
626                        };
627                        return self
628                            .finish_same_principal(transport, join, on_progress, cancel)
629                            .await;
630                    }
631                    let provider_ready = match transport
632                        .await_step::<ProviderReadyDeviceBootstrap>(timing)
633                        .await?
634                    {
635                        DeviceJoinStep::Continue(provider_ready) => provider_ready,
636                        DeviceJoinStep::Abandoned(abandonment) => {
637                            return self.accept_abandonment(transport, abandonment).await;
638                        }
639                    };
640                    let readiness = Box::pin(self.bootstrap_pending_device(
641                        provider_ready,
642                        on_progress,
643                        cancel,
644                    ))
645                    .await?;
646                    publish_once(
647                        transport,
648                        &mut published,
649                        DeviceJoinAction::TransferReadiness(readiness),
650                    )
651                    .await?;
652                }
653                Some(DeviceJoinStatus::AwaitingProviderCompletion { readiness }) => {
654                    if !matches!(
655                        &readiness.provider,
656                        coven_replication::sync::DeviceProviderReadiness::SamePrincipal
657                    ) {
658                        publish_once(
659                            transport,
660                            &mut published,
661                            DeviceJoinAction::TransferReadiness(readiness),
662                        )
663                        .await?;
664                    }
665                    on_progress(
666                        coven_replication::sync::JoiningDeviceJoinProgress::WaitingForActivation,
667                    );
668                    let activation = transport
669                        .await_artifact::<DeviceJoinActivation>(timing)
670                        .await?;
671                    return self.finish(transport, activation, on_progress).await;
672                }
673                Some(DeviceJoinStatus::AwaitingCompletion { activation }) => {
674                    return self.finish(transport, activation, on_progress).await;
675                }
676                Some(_) => {
677                    return Err(coven_replication::sync::DeviceJoinError::JournalConflict.into())
678                }
679            }
680        }
681    }
682
683    /// Save the store and clear the attempt's namespace: the join is complete,
684    /// so every artifact has been consumed and neither side has anything left
685    /// to read.
686    async fn finish(
687        &self,
688        transport: &DeviceJoinTransport<'_>,
689        activation: DeviceJoinActivation,
690        on_progress: &coven_replication::sync::JoiningDeviceJoinProgressObserver,
691    ) -> Result<DeviceJoinTransportOutcome, BootstrapError> {
692        let config = self.complete_device_join(activation, on_progress).await?;
693        transport.delete_attempt_slots().await?;
694        Ok(DeviceJoinTransportOutcome::Joined(config))
695    }
696
697    async fn finish_same_principal(
698        &self,
699        transport: &DeviceJoinTransport<'_>,
700        join: coven_replication::sync::SamePrincipalDeviceJoin,
701        on_progress: &coven_replication::sync::JoiningDeviceJoinProgressObserver,
702        cancel: &watch::Receiver<bool>,
703    ) -> Result<DeviceJoinTransportOutcome, BootstrapError> {
704        let config = self
705            .install_same_principal_device_join(join, on_progress, cancel)
706            .await?;
707        transport.delete_attempt_slots().await?;
708        Ok(DeviceJoinTransportOutcome::Joined(config))
709    }
710
711    /// Record the owner's abandonment and clear the attempt's namespace: the
712    /// abandonment is the last artifact either side publishes, and this device
713    /// has just read it.
714    async fn accept_abandonment(
715        &self,
716        transport: &DeviceJoinTransport<'_>,
717        abandonment: DeviceJoinAbandonment,
718    ) -> Result<DeviceJoinTransportOutcome, BootstrapError> {
719        let accepted = self.accept_device_join_abandonment(abandonment).await?;
720        transport.delete_attempt_slots().await?;
721        Ok(DeviceJoinTransportOutcome::Abandoned(accepted))
722    }
723}