1use crate::prefabs::ClientServerRemote;
64use crate::prelude::results::{PeerConnectSuccess, PeerRegisterStatus};
65use crate::prelude::*;
66use crate::remote_ext::remote_specialization::PeerRemote;
67use crate::remote_ext::results::LocalGroupPeerFullInfo;
68use std::ops::{Deref, DerefMut};
69
70use futures::StreamExt;
71use std::path::PathBuf;
72use std::time::Duration;
73
74pub(crate) mod user_ids {
75 use crate::prelude::*;
76 use std::ops::Deref;
77
78 #[derive(Debug)]
79 pub struct SymmetricIdentifierHandleRef<'a, R: Ratchet> {
81 pub(crate) user: VirtualTargetType,
82 pub(crate) remote: &'a NodeRemote<R>,
83 pub(crate) target_username: Option<String>,
84 }
85
86 impl<R: Ratchet> SymmetricIdentifierHandleRef<'_, R> {
87 pub fn into_owned(self) -> SymmetricIdentifierHandle<R> {
88 SymmetricIdentifierHandle {
89 user: self.user,
90 remote: self.remote.clone(),
91 target_username: self.target_username,
92 }
93 }
94 }
95
96 #[derive(Clone, Debug)]
97 pub struct SymmetricIdentifierHandle<R: Ratchet> {
99 user: VirtualTargetType,
100 remote: NodeRemote<R>,
101 target_username: Option<String>,
102 }
103
104 pub trait TargetLockedRemote<R: Ratchet>: Send + Sync {
105 fn user(&self) -> &VirtualTargetType;
106 fn remote(&self) -> &NodeRemote<R>;
107 fn target_username(&self) -> Option<&str>;
108 fn user_mut(&mut self) -> &mut VirtualTargetType;
109 fn session_security_settings(&self) -> Option<&SessionSecuritySettings>;
110 }
111
112 impl<R: Ratchet> TargetLockedRemote<R> for SymmetricIdentifierHandleRef<'_, R> {
113 fn user(&self) -> &VirtualTargetType {
114 &self.user
115 }
116 fn remote(&self) -> &NodeRemote<R> {
117 self.remote
118 }
119 fn target_username(&self) -> Option<&str> {
120 self.target_username.as_deref()
121 }
122 fn user_mut(&mut self) -> &mut VirtualTargetType {
123 &mut self.user
124 }
125
126 fn session_security_settings(&self) -> Option<&SessionSecuritySettings> {
127 None
128 }
129 }
130
131 impl<R: Ratchet> TargetLockedRemote<R> for SymmetricIdentifierHandle<R> {
132 fn user(&self) -> &VirtualTargetType {
133 &self.user
134 }
135 fn remote(&self) -> &NodeRemote<R> {
136 &self.remote
137 }
138 fn target_username(&self) -> Option<&str> {
139 self.target_username.as_deref()
140 }
141 fn user_mut(&mut self) -> &mut VirtualTargetType {
142 &mut self.user
143 }
144
145 fn session_security_settings(&self) -> Option<&SessionSecuritySettings> {
146 None
147 }
148 }
149
150 impl<R: Ratchet> From<SymmetricIdentifierHandleRef<'_, R>> for SymmetricIdentifierHandle<R> {
151 fn from(this: SymmetricIdentifierHandleRef<'_, R>) -> Self {
152 this.into_owned()
153 }
154 }
155
156 impl<R: Ratchet> Deref for SymmetricIdentifierHandle<R> {
157 type Target = NodeRemote<R>;
158
159 fn deref(&self) -> &Self::Target {
160 &self.remote
161 }
162 }
163
164 impl<R: Ratchet> Deref for SymmetricIdentifierHandleRef<'_, R> {
165 type Target = NodeRemote<R>;
166
167 fn deref(&self) -> &Self::Target {
168 self.remote
169 }
170 }
171}
172
173pub struct CitadelClientServerConnection<R: Ratchet> {
175 pub(crate) channel: Option<PeerChannel<R>>,
177 pub remote: ClientServerRemote<R>,
178 pub udp_channel_rx: Option<citadel_io::tokio::sync::oneshot::Receiver<UdpChannel<R>>>,
180 pub services: ServicesObject,
182 pub cid: u64,
183 pub session_security_settings: SessionSecuritySettings,
184}
185
186impl<R: Ratchet> CitadelClientServerConnection<R> {
187 pub fn split(self) -> (PeerChannelSendHalf<R>, PeerChannelRecvHalf<R>) {
193 self.channel.expect("Channel already taken").split()
194 }
195
196 pub fn take_channel(&mut self) -> Option<PeerChannel<R>> {
197 self.channel.take()
198 }
199}
200
201impl<R: Ratchet> Deref for CitadelClientServerConnection<R> {
202 type Target = ClientServerRemote<R>;
203
204 fn deref(&self) -> &Self::Target {
205 &self.remote
206 }
207}
208
209impl<R: Ratchet> DerefMut for CitadelClientServerConnection<R> {
210 fn deref_mut(&mut self) -> &mut Self::Target {
211 &mut self.remote
212 }
213}
214
215pub struct RegisterSuccess {
217 pub cid: u64,
218}
219
220#[async_trait]
221pub trait ProtocolRemoteExt<R: Ratchet>: Remote<R> {
223 async fn register<
226 T: std::net::ToSocketAddrs + Send,
227 P: Into<String> + Send,
228 V: Into<String> + Send,
229 K: Into<SecBuffer> + Send,
230 >(
231 &self,
232 addr: T,
233 full_name: P,
234 username: V,
235 proposed_password: K,
236 default_security_settings: SessionSecuritySettings,
237 server_password: Option<PreSharedKey>,
238 ) -> Result<RegisterSuccess, NetworkError> {
239 let creds =
240 ProposedCredentials::new_register(full_name, username, proposed_password.into())
241 .await?;
242 let register_request = NodeRequest::RegisterToHypernode(RegisterToHypernode {
243 remote_addr: addr.to_socket_addrs()?.next().ok_or(citadel_io::error!(
244 citadel_io::ErrorCode::RemoteInvalidSocketAddr
245 ))?,
246 proposed_credentials: creds,
247 static_security_settings: default_security_settings,
248 session_password: server_password.unwrap_or_default(),
249 });
250
251 let mut subscription = self.send_callback_subscription(register_request).await?;
252 while let Some(status) = subscription.next().await {
253 match status.into_result()? {
254 NodeResult::RegisterOkay(RegisterOkay { cid, .. }) => {
255 return Ok(RegisterSuccess { cid });
256 }
257 NodeResult::RegisterFailure(err) => {
258 return Err(citadel_io::error!(
259 citadel_io::ErrorCode::RemoteRegisterFailure,
260 err.error_message
261 ));
262 }
263 NodeResult::Disconnect(err) => {
264 return Err(citadel_io::error!(
265 citadel_io::ErrorCode::RemoteDisconnected,
266 err.message
267 ));
268 }
269 evt => {
270 log::warn!(target: "citadel", "Invalid NodeResult for Register request received: {evt:?}");
271 }
272 }
273 }
274
275 Err(citadel_io::error!(
276 citadel_io::ErrorCode::RemoteKernelStreamDied,
277 "register"
278 ))
279 }
280
281 async fn register_with_defaults<
284 T: std::net::ToSocketAddrs + Send,
285 P: Into<String> + Send,
286 V: Into<String> + Send,
287 K: Into<SecBuffer> + Send,
288 >(
289 &self,
290 addr: T,
291 full_name: P,
292 username: V,
293 proposed_password: K,
294 ) -> Result<RegisterSuccess, NetworkError> {
295 self.register(
296 addr,
297 full_name,
298 username,
299 proposed_password,
300 Default::default(),
301 Default::default(),
302 )
303 .await
304 }
305
306 async fn connect(
309 &self,
310 auth: AuthenticationRequest,
311 connect_mode: ConnectMode,
312 udp_mode: UdpMode,
313 keep_alive_timeout: Option<Duration>,
314 session_security_settings: SessionSecuritySettings,
315 server_password: Option<PreSharedKey>,
316 ) -> Result<CitadelClientServerConnection<R>, NetworkError> {
317 let connect_request = NodeRequest::ConnectToHypernode(ConnectToHypernode {
318 auth_request: auth,
319 connect_mode,
320 udp_mode,
321 keep_alive_timeout: keep_alive_timeout.map(|r| r.as_secs()),
322 session_security_settings,
323 session_password: server_password.unwrap_or_default(),
324 });
325
326 let mut subscription = self.send_callback_subscription(connect_request).await?;
327 let status = subscription.next().await.ok_or(citadel_io::error!(
328 citadel_io::ErrorCode::RemoteKernelStreamDied,
329 "connect"
330 ))?;
331
332 return match status.into_result()? {
333 NodeResult::ConnectSuccess(ConnectSuccess {
334 ticket: _,
335 session_cid: cid,
336 remote_addr: _,
337 is_personal: _,
338 v_conn_type,
339 services,
340 welcome_message: _,
341 channel,
342 udp_rx_opt: udp_channel_rx,
343 session_security_settings,
344 }) => Ok(CitadelClientServerConnection {
345 remote: ClientServerRemote::new(
346 v_conn_type,
347 self.remote_ref().clone(),
348 session_security_settings,
349 None,
350 None,
351 ),
352 channel: Some(*channel),
353 udp_channel_rx,
354 services,
355 cid,
356 session_security_settings,
357 }),
358 NodeResult::ConnectFail(ConnectFail {
359 ticket: _,
360 cid_opt: _,
361 error_message: err,
362 }) => Err(citadel_io::error!(
363 citadel_io::ErrorCode::RemoteConnectFailed,
364 err
365 )),
366 NodeResult::Disconnect(err) => {
367 return Err(citadel_io::error!(
368 citadel_io::ErrorCode::RemoteDisconnected,
369 err.message
370 ));
371 }
372 res => Err(citadel_io::error!(
373 citadel_io::ErrorCode::RemoteConnectUnexpectedResponse,
374 citadel_io::Dbg(res)
375 )),
376 };
377 }
378
379 async fn connect_with_defaults(
382 &self,
383 auth: AuthenticationRequest,
384 ) -> Result<CitadelClientServerConnection<R>, NetworkError> {
385 self.connect(
386 auth,
387 Default::default(),
388 Default::default(),
389 None,
390 Default::default(),
391 Default::default(),
392 )
393 .await
394 }
395
396 async fn find_target<T: Into<UserIdentifier> + Send, P: Into<UserIdentifier> + Send>(
409 &self,
410 local_user: T,
411 peer: P,
412 ) -> Result<SymmetricIdentifierHandleRef<'_, R>, NetworkError> {
413 let account_manager = self.account_manager();
414 account_manager
415 .find_target_information(local_user, peer)
416 .await?
417 .map(move |(cid, peer)| {
418 if peer.parent_icid != 0 {
419 SymmetricIdentifierHandleRef {
420 user: VirtualTargetType::ExternalGroupPeer {
421 session_cid: cid,
422 interserver_cid: peer.parent_icid,
423 peer_cid: peer.cid,
424 },
425 remote: self.remote_ref(),
426 target_username: None,
427 }
428 } else {
429 SymmetricIdentifierHandleRef {
430 user: VirtualTargetType::LocalGroupPeer {
431 session_cid: cid,
432 peer_cid: peer.cid,
433 },
434 remote: self.remote_ref(),
435 target_username: None,
436 }
437 }
438 })
439 .ok_or_else(|| citadel_io::error!(citadel_io::ErrorCode::RemoteTargetPairNotFound))
440 }
441
442 async fn propose_target<T: Into<UserIdentifier> + Send, P: Into<UserIdentifier> + Send>(
445 &self,
446 local_user: T,
447 peer: P,
448 ) -> Result<SymmetricIdentifierHandleRef<'_, R>, NetworkError> {
449 let local_cid = self.get_session_cid(local_user).await?;
450 match peer.into() {
451 UserIdentifier::ID(peer_cid) => Ok(SymmetricIdentifierHandleRef {
452 user: VirtualTargetType::LocalGroupPeer {
453 session_cid: local_cid,
454 peer_cid,
455 },
456 remote: self.remote_ref(),
457 target_username: None,
458 }),
459 UserIdentifier::Username(uname) => {
460 let peer_cid = self
461 .remote_ref()
462 .account_manager()
463 .find_target_information(local_cid, uname.clone())
464 .await?
465 .map(|r| r.1.cid)
466 .unwrap_or(0);
467 Ok(SymmetricIdentifierHandleRef {
468 user: VirtualTargetType::LocalGroupPeer {
469 session_cid: local_cid,
470 peer_cid,
471 },
472 remote: self.remote_ref(),
473 target_username: Some(uname),
474 })
475 }
476 }
477 }
478
479 async fn get_local_group_peers<T: Into<UserIdentifier> + Send>(
482 &self,
483 local_user: T,
484 limit: Option<usize>,
485 ) -> Result<Vec<LocalGroupPeerFullInfo>, NetworkError> {
486 let local_cid = self.get_session_cid(local_user).await?;
487 let command = NodeRequest::PeerCommand(PeerCommand {
488 session_cid: local_cid,
489 command: PeerSignal::GetRegisteredPeers {
490 peer_conn_type: ClientConnectionType::Server {
491 session_cid: local_cid,
492 },
493 response: None,
494 limit: limit.map(|r| r as i32),
495 },
496 });
497
498 let mut stream = self.send_callback_subscription(command).await?;
499
500 while let Some(status) = stream.next().await {
501 if let NodeResult::PeerEvent(PeerEvent {
502 event:
503 PeerSignal::GetRegisteredPeers {
504 peer_conn_type: _,
505 response: Some(PeerResponse::RegisteredCids(peer_info, is_onlines)),
506 limit: _,
507 },
508 ..
509 }) = status.into_result()?
510 {
511 return Ok(peer_info
512 .into_iter()
513 .zip(is_onlines)
514 .filter_map(|(peer_info, is_online)| {
515 peer_info.map(|info| LocalGroupPeerFullInfo {
516 cid: info.cid,
517 username: Some(info.username),
518 full_name: Some(info.full_name),
519 is_online,
520 })
521 })
522 .collect());
523 }
524 }
525
526 Err(citadel_io::error!(
527 citadel_io::ErrorCode::RemoteKernelStreamDied,
528 "get_local_group_peers"
529 ))
530 }
531
532 async fn get_local_group_mutual_peers<T: Into<UserIdentifier> + Send>(
534 &self,
535 local_user: T,
536 ) -> Result<Vec<LocalGroupPeerFullInfo>, NetworkError> {
537 let local_cid = self.get_session_cid(local_user).await?;
538 let command = NodeRequest::PeerCommand(PeerCommand {
539 session_cid: local_cid,
540 command: PeerSignal::GetMutuals {
541 v_conn_type: ClientConnectionType::Server {
542 session_cid: local_cid,
543 },
544 response: None,
545 },
546 });
547
548 let mut stream = self.send_callback_subscription(command).await?;
549
550 while let Some(status) = stream.next().await {
551 if let NodeResult::PeerEvent(PeerEvent {
552 event:
553 PeerSignal::GetMutuals {
554 v_conn_type: _,
555 response: Some(PeerResponse::RegisteredCids(peer_info, is_onlines)),
556 },
557 ..
558 }) = status.into_result()?
559 {
560 return Ok(peer_info
561 .into_iter()
562 .zip(is_onlines)
563 .filter_map(|(peer_info, is_online)| {
564 peer_info.map(|info| LocalGroupPeerFullInfo {
565 cid: info.cid,
566 username: Some(info.username),
567 full_name: Some(info.full_name),
568 is_online,
569 })
570 })
571 .collect());
572 }
573 }
574
575 Err(citadel_io::error!(
576 citadel_io::ErrorCode::RemoteSessionStreamDied
577 ))
578 }
579
580 async fn sessions(&self) -> Result<ActiveSessions, NetworkError> {
583 match self
584 .send_callback_subscription(NodeRequest::GetActiveSessions)
585 .await
586 {
587 Ok(mut stream) => {
588 while let Some(result) = stream.next().await {
589 match result.into_result()? {
590 NodeResult::SessionList(res) => return Ok(res.sessions),
591 res => {
592 citadel_logging::warn!("Received unexpected result: {res:?}");
593 }
594 }
595 }
596
597 citadel_logging::warn!("Failed to receive response from SDK (stream died)");
598 return Err(citadel_io::error!(
599 citadel_io::ErrorCode::RemoteKernelStreamDied,
600 "get_active_sessions"
601 ));
602 }
603 Err(e) => {
604 citadel_logging::warn!(target: "citadel", "Failed to query SDK sessions: {e}");
605 }
606 }
607
608 Err(citadel_io::error!(
609 citadel_io::ErrorCode::RemoteQuerySessionsFailed
610 ))
611 }
612
613 #[doc(hidden)]
614 fn remote_ref(&self) -> &NodeRemote<R>;
615
616 #[doc(hidden)]
617 async fn get_session_cid<T: Into<UserIdentifier> + Send>(
618 &self,
619 local_user: T,
620 ) -> Result<u64, NetworkError> {
621 let account_manager = self.account_manager();
622 Ok(account_manager
623 .find_local_user_information(local_user)
624 .await?
625 .ok_or(citadel_io::error!(
626 citadel_io::ErrorCode::RemoteUserDoesNotExist
627 ))?)
628 }
629}
630
631impl<R: Ratchet> ProtocolRemoteExt<R> for NodeRemote<R> {
632 fn remote_ref(&self) -> &NodeRemote<R> {
633 self
634 }
635}
636
637impl<R: Ratchet> ProtocolRemoteExt<R> for ClientServerRemote<R> {
638 fn remote_ref(&self) -> &NodeRemote<R> {
639 &self.inner
640 }
641}
642
643#[async_trait]
644pub trait ProtocolRemoteTargetExt<R: Ratchet>: TargetLockedRemote<R> {
646 async fn send_file_with_custom_opts<T: ObjectSource>(
648 &self,
649 source: T,
650 chunk_size: usize,
651 transfer_type: TransferType,
652 ) -> Result<(), NetworkError> {
653 let chunk_size = if chunk_size == 0 {
654 None
655 } else {
656 Some(chunk_size)
657 };
658 let session_cid = self.user().get_session_cid();
659 let user = *self.user();
660 let remote = self.remote();
661
662 let mut stream = remote
663 .send_callback_subscription(NodeRequest::SendObject(SendObject {
664 source: Box::new(source),
665 chunk_size,
666 session_cid,
667 v_conn_type: user,
668 transfer_type,
669 }))
670 .await?;
671
672 while let Some(event) = stream.next().await {
673 match event.into_result()? {
674 NodeResult::ObjectTransferHandle(ObjectTransferHandle { mut handle, .. }) => {
675 return handle.transfer_file().await.map_err(|err| {
676 citadel_io::error!(
677 citadel_io::ErrorCode::RemoteFileTransferFailed,
678 err.into_string()
679 )
680 });
681 }
682
683 NodeResult::PeerEvent(PeerEvent {
684 event: PeerSignal::SignalReceived { .. },
685 ..
686 }) => {}
687
688 res => {
689 log::warn!(target: "citadel", "Invalid NodeResult for FileTransfer request received: {res:?}")
690 }
691 }
692 }
693
694 Err(citadel_io::error!(
695 citadel_io::ErrorCode::RemoteFileTransferStreamDied
696 ))
697 }
698
699 async fn send_file<T: ObjectSource>(&self, source: T) -> Result<(), NetworkError> {
701 self.send_file_with_custom_opts(source, 0, TransferType::FileTransfer)
702 .await
703 }
704
705 async fn remote_encrypted_virtual_filesystem_push_custom_chunking<
708 T: ObjectSource,
709 P: Into<PathBuf> + Send,
710 >(
711 &self,
712 source: T,
713 virtual_directory: P,
714 chunk_size: usize,
715 security_level: SecurityLevel,
716 ) -> Result<(), NetworkError> {
717 self.can_use_revfs()?;
718 let mut virtual_path = virtual_directory.into();
719 virtual_path = prepare_virtual_path(virtual_path);
720 validate_virtual_path(&virtual_path).map_err(|err| {
721 citadel_io::error!(
722 citadel_io::ErrorCode::RemoteRevfsInvalidVirtualPath,
723 err.into_string()
724 )
725 })?;
726 let tx_type = TransferType::RemoteEncryptedVirtualFilesystem {
727 virtual_path,
728 security_level,
729 };
730 self.send_file_with_custom_opts(source, chunk_size, tx_type)
731 .await
732 }
733
734 async fn remote_encrypted_virtual_filesystem_push<T: ObjectSource, P: Into<PathBuf> + Send>(
737 &self,
738 source: T,
739 virtual_directory: P,
740 security_level: SecurityLevel,
741 ) -> Result<(), NetworkError> {
742 self.remote_encrypted_virtual_filesystem_push_custom_chunking(
743 source,
744 virtual_directory,
745 0,
746 security_level,
747 )
748 .await
749 }
750
751 async fn remote_encrypted_virtual_filesystem_pull<P: Into<PathBuf> + Send>(
754 &self,
755 virtual_directory: P,
756 transfer_security_level: SecurityLevel,
757 delete_on_pull: bool,
758 ) -> Result<PathBuf, NetworkError> {
759 self.can_use_revfs()?;
760 let request = NodeRequest::PullObject(PullObject {
761 v_conn: *self.user(),
762 virtual_dir: virtual_directory.into(),
763 delete_on_pull,
764 transfer_security_level,
765 });
766
767 let mut stream = self.remote().send_callback_subscription(request).await?;
768
769 while let Some(event) = stream.next().await {
770 match event.into_result()? {
771 NodeResult::ObjectTransferHandle(ObjectTransferHandle { mut handle, .. }) => {
772 return handle.receive_file().await.map_err(|err| {
773 citadel_io::error!(
774 citadel_io::ErrorCode::RemoteFileTransferFailed,
775 err.into_string()
776 )
777 });
778 }
779
780 NodeResult::PeerEvent(PeerEvent {
781 event: PeerSignal::SignalReceived { .. },
782 ..
783 }) => {}
784
785 res => {
786 log::error!(target: "citadel", "Invalid NodeResult for REVFS FileTransfer request received: {res:?}");
787 return Err(citadel_io::error!(
788 citadel_io::ErrorCode::RemoteRevfsInvalidResponse
789 ));
790 }
791 }
792 }
793
794 Err(citadel_io::error!(
795 citadel_io::ErrorCode::RemoteRevfsFileTransferStreamDied
796 ))
797 }
798
799 async fn remote_encrypted_virtual_filesystem_delete<P: Into<PathBuf> + Send>(
803 &self,
804 virtual_directory: P,
805 ) -> Result<(), NetworkError> {
806 self.can_use_revfs()?;
807 let request = NodeRequest::DeleteObject(DeleteObject {
808 v_conn: *self.user(),
809 virtual_dir: virtual_directory.into(),
810 security_level: Default::default(),
811 });
812
813 let mut stream = self.remote().send_callback_subscription(request).await?;
814 while let Some(event) = stream.next().await {
815 match event.into_result()? {
816 NodeResult::ReVFS(result) => {
817 return if let Some(error) = result.error_message {
818 Err(citadel_io::error!(
819 citadel_io::ErrorCode::RemoteFileTransferFailed,
820 error
821 ))
822 } else {
823 Ok(())
824 }
825 }
826
827 evt => {
828 log::error!(target: "citadel", "Invalid NodeResult for REVFS Delete request received: {evt:?}");
829 }
830 }
831 }
832
833 Err(citadel_io::error!(
834 citadel_io::ErrorCode::RemoteRevfsDeleteStreamDied
835 ))
836 }
837
838 async fn connect_to_peer_custom(
840 &self,
841 session_security_settings: SessionSecuritySettings,
842 udp_mode: UdpMode,
843 peer_session_password: Option<PreSharedKey>,
844 ) -> Result<PeerConnectSuccess<R>, NetworkError> {
845 use std::time::Duration;
846
847 const P2P_CONNECT_TIMEOUT: Duration = Duration::from_secs(60);
852
853 let session_cid = self.user().get_session_cid();
854 let peer_target = self.try_as_peer_connection().await?;
855
856 let mut stream = self
857 .remote()
858 .send_callback_subscription(NodeRequest::PeerCommand(PeerCommand {
859 session_cid,
860 command: PeerSignal::PostConnect {
861 peer_conn_type: peer_target,
862 ticket_opt: None,
863 invitee_response: None,
864 session_security_settings,
865 udp_mode,
866 session_password: peer_session_password,
867 },
868 }))
869 .await?;
870
871 let connect_task = async {
872 while let Some(status) = stream.next().await {
873 match status.into_result()? {
874 NodeResult::PeerChannelCreated(PeerChannelCreated {
875 ticket: _,
876 channel,
877 udp_rx_opt,
878 }) => {
879 let username = self.target_username().map(ToString::to_string);
880 let remote = PeerRemote {
881 inner: self.remote().clone(),
882 peer: peer_target.as_virtual_connection(),
883 username,
884 session_security_settings,
885 };
886
887 return Ok(PeerConnectSuccess {
888 remote,
889 channel: *channel,
890 udp_channel_rx: udp_rx_opt,
891 incoming_object_transfer_handles: None,
892 });
893 }
894
895 NodeResult::PeerEvent(PeerEvent {
896 event:
897 PeerSignal::PostConnect {
898 invitee_response, ..
899 },
900 ..
901 }) => match invitee_response {
902 Some(PeerResponse::Timeout) => {
903 return Err(citadel_io::error!(
904 citadel_io::ErrorCode::RemotePeerNoResponse
905 ))
906 }
907 Some(PeerResponse::Decline) => {
908 return Err(citadel_io::error!(
909 citadel_io::ErrorCode::RemotePeerDeclined
910 ))
911 }
912 _ => {}
913 },
914
915 NodeResult::PeerEvent(PeerEvent {
919 event: PeerSignal::SignalError { error, .. },
920 ..
921 }) => {
922 return Err(citadel_io::error!(
923 citadel_io::ErrorCode::RemoteP2pSignalRoutingFailed,
924 error
925 ))
926 }
927
928 _ => {}
929 }
930 }
931
932 Err(citadel_io::error!(
933 citadel_io::ErrorCode::RemoteKernelStreamDied,
934 "connect_to_peer_custom"
935 ))
936 };
937
938 match citadel_io::time::timeout(P2P_CONNECT_TIMEOUT, connect_task).await {
939 Ok(result) => result,
940 Err(_elapsed) => Err(citadel_io::error!(
941 citadel_io::ErrorCode::RemoteP2pConnectTimeout,
942 P2P_CONNECT_TIMEOUT.as_secs()
943 )),
944 }
945 }
946
947 async fn connect_to_peer(&self) -> Result<PeerConnectSuccess<R>, NetworkError> {
949 self.connect_to_peer_custom(Default::default(), Default::default(), Default::default())
950 .await
951 }
952
953 async fn register_to_peer(&self) -> Result<PeerRegisterStatus, NetworkError> {
955 let session_cid = self.user().get_session_cid();
956 let peer_target = self.try_as_peer_connection().await?;
957 let local_username = self
959 .remote()
960 .account_manager()
961 .get_username_by_cid(session_cid)
962 .await?
963 .ok_or_else(|| citadel_io::error!(citadel_io::ErrorCode::RemoteLocalUsernameMissing))?;
964 let peer_username_opt = self.target_username().map(ToString::to_string);
965
966 let mut stream = self
967 .remote()
968 .send_callback_subscription(NodeRequest::PeerCommand(PeerCommand {
969 session_cid,
970 command: PeerSignal::PostRegister {
971 peer_conn_type: peer_target,
972 inviter_username: local_username,
973 invitee_username: peer_username_opt,
974 ticket_opt: None,
975 invitee_response: None,
976 },
977 }))
978 .await?;
979
980 while let Some(status) = stream.next().await {
981 if let NodeResult::PeerEvent(PeerEvent {
982 event:
983 PeerSignal::PostRegister {
984 peer_conn_type: _,
985 inviter_username: _,
986 invitee_username: _,
987 ticket_opt: _,
988 invitee_response: Some(resp),
989 },
990 ..
991 }) = status.into_result()?
992 {
993 match resp {
994 PeerResponse::Accept(..) => return Ok(PeerRegisterStatus::Accepted),
995 PeerResponse::Decline => return Ok(PeerRegisterStatus::Declined),
996 PeerResponse::Timeout => return Ok(PeerRegisterStatus::Failed { reason: Some("Timeout on register request. Peer did not accept in time. Try again later".to_string()) }),
997 _ => {}
998 }
999 }
1000 }
1001
1002 Err(citadel_io::error!(
1003 citadel_io::ErrorCode::RemoteKernelStreamDied,
1004 format!("register_to_peer: {:?}", stream.callback_key())
1005 ))
1006 }
1007
1008 async fn deregister(&self) -> Result<(), NetworkError> {
1012 if let Ok(peer_conn) = self.try_as_peer_connection().await {
1013 let peer_request = PeerSignal::Deregister {
1014 peer_conn_type: peer_conn,
1015 };
1016 let session_cid = self.user().get_session_cid();
1017 let request = NodeRequest::PeerCommand(PeerCommand {
1018 session_cid,
1019 command: peer_request,
1020 });
1021
1022 let mut subscription = self.remote().send_callback_subscription(request).await?;
1023 while let Some(result) = subscription.next().await {
1024 if let NodeResult::PeerEvent(PeerEvent {
1025 event: PeerSignal::DeregistrationSuccess { .. },
1026 ..
1027 }) = result.into_result()?
1028 {
1029 return Ok(());
1030 }
1031 }
1032 } else {
1033 let cid = self.user().get_session_cid();
1035 let request = NodeRequest::DeregisterFromHypernode(DeregisterFromHypernode {
1036 session_cid: cid,
1037 v_conn_type: *self.user(),
1038 });
1039 let mut subscription = self.remote().send_callback_subscription(request).await?;
1040 while let Some(result) = subscription.next().await {
1041 match result.into_result()? {
1042 NodeResult::DeRegistration(DeRegistration {
1043 session_cid: _,
1044 ticket_opt: _,
1045 success: true,
1046 }) => return Ok(()),
1047 NodeResult::DeRegistration(DeRegistration {
1048 session_cid: _,
1049 ticket_opt: _,
1050 success: false,
1051 }) => {
1052 return Err(citadel_io::error!(
1053 citadel_io::ErrorCode::RemoteDeregisterFailed
1054 ))
1055 }
1056
1057 _ => {}
1058 }
1059 }
1060 }
1061
1062 Err(citadel_io::error!(
1063 citadel_io::ErrorCode::RemoteDeregisterEndedUnexpectedly
1064 ))
1065 }
1066
1067 async fn disconnect(&self) -> Result<(), NetworkError> {
1068 if let Ok(peer_conn) = self.try_as_peer_connection().await {
1069 if let PeerConnectionType::LocalGroupPeer {
1070 session_cid,
1071 peer_cid: _,
1072 } = peer_conn
1073 {
1074 let request = NodeRequest::PeerCommand(PeerCommand {
1075 session_cid,
1076 command: PeerSignal::Disconnect {
1077 peer_conn_type: peer_conn,
1078 disconnect_response: None,
1079 disconnect_token: None,
1080 },
1081 });
1082
1083 let mut subscription = self.remote().send_callback_subscription(request).await?;
1084
1085 while let Some(event) = subscription.next().await {
1086 if let NodeResult::PeerEvent(PeerEvent {
1087 event:
1088 PeerSignal::Disconnect {
1089 peer_conn_type: _,
1090 disconnect_response: Some(_),
1091 ..
1092 },
1093 ..
1094 }) = event.into_result()?
1095 {
1096 return Ok(());
1097 }
1098 }
1099
1100 Err(citadel_io::error!(
1101 citadel_io::ErrorCode::RemoteDisconnectEventMissing
1102 ))
1103 } else {
1104 Err(citadel_io::error!(
1105 citadel_io::ErrorCode::RemoteExternalGroupPeerUnsupported
1106 ))
1107 }
1108 } else {
1109 let cid = self.user().get_session_cid();
1111 let request =
1112 NodeRequest::DisconnectFromHypernode(DisconnectFromHypernode { session_cid: cid });
1113
1114 let mut subscription = self.remote().send_callback_subscription(request).await?;
1115 while let Some(event) = subscription.next().await {
1116 if let NodeResult::Disconnect(Disconnect {
1117 success, message, ..
1118 }) = event.into_result()?
1119 {
1120 return if success {
1121 Ok(())
1122 } else {
1123 Err(citadel_io::error!(
1124 citadel_io::ErrorCode::RemoteDisconnected,
1125 message
1126 ))
1127 };
1128 }
1129 }
1130
1131 Err(citadel_io::error!(
1132 citadel_io::ErrorCode::RemoteDisconnectEventMissing
1133 ))
1134 }
1135 }
1136
1137 async fn create_group(
1138 &self,
1139 initial_users_to_invite: Option<Vec<UserIdentifier>>,
1140 ) -> Result<GroupChannel, NetworkError> {
1141 self.create_group_with_options(initial_users_to_invite, MessageGroupOptions::default())
1142 .await
1143 }
1144
1145 async fn create_group_with_options(
1150 &self,
1151 initial_users_to_invite: Option<Vec<UserIdentifier>>,
1152 options: MessageGroupOptions,
1153 ) -> Result<GroupChannel, NetworkError> {
1154 let session_cid = self.user().get_session_cid();
1155
1156 let mut initial_users = vec![];
1157 if let Some(initial_users_to_invite) = initial_users_to_invite {
1160 for user in initial_users_to_invite {
1161 initial_users.push(
1162 self.remote()
1163 .account_manager()
1164 .find_target_information(session_cid, user.clone())
1165 .await?
1166 .ok_or_else(|| {
1167 citadel_io::error!(
1168 citadel_io::ErrorCode::RemoteGroupAccountNotFound,
1169 citadel_io::Dbg(user),
1170 citadel_io::Dbg(session_cid)
1171 )
1172 })
1173 .map(|r| r.1.cid)?,
1174 )
1175 }
1176 }
1177
1178 let group_request = GroupBroadcast::Create {
1179 initial_invitees: initial_users,
1180 options,
1181 };
1182 let request = NodeRequest::GroupBroadcastCommand(GroupBroadcastCommand {
1183 session_cid,
1184 command: group_request,
1185 });
1186 let mut subscription = self.remote().send_callback_subscription(request).await?;
1187 while let Some(evt) = subscription.next().await {
1188 if let NodeResult::GroupChannelCreated(GroupChannelCreated {
1189 ticket: _,
1190 channel,
1191 session_cid: _,
1192 }) = evt
1193 {
1194 return Ok(channel);
1195 }
1196 }
1197
1198 Err(citadel_io::error!(
1199 citadel_io::ErrorCode::RemoteCreateGroupEndedUnexpectedly
1200 ))
1201 }
1202
1203 async fn list_owned_groups(&self) -> Result<Vec<MessageGroupKey>, NetworkError> {
1205 let session_cid = self.user().get_session_cid();
1206 let cid_to_check_for = match self.try_as_peer_connection().await {
1207 Ok(res) => res.get_original_target_cid(),
1208 _ => session_cid,
1209 };
1210 let group_request = GroupBroadcast::ListGroupsFor {
1211 cid: cid_to_check_for,
1212 };
1213 let request = NodeRequest::GroupBroadcastCommand(GroupBroadcastCommand {
1214 session_cid,
1215 command: group_request,
1216 });
1217
1218 let mut subscription = self.remote().send_callback_subscription(request).await?;
1219
1220 while let Some(evt) = subscription.next().await {
1221 if let NodeResult::GroupEvent(GroupEvent {
1222 session_cid: _,
1223 ticket: _,
1224 event: GroupBroadcast::ListResponse { groups },
1225 }) = evt.into_result()?
1226 {
1227 return Ok(groups);
1228 }
1229 }
1230
1231 Err(citadel_io::error!(
1232 citadel_io::ErrorCode::RemoteListGroupsEndedUnexpectedly
1233 ))
1234 }
1235
1236 async fn list_sessions(&self) -> Result<ActiveSessions, NetworkError> {
1240 let request = NodeRequest::GetActiveSessions;
1241 let mut subscription = self.remote().send_callback_subscription(request).await?;
1242
1243 if let Some(NodeResult::SessionList(result)) = subscription.next().await {
1244 return Ok(result.sessions);
1245 }
1246
1247 Err(citadel_io::error!(
1248 citadel_io::ErrorCode::RemoteListSessionsEndedUnexpectedly
1249 ))
1250 }
1251
1252 async fn rekey(&self) -> Result<Option<u32>, NetworkError> {
1256 let request = NodeRequest::ReKey(ReKey {
1257 v_conn_type: *self.user(),
1258 });
1259 let mut subscription = self.remote().send_callback_subscription(request).await?;
1260
1261 while let Some(evt) = subscription.next().await {
1262 if let NodeResult::ReKeyResult(result) = evt {
1263 return match result.status {
1264 ReKeyReturnType::Success { version } => Ok(Some(version)),
1265 ReKeyReturnType::AlreadyInProgress => Ok(None),
1266 ReKeyReturnType::Failure { err } => Err(citadel_io::error!(
1267 citadel_io::ErrorCode::RemoteRekeyFailed,
1268 err
1269 )),
1270 };
1271 }
1272 }
1273
1274 Err(citadel_io::error!(
1275 citadel_io::ErrorCode::RemoteRekeyEndedUnexpectedly
1276 ))
1277 }
1278
1279 async fn is_peer_registered(&self) -> Result<bool, NetworkError> {
1281 let target = self.try_as_peer_connection().await?;
1282 if let PeerConnectionType::LocalGroupPeer {
1283 session_cid: local_cid,
1284 peer_cid,
1285 } = target
1286 {
1287 let peers = self.remote().get_local_group_peers(local_cid, None).await?;
1288 citadel_logging::info!(target: "citadel", "Checking to see if {target} is registered in {peers:?}");
1289 Ok(peers.iter().any(|p| p.cid == peer_cid))
1290 } else {
1291 Err(citadel_io::error!(
1292 citadel_io::ErrorCode::RemoteExternalGroupPeerUnsupportedYet
1293 ))
1294 }
1295 }
1296
1297 #[doc(hidden)]
1298 async fn try_as_peer_connection(&self) -> Result<PeerConnectionType, NetworkError> {
1299 let verified_return = |user: &VirtualTargetType| {
1300 user.try_as_peer_connection().ok_or(citadel_io::error!(
1301 citadel_io::ErrorCode::RemoteTargetNotPeer
1302 ))
1303 };
1304
1305 if self.user().get_target_cid() == 0 {
1306 let peer_username = self.target_username().ok_or_else(|| {
1309 citadel_io::error!(citadel_io::ErrorCode::RemoteTargetCidZeroNoUsername)
1310 })?;
1311 let session_cid = self.user().get_session_cid();
1312 let expected_peer_cid = self
1313 .remote()
1314 .account_manager()
1315 .get_persistence_handler()
1316 .get_cid_by_username(peer_username);
1317 let peer_cid = self
1320 .remote()
1321 .account_manager()
1322 .find_target_information(session_cid, peer_username)
1323 .await?
1324 .map(|r| r.1.cid)
1325 .unwrap_or(expected_peer_cid);
1326
1327 let mut user = *self.user();
1328 user.set_target_cid(peer_cid);
1329 verified_return(&user)
1330 } else {
1331 verified_return(self.user())
1332 }
1333 }
1334
1335 #[doc(hidden)]
1336 fn can_use_revfs(&self) -> Result<(), NetworkError> {
1337 if let Some(sess) = self.session_security_settings() {
1338 if sess.crypto_params.kem_algorithm == KemAlgorithm::MlKem {
1339 Ok(())
1340 } else {
1341 Err(citadel_io::error!(
1342 citadel_io::ErrorCode::RemoteRevfsRequiresKyber
1343 ))
1344 }
1345 } else {
1346 Err(citadel_io::error!(
1347 citadel_io::ErrorCode::RemoteRevfsUnsupportedRemote
1348 ))
1349 }
1350 }
1351}
1352
1353impl<T: TargetLockedRemote<R>, R: Ratchet> ProtocolRemoteTargetExt<R> for T {}
1354
1355pub mod results {
1356 use crate::prefabs::client::peer_connection::FileTransferHandleRx;
1357 use crate::prelude::{PeerChannel, UdpChannel};
1358 use crate::remote_ext::remote_specialization::PeerRemote;
1359 use citadel_io::tokio::sync::oneshot::Receiver;
1360 use citadel_proto::prelude::*;
1361 use std::fmt::Debug;
1362
1363 pub struct PeerConnectSuccess<R: Ratchet> {
1364 pub channel: PeerChannel<R>,
1365 pub udp_channel_rx: Option<Receiver<UdpChannel<R>>>,
1366 pub remote: PeerRemote<R>,
1367 pub(crate) incoming_object_transfer_handles: Option<FileTransferHandleRx>,
1370 }
1371
1372 impl<R: Ratchet> Debug for PeerConnectSuccess<R> {
1373 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1374 f.debug_struct("PeerConnectSuccess")
1375 .field("channel", &self.channel)
1376 .field("udp_channel_rx", &self.udp_channel_rx)
1377 .finish()
1378 }
1379 }
1380
1381 impl<R: Ratchet> PeerConnectSuccess<R> {
1382 pub fn get_incoming_file_transfer_handle(
1384 &mut self,
1385 ) -> Result<FileTransferHandleRx, NetworkError> {
1386 self.incoming_object_transfer_handles
1387 .take()
1388 .ok_or(citadel_io::error!(
1389 citadel_io::ErrorCode::RemoteFunctionAlreadyCalled
1390 ))
1391 }
1392 }
1393
1394 pub enum PeerRegisterStatus {
1395 Accepted,
1396 Declined,
1397 Failed { reason: Option<String> },
1398 }
1399
1400 #[derive(Clone, Debug)]
1401 pub struct LocalGroupPeer {
1402 pub cid: u64,
1403 pub is_online: bool,
1404 }
1405
1406 #[derive(Clone, Debug)]
1407 pub struct LocalGroupPeerFullInfo {
1408 pub cid: u64,
1409 pub username: Option<String>,
1410 pub full_name: Option<String>,
1411 pub is_online: bool,
1412 }
1413}
1414
1415pub mod remote_specialization {
1416 use crate::prelude::*;
1417 use std::ops::{Deref, DerefMut};
1418
1419 #[derive(Debug, Clone)]
1420 pub struct PeerRemote<R: Ratchet> {
1421 pub(crate) inner: NodeRemote<R>,
1422 pub(crate) peer: VirtualTargetType,
1423 pub(crate) username: Option<String>,
1424 pub(crate) session_security_settings: SessionSecuritySettings,
1425 }
1426
1427 impl<R: Ratchet> Deref for PeerRemote<R> {
1428 type Target = NodeRemote<R>;
1429 fn deref(&self) -> &Self::Target {
1430 &self.inner
1431 }
1432 }
1433
1434 impl<R: Ratchet> DerefMut for PeerRemote<R> {
1435 fn deref_mut(&mut self) -> &mut Self::Target {
1436 &mut self.inner
1437 }
1438 }
1439
1440 impl<R: Ratchet> TargetLockedRemote<R> for PeerRemote<R> {
1441 fn user(&self) -> &VirtualTargetType {
1442 &self.peer
1443 }
1444 fn remote(&self) -> &NodeRemote<R> {
1445 &self.inner
1446 }
1447 fn target_username(&self) -> Option<&str> {
1448 self.username.as_deref()
1449 }
1450 fn user_mut(&mut self) -> &mut VirtualTargetType {
1451 &mut self.peer
1452 }
1453
1454 fn session_security_settings(&self) -> Option<&SessionSecuritySettings> {
1455 Some(&self.session_security_settings)
1456 }
1457 }
1458}
1459
1460#[cfg(all(test, not(target_family = "wasm")))]
1461mod tests {
1462 use crate::prefabs::client::single_connection::SingleClientServerConnectionKernel;
1463 use crate::prefabs::client::DefaultServerConnectionSettingsBuilder;
1464 use crate::prelude::*;
1465 use citadel_io::tokio;
1466 use rstest::rstest;
1467 use std::net::SocketAddr;
1468 use std::sync::atomic::{AtomicBool, Ordering};
1469 use std::sync::Arc;
1470 use uuid::Uuid;
1471
1472 pub struct ReceiverFileTransferKernel<R: Ratchet>(
1473 pub Option<NodeRemote<R>>,
1474 pub Arc<AtomicBool>,
1475 );
1476
1477 #[async_trait]
1478 impl<R: Ratchet> NetKernel<R> for ReceiverFileTransferKernel<R> {
1479 fn load_remote(&mut self, node_remote: NodeRemote<R>) -> Result<(), NetworkError> {
1480 self.0 = Some(node_remote);
1481 Ok(())
1482 }
1483
1484 async fn on_start(&self) -> Result<(), NetworkError> {
1485 Ok(())
1486 }
1487
1488 async fn on_node_event_received(&self, message: NodeResult<R>) -> Result<(), NetworkError> {
1489 log::trace!(target: "citadel", "SERVER received {:?}", message);
1490 if let NodeResult::ObjectTransferHandle(ObjectTransferHandle { mut handle, .. }) =
1491 message.into_result()?
1492 {
1493 let mut path = None;
1494 handle
1496 .accept()
1497 .map_err(|err| NetworkError::msg(err.into_string()))?;
1498
1499 use citadel_types::proto::ObjectTransferStatus;
1500 use futures::StreamExt;
1501 while let Some(status) = handle.next().await {
1502 match status {
1503 ObjectTransferStatus::ReceptionComplete => {
1504 log::trace!(target: "citadel", "Server has finished receiving the file!");
1505 let cmp = include_bytes!("../../resources/TheBridge.pdf");
1506 let streamed_data = citadel_io::tokio::fs::read(path.clone().unwrap())
1507 .await
1508 .unwrap();
1509 assert_eq!(
1510 cmp,
1511 streamed_data.as_slice(),
1512 "Original data and streamed data does not match"
1513 );
1514
1515 self.1.store(true, Ordering::Relaxed);
1516 self.0.clone().unwrap().shutdown().await?;
1517 }
1518
1519 ObjectTransferStatus::ReceptionBeginning(file_path, vfm) => {
1520 path = Some(file_path);
1521 assert_eq!(vfm.name, "TheBridge.pdf")
1522 }
1523
1524 _ => {}
1525 }
1526 }
1527 }
1528
1529 Ok(())
1530 }
1531
1532 async fn on_stop(&mut self) -> Result<(), NetworkError> {
1533 Ok(())
1534 }
1535 }
1536
1537 pub fn server_info<'a, R: Ratchet>(
1538 switch: Arc<AtomicBool>,
1539 ) -> (NodeFuture<'a, ReceiverFileTransferKernel<R>>, SocketAddr) {
1540 crate::test_common::server_test_node(ReceiverFileTransferKernel(None, switch), |_| {})
1541 }
1542
1543 #[rstest]
1544 #[case(
1545 EncryptionAlgorithm::AES_GCM_256,
1546 KemAlgorithm::MlKem,
1547 SigAlgorithm::None
1548 )]
1549 #[case(
1550 EncryptionAlgorithm::MlKemHybrid,
1551 KemAlgorithm::MlKem,
1552 SigAlgorithm::MlDsa65
1553 )]
1554 #[timeout(std::time::Duration::from_secs(90))]
1555 #[tokio::test]
1556 async fn test_c2s_file_transfer(
1557 #[case] enx: EncryptionAlgorithm,
1558 #[case] kem: KemAlgorithm,
1559 #[case] sig: SigAlgorithm,
1560 ) {
1561 citadel_logging::setup_log();
1562 let client_success = &AtomicBool::new(false);
1563 let server_success = &Arc::new(AtomicBool::new(false));
1564 let (server, server_addr) = server_info::<StackedRatchet>(server_success.clone());
1565 let uuid = Uuid::new_v4();
1566
1567 let session_security_settings = SessionSecuritySettingsBuilder::default()
1568 .with_crypto_params(enx + kem + sig)
1569 .build()
1570 .unwrap();
1571
1572 let server_connection_settings =
1573 DefaultServerConnectionSettingsBuilder::transient_with_id(server_addr, uuid)
1574 .with_session_security_settings(session_security_settings)
1575 .disable_udp()
1576 .build()
1577 .unwrap();
1578
1579 let client_kernel = SingleClientServerConnectionKernel::new(
1580 server_connection_settings,
1581 |connection| async move {
1582 log::trace!(target: "citadel", "***CLIENT LOGIN SUCCESS :: File transfer next ***");
1583 connection
1584 .send_file_with_custom_opts(
1585 "../resources/TheBridge.pdf",
1586 32 * 1024,
1587 TransferType::FileTransfer,
1588 )
1589 .await
1590 .unwrap();
1591 log::trace!(target: "citadel", "***CLIENT FILE TRANSFER SUCCESS***");
1592 client_success.store(true, Ordering::Relaxed);
1593 connection.shutdown_kernel().await
1594 },
1595 );
1596
1597 let client = DefaultNodeBuilder::default().build(client_kernel).unwrap();
1598
1599 let joined = futures::future::try_join(server, client);
1600
1601 let _ = joined.await.unwrap();
1602
1603 assert!(client_success.load(Ordering::Relaxed));
1604 assert!(server_success.load(Ordering::Relaxed));
1605 }
1606}