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            // Bounded, because this wait is otherwise the end of the line.
1116            //
1117            // The stream closing is handled below, but a subscription that
1118            // stays OPEN and simply never carries the matching Disconnect
1119            // parks here for ever. That is the reconnection suite's 240s
1120            // timeout: with the phase markers reaching CI, both peers finish
1121            // phase one and the disconnecting side never logs the line after
1122            // this call, while the other blocks on the barrier it never
1123            // reaches. The only await between those two markers is this one.
1124            //
1125            // The session is being torn down either way, so an unconfirmed
1126            // disconnect is reported rather than waited on: a caller can retry
1127            // or proceed, and neither is possible from inside an indefinite
1128            // await.
1129            const DISCONNECT_CONFIRMATION_TIMEOUT: Duration = Duration::from_secs(30);
1130            let deadline = citadel_io::time::Instant::now() + DISCONNECT_CONFIRMATION_TIMEOUT;
1131            // What the subscription DID carry, for the occasion when it does not
1132            // carry a Disconnect.
1133            //
1134            // `RemoteDisconnectEventMissing` says only that nothing matched in
1135            // thirty seconds, which cannot distinguish the two things it might
1136            // mean. A NON-ZERO count says the subscription was alive and other
1137            // events reached it, so the Disconnect was never emitted or was
1138            // emitted under a key this subscription does not hold. ZERO says the
1139            // subscription heard nothing at all, which points at the session
1140            // rather than at the routing. Those want different fixes, and the
1141            // error tells them apart not at all.
1142            //
1143            // This is the same instrument that localised the UDP one-shot race:
1144            // an intermittent CI failure that no local run reproduces is only
1145            // ever going to be solved by what it says about itself when it next
1146            // happens.
1147            let mut carried = 0usize;
1148            loop {
1149                let remaining =
1150                    deadline.saturating_duration_since(citadel_io::time::Instant::now());
1151                if remaining.is_zero() {
1152                    log::warn!(target: "citadel", "[dc-wait] no Disconnect for cid {cid} within {DISCONNECT_CONFIRMATION_TIMEOUT:?}; the subscription carried {carried} other event(s)");
1153                    return Err(citadel_io::error!(
1154                        citadel_io::ErrorCode::RemoteDisconnectEventMissing
1155                    ));
1156                }
1157                match citadel_io::time::timeout(remaining, subscription.next()).await {
1158                    Ok(Some(event)) => {
1159                        carried += 1;
1160                        if let NodeResult::Disconnect(Disconnect {
1161                            success, message, ..
1162                        }) = event.into_result()?
1163                        {
1164                            return if success {
1165                                Ok(())
1166                            } else {
1167                                Err(citadel_io::error!(
1168                                    citadel_io::ErrorCode::RemoteDisconnected,
1169                                    message
1170                                ))
1171                            };
1172                        }
1173                    }
1174                    // The stream ended without ever carrying the event.
1175                    Ok(None) => {
1176                        log::warn!(target: "citadel", "[dc-wait] the subscription for cid {cid} closed without a Disconnect; it carried {carried} other event(s)");
1177                        break;
1178                    }
1179                    Err(_elapsed) => {
1180                        log::warn!(target: "citadel", "[dc-wait] timed out waiting for a Disconnect for cid {cid}; the subscription carried {carried} other event(s)");
1181                        return Err(citadel_io::error!(
1182                            citadel_io::ErrorCode::RemoteDisconnectEventMissing
1183                        ));
1184                    }
1185                }
1186            }
1187
1188            Err(citadel_io::error!(
1189                citadel_io::ErrorCode::RemoteDisconnectEventMissing
1190            ))
1191        }
1192    }
1193
1194    async fn create_group(
1195        &self,
1196        initial_users_to_invite: Option<Vec<UserIdentifier>>,
1197    ) -> Result<GroupChannel, NetworkError> {
1198        self.create_group_with_options(initial_users_to_invite, MessageGroupOptions::default())
1199            .await
1200    }
1201
1202    /// Create a group with explicit [`MessageGroupOptions`] — e.g. a zero-trust
1203    /// [`GroupHierarchyMode::CommandHierarchy`](citadel_types::proto::GroupHierarchyMode) where a
1204    /// superior can read its subordinates' messages. `options.hierarchy` carries the initial rank
1205    /// assignment (member cid → command path); it is consumed owner-locally and never reaches the relay.
1206    async fn create_group_with_options(
1207        &self,
1208        initial_users_to_invite: Option<Vec<UserIdentifier>>,
1209        options: MessageGroupOptions,
1210    ) -> Result<GroupChannel, NetworkError> {
1211        let session_cid = self.user().get_session_cid();
1212
1213        let mut initial_users = vec![];
1214        // NOTE: default is PRIVATE mode, meaning all users in group must be registered to the owner.
1215        // Initial users are UserIdentifiers resolved to cids below.
1216        if let Some(initial_users_to_invite) = initial_users_to_invite {
1217            for user in initial_users_to_invite {
1218                initial_users.push(
1219                    self.remote()
1220                        .account_manager()
1221                        .find_target_information(session_cid, user.clone())
1222                        .await?
1223                        .ok_or_else(|| {
1224                            citadel_io::error!(
1225                                citadel_io::ErrorCode::RemoteGroupAccountNotFound,
1226                                citadel_io::Dbg(user),
1227                                citadel_io::Dbg(session_cid)
1228                            )
1229                        })
1230                        .map(|r| r.1.cid)?,
1231                )
1232            }
1233        }
1234
1235        let group_request = GroupBroadcast::Create {
1236            initial_invitees: initial_users,
1237            options,
1238        };
1239        let request = NodeRequest::GroupBroadcastCommand(GroupBroadcastCommand {
1240            session_cid,
1241            command: group_request,
1242        });
1243        let mut subscription = self.remote().send_callback_subscription(request).await?;
1244
1245        // The loop lives in `group_create_wait` so its refusal paths can be
1246        // tested without a running node: the only way to make a real server
1247        // refuse a Create is to give the owner 256 groups first.
1248        //
1249        // It matched GroupChannelCreated and nothing else, and — alone among its
1250        // neighbours — never called into_result(). So SignalError,
1251        // OutboundRequestRejected and InternalServerError were discarded, and so
1252        // was the server's own answer to a Create it could not perform:
1253        // CreateResponse { key: None }, delivered on this very ticket. Every one
1254        // of those left the caller parked on a stream that would never speak
1255        // again. The broadcast prefab handles that same event as
1256        // BroadcastCreateGroupFailed — the guard existed, in the prefab only.
1257        match crate::group_create_wait::await_group_creation(&mut subscription).await? {
1258            crate::group_create_wait::GroupCreation::Created(channel) => Ok(channel),
1259            crate::group_create_wait::GroupCreation::Refused => Err(citadel_io::error!(
1260                citadel_io::ErrorCode::Generic,
1261                "The server refused to create the group"
1262            )),
1263            crate::group_create_wait::GroupCreation::Ended => Err(citadel_io::error!(
1264                citadel_io::ErrorCode::RemoteCreateGroupEndedUnexpectedly
1265            )),
1266        }
1267    }
1268
1269    /// Lists all groups that which the current peer owns
1270    async fn list_owned_groups(&self) -> Result<Vec<MessageGroupKey>, NetworkError> {
1271        let session_cid = self.user().get_session_cid();
1272        let cid_to_check_for = match self.try_as_peer_connection().await {
1273            Ok(res) => res.get_original_target_cid(),
1274            _ => session_cid,
1275        };
1276        let group_request = GroupBroadcast::ListGroupsFor {
1277            cid: cid_to_check_for,
1278        };
1279        let request = NodeRequest::GroupBroadcastCommand(GroupBroadcastCommand {
1280            session_cid,
1281            command: group_request,
1282        });
1283
1284        let mut subscription = self.remote().send_callback_subscription(request).await?;
1285
1286        while let Some(evt) = subscription.next().await {
1287            if let NodeResult::GroupEvent(GroupEvent {
1288                session_cid: _,
1289                ticket: _,
1290                event: GroupBroadcast::ListResponse { groups },
1291            }) = evt.into_result()?
1292            {
1293                return Ok(groups);
1294            }
1295        }
1296
1297        Err(citadel_io::error!(
1298            citadel_io::ErrorCode::RemoteListGroupsEndedUnexpectedly
1299        ))
1300    }
1301
1302    /// Lists all active sessions, including the local nat type. For each active session,
1303    /// lists each connection info which includes the remote nat type, connection status,
1304    /// peer id, and latest ratchet version.
1305    async fn list_sessions(&self) -> Result<ActiveSessions, NetworkError> {
1306        let request = NodeRequest::GetActiveSessions;
1307        let mut subscription = self.remote().send_callback_subscription(request).await?;
1308
1309        if let Some(NodeResult::SessionList(result)) = subscription.next().await {
1310            return Ok(result.sessions);
1311        }
1312
1313        Err(citadel_io::error!(
1314            citadel_io::ErrorCode::RemoteListSessionsEndedUnexpectedly
1315        ))
1316    }
1317
1318    /// Begins a re-key, updating the container in the process.
1319    /// Returns the new key matrix version. Does not return the new key version
1320    /// if the rekey fails, or, if a current rekey is already executing
1321    async fn rekey(&self) -> Result<Option<u32>, NetworkError> {
1322        let request = NodeRequest::ReKey(ReKey {
1323            v_conn_type: *self.user(),
1324        });
1325        let mut subscription = self.remote().send_callback_subscription(request).await?;
1326
1327        while let Some(evt) = subscription.next().await {
1328            if let NodeResult::ReKeyResult(result) = evt {
1329                return match result.status {
1330                    ReKeyReturnType::Success { version } => Ok(Some(version)),
1331                    ReKeyReturnType::AlreadyInProgress => Ok(None),
1332                    ReKeyReturnType::Failure { err } => Err(citadel_io::error!(
1333                        citadel_io::ErrorCode::RemoteRekeyFailed,
1334                        err
1335                    )),
1336                };
1337            }
1338        }
1339
1340        Err(citadel_io::error!(
1341            citadel_io::ErrorCode::RemoteRekeyEndedUnexpectedly
1342        ))
1343    }
1344
1345    /// Checks if the locked target is registered
1346    async fn is_peer_registered(&self) -> Result<bool, NetworkError> {
1347        let target = self.try_as_peer_connection().await?;
1348        if let PeerConnectionType::LocalGroupPeer {
1349            session_cid: local_cid,
1350            peer_cid,
1351        } = target
1352        {
1353            let peers = self.remote().get_local_group_peers(local_cid, None).await?;
1354            citadel_logging::info!(target: "citadel", "Checking to see if {target} is registered in {peers:?}");
1355            Ok(peers.iter().any(|p| p.cid == peer_cid))
1356        } else {
1357            Err(citadel_io::error!(
1358                citadel_io::ErrorCode::RemoteExternalGroupPeerUnsupportedYet
1359            ))
1360        }
1361    }
1362
1363    #[doc(hidden)]
1364    async fn try_as_peer_connection(&self) -> Result<PeerConnectionType, NetworkError> {
1365        let verified_return = |user: &VirtualTargetType| {
1366            user.try_as_peer_connection().ok_or(citadel_io::error!(
1367                citadel_io::ErrorCode::RemoteTargetNotPeer
1368            ))
1369        };
1370
1371        if self.user().get_target_cid() == 0 {
1372            // in this case, the user re-used a remote locked to a registration target
1373            // where the username was provided, but the cid was 0 (unknown).
1374            let peer_username = self.target_username().ok_or_else(|| {
1375                citadel_io::error!(citadel_io::ErrorCode::RemoteTargetCidZeroNoUsername)
1376            })?;
1377            let session_cid = self.user().get_session_cid();
1378            let expected_peer_cid = self
1379                .remote()
1380                .account_manager()
1381                .get_persistence_handler()
1382                .get_cid_by_username(peer_username);
1383            // get the peer cid from the account manager (implying the peers are already registered).
1384            // fallback to the mapped cid if the peer is not registered
1385            let peer_cid = self
1386                .remote()
1387                .account_manager()
1388                .find_target_information(session_cid, peer_username)
1389                .await?
1390                .map(|r| r.1.cid)
1391                .unwrap_or(expected_peer_cid);
1392
1393            let mut user = *self.user();
1394            user.set_target_cid(peer_cid);
1395            verified_return(&user)
1396        } else {
1397            verified_return(self.user())
1398        }
1399    }
1400
1401    #[doc(hidden)]
1402    fn can_use_revfs(&self) -> Result<(), NetworkError> {
1403        if let Some(sess) = self.session_security_settings() {
1404            if sess.crypto_params.kem_algorithm == KemAlgorithm::MlKem {
1405                Ok(())
1406            } else {
1407                Err(citadel_io::error!(
1408                    citadel_io::ErrorCode::RemoteRevfsRequiresKyber
1409                ))
1410            }
1411        } else {
1412            Err(citadel_io::error!(
1413                citadel_io::ErrorCode::RemoteRevfsUnsupportedRemote
1414            ))
1415        }
1416    }
1417}
1418
1419impl<T: TargetLockedRemote<R>, R: Ratchet> ProtocolRemoteTargetExt<R> for T {}
1420
1421pub mod results {
1422    use crate::prefabs::client::peer_connection::FileTransferHandleRx;
1423    use crate::prelude::{PeerChannel, UdpChannel};
1424    use crate::remote_ext::remote_specialization::PeerRemote;
1425    use citadel_io::tokio::sync::oneshot::Receiver;
1426    use citadel_proto::prelude::*;
1427    use std::fmt::Debug;
1428
1429    pub struct PeerConnectSuccess<R: Ratchet> {
1430        pub channel: PeerChannel<R>,
1431        pub udp_channel_rx: Option<Receiver<UdpChannel<R>>>,
1432        pub remote: PeerRemote<R>,
1433        /// Receives incoming file/object transfer requests. The handles must be
1434        /// .accepted() before the file/object transfer is allowed to proceed
1435        pub(crate) incoming_object_transfer_handles: Option<FileTransferHandleRx>,
1436    }
1437
1438    impl<R: Ratchet> Debug for PeerConnectSuccess<R> {
1439        fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1440            f.debug_struct("PeerConnectSuccess")
1441                .field("channel", &self.channel)
1442                .field("udp_channel_rx", &self.udp_channel_rx)
1443                .finish()
1444        }
1445    }
1446
1447    impl<R: Ratchet> PeerConnectSuccess<R> {
1448        /// Obtains a receiver which yields incoming file/object transfer handles
1449        pub fn get_incoming_file_transfer_handle(
1450            &mut self,
1451        ) -> Result<FileTransferHandleRx, NetworkError> {
1452            self.incoming_object_transfer_handles
1453                .take()
1454                .ok_or(citadel_io::error!(
1455                    citadel_io::ErrorCode::RemoteFunctionAlreadyCalled
1456                ))
1457        }
1458    }
1459
1460    /// The peer's answer to a registration request.
1461    ///
1462    /// `Ok` means the request completed, not that it was accepted — a decline is
1463    /// a successful round trip with the answer "no". This type derived nothing,
1464    /// so a caller could neither compare it nor log it, and all three real
1465    /// callers wrote `Ok(_)` / `let _ =` and carried on as though the peer had
1466    /// said yes. `peer_connection.rs` then issued a PostConnect to a peer that
1467    /// had just declined and waited out a 60s RemoteP2pConnectTimeout, reporting
1468    /// that instead of the decline the SDK had in hand.
1469    ///
1470    /// The derives are the fix for the type; `is_accepted` is the fix for the
1471    /// call sites, which needed something shorter to write than the mistake.
1472    #[derive(Clone, Debug, PartialEq, Eq)]
1473    pub enum PeerRegisterStatus {
1474        Accepted,
1475        Declined,
1476        Failed { reason: Option<String> },
1477    }
1478
1479    impl PeerRegisterStatus {
1480        /// Did the peer say yes?
1481        pub fn is_accepted(&self) -> bool {
1482            matches!(self, PeerRegisterStatus::Accepted)
1483        }
1484
1485        /// The reason this was not an acceptance, or `None` if it was.
1486        pub fn refusal_reason(&self) -> Option<String> {
1487            match self {
1488                PeerRegisterStatus::Accepted => None,
1489                PeerRegisterStatus::Declined => Some("The peer declined the request".to_string()),
1490                PeerRegisterStatus::Failed { reason } => Some(
1491                    reason
1492                        .clone()
1493                        .unwrap_or_else(|| "The registration request failed".to_string()),
1494                ),
1495            }
1496        }
1497    }
1498
1499    #[derive(Clone, Debug)]
1500    pub struct LocalGroupPeer {
1501        pub cid: u64,
1502        pub is_online: bool,
1503    }
1504
1505    #[derive(Clone, Debug)]
1506    pub struct LocalGroupPeerFullInfo {
1507        pub cid: u64,
1508        pub username: Option<String>,
1509        pub full_name: Option<String>,
1510        pub is_online: bool,
1511    }
1512}
1513
1514pub mod remote_specialization {
1515    use crate::prelude::*;
1516    use std::ops::{Deref, DerefMut};
1517
1518    #[derive(Debug, Clone)]
1519    pub struct PeerRemote<R: Ratchet> {
1520        pub(crate) inner: NodeRemote<R>,
1521        pub(crate) peer: VirtualTargetType,
1522        pub(crate) username: Option<String>,
1523        pub(crate) session_security_settings: SessionSecuritySettings,
1524    }
1525
1526    impl<R: Ratchet> Deref for PeerRemote<R> {
1527        type Target = NodeRemote<R>;
1528        fn deref(&self) -> &Self::Target {
1529            &self.inner
1530        }
1531    }
1532
1533    impl<R: Ratchet> DerefMut for PeerRemote<R> {
1534        fn deref_mut(&mut self) -> &mut Self::Target {
1535            &mut self.inner
1536        }
1537    }
1538
1539    impl<R: Ratchet> TargetLockedRemote<R> for PeerRemote<R> {
1540        fn user(&self) -> &VirtualTargetType {
1541            &self.peer
1542        }
1543        fn remote(&self) -> &NodeRemote<R> {
1544            &self.inner
1545        }
1546        fn target_username(&self) -> Option<&str> {
1547            self.username.as_deref()
1548        }
1549        fn user_mut(&mut self) -> &mut VirtualTargetType {
1550            &mut self.peer
1551        }
1552
1553        fn session_security_settings(&self) -> Option<&SessionSecuritySettings> {
1554            Some(&self.session_security_settings)
1555        }
1556    }
1557}
1558
1559#[cfg(all(test, not(target_family = "wasm")))]
1560mod tests {
1561    use crate::prefabs::client::single_connection::SingleClientServerConnectionKernel;
1562    use crate::prefabs::client::DefaultServerConnectionSettingsBuilder;
1563    use crate::prelude::*;
1564    use citadel_io::tokio;
1565    use rstest::rstest;
1566    use std::net::SocketAddr;
1567    use std::sync::atomic::{AtomicBool, Ordering};
1568    use std::sync::Arc;
1569    use uuid::Uuid;
1570
1571    pub struct ReceiverFileTransferKernel<R: Ratchet>(
1572        pub Option<NodeRemote<R>>,
1573        pub Arc<AtomicBool>,
1574    );
1575
1576    #[async_trait]
1577    impl<R: Ratchet> NetKernel<R> for ReceiverFileTransferKernel<R> {
1578        fn load_remote(&mut self, node_remote: NodeRemote<R>) -> Result<(), NetworkError> {
1579            self.0 = Some(node_remote);
1580            Ok(())
1581        }
1582
1583        async fn on_start(&self) -> Result<(), NetworkError> {
1584            Ok(())
1585        }
1586
1587        async fn on_node_event_received(&self, message: NodeResult<R>) -> Result<(), NetworkError> {
1588            log::trace!(target: "citadel", "SERVER received {:?}", message);
1589            if let NodeResult::ObjectTransferHandle(ObjectTransferHandle { mut handle, .. }) =
1590                message.into_result()?
1591            {
1592                let mut path = None;
1593                // accept the transfer
1594                handle
1595                    .accept()
1596                    .map_err(|err| NetworkError::msg(err.into_string()))?;
1597
1598                use citadel_types::proto::ObjectTransferStatus;
1599                use futures::StreamExt;
1600                while let Some(status) = handle.next().await {
1601                    match status {
1602                        ObjectTransferStatus::ReceptionComplete => {
1603                            log::trace!(target: "citadel", "Server has finished receiving the file!");
1604                            let cmp = include_bytes!("../../resources/TheBridge.pdf");
1605                            let streamed_data = citadel_io::tokio::fs::read(path.clone().unwrap())
1606                                .await
1607                                .unwrap();
1608                            assert_eq!(
1609                                cmp,
1610                                streamed_data.as_slice(),
1611                                "Original data and streamed data does not match"
1612                            );
1613
1614                            self.1.store(true, Ordering::Relaxed);
1615                            self.0.clone().unwrap().shutdown().await?;
1616                        }
1617
1618                        ObjectTransferStatus::ReceptionBeginning(file_path, vfm) => {
1619                            path = Some(file_path);
1620                            assert_eq!(vfm.name, "TheBridge.pdf")
1621                        }
1622
1623                        _ => {}
1624                    }
1625                }
1626            }
1627
1628            Ok(())
1629        }
1630
1631        async fn on_stop(&mut self) -> Result<(), NetworkError> {
1632            Ok(())
1633        }
1634    }
1635
1636    pub fn server_info<'a, R: Ratchet>(
1637        switch: Arc<AtomicBool>,
1638    ) -> (NodeFuture<'a, ReceiverFileTransferKernel<R>>, SocketAddr) {
1639        crate::test_common::server_test_node(ReceiverFileTransferKernel(None, switch), |_| {})
1640    }
1641
1642    #[rstest]
1643    #[case(
1644        EncryptionAlgorithm::AES_GCM_256,
1645        KemAlgorithm::MlKem,
1646        SigAlgorithm::None
1647    )]
1648    #[case(
1649        EncryptionAlgorithm::MlKemHybrid,
1650        KemAlgorithm::MlKem,
1651        SigAlgorithm::MlDsa65
1652    )]
1653    #[timeout(std::time::Duration::from_secs(90))]
1654    #[tokio::test]
1655    async fn test_c2s_file_transfer(
1656        #[case] enx: EncryptionAlgorithm,
1657        #[case] kem: KemAlgorithm,
1658        #[case] sig: SigAlgorithm,
1659    ) {
1660        citadel_logging::setup_log();
1661        let client_success = &AtomicBool::new(false);
1662        let server_success = &Arc::new(AtomicBool::new(false));
1663        let (server, server_addr) = server_info::<StackedRatchet>(server_success.clone());
1664        let uuid = Uuid::new_v4();
1665
1666        let session_security_settings = SessionSecuritySettingsBuilder::default()
1667            .with_crypto_params(enx + kem + sig)
1668            .build()
1669            .unwrap();
1670
1671        let server_connection_settings =
1672            DefaultServerConnectionSettingsBuilder::transient_with_id(server_addr, uuid)
1673                .with_session_security_settings(session_security_settings)
1674                .disable_udp()
1675                .build()
1676                .unwrap();
1677
1678        let client_kernel = SingleClientServerConnectionKernel::new(
1679            server_connection_settings,
1680            |connection| async move {
1681                log::trace!(target: "citadel", "***CLIENT LOGIN SUCCESS :: File transfer next ***");
1682                connection
1683                    .send_file_with_custom_opts(
1684                        "../resources/TheBridge.pdf",
1685                        32 * 1024,
1686                        TransferType::FileTransfer,
1687                    )
1688                    .await
1689                    .unwrap();
1690                log::trace!(target: "citadel", "***CLIENT FILE TRANSFER SUCCESS***");
1691                client_success.store(true, Ordering::Relaxed);
1692                connection.shutdown_kernel().await
1693            },
1694        );
1695
1696        let client = DefaultNodeBuilder::default().build(client_kernel).unwrap();
1697
1698        let joined = futures::future::try_join(server, client);
1699
1700        let _ = joined.await.unwrap();
1701
1702        assert!(client_success.load(Ordering::Relaxed));
1703        assert!(server_success.load(Ordering::Relaxed));
1704    }
1705}