Skip to main content

coven/
device_pairing.rs

1use coven_domain::joining::{
2    DevicePairingHost, DevicePairingOffer, DevicePairingRequest, DevicePairingTransportError,
3};
4use coven_protocol::membership::MemberRole;
5use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
6
7use crate::store_joining::StoreJoining;
8use crate::store_sync::{ConfigProvider, StoreSync};
9
10const DEVICE_PAIRING_PORT: u16 = 24_821;
11const DEVICE_PAIRING_LIFETIME: chrono::Duration = chrono::Duration::minutes(15);
12
13#[derive(Debug, thiserror::Error)]
14pub enum StartDevicePairingError {
15    #[error("this Store has no cloud provider")]
16    NoCloudProvider,
17    #[error("no active local network interface can receive the joining device")]
18    NoLocalInterface,
19    #[error("local interfaces: {0}")]
20    Interfaces(#[source] std::io::Error),
21    #[error("pairing listener: {0}")]
22    Listen(#[from] std::io::Error),
23    #[error("pairing offer: {0}")]
24    Pairing(#[from] coven_domain::joining::DevicePairingError),
25    #[error("pairing host: {0}")]
26    Host(#[from] DevicePairingTransportError),
27}
28
29#[derive(Debug, thiserror::Error)]
30pub enum ApproveDevicePairingError {
31    #[error("device pairing was cancelled")]
32    Cancelled,
33    #[error("pairing transport: {0}")]
34    Pairing(#[from] DevicePairingTransportError),
35    #[error("device invitation: {0}")]
36    Invitation(#[from] crate::DeviceAdmissionError),
37    #[error("persisted device invitation: {0}")]
38    PersistedInvitation(#[from] coven_domain::joining::BootstrapError),
39    #[error("device join: {0}")]
40    Join(#[from] crate::SyncError),
41    #[error("the invitation was created for another pairing request")]
42    RequestMismatch,
43}
44
45#[derive(Clone)]
46pub(crate) struct StoreDevicePairing {
47    config_provider: ConfigProvider,
48    journal_path: std::path::PathBuf,
49    clock: coven_foundation::clock::ClockRef,
50    joining: StoreJoining,
51    sync: StoreSync,
52}
53
54impl StoreDevicePairing {
55    pub(crate) fn new(
56        config_provider: ConfigProvider,
57        journal_path: std::path::PathBuf,
58        clock: coven_foundation::clock::ClockRef,
59        joining: StoreJoining,
60        sync: StoreSync,
61    ) -> Self {
62        Self {
63            config_provider,
64            journal_path,
65            clock,
66            joining,
67            sync,
68        }
69    }
70
71    /// Start the one pairing session this process can present on the LAN. The
72    /// returned code is the only value the joining device scans.
73    pub(crate) async fn start(&self) -> Result<DevicePairingHost, StartDevicePairingError> {
74        let config = (self.config_provider)();
75        let cloud_provider = config
76            .cloud_home
77            .provider
78            .clone()
79            .ok_or(StartDevicePairingError::NoCloudProvider)?;
80        let endpoints = local_pairing_endpoints(DEVICE_PAIRING_PORT)?;
81        let bind_address = if endpoints[0].is_ipv4() {
82            SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), DEVICE_PAIRING_PORT)
83        } else {
84            SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), DEVICE_PAIRING_PORT)
85        };
86        let listener = tokio::net::TcpListener::bind(bind_address).await?;
87        let pairing_key = coven_keys::keys::UserKeypair::generate();
88        let offer = DevicePairingOffer::new(
89            &pairing_key,
90            endpoints,
91            config.store_name,
92            cloud_provider,
93            (self.clock.now() + DEVICE_PAIRING_LIFETIME).timestamp(),
94        )?;
95        Ok(DevicePairingHost::start_or_resume(
96            listener,
97            offer,
98            pairing_key,
99            self.journal_path.clone(),
100            self.clock.clone(),
101        )
102        .await?)
103    }
104
105    /// Admit the exact signed request the owner reviewed, return its sealed
106    /// invitation over the local pairing session, and drive the Store
107    /// registration protocol to its terminal outcome.
108    pub(crate) async fn approve(
109        &self,
110        host: &DevicePairingHost,
111        request: &DevicePairingRequest,
112        role: MemberRole,
113        policy: crate::DeviceJoinApprovalPolicy<'_>,
114        access_administrator: Option<&dyn crate::DeviceProviderAccessAdministrator>,
115        on_progress: &(dyn Fn(crate::AdmittingDeviceJoinProgress) + Send + Sync),
116        cancel: tokio::sync::watch::Receiver<bool>,
117    ) -> Result<crate::DeviceJoinDriveOutcome, ApproveDevicePairingError> {
118        let timing = crate::DeviceJoinTransportTiming::interactive();
119        if let Some(bytes) = host.cancellation_invitation(request)? {
120            let invitation = coven_domain::joining::DeviceJoinInvite::from_bytes(&bytes)?;
121            self.sync
122                .abort_device_join_transport(&invitation.bundle)
123                .await?;
124            host.finish()?;
125            return Err(ApproveDevicePairingError::Cancelled);
126        }
127        on_progress(crate::AdmittingDeviceJoinProgress::PreparingInvitation);
128        let invitation = match host.invitation(request)? {
129            Some(bytes) => coven_domain::joining::DeviceJoinInvite::from_bytes(&bytes)?,
130            None => self.joining.begin_invite(request, role).await?,
131        };
132        if invitation.bundle.offer.member_pubkey != request.public_key() {
133            return Err(ApproveDevicePairingError::RequestMismatch);
134        }
135        host.deliver_invitation(request, invitation.to_bytes())?;
136        let drive = self.sync.drive_device_join(
137            &invitation.bundle,
138            policy,
139            access_administrator,
140            on_progress,
141            timing,
142        );
143        let cancellation = cancellation_requested(cancel);
144        tokio::pin!(drive);
145        tokio::pin!(cancellation);
146        let outcome = tokio::select! {
147            outcome = &mut drive => outcome?,
148            () = &mut cancellation => {
149                host.cancel()?;
150                self.sync
151                    .abort_device_join_transport(&invitation.bundle)
152                    .await?;
153                host.finish()?;
154                return Err(ApproveDevicePairingError::Cancelled);
155            }
156        };
157        host.finish()?;
158        Ok(outcome)
159    }
160
161    /// Persist cancellation, unwind the exact Store attempt retained by the
162    /// pairing journal, and close the local pairing session.
163    pub(crate) async fn cancel(
164        &self,
165        host: &DevicePairingHost,
166    ) -> Result<(), ApproveDevicePairingError> {
167        if let Some(bytes) = host.cancel()? {
168            let invitation = coven_domain::joining::DeviceJoinInvite::from_bytes(&bytes)?;
169            self.sync
170                .abort_device_join_transport(&invitation.bundle)
171                .await?;
172        }
173        host.finish()?;
174        Ok(())
175    }
176}
177
178async fn cancellation_requested(mut cancel: tokio::sync::watch::Receiver<bool>) {
179    while !*cancel.borrow() {
180        if cancel.changed().await.is_err() {
181            std::future::pending::<()>().await;
182        }
183    }
184}
185
186fn local_pairing_endpoints(port: u16) -> Result<Vec<SocketAddr>, StartDevicePairingError> {
187    let addresses = if_addrs::get_if_addrs()
188        .map_err(StartDevicePairingError::Interfaces)?
189        .into_iter()
190        .filter(|interface| interface.is_oper_up() && !interface.is_loopback())
191        .filter(|interface| !interface.is_link_local())
192        .map(|interface| interface.ip());
193    select_pairing_endpoints(addresses, port)
194}
195
196fn select_pairing_endpoints(
197    addresses: impl IntoIterator<Item = IpAddr>,
198    port: u16,
199) -> Result<Vec<SocketAddr>, StartDevicePairingError> {
200    let mut ipv4 = Vec::new();
201    let mut ipv6 = Vec::new();
202    for address in addresses {
203        match address {
204            IpAddr::V4(address) if !address.is_loopback() => ipv4.push(IpAddr::V4(address)),
205            IpAddr::V6(address) if !address.is_loopback() => ipv6.push(IpAddr::V6(address)),
206            _ => {}
207        }
208    }
209    let selected = if ipv4.is_empty() { ipv6 } else { ipv4 };
210    let mut endpoints: Vec<_> = selected
211        .into_iter()
212        .map(|address| SocketAddr::new(address, port))
213        .collect();
214    endpoints.sort();
215    endpoints.dedup();
216    if endpoints.is_empty() {
217        return Err(StartDevicePairingError::NoLocalInterface);
218    }
219    Ok(endpoints)
220}
221
222#[cfg(test)]
223mod tests {
224    use super::*;
225
226    #[test]
227    fn endpoint_discovery_never_advertises_loopback() {
228        match local_pairing_endpoints(DEVICE_PAIRING_PORT) {
229            Ok(endpoints) => assert!(endpoints
230                .iter()
231                .all(|endpoint| !endpoint.ip().is_loopback())),
232            Err(StartDevicePairingError::NoLocalInterface) => {}
233            Err(error) => panic!("interface discovery failed: {error}"),
234        }
235    }
236
237    #[test]
238    fn endpoint_selection_uses_one_listener_family_and_supports_ipv6_only_networks() {
239        let ipv4 = "192.0.2.4".parse().expect("IPv4 address");
240        let ipv6 = "2001:db8::4".parse().expect("IPv6 address");
241        assert_eq!(
242            select_pairing_endpoints([ipv6, ipv4], 7).expect("select IPv4 endpoints"),
243            vec!["192.0.2.4:7".parse().expect("IPv4 endpoint")]
244        );
245        assert_eq!(
246            select_pairing_endpoints([ipv6], 7).expect("select IPv6 endpoints"),
247            vec!["[2001:db8::4]:7".parse().expect("IPv6 endpoint")]
248        );
249    }
250}