1use crate::{
2 AuthorizedDevice, DevicePubkey, DeviceRoster, DomainError, Error, Invite, InviteResponse,
3 InviteResponseEnvelope, MessageEnvelope, OwnerPubkey, ProtocolContext, Result,
4 RosterSnapshotDecision, Session, SessionState, UnixSeconds,
5};
6use rand::{CryptoRng, RngCore};
7use serde::{Deserialize, Serialize};
8use std::collections::{BTreeMap, BTreeSet};
9
10const MAX_INACTIVE_SESSIONS: usize = 10;
11#[derive(Debug, Clone)]
12pub struct SessionManager {
13 local_owner_pubkey: OwnerPubkey,
14 local_device_pubkey: DevicePubkey,
15 local_device_secret_key: [u8; 32],
16 local_invite: Option<Invite>,
17 users: BTreeMap<OwnerPubkey, UserRecord>,
18}
19
20#[derive(Debug, Clone)]
21struct UserRecord {
22 owner_pubkey: OwnerPubkey,
23 roster: Option<DeviceRoster>,
24 devices: BTreeMap<DevicePubkey, DeviceRecord>,
25}
26
27#[derive(Debug, Clone)]
28struct DeviceRecord {
29 device_pubkey: DevicePubkey,
30 authorized: bool,
31 is_stale: bool,
32 stale_since: Option<UnixSeconds>,
33 claimed_owner_pubkey: Option<OwnerPubkey>,
34 public_invite: Option<Invite>,
35 invite_response_generated: bool,
36 active_session: Option<Session>,
37 inactive_sessions: Vec<Session>,
38 last_activity: Option<UnixSeconds>,
39 created_at: UnixSeconds,
40}
41
42#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
43pub struct SessionManagerSnapshot {
44 pub local_owner_pubkey: OwnerPubkey,
45 pub local_device_pubkey: DevicePubkey,
46 pub local_invite: Option<Invite>,
47 pub users: Vec<UserRecordSnapshot>,
48}
49
50#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
51pub struct UserRecordSnapshot {
52 pub owner_pubkey: OwnerPubkey,
53 pub roster: Option<DeviceRoster>,
54 pub devices: Vec<DeviceRecordSnapshot>,
55}
56
57#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
58pub struct DeviceRecordSnapshot {
59 pub device_pubkey: DevicePubkey,
60 pub authorized: bool,
61 pub is_stale: bool,
62 pub stale_since: Option<UnixSeconds>,
63 #[serde(default, skip_serializing_if = "Option::is_none")]
64 pub claimed_owner_pubkey: Option<OwnerPubkey>,
65 pub public_invite: Option<Invite>,
66 #[serde(default)]
67 pub invite_response_generated: bool,
68 pub active_session: Option<SessionState>,
69 pub inactive_sessions: Vec<SessionState>,
70 pub last_activity: Option<UnixSeconds>,
71 pub created_at: UnixSeconds,
72}
73
74#[derive(Debug, Clone, PartialEq, Eq)]
75pub struct PreparedSend {
76 pub recipient_owner: OwnerPubkey,
77 pub payload: Vec<u8>,
78 pub deliveries: Vec<Delivery>,
79 pub invite_responses: Vec<InviteResponseEnvelope>,
80 pub relay_gaps: Vec<RelayGap>,
81}
82
83#[derive(Debug, Clone, PartialEq, Eq)]
84pub struct Delivery {
85 pub owner_pubkey: OwnerPubkey,
86 pub device_pubkey: DevicePubkey,
87 pub envelope: MessageEnvelope,
88}
89
90#[derive(Debug, Clone, PartialEq, Eq)]
91pub struct ProcessedInviteResponse {
92 pub owner_pubkey: OwnerPubkey,
93 pub device_pubkey: DevicePubkey,
94 pub claimed_owner_pubkey: Option<OwnerPubkey>,
95}
96
97#[derive(Debug, Clone, PartialEq, Eq)]
98pub struct ReceivedMessage {
99 pub owner_pubkey: OwnerPubkey,
100 pub device_pubkey: DevicePubkey,
101 pub payload: Vec<u8>,
102}
103
104#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord)]
105pub enum RelayGap {
106 MissingRoster {
107 owner_pubkey: OwnerPubkey,
108 },
109 MissingDeviceInvite {
110 owner_pubkey: OwnerPubkey,
111 device_pubkey: DevicePubkey,
112 },
113}
114
115#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
116pub struct PruneReport {
117 pub removed_devices: Vec<(OwnerPubkey, DevicePubkey)>,
118 pub removed_users: Vec<OwnerPubkey>,
119}
120
121#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
122struct TargetDevice {
123 owner_pubkey: OwnerPubkey,
124 device_pubkey: DevicePubkey,
125}
126
127#[derive(Debug, Clone, Copy, PartialEq, Eq)]
128enum SendSessionSource {
129 Active,
130 Inactive(usize),
131}
132
133impl SessionManager {
134 pub fn new(local_owner_pubkey: OwnerPubkey, local_device_secret_key: [u8; 32]) -> Self {
135 let local_device_pubkey = crate::device_pubkey_from_secret_bytes(&local_device_secret_key)
136 .expect("local device secret key must derive a valid device public key");
137
138 Self {
139 local_owner_pubkey,
140 local_device_pubkey,
141 local_device_secret_key,
142 local_invite: None,
143 users: BTreeMap::new(),
144 }
145 }
146
147 pub fn from_snapshot(
148 snapshot: SessionManagerSnapshot,
149 local_device_secret_key: [u8; 32],
150 ) -> Result<Self> {
151 let derived_local_device_pubkey =
152 crate::device_pubkey_from_secret_bytes(&local_device_secret_key)?;
153 if derived_local_device_pubkey != snapshot.local_device_pubkey {
154 return Err(DomainError::InvalidState(
155 "snapshot local device pubkey does not match provided secret key".to_string(),
156 )
157 .into());
158 }
159
160 let users = snapshot
161 .users
162 .into_iter()
163 .map(UserRecord::from_snapshot)
164 .map(|record| (record.owner_pubkey, record))
165 .collect();
166
167 Ok(Self {
168 local_owner_pubkey: snapshot.local_owner_pubkey,
169 local_device_pubkey: snapshot.local_device_pubkey,
170 local_device_secret_key,
171 local_invite: snapshot.local_invite,
172 users,
173 })
174 }
175
176 pub fn snapshot(&self) -> SessionManagerSnapshot {
177 SessionManagerSnapshot {
178 local_owner_pubkey: self.local_owner_pubkey,
179 local_device_pubkey: self.local_device_pubkey,
180 local_invite: self.local_invite.clone(),
181 users: self.users.values().map(UserRecord::snapshot).collect(),
182 }
183 }
184
185 pub fn local_device_pubkey(&self) -> DevicePubkey {
186 self.local_device_pubkey
187 }
188
189 pub fn replace_local_invite(&mut self, invite: Invite) {
190 self.local_invite = Some(invite);
191 }
192
193 pub fn ensure_local_invite<R>(&mut self, ctx: &mut ProtocolContext<'_, R>) -> Result<&Invite>
194 where
195 R: RngCore + CryptoRng,
196 {
197 if self.local_invite.is_none() {
198 let invite = Invite::create_new_with_context(
199 ctx,
200 self.local_device_pubkey,
201 Some(self.local_owner_pubkey),
202 None,
203 )?;
204 self.observe_public_invite(self.local_owner_pubkey, invite.clone())?;
205 self.local_invite = Some(invite);
206 }
207
208 Ok(self.local_invite.as_ref().expect("local invite must exist"))
209 }
210
211 pub fn apply_local_roster(&mut self, roster: DeviceRoster) -> RosterSnapshotDecision {
212 self.apply_roster_for_owner(self.local_owner_pubkey, roster)
213 }
214
215 pub fn replace_local_roster(&mut self, roster: DeviceRoster) -> RosterSnapshotDecision {
216 self.apply_roster_for_owner_inner(self.local_owner_pubkey, roster, true)
217 }
218
219 pub fn observe_peer_roster(
220 &mut self,
221 owner_pubkey: OwnerPubkey,
222 roster: DeviceRoster,
223 ) -> RosterSnapshotDecision {
224 self.apply_roster_for_owner(owner_pubkey, roster)
225 }
226
227 pub fn observe_device_invite(
228 &mut self,
229 owner_pubkey: OwnerPubkey,
230 invite: Invite,
231 ) -> Result<()> {
232 self.observe_public_invite(owner_pubkey, invite)
233 }
234
235 pub fn observe_invite_response<R>(
236 &mut self,
237 ctx: &mut ProtocolContext<'_, R>,
238 envelope: &InviteResponseEnvelope,
239 ) -> Result<Option<ProcessedInviteResponse>>
240 where
241 R: RngCore + CryptoRng,
242 {
243 let Some(invite) = self.local_invite.clone() else {
244 return Ok(None);
245 };
246
247 let mut owned_invite = invite;
248 let InviteResponse {
249 session,
250 invitee_device_pubkey,
251 invitee_owner_pubkey,
252 ..
253 } = owned_invite.process_response(ctx, envelope, self.local_device_secret_key)?;
254
255 self.local_invite = Some(owned_invite);
256
257 let device_owner_pubkey = crate::owner_pubkey_from_device_pubkey(invitee_device_pubkey);
258 let invitee_owner_pubkey = invitee_owner_pubkey.ok_or_else(|| {
259 DomainError::InvalidState("invite response missing owner claim".to_string())
260 })?;
261 let claimed_owner_pubkey =
262 (invitee_owner_pubkey != device_owner_pubkey).then_some(invitee_owner_pubkey);
263 let owner_pubkey = claimed_owner_pubkey
264 .filter(|claimed_owner_pubkey| {
265 self.users
266 .get(claimed_owner_pubkey)
267 .and_then(|user| user.roster.as_ref())
268 .and_then(|roster| roster.get_device(&invitee_device_pubkey))
269 .is_some()
270 })
271 .unwrap_or(device_owner_pubkey);
272 let should_seed_single_device_roster = owner_pubkey == device_owner_pubkey
273 && self
274 .users
275 .get(&owner_pubkey)
276 .and_then(|user| user.roster.as_ref())
277 .is_none();
278 let user = self.user_record_mut(owner_pubkey);
279 if should_seed_single_device_roster {
280 user.roster = Some(DeviceRoster::new(
281 ctx.now,
282 vec![AuthorizedDevice::new(invitee_device_pubkey, ctx.now)],
283 ));
284 }
285 let record = user.device_record_mut(invitee_device_pubkey, ctx.now);
286 record.claimed_owner_pubkey = claimed_owner_pubkey
287 .filter(|claimed_owner_pubkey| *claimed_owner_pubkey != owner_pubkey);
288 if should_seed_single_device_roster {
289 record.authorized = true;
290 record.is_stale = false;
291 }
292 record.invite_response_generated = true;
293 record.upsert_session(session, ctx.now);
294
295 Ok(Some(ProcessedInviteResponse {
296 owner_pubkey,
297 device_pubkey: invitee_device_pubkey,
298 claimed_owner_pubkey: record.claimed_owner_pubkey,
299 }))
300 }
301
302 pub fn prepare_send<R>(
303 &mut self,
304 ctx: &mut ProtocolContext<'_, R>,
305 recipient_owner: OwnerPubkey,
306 payload: Vec<u8>,
307 ) -> Result<PreparedSend>
308 where
309 R: RngCore + CryptoRng,
310 {
311 self.prepare_send_inner(ctx, recipient_owner, payload, true)
312 }
313
314 pub fn prepare_remote_send<R>(
321 &mut self,
322 ctx: &mut ProtocolContext<'_, R>,
323 recipient_owner: OwnerPubkey,
324 payload: Vec<u8>,
325 ) -> Result<PreparedSend>
326 where
327 R: RngCore + CryptoRng,
328 {
329 self.prepare_send_inner(ctx, recipient_owner, payload, false)
330 }
331
332 pub fn prepare_remote_send_to_devices<R>(
333 &mut self,
334 ctx: &mut ProtocolContext<'_, R>,
335 recipient_owner: OwnerPubkey,
336 device_pubkeys: impl IntoIterator<Item = DevicePubkey>,
337 payload: Vec<u8>,
338 ) -> Result<PreparedSend>
339 where
340 R: RngCore + CryptoRng,
341 {
342 let targets = device_pubkeys
343 .into_iter()
344 .map(|device_pubkey| TargetDevice {
345 owner_pubkey: recipient_owner,
346 device_pubkey,
347 })
348 .collect();
349 self.prepare_explicit_send(ctx, recipient_owner, targets, payload, false)
350 }
351
352 pub fn prepare_local_sibling_send<R>(
353 &mut self,
354 ctx: &mut ProtocolContext<'_, R>,
355 payload: Vec<u8>,
356 ) -> Result<PreparedSend>
357 where
358 R: RngCore + CryptoRng,
359 {
360 self.prepare_local_sibling_send_inner(ctx, payload, false)
361 }
362
363 pub fn prepare_local_sibling_send_to_devices<R>(
364 &mut self,
365 ctx: &mut ProtocolContext<'_, R>,
366 device_pubkeys: impl IntoIterator<Item = DevicePubkey>,
367 payload: Vec<u8>,
368 ) -> Result<PreparedSend>
369 where
370 R: RngCore + CryptoRng,
371 {
372 let targets = device_pubkeys
373 .into_iter()
374 .map(|device_pubkey| TargetDevice {
375 owner_pubkey: self.local_owner_pubkey,
376 device_pubkey,
377 })
378 .collect();
379 self.prepare_explicit_send(ctx, self.local_owner_pubkey, targets, payload, false)
380 }
381
382 pub fn prepare_local_sibling_send_reusing_sessions<R>(
383 &mut self,
384 ctx: &mut ProtocolContext<'_, R>,
385 payload: Vec<u8>,
386 ) -> Result<PreparedSend>
387 where
388 R: RngCore + CryptoRng,
389 {
390 self.prepare_local_sibling_send_inner(ctx, payload, false)
391 }
392
393 pub fn prepare_local_sibling_send_refreshing_one_way_sessions<R>(
394 &mut self,
395 ctx: &mut ProtocolContext<'_, R>,
396 payload: Vec<u8>,
397 ) -> Result<PreparedSend>
398 where
399 R: RngCore + CryptoRng,
400 {
401 self.prepare_local_sibling_send_inner(ctx, payload, true)
402 }
403
404 pub fn prepare_local_sibling_send_reusing_all_sessions<R>(
405 &mut self,
406 ctx: &mut ProtocolContext<'_, R>,
407 payload: Vec<u8>,
408 ) -> Result<PreparedSend>
409 where
410 R: RngCore + CryptoRng,
411 {
412 let mut targets = BTreeSet::new();
413 self.collect_local_sibling_targets(&mut targets);
414
415 let mut deliveries = Vec::new();
416 let mut invite_responses = Vec::new();
417 let mut relay_gaps = Vec::new();
418
419 for target in targets {
420 let mut target_deliveries = self.prepare_device_deliveries_for_all_send_sessions(
421 ctx,
422 target.owner_pubkey,
423 target.device_pubkey,
424 &payload,
425 )?;
426 if target_deliveries.is_empty() {
427 match self.prepare_device_delivery(
428 ctx,
429 target.owner_pubkey,
430 target.device_pubkey,
431 &payload,
432 false,
433 )? {
434 Some((delivery, maybe_response)) => {
435 target_deliveries.push(delivery);
436 if let Some(response) = maybe_response {
437 invite_responses.push(response);
438 }
439 }
440 None => {
441 relay_gaps.push(RelayGap::MissingDeviceInvite {
442 owner_pubkey: target.owner_pubkey,
443 device_pubkey: target.device_pubkey,
444 });
445 }
446 }
447 }
448 deliveries.extend(target_deliveries);
449 }
450
451 relay_gaps.sort();
452
453 Ok(PreparedSend {
454 recipient_owner: self.local_owner_pubkey,
455 payload,
456 deliveries,
457 invite_responses,
458 relay_gaps,
459 })
460 }
461
462 fn prepare_local_sibling_send_inner<R>(
463 &mut self,
464 ctx: &mut ProtocolContext<'_, R>,
465 payload: Vec<u8>,
466 refresh_one_way_bootstrap: bool,
467 ) -> Result<PreparedSend>
468 where
469 R: RngCore + CryptoRng,
470 {
471 let mut targets = BTreeSet::new();
472 self.collect_local_sibling_targets(&mut targets);
473
474 let mut deliveries = Vec::new();
475 let mut invite_responses = Vec::new();
476 let mut relay_gaps = Vec::new();
477
478 for target in targets {
479 match self.prepare_device_delivery(
480 ctx,
481 target.owner_pubkey,
482 target.device_pubkey,
483 &payload,
484 refresh_one_way_bootstrap,
485 )? {
486 Some((delivery, maybe_response)) => {
487 deliveries.push(delivery);
488 if let Some(response) = maybe_response {
489 invite_responses.push(response);
490 }
491 }
492 None => {
493 relay_gaps.push(RelayGap::MissingDeviceInvite {
494 owner_pubkey: target.owner_pubkey,
495 device_pubkey: target.device_pubkey,
496 });
497 }
498 }
499 }
500
501 relay_gaps.sort();
502
503 Ok(PreparedSend {
504 recipient_owner: self.local_owner_pubkey,
505 payload,
506 deliveries,
507 invite_responses,
508 relay_gaps,
509 })
510 }
511
512 pub(crate) fn has_authorized_local_siblings(&self) -> bool {
513 let Some(user) = self.users.get(&self.local_owner_pubkey) else {
514 return false;
515 };
516 if user.roster.is_none() {
517 return false;
518 }
519 user.authorized_non_stale_devices()
520 .into_iter()
521 .any(|device_pubkey| device_pubkey != self.local_device_pubkey)
522 }
523
524 pub fn receive<R>(
525 &mut self,
526 ctx: &mut ProtocolContext<'_, R>,
527 sender_owner: OwnerPubkey,
528 envelope: &MessageEnvelope,
529 ) -> Result<Option<ReceivedMessage>>
530 where
531 R: RngCore + CryptoRng,
532 {
533 let Some(user) = self.users.get_mut(&sender_owner) else {
534 return Ok(None);
535 };
536
537 let device_pubkeys: Vec<DevicePubkey> = user.devices.keys().copied().collect();
538 for device_pubkey in device_pubkeys {
539 let record = user
540 .devices
541 .get_mut(&device_pubkey)
542 .expect("device key collected from map");
543
544 if let Some(active_session) = record.active_session.as_ref() {
545 if active_session.matches_sender(envelope.sender) {
546 let plan = active_session.plan_receive(ctx, envelope)?;
547 let outcome = record
548 .active_session
549 .as_mut()
550 .expect("active session must still exist")
551 .apply_receive(plan);
552 record.last_activity = Some(ctx.now);
553 return Ok(Some(ReceivedMessage {
554 owner_pubkey: sender_owner,
555 device_pubkey,
556 payload: outcome.payload,
557 }));
558 }
559 }
560
561 let mut matched_inactive = None;
562 for (index, session) in record.inactive_sessions.iter().enumerate() {
563 if !session.matches_sender(envelope.sender) {
564 continue;
565 }
566 let plan = session.plan_receive(ctx, envelope)?;
567 matched_inactive = Some((index, plan));
568 break;
569 }
570
571 if let Some((index, plan)) = matched_inactive {
572 let mut session = record.inactive_sessions.remove(index);
573 let outcome = session.apply_receive(plan);
574 record.promote_inactive_session(session);
575 record.last_activity = Some(ctx.now);
576 return Ok(Some(ReceivedMessage {
577 owner_pubkey: sender_owner,
578 device_pubkey,
579 payload: outcome.payload,
580 }));
581 }
582 }
583
584 Ok(None)
585 }
586
587 pub fn prune_stale(&mut self, _now: UnixSeconds) -> PruneReport {
588 let mut removed_devices = Vec::new();
589 let mut removed_users = Vec::new();
590
591 self.users.retain(|owner_pubkey, user| {
592 user.devices.retain(|device_pubkey, record| {
593 let keep = !record.is_stale;
594 if !keep {
595 removed_devices.push((*owner_pubkey, *device_pubkey));
596 }
597 keep
598 });
599
600 let keep_user = !user.devices.is_empty() || user.roster.is_some();
601 if !keep_user {
602 removed_users.push(*owner_pubkey);
603 }
604 keep_user
605 });
606
607 removed_devices.sort();
608 removed_users.sort();
609
610 PruneReport {
611 removed_devices,
612 removed_users,
613 }
614 }
615
616 pub fn delete_user(&mut self, owner_pubkey: OwnerPubkey) {
617 if owner_pubkey != self.local_owner_pubkey {
618 self.users.remove(&owner_pubkey);
619 }
620 }
621
622 pub fn import_session_state(
623 &mut self,
624 owner_pubkey: OwnerPubkey,
625 device_pubkey: DevicePubkey,
626 state: SessionState,
627 now: UnixSeconds,
628 ) {
629 let user = self.user_record_mut(owner_pubkey);
630 let record = user.device_record_mut(device_pubkey, now);
631 record.authorized = true;
632 record.is_stale = false;
633 record.invite_response_generated = true;
634 record.upsert_session(Session::from_state(state), now);
635 }
636
637 fn prepare_device_delivery<R>(
638 &mut self,
639 ctx: &mut ProtocolContext<'_, R>,
640 owner_pubkey: OwnerPubkey,
641 device_pubkey: DevicePubkey,
642 payload: &[u8],
643 refresh_one_way_bootstrap: bool,
644 ) -> Result<Option<(Delivery, Option<InviteResponseEnvelope>)>>
645 where
646 R: RngCore + CryptoRng,
647 {
648 let claimed_owner = Some(self.local_owner_pubkey);
649 let local_owner_pubkey = self.local_owner_pubkey;
650 let local_device_pubkey = self.local_device_pubkey;
651 let local_device_secret_key = self.local_device_secret_key;
652 let user = self.user_record_mut(owner_pubkey);
653 let record = user.device_record_mut(device_pubkey, ctx.now);
654
655 if !record.authorized || record.is_stale {
656 return Ok(None);
657 }
658
659 let source = record.best_send_session_source();
660 let should_refresh_local_sibling_bootstrap = refresh_one_way_bootstrap
661 && owner_pubkey == local_owner_pubkey
662 && device_pubkey != local_device_pubkey
663 && record.public_invite.is_some()
664 && source
665 .as_ref()
666 .and_then(|source| record.session_for_send_source(source))
667 .is_some_and(is_one_way_bootstrap_session);
668
669 if should_refresh_local_sibling_bootstrap {
670 let public_invite = record
671 .public_invite
672 .clone()
673 .expect("checked public invite presence");
674 match public_invite.accept_with_owner_context(
675 ctx,
676 local_device_pubkey,
677 local_device_secret_key,
678 claimed_owner,
679 ) {
680 Ok((mut session, invite_response)) => {
681 let mut envelope = session
682 .apply_send(session.plan_send(payload, ctx.now)?)
683 .envelope;
684 envelope.recipient = Some(device_pubkey);
685 record.invite_response_generated = true;
686 record.upsert_session(session, ctx.now);
687
688 return Ok(Some((
689 Delivery {
690 owner_pubkey,
691 device_pubkey,
692 envelope,
693 },
694 Some(invite_response),
695 )));
696 }
697 Err(Error::Domain(
698 DomainError::InviteAlreadyUsed | DomainError::InviteExhausted,
699 )) => {}
700 Err(error) => return Err(error),
701 }
702 }
703
704 if let Some(source) = source {
705 let plan = match source {
706 SendSessionSource::Active => record
707 .active_session
708 .as_ref()
709 .expect("active session must exist")
710 .plan_send(payload, ctx.now)?,
711 SendSessionSource::Inactive(index) => {
712 record.inactive_sessions[index].plan_send(payload, ctx.now)?
713 }
714 };
715
716 let mut envelope = match source {
717 SendSessionSource::Active => {
718 record
719 .active_session
720 .as_mut()
721 .expect("active session must exist")
722 .apply_send(plan)
723 .envelope
724 }
725 SendSessionSource::Inactive(index) => {
726 let mut session = record.inactive_sessions.remove(index);
727 let outcome = session.apply_send(plan);
728 record.upsert_session(session, ctx.now);
729 outcome.envelope
730 }
731 };
732 envelope.recipient = Some(device_pubkey);
733
734 record.last_activity = Some(ctx.now);
735 return Ok(Some((
736 Delivery {
737 owner_pubkey,
738 device_pubkey,
739 envelope,
740 },
741 None,
742 )));
743 }
744
745 let Some(public_invite) = record.public_invite.clone() else {
746 return Ok(None);
747 };
748
749 let (mut session, invite_response) = match public_invite.accept_with_owner_context(
750 ctx,
751 local_device_pubkey,
752 local_device_secret_key,
753 claimed_owner,
754 ) {
755 Ok(result) => result,
756 Err(Error::Domain(DomainError::InviteAlreadyUsed | DomainError::InviteExhausted)) => {
757 return Ok(None)
758 }
759 Err(error) => return Err(error),
760 };
761 let mut envelope = session
762 .apply_send(session.plan_send(payload, ctx.now)?)
763 .envelope;
764 envelope.recipient = Some(device_pubkey);
765 record.invite_response_generated = true;
766 record.upsert_session(session, ctx.now);
767
768 Ok(Some((
769 Delivery {
770 owner_pubkey,
771 device_pubkey,
772 envelope,
773 },
774 Some(invite_response),
775 )))
776 }
777
778 fn prepare_device_deliveries_for_all_send_sessions<R>(
779 &mut self,
780 ctx: &mut ProtocolContext<'_, R>,
781 owner_pubkey: OwnerPubkey,
782 device_pubkey: DevicePubkey,
783 payload: &[u8],
784 ) -> Result<Vec<Delivery>>
785 where
786 R: RngCore + CryptoRng,
787 {
788 let user = self.user_record_mut(owner_pubkey);
789 let record = user.device_record_mut(device_pubkey, ctx.now);
790
791 if !record.authorized || record.is_stale {
792 return Ok(Vec::new());
793 }
794
795 let mut deliveries = Vec::new();
796
797 if let Some(active_session) = record.active_session.as_mut() {
798 if active_session.can_send() {
799 let plan = active_session.plan_send(payload, ctx.now)?;
800 let mut envelope = active_session.apply_send(plan).envelope;
801 envelope.recipient = Some(device_pubkey);
802 deliveries.push(Delivery {
803 owner_pubkey,
804 device_pubkey,
805 envelope,
806 });
807 }
808 }
809
810 let inactive_sessions = std::mem::take(&mut record.inactive_sessions);
811 for mut session in inactive_sessions {
812 if session.can_send() {
813 let plan = session.plan_send(payload, ctx.now)?;
814 let mut envelope = session.apply_send(plan).envelope;
815 envelope.recipient = Some(device_pubkey);
816 deliveries.push(Delivery {
817 owner_pubkey,
818 device_pubkey,
819 envelope,
820 });
821 }
822 record.upsert_session(session, ctx.now);
823 }
824
825 if !deliveries.is_empty() {
826 record.last_activity = Some(ctx.now);
827 }
828
829 Ok(deliveries)
830 }
831
832 fn prepare_send_inner<R>(
833 &mut self,
834 ctx: &mut ProtocolContext<'_, R>,
835 recipient_owner: OwnerPubkey,
836 payload: Vec<u8>,
837 include_local_siblings: bool,
838 ) -> Result<PreparedSend>
839 where
840 R: RngCore + CryptoRng,
841 {
842 let mut relay_gaps = Vec::new();
843 let mut targets = BTreeSet::new();
844
845 self.collect_recipient_targets(recipient_owner, &mut targets, &mut relay_gaps);
846 if include_local_siblings {
847 self.collect_local_sibling_targets(&mut targets);
848 }
849
850 let mut deliveries = Vec::new();
851 let mut invite_responses = Vec::new();
852
853 for target in targets {
854 match self.prepare_device_delivery(
855 ctx,
856 target.owner_pubkey,
857 target.device_pubkey,
858 &payload,
859 false,
860 )? {
861 Some((delivery, maybe_response)) => {
862 deliveries.push(delivery);
863 if let Some(response) = maybe_response {
864 invite_responses.push(response);
865 }
866 }
867 None => {
868 relay_gaps.push(RelayGap::MissingDeviceInvite {
869 owner_pubkey: target.owner_pubkey,
870 device_pubkey: target.device_pubkey,
871 });
872 }
873 }
874 }
875
876 relay_gaps.sort();
877
878 Ok(PreparedSend {
879 recipient_owner,
880 payload,
881 deliveries,
882 invite_responses,
883 relay_gaps,
884 })
885 }
886
887 fn prepare_explicit_send<R>(
888 &mut self,
889 ctx: &mut ProtocolContext<'_, R>,
890 recipient_owner: OwnerPubkey,
891 targets: BTreeSet<TargetDevice>,
892 payload: Vec<u8>,
893 refresh_one_way_bootstrap: bool,
894 ) -> Result<PreparedSend>
895 where
896 R: RngCore + CryptoRng,
897 {
898 let mut deliveries = Vec::new();
899 let mut invite_responses = Vec::new();
900 let mut relay_gaps = Vec::new();
901
902 for target in targets {
903 match self.prepare_device_delivery(
904 ctx,
905 target.owner_pubkey,
906 target.device_pubkey,
907 &payload,
908 refresh_one_way_bootstrap,
909 )? {
910 Some((delivery, maybe_response)) => {
911 deliveries.push(delivery);
912 if let Some(response) = maybe_response {
913 invite_responses.push(response);
914 }
915 }
916 None => {
917 relay_gaps.push(RelayGap::MissingDeviceInvite {
918 owner_pubkey: target.owner_pubkey,
919 device_pubkey: target.device_pubkey,
920 });
921 }
922 }
923 }
924
925 relay_gaps.sort();
926
927 Ok(PreparedSend {
928 recipient_owner,
929 payload,
930 deliveries,
931 invite_responses,
932 relay_gaps,
933 })
934 }
935
936 fn collect_recipient_targets(
937 &self,
938 recipient_owner: OwnerPubkey,
939 targets: &mut BTreeSet<TargetDevice>,
940 relay_gaps: &mut Vec<RelayGap>,
941 ) {
942 let Some(user) = self.users.get(&recipient_owner) else {
943 relay_gaps.push(RelayGap::MissingRoster {
944 owner_pubkey: recipient_owner,
945 });
946 return;
947 };
948
949 if user
950 .roster
951 .as_ref()
952 .is_none_or(|roster| roster.devices().is_empty())
953 {
954 relay_gaps.push(RelayGap::MissingRoster {
955 owner_pubkey: recipient_owner,
956 });
957 return;
958 }
959
960 for device_pubkey in user.authorized_non_stale_devices() {
961 targets.insert(TargetDevice {
962 owner_pubkey: recipient_owner,
963 device_pubkey,
964 });
965 }
966 }
967
968 fn collect_local_sibling_targets(&self, targets: &mut BTreeSet<TargetDevice>) {
969 let Some(user) = self.users.get(&self.local_owner_pubkey) else {
970 return;
971 };
972
973 if user.roster.is_none() {
974 return;
975 }
976
977 for device_pubkey in user.authorized_non_stale_devices() {
978 if device_pubkey == self.local_device_pubkey {
979 continue;
980 }
981 targets.insert(TargetDevice {
982 owner_pubkey: self.local_owner_pubkey,
983 device_pubkey,
984 });
985 }
986 }
987
988 fn observe_public_invite(&mut self, owner_pubkey: OwnerPubkey, invite: Invite) -> Result<()> {
989 if let Some(inviter_owner_pubkey) = invite.inviter_owner_pubkey {
990 if inviter_owner_pubkey != owner_pubkey {
991 return Err(DomainError::InvalidState(format!(
992 "invite owner mismatch: expected {owner_pubkey}, got {inviter_owner_pubkey}"
993 ))
994 .into());
995 }
996 }
997
998 let device_pubkey = invite.inviter_device_pubkey;
999 let mut public_invite = invite;
1000 public_invite.inviter_ephemeral_private_key = None;
1001
1002 let user = self.user_record_mut(owner_pubkey);
1003 let record = user.device_record_mut(device_pubkey, public_invite.created_at);
1004
1005 let should_replace_invite = record
1006 .public_invite
1007 .as_ref()
1008 .is_none_or(|existing| public_invite.created_at >= existing.created_at);
1009
1010 record.created_at = merge_created_at(record.created_at, public_invite.created_at);
1011 if should_replace_invite {
1012 record.public_invite = Some(public_invite);
1013 }
1014 Ok(())
1015 }
1016
1017 fn apply_roster_for_owner(
1018 &mut self,
1019 owner_pubkey: OwnerPubkey,
1020 incoming_roster: DeviceRoster,
1021 ) -> RosterSnapshotDecision {
1022 self.apply_roster_for_owner_inner(owner_pubkey, incoming_roster, false)
1023 }
1024
1025 fn apply_roster_for_owner_inner(
1026 &mut self,
1027 owner_pubkey: OwnerPubkey,
1028 incoming_roster: DeviceRoster,
1029 replace_existing: bool,
1030 ) -> RosterSnapshotDecision {
1031 let user = self.user_record_mut(owner_pubkey);
1032 let current_roster = user.roster.as_ref();
1033 let (decision, next_roster) = if replace_existing {
1034 (RosterSnapshotDecision::Advanced, incoming_roster)
1035 } else {
1036 apply_roster_snapshot(current_roster, &incoming_roster)
1037 };
1038
1039 let previous_authorized = current_roster
1040 .map(authorized_device_set)
1041 .unwrap_or_default();
1042 let next_authorized = authorized_device_set(&next_roster);
1043
1044 user.roster = Some(next_roster.clone());
1045
1046 for device in next_roster.devices() {
1047 let record = user.device_record_mut(device.device_pubkey, device.created_at);
1048 record.authorized = true;
1049 record.is_stale = false;
1050 record.stale_since = None;
1051 record.created_at = merge_created_at(record.created_at, device.created_at);
1052 }
1053
1054 for removed in previous_authorized.difference(&next_authorized) {
1055 let record = user.device_record_mut(*removed, next_roster.created_at);
1056 record.authorized = false;
1057 record.is_stale = true;
1058 if record.stale_since.is_none() {
1059 record.stale_since = Some(next_roster.created_at);
1060 }
1061 }
1062
1063 self.reconcile_verified_claimed_devices(owner_pubkey, &next_roster, next_roster.created_at);
1064
1065 decision
1066 }
1067
1068 fn reconcile_verified_claimed_devices(
1069 &mut self,
1070 owner_pubkey: OwnerPubkey,
1071 roster: &DeviceRoster,
1072 now: UnixSeconds,
1073 ) {
1074 let roster_devices = authorized_device_set(roster);
1075 if roster_devices.is_empty() {
1076 return;
1077 }
1078
1079 let source_owners: Vec<OwnerPubkey> = self
1080 .users
1081 .keys()
1082 .copied()
1083 .filter(|candidate_owner_pubkey| *candidate_owner_pubkey != owner_pubkey)
1084 .collect();
1085
1086 let mut migrated = Vec::new();
1087 let mut empty_sources = Vec::new();
1088
1089 for source_owner_pubkey in source_owners {
1090 let matching_devices = self
1091 .users
1092 .get(&source_owner_pubkey)
1093 .map(|user| {
1094 user.devices
1095 .values()
1096 .filter(|record| {
1097 if !roster_devices.contains(&record.device_pubkey) {
1098 return false;
1099 }
1100 if record.claimed_owner_pubkey == Some(owner_pubkey) {
1101 return true;
1102 }
1103 user.roster
1104 .as_ref()
1105 .and_then(|roster| roster.get_device(&record.device_pubkey))
1106 .is_none()
1107 })
1108 .map(|record| record.device_pubkey)
1109 .collect::<Vec<_>>()
1110 })
1111 .unwrap_or_default();
1112
1113 if matching_devices.is_empty() {
1114 continue;
1115 }
1116
1117 if let Some(user) = self.users.get_mut(&source_owner_pubkey) {
1118 let source_roster_is_provisional = user.roster.as_ref().is_some_and(|roster| {
1119 roster.devices().iter().all(|device| {
1120 matching_devices.contains(&device.device_pubkey)
1121 && crate::owner_pubkey_from_device_pubkey(device.device_pubkey)
1122 == source_owner_pubkey
1123 })
1124 });
1125
1126 for device_pubkey in matching_devices {
1127 if let Some(mut record) = user.devices.remove(&device_pubkey) {
1128 record.claimed_owner_pubkey = None;
1129 migrated.push(record);
1130 }
1131 }
1132
1133 if user.devices.is_empty()
1134 && (user.roster.is_none() || source_roster_is_provisional)
1135 {
1136 empty_sources.push(source_owner_pubkey);
1137 }
1138 }
1139 }
1140
1141 for source_owner_pubkey in empty_sources {
1142 self.users.remove(&source_owner_pubkey);
1143 }
1144
1145 if migrated.is_empty() {
1146 return;
1147 }
1148
1149 let user = self.user_record_mut(owner_pubkey);
1150 for record in migrated {
1151 let device_pubkey = record.device_pubkey;
1152 user.device_record_mut(device_pubkey, record.created_at)
1153 .absorb(record, now);
1154 }
1155 }
1156
1157 fn user_record_mut(&mut self, owner_pubkey: OwnerPubkey) -> &mut UserRecord {
1158 self.users
1159 .entry(owner_pubkey)
1160 .or_insert_with(|| UserRecord::new(owner_pubkey))
1161 }
1162}
1163
1164impl UserRecord {
1165 fn new(owner_pubkey: OwnerPubkey) -> Self {
1166 Self {
1167 owner_pubkey,
1168 roster: None,
1169 devices: BTreeMap::new(),
1170 }
1171 }
1172
1173 fn from_snapshot(snapshot: UserRecordSnapshot) -> Self {
1174 Self {
1175 owner_pubkey: snapshot.owner_pubkey,
1176 roster: snapshot.roster,
1177 devices: snapshot
1178 .devices
1179 .into_iter()
1180 .map(DeviceRecord::from_snapshot)
1181 .map(|record| (record.device_pubkey, record))
1182 .collect(),
1183 }
1184 }
1185
1186 fn snapshot(&self) -> UserRecordSnapshot {
1187 UserRecordSnapshot {
1188 owner_pubkey: self.owner_pubkey,
1189 roster: self.roster.clone(),
1190 devices: self.devices.values().map(DeviceRecord::snapshot).collect(),
1191 }
1192 }
1193
1194 fn device_record_mut(
1195 &mut self,
1196 device_pubkey: DevicePubkey,
1197 created_at: UnixSeconds,
1198 ) -> &mut DeviceRecord {
1199 self.devices
1200 .entry(device_pubkey)
1201 .or_insert_with(|| DeviceRecord::new(device_pubkey, created_at))
1202 }
1203
1204 fn authorized_non_stale_devices(&self) -> Vec<DevicePubkey> {
1205 self.devices
1206 .values()
1207 .filter(|record| record.authorized && !record.is_stale)
1208 .map(|record| record.device_pubkey)
1209 .collect()
1210 }
1211}
1212
1213impl DeviceRecord {
1214 fn new(device_pubkey: DevicePubkey, created_at: UnixSeconds) -> Self {
1215 Self {
1216 device_pubkey,
1217 authorized: false,
1218 is_stale: false,
1219 stale_since: None,
1220 claimed_owner_pubkey: None,
1221 public_invite: None,
1222 invite_response_generated: false,
1223 active_session: None,
1224 inactive_sessions: Vec::new(),
1225 last_activity: None,
1226 created_at,
1227 }
1228 }
1229
1230 fn from_snapshot(snapshot: DeviceRecordSnapshot) -> Self {
1231 Self {
1232 device_pubkey: snapshot.device_pubkey,
1233 authorized: snapshot.authorized,
1234 is_stale: snapshot.is_stale,
1235 stale_since: snapshot.stale_since,
1236 claimed_owner_pubkey: snapshot.claimed_owner_pubkey,
1237 public_invite: snapshot.public_invite,
1238 invite_response_generated: snapshot.invite_response_generated,
1239 active_session: snapshot.active_session.map(Session::from_state),
1240 inactive_sessions: snapshot
1241 .inactive_sessions
1242 .into_iter()
1243 .map(Session::from_state)
1244 .collect(),
1245 last_activity: snapshot.last_activity,
1246 created_at: snapshot.created_at,
1247 }
1248 }
1249
1250 fn snapshot(&self) -> DeviceRecordSnapshot {
1251 DeviceRecordSnapshot {
1252 device_pubkey: self.device_pubkey,
1253 authorized: self.authorized,
1254 is_stale: self.is_stale,
1255 stale_since: self.stale_since,
1256 claimed_owner_pubkey: self.claimed_owner_pubkey,
1257 public_invite: self.public_invite.clone(),
1258 invite_response_generated: self.invite_response_generated,
1259 active_session: self
1260 .active_session
1261 .as_ref()
1262 .map(|session| session.state.clone()),
1263 inactive_sessions: self
1264 .inactive_sessions
1265 .iter()
1266 .map(|session| session.state.clone())
1267 .collect(),
1268 last_activity: self.last_activity,
1269 created_at: self.created_at,
1270 }
1271 }
1272
1273 fn best_send_session_source(&self) -> Option<SendSessionSource> {
1274 let mut best: Option<(SendSessionSource, (u8, u32, u32))> = None;
1275
1276 if let Some(active_session) = self.active_session.as_ref() {
1277 if active_session.can_send() {
1278 best = Some((SendSessionSource::Active, session_priority(active_session)));
1279 }
1280 }
1281
1282 for (index, session) in self.inactive_sessions.iter().enumerate() {
1283 if !session.can_send() {
1284 continue;
1285 }
1286 let priority = session_priority(session);
1287 if best
1288 .as_ref()
1289 .is_none_or(|(_, current_priority)| priority > *current_priority)
1290 {
1291 best = Some((SendSessionSource::Inactive(index), priority));
1292 }
1293 }
1294
1295 best.map(|(source, _)| source)
1296 }
1297
1298 fn session_for_send_source(&self, source: &SendSessionSource) -> Option<&Session> {
1299 match source {
1300 SendSessionSource::Active => self.active_session.as_ref(),
1301 SendSessionSource::Inactive(index) => self.inactive_sessions.get(*index),
1302 }
1303 }
1304
1305 fn upsert_session(&mut self, session: Session, now: UnixSeconds) {
1306 if self.contains_state(&session.state) {
1307 self.compact_duplicate_sessions();
1308 self.last_activity = Some(now);
1309 return;
1310 }
1311
1312 let new_priority = session_priority(&session);
1313 let old_priority = self
1314 .active_session
1315 .as_ref()
1316 .map(session_priority)
1317 .unwrap_or((0, 0, 0));
1318
1319 if let Some(old_active) = self.active_session.take() {
1320 if old_priority >= new_priority {
1321 self.inactive_sessions.push(session);
1322 self.active_session = Some(old_active);
1323 } else {
1324 self.inactive_sessions.push(old_active);
1325 self.active_session = Some(session);
1326 }
1327 } else {
1328 self.active_session = Some(session);
1329 }
1330
1331 self.compact_duplicate_sessions();
1332 if self.inactive_sessions.len() > MAX_INACTIVE_SESSIONS {
1333 self.inactive_sessions.truncate(MAX_INACTIVE_SESSIONS);
1334 }
1335 self.last_activity = Some(now);
1336 }
1337
1338 fn absorb(&mut self, mut other: DeviceRecord, now: UnixSeconds) {
1339 self.authorized |= other.authorized;
1340 self.is_stale &= other.is_stale;
1341 self.stale_since = match (self.stale_since, other.stale_since) {
1342 (Some(existing), Some(incoming)) => Some(existing.min(incoming)),
1343 (None, incoming) => incoming,
1344 (existing, None) => existing,
1345 };
1346 self.claimed_owner_pubkey = self
1347 .claimed_owner_pubkey
1348 .or(other.claimed_owner_pubkey.take());
1349 self.created_at = merge_created_at(self.created_at, other.created_at);
1350
1351 if let Some(public_invite) = other.public_invite.take() {
1352 let should_replace_invite = self
1353 .public_invite
1354 .as_ref()
1355 .is_none_or(|existing| public_invite.created_at >= existing.created_at);
1356 if should_replace_invite {
1357 self.public_invite = Some(public_invite);
1358 }
1359 }
1360 self.invite_response_generated |= other.invite_response_generated;
1361
1362 if let Some(session) = other.active_session.take() {
1363 self.upsert_session(session, now);
1364 }
1365
1366 for session in other.inactive_sessions.drain(..) {
1367 self.upsert_session(session, now);
1368 }
1369
1370 self.last_activity = match (self.last_activity, other.last_activity) {
1371 (Some(existing), Some(incoming)) => Some(existing.max(incoming)),
1372 (None, incoming) => incoming,
1373 (existing, None) => existing,
1374 };
1375 }
1376
1377 fn promote_inactive_session(&mut self, session: Session) {
1378 let new_priority = session_priority(&session);
1379 if let Some(old_active) = self.active_session.take() {
1380 let old_priority = session_priority(&old_active);
1381 if new_priority > old_priority {
1382 if old_active.state != session.state {
1383 self.inactive_sessions.push(old_active);
1384 }
1385 self.active_session = Some(session);
1386 } else {
1387 self.inactive_sessions.push(session);
1388 self.active_session = Some(old_active);
1389 }
1390 } else {
1391 self.active_session = Some(session);
1392 }
1393 self.compact_duplicate_sessions();
1394 if self.inactive_sessions.len() > MAX_INACTIVE_SESSIONS {
1395 self.inactive_sessions.truncate(MAX_INACTIVE_SESSIONS);
1396 }
1397 }
1398
1399 fn contains_state(&self, state: &SessionState) -> bool {
1400 self.active_session
1401 .as_ref()
1402 .is_some_and(|session| session.state == *state)
1403 || self
1404 .inactive_sessions
1405 .iter()
1406 .any(|session| session.state == *state)
1407 }
1408
1409 fn compact_duplicate_sessions(&mut self) {
1410 let active_state = self
1411 .active_session
1412 .as_ref()
1413 .map(|session| session.state.clone());
1414 let mut unique_states = Vec::new();
1415 let mut inactive_sessions = Vec::with_capacity(self.inactive_sessions.len());
1416
1417 for session in self.inactive_sessions.drain(..) {
1418 let is_duplicate = active_state
1419 .as_ref()
1420 .is_some_and(|state| *state == session.state)
1421 || unique_states.contains(&session.state);
1422 if is_duplicate {
1423 continue;
1424 }
1425 unique_states.push(session.state.clone());
1426 inactive_sessions.push(session);
1427 }
1428
1429 self.inactive_sessions = inactive_sessions;
1430 }
1431}
1432
1433fn apply_roster_snapshot(
1434 current_roster: Option<&DeviceRoster>,
1435 incoming_roster: &DeviceRoster,
1436) -> (RosterSnapshotDecision, DeviceRoster) {
1437 let Some(current_roster) = current_roster else {
1438 return (RosterSnapshotDecision::Advanced, incoming_roster.clone());
1439 };
1440
1441 if incoming_roster.created_at > current_roster.created_at {
1442 return (RosterSnapshotDecision::Advanced, incoming_roster.clone());
1443 }
1444
1445 if incoming_roster.created_at < current_roster.created_at {
1446 return (RosterSnapshotDecision::Stale, current_roster.clone());
1447 }
1448
1449 (
1450 RosterSnapshotDecision::MergedEqualTimestamp,
1451 current_roster.merge(incoming_roster),
1452 )
1453}
1454
1455fn authorized_device_set(roster: &DeviceRoster) -> BTreeSet<DevicePubkey> {
1456 roster
1457 .devices()
1458 .iter()
1459 .map(|device| device.device_pubkey)
1460 .collect()
1461}
1462
1463fn session_priority(session: &Session) -> (u8, u32, u32) {
1464 let can_send = session.can_send();
1465 let can_receive = session.state.receiving_chain_key.is_some()
1466 || session.state.their_current_nostr_public_key.is_some()
1467 || session.state.receiving_chain_message_number > 0;
1468
1469 let directionality = match (can_send, can_receive) {
1470 (true, true) => 3,
1471 (true, false) => 2,
1472 (false, true) => 1,
1473 (false, false) => 0,
1474 };
1475
1476 (
1477 directionality,
1478 session.state.receiving_chain_message_number,
1479 session.state.sending_chain_message_number,
1480 )
1481}
1482
1483fn is_one_way_bootstrap_session(session: &Session) -> bool {
1484 session.state.receiving_chain_key.is_none()
1485 && session.state.their_current_nostr_public_key.is_none()
1486}
1487
1488fn merge_created_at(current: UnixSeconds, observed: UnixSeconds) -> UnixSeconds {
1489 match (current.get(), observed.get()) {
1490 (0, _) => observed,
1491 (_, 0) => current,
1492 _ => current.min(observed),
1493 }
1494}