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 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 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 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}