Skip to main content

citadel_sdk/
remote_ext.rs

1//! Remote Protocol Extensions
2//!
3//! This module extends the core NodeRemote functionality with high-level operations
4//! for managing connections, file transfers, and peer interactions in the Citadel
5//! Protocol network.
6//!
7//! # Features
8//! - User registration and authentication
9//! - Connection management
10//! - File transfer operations
11//! - Virtual filesystem support
12//! - Peer discovery and management
13//! - Group communication
14//! - Security settings configuration
15//!
16//! # Example
17//! ```rust
18//! use citadel_sdk::prelude::*;
19//!
20//! async fn example<R: Ratchet>(remote: NodeRemote<R>) -> Result<(), NetworkError> {
21//!     // Register a new user
22//!     let reg = remote.register_with_defaults(
23//!         "127.0.0.1:25021",
24//!         "John Doe",
25//!         "john.doe",
26//!         "password123"
27//!     ).await?;
28//!
29//!     // Connect to a peer
30//!     let auth = AuthenticationRequest::credentialed("john.doe", "password123");
31//!     let conn = remote.connect_with_defaults(auth).await?;
32//!
33//!     // Send a file to a peer
34//!     remote.find_target("john.doe", "peer.name")
35//!         .await?
36//!         .send_file("/path/to/file.txt")
37//!         .await?;
38//!
39//!     Ok(())
40//! }
41//! ```
42//!
43//! # Important Notes
44//! - All operations are asynchronous
45//! - Connections are automatically managed
46//! - File transfers support chunking
47//! - Virtual filesystem is encrypted
48//! - Peer connections require mutual registration
49//!
50//! # Related Components
51//! - [`NodeRemote`]: Core remote interface
52//! - [`ClientServerRemote`]: Client-server communication
53//! - [`PeerRemote`]: Peer-to-peer communication
54//! - [`CitadelClientServerConnection`]: Connection management
55//! - [`RegisterSuccess`]: Registration handling
56//!
57//! [`NodeRemote`]: crate::prelude::NodeRemote
58//! [`ClientServerRemote`]: crate::prelude::ClientServerRemote
59//! [`PeerRemote`]: crate::prelude::PeerRemote
60//! [`CitadelClientServerConnection`]: crate::prelude::CitadelClientServerConnection
61//! [`RegisterSuccess`]: crate::prelude::RegisterSuccess
62
63use 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    /// A reference to a user identifier
80    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    /// A convenience structure for executing commands that depend on a specific registered user
98    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
173/// Contains the elements required to communicate with the adjacent node
174pub struct CitadelClientServerConnection<R: Ratchet> {
175    /// An interface to send ordered, reliable, and encrypted messages
176    pub(crate) channel: Option<PeerChannel<R>>,
177    pub remote: ClientServerRemote<R>,
178    /// Only available if UdpMode was enabled at the beginning of a session
179    pub udp_channel_rx: Option<citadel_io::tokio::sync::oneshot::Receiver<UdpChannel<R>>>,
180    /// Contains the Google auth minted at the central server (if the central server enabled it), as well as any other services enabled by the central server
181    pub services: ServicesObject,
182    pub cid: u64,
183    pub session_security_settings: SessionSecuritySettings,
184}
185
186impl<R: Ratchet> CitadelClientServerConnection<R> {
187    /// Splits the channel into a send and receive half. This will render
188    /// the other fields of this connection object innaccessible
189    ///
190    /// # Panics
191    ///  - If the channel has already been taken
192    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
215/// Contains the elements entailed by a successful registration
216pub struct RegisterSuccess {
217    pub cid: u64,
218}
219
220#[async_trait]
221/// Endows the [NodeRemote](NodeRemote) with additional functions
222pub trait ProtocolRemoteExt<R: Ratchet>: Remote<R> {
223    /// Registers with custom settings
224    /// Returns a ticket which is used to uniquely identify the request in the protocol
225    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    /// Registers using the default settings. The default uses No Google FCM keys and the default session security settings
282    /// Returns a ticket which is used to uniquely identify the request in the protocol
283    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    /// Connects with custom settings
307    /// Returns a ticket which is used to uniquely identify the request in the protocol
308    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    /// Connects with the default settings
380    /// If FCM keys were created during the registration phase, then those keys will be used for the session. If new FCM keys need to be used, consider using [`Self::connect`]
381    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    /// Creates a valid target identifier used to make protocol requests. Raw user IDs or usernames can be used
397    /// ```
398    /// use citadel_sdk::prelude::*;
399    /// # use citadel_sdk::prefabs::client::single_connection::SingleClientServerConnectionKernel;
400    ///
401    /// let server_connection_settings = DefaultServerConnectionSettingsBuilder::credentialed_login("127.0.0.1:25021", "john.doe", "password").build().unwrap();
402    ///
403    /// # SingleClientServerConnectionKernel::new(server_connection_settings, |conn| async move {
404    /// conn.find_target("my_account", "my_peer").await?.send_file("/path/to/file.pdf").await
405    /// // or: conn.find_target(1234, "my_peer").await? [...]
406    /// # });
407    /// ```
408    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    /// Creates a proposed target from the valid local user to an unregistered peer in the network. Used when creating registration requests for peers.
443    /// Currently only supports LocalGroup <-> LocalGroup peer connections
444    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    /// Returns a list of local group peers on the network for local_user. May or may not be registered to the user. To get a list of registered users to local_user, run [`Self::get_local_group_mutual_peers`]
480    /// - limit: if None, all peers are obtained. If Some, at most the specified number of peers will be obtained
481    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    /// Returns a list of mutually-registered peers with the local_user
533    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    /// Returns all the active sessions in the protocol, including all P2P connections hierarchically placed as children to C2S
581    /// connections
582    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]
644/// Some functions require that a target exists
645pub trait ProtocolRemoteTargetExt<R: Ratchet>: TargetLockedRemote<R> {
646    /// Sends a file with a custom size. The smaller the chunks, the higher the degree of scrambling, but the higher the performance cost. A chunk size of zero will use the default
647    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    /// Sends a file to the provided target using the default chunking size
700    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    /// Sends a file to the provided target using custom chunking size with local encryption.
706    /// Only this local node may decrypt the information send to the adjacent node.
707    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    /// Sends a file to the provided target using the default chunking size with local encryption.
735    /// Only this local node may decrypt the information send to the adjacent node.
736    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    /// Pulls a virtual file from the RE-VFS. If `delete_on_pull` is true, then, the virtual file
752    /// will be taken from the RE-VFS
753    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    /// Deletes the file from the RE-VFS. If the contents are desired on delete,
800    /// consider calling `Self::remote_encrypted_virtual_filesystem_pull` with the delete
801    /// parameter set to true
802    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    /// Connects to the peer with custom settings
839    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        // Timeout for the entire P2P connection process.
848        // This prevents indefinite hangs when the server never responds with PeerChannelCreated.
849        // The timeout should be long enough for normal hole punching (which has its own 30s timeout)
850        // plus key exchange, but short enough to fail fast when something is stuck.
851        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                    // A routing failure for this ticket (e.g. the server could not deliver our
916                    // accept or request). Fail fast with the reason instead of waiting out the
917                    // full 60s watchdog — the caller can retry immediately.
918                    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    /// Connects to the target peer with default settings
948    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    /// Posts a registration request to a peer
954    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        // TODO: Get rid of this step. Should be handled by the protocol
958        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    /// Deregisters the currently locked target. If the target is a client to server
1009    /// connection, deregisters from the server. If the target is a p2p connection,
1010    /// deregisters the p2p
1011    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            // c2s conn
1034            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            //c2s conn
1110            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    /// Create a group with explicit [`MessageGroupOptions`] — e.g. a zero-trust
1146    /// [`GroupHierarchyMode::CommandHierarchy`](citadel_types::proto::GroupHierarchyMode) where a
1147    /// superior can read its subordinates' messages. `options.hierarchy` carries the initial rank
1148    /// assignment (member cid → command path); it is consumed owner-locally and never reaches the relay.
1149    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        // NOTE: default is PRIVATE mode, meaning all users in group must be registered to the owner.
1158        // Initial users are UserIdentifiers resolved to cids below.
1159        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    /// Lists all groups that which the current peer owns
1204    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    /// Lists all active sessions, including the local nat type. For each active session,
1237    /// lists each connection info which includes the remote nat type, connection status,
1238    /// peer id, and latest ratchet version.
1239    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    /// Begins a re-key, updating the container in the process.
1253    /// Returns the new key matrix version. Does not return the new key version
1254    /// if the rekey fails, or, if a current rekey is already executing
1255    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    /// Checks if the locked target is registered
1280    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            // in this case, the user re-used a remote locked to a registration target
1307            // where the username was provided, but the cid was 0 (unknown).
1308            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            // get the peer cid from the account manager (implying the peers are already registered).
1318            // fallback to the mapped cid if the peer is not registered
1319            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        /// Receives incoming file/object transfer requests. The handles must be
1368        /// .accepted() before the file/object transfer is allowed to proceed
1369        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        /// Obtains a receiver which yields incoming file/object transfer handles
1383        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                // accept the transfer
1495                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}