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 _ => {}
916 }
917 }
918
919 Err(citadel_io::error!(
920 citadel_io::ErrorCode::RemoteKernelStreamDied,
921 "connect_to_peer_custom"
922 ))
923 };
924
925 match citadel_io::time::timeout(P2P_CONNECT_TIMEOUT, connect_task).await {
926 Ok(result) => result,
927 Err(_elapsed) => Err(citadel_io::error!(
928 citadel_io::ErrorCode::RemoteP2pConnectTimeout,
929 P2P_CONNECT_TIMEOUT.as_secs()
930 )),
931 }
932 }
933
934 async fn connect_to_peer(&self) -> Result<PeerConnectSuccess<R>, NetworkError> {
936 self.connect_to_peer_custom(Default::default(), Default::default(), Default::default())
937 .await
938 }
939
940 async fn register_to_peer(&self) -> Result<PeerRegisterStatus, NetworkError> {
942 let session_cid = self.user().get_session_cid();
943 let peer_target = self.try_as_peer_connection().await?;
944 let local_username = self
946 .remote()
947 .account_manager()
948 .get_username_by_cid(session_cid)
949 .await?
950 .ok_or_else(|| citadel_io::error!(citadel_io::ErrorCode::RemoteLocalUsernameMissing))?;
951 let peer_username_opt = self.target_username().map(ToString::to_string);
952
953 let mut stream = self
954 .remote()
955 .send_callback_subscription(NodeRequest::PeerCommand(PeerCommand {
956 session_cid,
957 command: PeerSignal::PostRegister {
958 peer_conn_type: peer_target,
959 inviter_username: local_username,
960 invitee_username: peer_username_opt,
961 ticket_opt: None,
962 invitee_response: None,
963 },
964 }))
965 .await?;
966
967 while let Some(status) = stream.next().await {
968 if let NodeResult::PeerEvent(PeerEvent {
969 event:
970 PeerSignal::PostRegister {
971 peer_conn_type: _,
972 inviter_username: _,
973 invitee_username: _,
974 ticket_opt: _,
975 invitee_response: Some(resp),
976 },
977 ..
978 }) = status.into_result()?
979 {
980 match resp {
981 PeerResponse::Accept(..) => return Ok(PeerRegisterStatus::Accepted),
982 PeerResponse::Decline => return Ok(PeerRegisterStatus::Declined),
983 PeerResponse::Timeout => return Ok(PeerRegisterStatus::Failed { reason: Some("Timeout on register request. Peer did not accept in time. Try again later".to_string()) }),
984 _ => {}
985 }
986 }
987 }
988
989 Err(citadel_io::error!(
990 citadel_io::ErrorCode::RemoteKernelStreamDied,
991 format!("register_to_peer: {:?}", stream.callback_key())
992 ))
993 }
994
995 async fn deregister(&self) -> Result<(), NetworkError> {
999 if let Ok(peer_conn) = self.try_as_peer_connection().await {
1000 let peer_request = PeerSignal::Deregister {
1001 peer_conn_type: peer_conn,
1002 };
1003 let session_cid = self.user().get_session_cid();
1004 let request = NodeRequest::PeerCommand(PeerCommand {
1005 session_cid,
1006 command: peer_request,
1007 });
1008
1009 let mut subscription = self.remote().send_callback_subscription(request).await?;
1010 while let Some(result) = subscription.next().await {
1011 if let NodeResult::PeerEvent(PeerEvent {
1012 event: PeerSignal::DeregistrationSuccess { .. },
1013 ..
1014 }) = result.into_result()?
1015 {
1016 return Ok(());
1017 }
1018 }
1019 } else {
1020 let cid = self.user().get_session_cid();
1022 let request = NodeRequest::DeregisterFromHypernode(DeregisterFromHypernode {
1023 session_cid: cid,
1024 v_conn_type: *self.user(),
1025 });
1026 let mut subscription = self.remote().send_callback_subscription(request).await?;
1027 while let Some(result) = subscription.next().await {
1028 match result.into_result()? {
1029 NodeResult::DeRegistration(DeRegistration {
1030 session_cid: _,
1031 ticket_opt: _,
1032 success: true,
1033 }) => return Ok(()),
1034 NodeResult::DeRegistration(DeRegistration {
1035 session_cid: _,
1036 ticket_opt: _,
1037 success: false,
1038 }) => {
1039 return Err(citadel_io::error!(
1040 citadel_io::ErrorCode::RemoteDeregisterFailed
1041 ))
1042 }
1043
1044 _ => {}
1045 }
1046 }
1047 }
1048
1049 Err(citadel_io::error!(
1050 citadel_io::ErrorCode::RemoteDeregisterEndedUnexpectedly
1051 ))
1052 }
1053
1054 async fn disconnect(&self) -> Result<(), NetworkError> {
1055 if let Ok(peer_conn) = self.try_as_peer_connection().await {
1056 if let PeerConnectionType::LocalGroupPeer {
1057 session_cid,
1058 peer_cid: _,
1059 } = peer_conn
1060 {
1061 let request = NodeRequest::PeerCommand(PeerCommand {
1062 session_cid,
1063 command: PeerSignal::Disconnect {
1064 peer_conn_type: peer_conn,
1065 disconnect_response: None,
1066 disconnect_token: None,
1067 },
1068 });
1069
1070 let mut subscription = self.remote().send_callback_subscription(request).await?;
1071
1072 while let Some(event) = subscription.next().await {
1073 if let NodeResult::PeerEvent(PeerEvent {
1074 event:
1075 PeerSignal::Disconnect {
1076 peer_conn_type: _,
1077 disconnect_response: Some(_),
1078 ..
1079 },
1080 ..
1081 }) = event.into_result()?
1082 {
1083 return Ok(());
1084 }
1085 }
1086
1087 Err(citadel_io::error!(
1088 citadel_io::ErrorCode::RemoteDisconnectEventMissing
1089 ))
1090 } else {
1091 Err(citadel_io::error!(
1092 citadel_io::ErrorCode::RemoteExternalGroupPeerUnsupported
1093 ))
1094 }
1095 } else {
1096 let cid = self.user().get_session_cid();
1098 let request =
1099 NodeRequest::DisconnectFromHypernode(DisconnectFromHypernode { session_cid: cid });
1100
1101 let mut subscription = self.remote().send_callback_subscription(request).await?;
1102 while let Some(event) = subscription.next().await {
1103 if let NodeResult::Disconnect(Disconnect {
1104 success, message, ..
1105 }) = event.into_result()?
1106 {
1107 return if success {
1108 Ok(())
1109 } else {
1110 Err(citadel_io::error!(
1111 citadel_io::ErrorCode::RemoteDisconnected,
1112 message
1113 ))
1114 };
1115 }
1116 }
1117
1118 Err(citadel_io::error!(
1119 citadel_io::ErrorCode::RemoteDisconnectEventMissing
1120 ))
1121 }
1122 }
1123
1124 async fn create_group(
1125 &self,
1126 initial_users_to_invite: Option<Vec<UserIdentifier>>,
1127 ) -> Result<GroupChannel, NetworkError> {
1128 self.create_group_with_options(initial_users_to_invite, MessageGroupOptions::default())
1129 .await
1130 }
1131
1132 async fn create_group_with_options(
1137 &self,
1138 initial_users_to_invite: Option<Vec<UserIdentifier>>,
1139 options: MessageGroupOptions,
1140 ) -> Result<GroupChannel, NetworkError> {
1141 let session_cid = self.user().get_session_cid();
1142
1143 let mut initial_users = vec![];
1144 if let Some(initial_users_to_invite) = initial_users_to_invite {
1147 for user in initial_users_to_invite {
1148 initial_users.push(
1149 self.remote()
1150 .account_manager()
1151 .find_target_information(session_cid, user.clone())
1152 .await?
1153 .ok_or_else(|| {
1154 citadel_io::error!(
1155 citadel_io::ErrorCode::RemoteGroupAccountNotFound,
1156 citadel_io::Dbg(user),
1157 citadel_io::Dbg(session_cid)
1158 )
1159 })
1160 .map(|r| r.1.cid)?,
1161 )
1162 }
1163 }
1164
1165 let group_request = GroupBroadcast::Create {
1166 initial_invitees: initial_users,
1167 options,
1168 };
1169 let request = NodeRequest::GroupBroadcastCommand(GroupBroadcastCommand {
1170 session_cid,
1171 command: group_request,
1172 });
1173 let mut subscription = self.remote().send_callback_subscription(request).await?;
1174 while let Some(evt) = subscription.next().await {
1175 if let NodeResult::GroupChannelCreated(GroupChannelCreated {
1176 ticket: _,
1177 channel,
1178 session_cid: _,
1179 }) = evt
1180 {
1181 return Ok(channel);
1182 }
1183 }
1184
1185 Err(citadel_io::error!(
1186 citadel_io::ErrorCode::RemoteCreateGroupEndedUnexpectedly
1187 ))
1188 }
1189
1190 async fn list_owned_groups(&self) -> Result<Vec<MessageGroupKey>, NetworkError> {
1192 let session_cid = self.user().get_session_cid();
1193 let cid_to_check_for = match self.try_as_peer_connection().await {
1194 Ok(res) => res.get_original_target_cid(),
1195 _ => session_cid,
1196 };
1197 let group_request = GroupBroadcast::ListGroupsFor {
1198 cid: cid_to_check_for,
1199 };
1200 let request = NodeRequest::GroupBroadcastCommand(GroupBroadcastCommand {
1201 session_cid,
1202 command: group_request,
1203 });
1204
1205 let mut subscription = self.remote().send_callback_subscription(request).await?;
1206
1207 while let Some(evt) = subscription.next().await {
1208 if let NodeResult::GroupEvent(GroupEvent {
1209 session_cid: _,
1210 ticket: _,
1211 event: GroupBroadcast::ListResponse { groups },
1212 }) = evt.into_result()?
1213 {
1214 return Ok(groups);
1215 }
1216 }
1217
1218 Err(citadel_io::error!(
1219 citadel_io::ErrorCode::RemoteListGroupsEndedUnexpectedly
1220 ))
1221 }
1222
1223 async fn list_sessions(&self) -> Result<ActiveSessions, NetworkError> {
1227 let request = NodeRequest::GetActiveSessions;
1228 let mut subscription = self.remote().send_callback_subscription(request).await?;
1229
1230 if let Some(NodeResult::SessionList(result)) = subscription.next().await {
1231 return Ok(result.sessions);
1232 }
1233
1234 Err(citadel_io::error!(
1235 citadel_io::ErrorCode::RemoteListSessionsEndedUnexpectedly
1236 ))
1237 }
1238
1239 async fn rekey(&self) -> Result<Option<u32>, NetworkError> {
1243 let request = NodeRequest::ReKey(ReKey {
1244 v_conn_type: *self.user(),
1245 });
1246 let mut subscription = self.remote().send_callback_subscription(request).await?;
1247
1248 while let Some(evt) = subscription.next().await {
1249 if let NodeResult::ReKeyResult(result) = evt {
1250 return match result.status {
1251 ReKeyReturnType::Success { version } => Ok(Some(version)),
1252 ReKeyReturnType::AlreadyInProgress => Ok(None),
1253 ReKeyReturnType::Failure { err } => Err(citadel_io::error!(
1254 citadel_io::ErrorCode::RemoteRekeyFailed,
1255 err
1256 )),
1257 };
1258 }
1259 }
1260
1261 Err(citadel_io::error!(
1262 citadel_io::ErrorCode::RemoteRekeyEndedUnexpectedly
1263 ))
1264 }
1265
1266 async fn is_peer_registered(&self) -> Result<bool, NetworkError> {
1268 let target = self.try_as_peer_connection().await?;
1269 if let PeerConnectionType::LocalGroupPeer {
1270 session_cid: local_cid,
1271 peer_cid,
1272 } = target
1273 {
1274 let peers = self.remote().get_local_group_peers(local_cid, None).await?;
1275 citadel_logging::info!(target: "citadel", "Checking to see if {target} is registered in {peers:?}");
1276 Ok(peers.iter().any(|p| p.cid == peer_cid))
1277 } else {
1278 Err(citadel_io::error!(
1279 citadel_io::ErrorCode::RemoteExternalGroupPeerUnsupportedYet
1280 ))
1281 }
1282 }
1283
1284 #[doc(hidden)]
1285 async fn try_as_peer_connection(&self) -> Result<PeerConnectionType, NetworkError> {
1286 let verified_return = |user: &VirtualTargetType| {
1287 user.try_as_peer_connection().ok_or(citadel_io::error!(
1288 citadel_io::ErrorCode::RemoteTargetNotPeer
1289 ))
1290 };
1291
1292 if self.user().get_target_cid() == 0 {
1293 let peer_username = self.target_username().ok_or_else(|| {
1296 citadel_io::error!(citadel_io::ErrorCode::RemoteTargetCidZeroNoUsername)
1297 })?;
1298 let session_cid = self.user().get_session_cid();
1299 let expected_peer_cid = self
1300 .remote()
1301 .account_manager()
1302 .get_persistence_handler()
1303 .get_cid_by_username(peer_username);
1304 let peer_cid = self
1307 .remote()
1308 .account_manager()
1309 .find_target_information(session_cid, peer_username)
1310 .await?
1311 .map(|r| r.1.cid)
1312 .unwrap_or(expected_peer_cid);
1313
1314 let mut user = *self.user();
1315 user.set_target_cid(peer_cid);
1316 verified_return(&user)
1317 } else {
1318 verified_return(self.user())
1319 }
1320 }
1321
1322 #[doc(hidden)]
1323 fn can_use_revfs(&self) -> Result<(), NetworkError> {
1324 if let Some(sess) = self.session_security_settings() {
1325 if sess.crypto_params.kem_algorithm == KemAlgorithm::MlKem {
1326 Ok(())
1327 } else {
1328 Err(citadel_io::error!(
1329 citadel_io::ErrorCode::RemoteRevfsRequiresKyber
1330 ))
1331 }
1332 } else {
1333 Err(citadel_io::error!(
1334 citadel_io::ErrorCode::RemoteRevfsUnsupportedRemote
1335 ))
1336 }
1337 }
1338}
1339
1340impl<T: TargetLockedRemote<R>, R: Ratchet> ProtocolRemoteTargetExt<R> for T {}
1341
1342pub mod results {
1343 use crate::prefabs::client::peer_connection::FileTransferHandleRx;
1344 use crate::prelude::{PeerChannel, UdpChannel};
1345 use crate::remote_ext::remote_specialization::PeerRemote;
1346 use citadel_io::tokio::sync::oneshot::Receiver;
1347 use citadel_proto::prelude::*;
1348 use std::fmt::Debug;
1349
1350 pub struct PeerConnectSuccess<R: Ratchet> {
1351 pub channel: PeerChannel<R>,
1352 pub udp_channel_rx: Option<Receiver<UdpChannel<R>>>,
1353 pub remote: PeerRemote<R>,
1354 pub(crate) incoming_object_transfer_handles: Option<FileTransferHandleRx>,
1357 }
1358
1359 impl<R: Ratchet> Debug for PeerConnectSuccess<R> {
1360 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1361 f.debug_struct("PeerConnectSuccess")
1362 .field("channel", &self.channel)
1363 .field("udp_channel_rx", &self.udp_channel_rx)
1364 .finish()
1365 }
1366 }
1367
1368 impl<R: Ratchet> PeerConnectSuccess<R> {
1369 pub fn get_incoming_file_transfer_handle(
1371 &mut self,
1372 ) -> Result<FileTransferHandleRx, NetworkError> {
1373 self.incoming_object_transfer_handles
1374 .take()
1375 .ok_or(citadel_io::error!(
1376 citadel_io::ErrorCode::RemoteFunctionAlreadyCalled
1377 ))
1378 }
1379 }
1380
1381 pub enum PeerRegisterStatus {
1382 Accepted,
1383 Declined,
1384 Failed { reason: Option<String> },
1385 }
1386
1387 #[derive(Clone, Debug)]
1388 pub struct LocalGroupPeer {
1389 pub cid: u64,
1390 pub is_online: bool,
1391 }
1392
1393 #[derive(Clone, Debug)]
1394 pub struct LocalGroupPeerFullInfo {
1395 pub cid: u64,
1396 pub username: Option<String>,
1397 pub full_name: Option<String>,
1398 pub is_online: bool,
1399 }
1400}
1401
1402pub mod remote_specialization {
1403 use crate::prelude::*;
1404 use std::ops::{Deref, DerefMut};
1405
1406 #[derive(Debug, Clone)]
1407 pub struct PeerRemote<R: Ratchet> {
1408 pub(crate) inner: NodeRemote<R>,
1409 pub(crate) peer: VirtualTargetType,
1410 pub(crate) username: Option<String>,
1411 pub(crate) session_security_settings: SessionSecuritySettings,
1412 }
1413
1414 impl<R: Ratchet> Deref for PeerRemote<R> {
1415 type Target = NodeRemote<R>;
1416 fn deref(&self) -> &Self::Target {
1417 &self.inner
1418 }
1419 }
1420
1421 impl<R: Ratchet> DerefMut for PeerRemote<R> {
1422 fn deref_mut(&mut self) -> &mut Self::Target {
1423 &mut self.inner
1424 }
1425 }
1426
1427 impl<R: Ratchet> TargetLockedRemote<R> for PeerRemote<R> {
1428 fn user(&self) -> &VirtualTargetType {
1429 &self.peer
1430 }
1431 fn remote(&self) -> &NodeRemote<R> {
1432 &self.inner
1433 }
1434 fn target_username(&self) -> Option<&str> {
1435 self.username.as_deref()
1436 }
1437 fn user_mut(&mut self) -> &mut VirtualTargetType {
1438 &mut self.peer
1439 }
1440
1441 fn session_security_settings(&self) -> Option<&SessionSecuritySettings> {
1442 Some(&self.session_security_settings)
1443 }
1444 }
1445}
1446
1447#[cfg(all(test, not(target_family = "wasm")))]
1448mod tests {
1449 use crate::prefabs::client::single_connection::SingleClientServerConnectionKernel;
1450 use crate::prefabs::client::DefaultServerConnectionSettingsBuilder;
1451 use crate::prelude::*;
1452 use citadel_io::tokio;
1453 use rstest::rstest;
1454 use std::net::SocketAddr;
1455 use std::sync::atomic::{AtomicBool, Ordering};
1456 use std::sync::Arc;
1457 use uuid::Uuid;
1458
1459 pub struct ReceiverFileTransferKernel<R: Ratchet>(
1460 pub Option<NodeRemote<R>>,
1461 pub Arc<AtomicBool>,
1462 );
1463
1464 #[async_trait]
1465 impl<R: Ratchet> NetKernel<R> for ReceiverFileTransferKernel<R> {
1466 fn load_remote(&mut self, node_remote: NodeRemote<R>) -> Result<(), NetworkError> {
1467 self.0 = Some(node_remote);
1468 Ok(())
1469 }
1470
1471 async fn on_start(&self) -> Result<(), NetworkError> {
1472 Ok(())
1473 }
1474
1475 async fn on_node_event_received(&self, message: NodeResult<R>) -> Result<(), NetworkError> {
1476 log::trace!(target: "citadel", "SERVER received {:?}", message);
1477 if let NodeResult::ObjectTransferHandle(ObjectTransferHandle { mut handle, .. }) =
1478 message.into_result()?
1479 {
1480 let mut path = None;
1481 handle
1483 .accept()
1484 .map_err(|err| NetworkError::msg(err.into_string()))?;
1485
1486 use citadel_types::proto::ObjectTransferStatus;
1487 use futures::StreamExt;
1488 while let Some(status) = handle.next().await {
1489 match status {
1490 ObjectTransferStatus::ReceptionComplete => {
1491 log::trace!(target: "citadel", "Server has finished receiving the file!");
1492 let cmp = include_bytes!("../../resources/TheBridge.pdf");
1493 let streamed_data = citadel_io::tokio::fs::read(path.clone().unwrap())
1494 .await
1495 .unwrap();
1496 assert_eq!(
1497 cmp,
1498 streamed_data.as_slice(),
1499 "Original data and streamed data does not match"
1500 );
1501
1502 self.1.store(true, Ordering::Relaxed);
1503 self.0.clone().unwrap().shutdown().await?;
1504 }
1505
1506 ObjectTransferStatus::ReceptionBeginning(file_path, vfm) => {
1507 path = Some(file_path);
1508 assert_eq!(vfm.name, "TheBridge.pdf")
1509 }
1510
1511 _ => {}
1512 }
1513 }
1514 }
1515
1516 Ok(())
1517 }
1518
1519 async fn on_stop(&mut self) -> Result<(), NetworkError> {
1520 Ok(())
1521 }
1522 }
1523
1524 pub fn server_info<'a, R: Ratchet>(
1525 switch: Arc<AtomicBool>,
1526 ) -> (NodeFuture<'a, ReceiverFileTransferKernel<R>>, SocketAddr) {
1527 crate::test_common::server_test_node(ReceiverFileTransferKernel(None, switch), |_| {})
1528 }
1529
1530 #[rstest]
1531 #[case(
1532 EncryptionAlgorithm::AES_GCM_256,
1533 KemAlgorithm::MlKem,
1534 SigAlgorithm::None
1535 )]
1536 #[case(
1537 EncryptionAlgorithm::MlKemHybrid,
1538 KemAlgorithm::MlKem,
1539 SigAlgorithm::MlDsa65
1540 )]
1541 #[timeout(std::time::Duration::from_secs(90))]
1542 #[tokio::test]
1543 async fn test_c2s_file_transfer(
1544 #[case] enx: EncryptionAlgorithm,
1545 #[case] kem: KemAlgorithm,
1546 #[case] sig: SigAlgorithm,
1547 ) {
1548 citadel_logging::setup_log();
1549 let client_success = &AtomicBool::new(false);
1550 let server_success = &Arc::new(AtomicBool::new(false));
1551 let (server, server_addr) = server_info::<StackedRatchet>(server_success.clone());
1552 let uuid = Uuid::new_v4();
1553
1554 let session_security_settings = SessionSecuritySettingsBuilder::default()
1555 .with_crypto_params(enx + kem + sig)
1556 .build()
1557 .unwrap();
1558
1559 let server_connection_settings =
1560 DefaultServerConnectionSettingsBuilder::transient_with_id(server_addr, uuid)
1561 .with_session_security_settings(session_security_settings)
1562 .disable_udp()
1563 .build()
1564 .unwrap();
1565
1566 let client_kernel = SingleClientServerConnectionKernel::new(
1567 server_connection_settings,
1568 |connection| async move {
1569 log::trace!(target: "citadel", "***CLIENT LOGIN SUCCESS :: File transfer next ***");
1570 connection
1571 .send_file_with_custom_opts(
1572 "../resources/TheBridge.pdf",
1573 32 * 1024,
1574 TransferType::FileTransfer,
1575 )
1576 .await
1577 .unwrap();
1578 log::trace!(target: "citadel", "***CLIENT FILE TRANSFER SUCCESS***");
1579 client_success.store(true, Ordering::Relaxed);
1580 connection.shutdown_kernel().await
1581 },
1582 );
1583
1584 let client = DefaultNodeBuilder::default().build(client_kernel).unwrap();
1585
1586 let joined = futures::future::try_join(server, client);
1587
1588 let _ = joined.await.unwrap();
1589
1590 assert!(client_success.load(Ordering::Relaxed));
1591 assert!(server_success.load(Ordering::Relaxed));
1592 }
1593}