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                    _ => {}
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    /// Connects to the target peer with default settings
935    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    /// Posts a registration request to a peer
941    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        // TODO: Get rid of this step. Should be handled by the protocol
945        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    /// Deregisters the currently locked target. If the target is a client to server
996    /// connection, deregisters from the server. If the target is a p2p connection,
997    /// deregisters the p2p
998    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            // c2s conn
1021            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            //c2s conn
1097            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    /// Create a group with explicit [`MessageGroupOptions`] — e.g. a zero-trust
1133    /// [`GroupHierarchyMode::CommandHierarchy`](citadel_types::proto::GroupHierarchyMode) where a
1134    /// superior can read its subordinates' messages. `options.hierarchy` carries the initial rank
1135    /// assignment (member cid → command path); it is consumed owner-locally and never reaches the relay.
1136    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        // NOTE: default is PRIVATE mode, meaning all users in group must be registered to the owner.
1145        // Initial users are UserIdentifiers resolved to cids below.
1146        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    /// Lists all groups that which the current peer owns
1191    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    /// Lists all active sessions, including the local nat type. For each active session,
1224    /// lists each connection info which includes the remote nat type, connection status,
1225    /// peer id, and latest ratchet version.
1226    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    /// Begins a re-key, updating the container in the process.
1240    /// Returns the new key matrix version. Does not return the new key version
1241    /// if the rekey fails, or, if a current rekey is already executing
1242    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    /// Checks if the locked target is registered
1267    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            // in this case, the user re-used a remote locked to a registration target
1294            // where the username was provided, but the cid was 0 (unknown).
1295            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            // get the peer cid from the account manager (implying the peers are already registered).
1305            // fallback to the mapped cid if the peer is not registered
1306            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        /// Receives incoming file/object transfer requests. The handles must be
1355        /// .accepted() before the file/object transfer is allowed to proceed
1356        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        /// Obtains a receiver which yields incoming file/object transfer handles
1370        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                // accept the transfer
1482                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}