Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion api/snapshots/rns-runtime.txt
Original file line number Diff line number Diff line change
Expand Up @@ -587,7 +587,7 @@ pub fn rns_runtime::link_manager::LinkManager::set_link_packet_channel(&mut self
pub fn rns_runtime::link_manager::LinkManager::set_link_packet_proof_channel(&mut self, tokio::sync::mpsc::bounded::Sender<rns_runtime::link_manager::LinkPacketProof>)
pub fn rns_runtime::link_manager::LinkManager::set_outbound_resource_proof_channel(&mut self, tokio::sync::mpsc::bounded::Sender<rns_runtime::link_manager::LinkResourceProof>)
pub fn rns_runtime::link_manager::LinkManager::set_request_handler<F>(&mut self, F) where F: core::ops::function::Fn([u8; 16], [u8; 16], alloc::vec::Vec<u8>) -> core::option::Option<alloc::vec::Vec<u8>> + core::marker::Send + 'static
pub fn rns_runtime::link_manager::LinkManager::set_request_handler_ex<F>(&mut self, F) where F: core::ops::function::Fn([u8; 16], [u8; 16], alloc::vec::Vec<u8>) -> rns_runtime::link_manager::RequestOutcome + core::marker::Send + 'static
pub fn rns_runtime::link_manager::LinkManager::set_request_handler_ex<F>(&mut self, F) where F: core::ops::function::Fn([u8; 16], [u8; 16], alloc::vec::Vec<u8>, core::option::Option<rns_identity::identity::Identity>) -> rns_runtime::link_manager::RequestOutcome + core::marker::Send + 'static
pub fn rns_runtime::link_manager::LinkManager::set_resource_accept_handler<F>(&mut self, F) where F: core::ops::function::Fn([u8; 16], &rns_protocol::resource_adv::ResourceAdvertisement) -> bool + core::marker::Send + 'static
pub fn rns_runtime::link_manager::LinkManager::set_resource_completed_channel(&mut self, tokio::sync::mpsc::bounded::Sender<(alloc::vec::Vec<u8>, [u8; 16])>)
pub fn rns_runtime::link_manager::LinkManager::set_resource_completion_channel(&mut self, tokio::sync::mpsc::bounded::Sender<rns_runtime::link_manager::ResourceCompletion>)
Expand Down
4 changes: 3 additions & 1 deletion crates/rns-runtime/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,9 @@ pub mod prelude {
RegisteredDestination, ResourceAcceptPolicy,
};
pub use crate::lifecycle::ShutdownSignal;
pub use crate::link_manager::{DestinationAnnounceOptions, DestinationRequest, RequestOutcome};
pub use crate::link_manager::{
DestinationAnnounceOptions, DestinationRequest, RequestOutcome, pack_file_name_metadata,
};
pub use crate::link_session::{
LinkSession, LinkSessionChannelError, LinkSessionChannelHandle, LinkSessionCloseReason,
LinkSessionError, LinkSessionEvent, LinkSessionHandle, LinkSessionResourceError,
Expand Down
49 changes: 38 additions & 11 deletions crates/rns-runtime/src/link_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,17 @@ pub enum LinkClientError {
UnexpectedResponse(String),
}

/// Successful Link request response.
///
/// Ordinary replies carry packed response bytes in [`Self::data`] with
/// [`Self::metadata`] unset. NomadNet `/file/...` replies are response Resources
/// whose payload is raw file bytes plus optional msgpack filename metadata.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LinkQueryResponse {
pub data: Vec<u8>,
pub metadata: Option<Vec<u8>>,
}

#[derive(Clone)]
pub struct LinkClient {
transport_tx: mpsc::Sender<TransportMessage>,
Expand All @@ -66,7 +77,7 @@ impl LinkClient {
}

/// Open a Link to `app_name` on `remote_transport_hash`, send one
/// request, return the response.
/// request, return the response (and optional Resource metadata).
pub async fn query(
&self,
remote_transport_hash: [u8; 16],
Expand All @@ -75,7 +86,7 @@ impl LinkClient {
payload: Vec<u8>,
hops: u8,
overall_timeout: Duration,
) -> Result<Vec<u8>, LinkClientError> {
) -> Result<LinkQueryResponse, LinkClientError> {
let started = Instant::now();
let deadline = started + overall_timeout;

Expand Down Expand Up @@ -280,7 +291,7 @@ async fn wait_for_response(
link_id: [u8; 16],
request_id: [u8; 16],
deadline: Duration,
) -> Result<Vec<u8>, LinkClientError> {
) -> Result<LinkQueryResponse, LinkClientError> {
let fut = async {
let mut inbound_resources: HashMap<[u8; 32], InboundTransfer> = HashMap::new();

Expand Down Expand Up @@ -318,7 +329,10 @@ async fn wait_for_response(
match link.handle_response(body) {
Ok((id, response_data)) => {
if id == request_id {
return Ok(response_data);
return Ok(LinkQueryResponse {
data: response_data,
metadata: None,
});
}
}
Err(e) => {
Expand Down Expand Up @@ -432,7 +446,7 @@ async fn wait_for_response(
}

if let Some(rh) = completed_rh {
let (assembled, proof) = {
let (assembled, proof, metadata) = {
let transfer =
inbound_resources.get_mut(&rh).ok_or_else(|| {
LinkClientError::UnexpectedResponse(
Expand All @@ -451,19 +465,32 @@ async fn wait_for_response(
},
)
};
transfer.complete(Some(&decrypt_fn)).map_err(|e| {
LinkClientError::UnexpectedResponse(format!(
"resource assemble: {e:?}"
))
})?
let (assembled, proof) =
transfer.complete(Some(&decrypt_fn)).map_err(|e| {
LinkClientError::UnexpectedResponse(format!(
"resource assemble: {e:?}"
))
})?;
let metadata = transfer.resource.metadata.clone();
(assembled, proof, metadata)
};

send_link_proof(transport_tx, link_id, &proof).await?;
inbound_resources.remove(&rh);
if metadata.is_some() {
// NomadNet file response: raw payload + Resource metadata.
return Ok(LinkQueryResponse {
data: assembled,
metadata,
});
}
match link.handle_response_plaintext(&assembled) {
Ok((id, response_data)) => {
if id == request_id {
return Ok(response_data);
return Ok(LinkQueryResponse {
data: response_data,
metadata: None,
});
}
}
Err(e) => {
Expand Down
Loading