1use serde::{Deserialize, Serialize};
7use tracing::{debug, info};
8
9use super::peer_types::Address;
10use crate::{
11 Error, Result, Transaction,
12 auth::{Permission, crypto::PublicKey},
13 crdt::Doc,
14 entry::ID,
15 store::{StoreError, Table},
16};
17
18pub(super) const BOOTSTRAP_REQUESTS_SUBTREE: &str = "bootstrap_requests";
20
21pub(super) struct BootstrapRequestManager<'a> {
26 txn: &'a Transaction,
27}
28
29#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
31pub struct BootstrapRequest {
32 pub tree_id: ID,
34 pub requesting_pubkey: PublicKey,
36 pub requesting_key_name: String,
38 pub requested_permission: Permission,
40 pub timestamp: String,
42 pub status: RequestStatus,
44 pub peer_address: Address,
46 #[serde(default)]
49 pub metadata: Option<Doc>,
50}
51
52#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
54pub enum RequestStatus {
55 Pending,
57 Approved {
59 approved_by: String,
61 approval_time: String,
63 },
64 Rejected {
66 rejected_by: String,
68 rejection_time: String,
70 },
71}
72
73pub(super) fn request_id_for(
75 tree_id: &ID,
76 requesting_pubkey: &PublicKey,
77 requested_permission: &Permission,
78) -> String {
79 let identity = serde_ipld_dagcbor::to_vec(&(tree_id, requesting_pubkey, requested_permission))
80 .expect("bootstrap request identity is serializable");
81 ID::from_dagcbor_bytes(identity).to_string()
82}
83
84impl<'a> BootstrapRequestManager<'a> {
85 pub(super) fn new(txn: &'a Transaction) -> Self {
87 Self { txn }
88 }
89
90 pub(super) async fn store_request(&self, request: BootstrapRequest) -> Result<String> {
98 let requests = self
99 .txn
100 .get_store::<Table<BootstrapRequest>>(BOOTSTRAP_REQUESTS_SUBTREE)
101 .await?;
102
103 debug!(tree_id = %request.tree_id, "Storing bootstrap request");
104
105 let request_id = request_id_for(
106 &request.tree_id,
107 &request.requesting_pubkey,
108 &request.requested_permission,
109 );
110 requests.set(&request_id, request.clone()).await?;
111
112 info!(request_id = %request_id, tree_id = %request.tree_id, "Successfully stored bootstrap request");
113 Ok(request_id)
114 }
115
116 pub(super) async fn get_request(&self, request_id: &str) -> Result<Option<BootstrapRequest>> {
124 let requests = self
125 .txn
126 .get_store::<Table<BootstrapRequest>>(BOOTSTRAP_REQUESTS_SUBTREE)
127 .await?;
128
129 match requests.get(request_id).await {
130 Ok(request) => Ok(Some(request)),
131 Err(Error::Store(ref e)) if matches!(**e, StoreError::KeyNotFound { .. }) => Ok(None),
132 Err(e) => Err(e),
133 }
134 }
135
136 async fn filter_requests(
138 &self,
139 status_filter: &RequestStatus,
140 ) -> Result<Vec<(String, BootstrapRequest)>> {
141 let requests = self
142 .txn
143 .get_store::<Table<BootstrapRequest>>(BOOTSTRAP_REQUESTS_SUBTREE)
144 .await?;
145
146 let results = requests
147 .search(|request| {
148 std::mem::discriminant(status_filter) == std::mem::discriminant(&request.status)
149 })
150 .await?;
151
152 Ok(results)
153 }
154
155 pub(super) async fn pending_requests(&self) -> Result<Vec<(String, BootstrapRequest)>> {
160 self.filter_requests(&RequestStatus::Pending).await
161 }
162
163 pub(super) async fn approved_requests(&self) -> Result<Vec<(String, BootstrapRequest)>> {
168 self.filter_requests(&RequestStatus::Approved {
169 approved_by: String::new(),
170 approval_time: String::new(),
171 })
172 .await
173 }
174
175 pub(super) async fn rejected_requests(&self) -> Result<Vec<(String, BootstrapRequest)>> {
180 self.filter_requests(&RequestStatus::Rejected {
181 rejected_by: String::new(),
182 rejection_time: String::new(),
183 })
184 .await
185 }
186
187 pub(super) async fn find_existing_request(
189 &self,
190 tree_id: &ID,
191 requesting_pubkey: &PublicKey,
192 requested_permission: &Permission,
193 ) -> Result<Option<(String, BootstrapRequest)>> {
194 let requests = self
195 .txn
196 .get_store::<Table<BootstrapRequest>>(BOOTSTRAP_REQUESTS_SUBTREE)
197 .await?;
198 let matches = requests
199 .search(|request| {
200 &request.tree_id == tree_id
201 && &request.requesting_pubkey == requesting_pubkey
202 && &request.requested_permission == requested_permission
203 })
204 .await?;
205 Ok(matches
206 .iter()
207 .find(|(_, request)| matches!(request.status, RequestStatus::Rejected { .. }))
208 .cloned()
209 .or_else(|| {
210 matches
211 .iter()
212 .find(|(_, request)| matches!(request.status, RequestStatus::Pending))
213 .cloned()
214 })
215 .or_else(|| matches.into_iter().next()))
216 }
217
218 pub(super) async fn update_status(
227 &self,
228 request_id: &str,
229 new_status: RequestStatus,
230 ) -> Result<()> {
231 let requests = self
232 .txn
233 .get_store::<Table<BootstrapRequest>>(BOOTSTRAP_REQUESTS_SUBTREE)
234 .await?;
235
236 let mut request = requests.get(request_id).await?;
238
239 request.status = new_status;
241
242 requests.set(request_id, request).await?;
244
245 debug!(request_id = %request_id, "Updated bootstrap request status");
246 Ok(())
247 }
248}
249
250#[cfg(test)]
251mod tests {
252 use super::*;
253 use crate::{
254 Clock, Database, Instance, auth::types::Permission, backend::database::InMemory,
255 clock::FixedClock, crdt::Doc,
256 };
257 use std::sync::Arc;
258
259 async fn create_test_sync_tree() -> (Instance, Database, Arc<FixedClock>) {
260 let clock = Arc::new(FixedClock::default());
261 let (instance, mut user) = Instance::create_backend_with_clock(
262 Box::new(InMemory::new()),
263 clock.clone(),
264 crate::NewUser::passwordless("test"),
265 )
266 .await
267 .expect("Failed to create test instance");
268
269 let mut sync_settings = Doc::new();
270 sync_settings.set("name", "_sync");
271 sync_settings.set("type", "sync_settings");
272
273 let (database, _) = user
274 .new_database()
275 .settings(sync_settings)
276 .build()
277 .await
278 .unwrap();
279
280 (instance, database, clock)
281 }
282
283 fn create_test_request(clock: &FixedClock) -> BootstrapRequest {
284 BootstrapRequest {
285 tree_id: ID::from_bytes("test_tree_id"),
287 requesting_pubkey: PublicKey::random(),
288 requesting_key_name: "laptop_key".to_string(),
289 requested_permission: Permission::Write(5),
290 timestamp: clock.now_rfc3339(),
291 status: RequestStatus::Pending,
292 peer_address: Address {
293 transport_type: "http".to_string(),
294 address: "127.0.0.1:8080".to_string(),
295 },
296 metadata: None,
297 }
298 }
299
300 #[tokio::test]
301 async fn test_store_and_get_request() {
302 let (_instance, sync_tree, clock) = create_test_sync_tree().await;
303 let txn = sync_tree.new_transaction().await.unwrap();
304 let manager = BootstrapRequestManager::new(&txn);
305
306 let request = create_test_request(&clock);
307
308 let request_id = manager.store_request(request.clone()).await.unwrap();
310
311 let retrieved = manager.get_request(&request_id).await.unwrap().unwrap();
313 assert_eq!(retrieved.tree_id, request.tree_id);
314 assert_eq!(retrieved.requesting_pubkey, request.requesting_pubkey);
315 assert_eq!(retrieved.requesting_key_name, request.requesting_key_name);
316 assert_eq!(retrieved.requested_permission, request.requested_permission);
317 assert_eq!(retrieved.status, request.status);
318 assert_eq!(retrieved.peer_address, request.peer_address);
319 }
320
321 #[tokio::test]
322 async fn test_list_requests() {
323 let (_instance, sync_tree, clock) = create_test_sync_tree().await;
324 let txn = sync_tree.new_transaction().await.unwrap();
325 let manager = BootstrapRequestManager::new(&txn);
326
327 let request1 = create_test_request(&clock);
329
330 let mut request2 = create_test_request(&clock);
331 request2.status = RequestStatus::Approved {
332 approved_by: "admin".to_string(),
333 approval_time: clock.now_rfc3339(),
334 };
335
336 manager.store_request(request1).await.unwrap();
337 manager.store_request(request2).await.unwrap();
338
339 let pending_requests = manager.pending_requests().await.unwrap();
341 assert_eq!(pending_requests.len(), 1);
342
343 let approved_requests = manager.approved_requests().await.unwrap();
345 assert_eq!(approved_requests.len(), 1);
346
347 assert!(matches!(
349 pending_requests[0].1.status,
350 RequestStatus::Pending
351 ));
352 assert!(matches!(
353 approved_requests[0].1.status,
354 RequestStatus::Approved { .. }
355 ));
356 }
357
358 #[tokio::test]
359 async fn test_update_status() {
360 let (_instance, sync_tree, clock) = create_test_sync_tree().await;
361 let txn = sync_tree.new_transaction().await.unwrap();
362 let manager = BootstrapRequestManager::new(&txn);
363
364 let request = create_test_request(&clock);
365
366 let request_id = manager.store_request(request).await.unwrap();
368
369 let new_status = RequestStatus::Approved {
371 approved_by: "admin".to_string(),
372 approval_time: clock.now_rfc3339(),
373 };
374 manager
375 .update_status(&request_id, new_status.clone())
376 .await
377 .unwrap();
378
379 let updated_request = manager.get_request(&request_id).await.unwrap().unwrap();
381 assert_eq!(updated_request.status, new_status);
382 }
383
384 #[test]
385 fn request_identity_is_deterministic_and_distinguishes_semantic_fields() {
386 let tree = ID::from_bytes("test_tree_id");
387 let other_tree = ID::from_bytes("other_tree_id");
388 let key = PublicKey::random();
389 let other_key = PublicKey::random();
390 let request_id = request_id_for(&tree, &key, &Permission::Write(5));
391 let storage_id = ID::parse(&request_id).expect("request ID uses the project ID encoding");
392
393 assert_eq!(storage_id.as_cid().unwrap().codec(), 0x71);
394 assert_eq!(storage_id.hash_code(), Some(0x1e));
395 assert_eq!(
396 request_id,
397 request_id_for(&tree, &key, &Permission::Write(5))
398 );
399 assert_ne!(
400 request_id,
401 request_id_for(&other_tree, &key, &Permission::Write(5))
402 );
403 assert_ne!(
404 request_id,
405 request_id_for(&tree, &other_key, &Permission::Write(5))
406 );
407 assert_ne!(
408 request_id,
409 request_id_for(&tree, &key, &Permission::Admin(5))
410 );
411 assert_ne!(
412 request_id,
413 request_id_for(&tree, &key, &Permission::Write(6))
414 );
415 assert_ne!(request_id, request_id_for(&tree, &key, &Permission::Read));
416 }
417
418 #[test]
419 fn request_identity_ignores_mutable_record_fields() {
420 let clock = FixedClock::default();
421 let request = create_test_request(&clock);
422 let request_id = request_id_for(
423 &request.tree_id,
424 &request.requesting_pubkey,
425 &request.requested_permission,
426 );
427 let mut changed = request;
428 changed.requesting_key_name = "renamed key".to_string();
429 changed.timestamp = "2026-09-16T12:34:56Z".to_string();
430 changed.status = RequestStatus::Rejected {
431 rejected_by: "admin".to_string(),
432 rejection_time: "2026-09-16T12:35:00Z".to_string(),
433 };
434 changed.peer_address = Address {
435 transport_type: "iroh".to_string(),
436 address: "new-address".to_string(),
437 };
438 let mut metadata = Doc::new();
439 metadata.set("note", "changed");
440 changed.metadata = Some(metadata);
441
442 assert_eq!(
443 request_id,
444 request_id_for(
445 &changed.tree_id,
446 &changed.requesting_pubkey,
447 &changed.requested_permission,
448 )
449 );
450 }
451
452 #[tokio::test]
453 async fn concurrent_writers_converge_on_one_request() {
454 let (_instance, sync_tree, clock) = create_test_sync_tree().await;
455 let request = create_test_request(&clock);
456 let first = sync_tree.new_transaction().await.unwrap();
457 let second = sync_tree.new_transaction().await.unwrap();
458 let first_id = BootstrapRequestManager::new(&first)
459 .store_request(request.clone())
460 .await
461 .unwrap();
462 let mut later = request;
463 later.timestamp = "2026-09-15T23:00:01Z".to_string();
464 let second_id = BootstrapRequestManager::new(&second)
465 .store_request(later)
466 .await
467 .unwrap();
468 first.commit().await.unwrap();
469 second.commit().await.unwrap();
470 assert_eq!(first_id, second_id);
471 let txn = sync_tree.new_transaction().await.unwrap();
472 let pending = BootstrapRequestManager::new(&txn)
473 .pending_requests()
474 .await
475 .unwrap();
476 assert_eq!(pending.len(), 1, "racing writers left {pending:#?}");
477 }
478
479 #[tokio::test]
480 async fn test_get_nonexistent_request() {
481 let (_instance, sync_tree, _clock) = create_test_sync_tree().await;
482 let txn = sync_tree.new_transaction().await.unwrap();
483 let manager = BootstrapRequestManager::new(&txn);
484
485 let result = manager.get_request("nonexistent").await.unwrap();
486 assert!(result.is_none());
487 }
488}