From aa2f506996840bd9b667709f2c44807d9bc62ea2 Mon Sep 17 00:00:00 2001 From: timedout Date: Fri, 26 Jun 2026 01:20:08 +0100 Subject: [PATCH] style: Document & refactor sending service --- src/service/sending/appservice.rs | 192 +++++++++++++++--------------- src/service/sending/dest.rs | 75 ++++++------ 2 files changed, 136 insertions(+), 131 deletions(-) diff --git a/src/service/sending/appservice.rs b/src/service/sending/appservice.rs index a9bd50d43..e7242af42 100644 --- a/src/service/sending/appservice.rs +++ b/src/service/sending/appservice.rs @@ -2,7 +2,7 @@ use bytes::BytesMut; use conduwuit::{ - Err, Result, debug_error, err, implement, trace, utils, utils::response::LimitReadExt, warn, + Err, Result, debug_error, err, trace, utils, utils::response::LimitReadExt, warn, }; use ruma::api::{ IncomingResponse, OutgoingRequest, @@ -11,99 +11,103 @@ path_builder::SinglePath, }; -/// Sends a request to an appservice -/// -/// Only returns Ok(None) if there is no url specified in the appservice -/// registration file -#[implement(super::Service)] -pub async fn send_appservice_request( - &self, - registration: Registration, - request: T, -) -> Result> -where - T: OutgoingRequest + Debug + Send, -{ - let Some(dest) = registration.url else { - return Ok(None); - }; +impl super::Service { + /// Sends a request to an appservice + /// + /// Only returns Ok(None) if there is no url specified in the appservice + /// registration file + pub async fn send_appservice_request( + &self, + registration: Registration, + request: T, + ) -> Result> + where + T: OutgoingRequest + Debug + Send, + { + let Some(dest) = registration.url else { + return Ok(None); + }; - if dest == *"null" || dest.is_empty() { - return Ok(None); + if dest == *"null" || dest.is_empty() { + return Ok(None); + } + + trace!("Appservice URL \"{dest}\", Appservice ID: {}", registration.id); + + let hs_token = registration.hs_token.as_str(); + let mut http_request = request + .try_into_http_request::(&dest, SendAccessToken::Appservice(hs_token), ()) + .map_err(|e| { + err!(BadServerResponse( + warn!(appservice = %registration.id, "Failed to find destination {dest}: {e:?}") + )) + })? + .map(BytesMut::freeze); + + let mut parts = http_request.uri().clone().into_parts(); + let old_path_and_query = parts.path_and_query.unwrap().as_str().to_owned(); + let symbol = if old_path_and_query.contains('?') { "&" } else { "?" }; + + parts.path_and_query = Some( + (old_path_and_query + symbol + "access_token=" + hs_token) + .parse() + .unwrap(), + ); + *http_request.uri_mut() = parts.try_into().expect("our manipulation is always valid"); + + let reqwest_request = reqwest::Request::try_from(http_request)?; + + let client = &self.services.client.appservice; + + let mut response = client.execute(reqwest_request).await.map_err(|e| { + warn!( + "Could not send request to appservice \"{}\" at {dest}: {e:?}", + registration.id + ); + e + })?; + + // reqwest::Response -> http::Response conversion + let status = response.status(); + let mut http_response_builder = http::Response::builder() + .status(status) + .version(response.version()); + mem::swap( + response.headers_mut(), + http_response_builder + .headers_mut() + .expect("http::response::Builder is usable"), + ); + + let body = response + .limit_read( + self.server + .config + .max_request_size + .try_into() + .expect("usize fits into u64"), + ) + .await?; + + if !status.is_success() { + debug_error!("Appservice response bytes: {:?}", utils::string_from_bytes(&body)); + return Err!(BadServerResponse(warn!( + "Appservice \"{}\" returned unsuccessful HTTP response {status} at {dest}", + registration.id + ))); + } + + let response = T::IncomingResponse::try_from_http_response( + http_response_builder + .body(body) + .expect("reqwest body is valid http body"), + ); + + response.map(Some).map_err(|e| { + err!(BadServerResponse(warn!( + "Appservice \"{}\" returned invalid/malformed response bytes {dest}: {e}", + registration.id + ))) + }) } - - trace!("Appservice URL \"{dest}\", Appservice ID: {}", registration.id); - - let hs_token = registration.hs_token.as_str(); - let mut http_request = request - .try_into_http_request::(&dest, SendAccessToken::Appservice(hs_token), ()) - .map_err(|e| { - err!(BadServerResponse( - warn!(appservice = %registration.id, "Failed to find destination {dest}: {e:?}") - )) - })? - .map(BytesMut::freeze); - - let mut parts = http_request.uri().clone().into_parts(); - let old_path_and_query = parts.path_and_query.unwrap().as_str().to_owned(); - let symbol = if old_path_and_query.contains('?') { "&" } else { "?" }; - - parts.path_and_query = Some( - (old_path_and_query + symbol + "access_token=" + hs_token) - .parse() - .unwrap(), - ); - *http_request.uri_mut() = parts.try_into().expect("our manipulation is always valid"); - - let reqwest_request = reqwest::Request::try_from(http_request)?; - - let client = &self.services.client.appservice; - - let mut response = client.execute(reqwest_request).await.map_err(|e| { - warn!("Could not send request to appservice \"{}\" at {dest}: {e:?}", registration.id); - e - })?; - - // reqwest::Response -> http::Response conversion - let status = response.status(); - let mut http_response_builder = http::Response::builder() - .status(status) - .version(response.version()); - mem::swap( - response.headers_mut(), - http_response_builder - .headers_mut() - .expect("http::response::Builder is usable"), - ); - - let body = response - .limit_read( - self.server - .config - .max_request_size - .try_into() - .expect("usize fits into u64"), - ) - .await?; - - if !status.is_success() { - debug_error!("Appservice response bytes: {:?}", utils::string_from_bytes(&body)); - return Err!(BadServerResponse(warn!( - "Appservice \"{}\" returned unsuccessful HTTP response {status} at {dest}", - registration.id - ))); - } - - let response = T::IncomingResponse::try_from_http_response( - http_response_builder - .body(body) - .expect("reqwest body is valid http body"), - ); - - response.map(Some).map_err(|e| { - err!(BadServerResponse(warn!( - "Appservice \"{}\" returned invalid/malformed response bytes {dest}: {e}", - registration.id - ))) - }) } diff --git a/src/service/sending/dest.rs b/src/service/sending/dest.rs index 4099d3722..878b0fc41 100644 --- a/src/service/sending/dest.rs +++ b/src/service/sending/dest.rs @@ -1,6 +1,5 @@ use std::fmt::Debug; -use conduwuit::implement; use ruma::{OwnedServerName, OwnedUserId}; #[derive(Clone, Debug, PartialEq, Eq, Hash)] @@ -10,44 +9,46 @@ pub enum Destination { Federation(OwnedServerName), } -#[implement(Destination)] -#[must_use] -pub(super) fn get_prefix(&self) -> Vec { - match self { - | Self::Federation(server) => { - let len = server.as_bytes().len().saturating_add(1); +impl Destination { + /// Gets the prefix for this destination. + #[must_use] + pub(super) fn get_prefix(&self) -> Vec { + match self { + | Self::Federation(server) => { + let len = server.as_bytes().len().saturating_add(1); - let mut p = Vec::with_capacity(len); - p.extend_from_slice(server.as_bytes()); - p.push(0xFF); - p - }, - | Self::Appservice(server) => { - let sigil = b"+"; - let len = sigil.len().saturating_add(server.len()).saturating_add(1); + let mut p = Vec::with_capacity(len); + p.extend_from_slice(server.as_bytes()); + p.push(0xFF); + p + }, + | Self::Appservice(server) => { + let sigil = b"+"; + let len = sigil.len().saturating_add(server.len()).saturating_add(1); - let mut p = Vec::with_capacity(len); - p.extend_from_slice(sigil); - p.extend_from_slice(server.as_bytes()); - p.push(0xFF); - p - }, - | Self::Push(user, pushkey) => { - let sigil = b"$"; - let len = sigil - .len() - .saturating_add(user.as_bytes().len()) - .saturating_add(1) - .saturating_add(pushkey.len()) - .saturating_add(1); + let mut p = Vec::with_capacity(len); + p.extend_from_slice(sigil); + p.extend_from_slice(server.as_bytes()); + p.push(0xFF); + p + }, + | Self::Push(user, pushkey) => { + let sigil = b"$"; + let len = sigil + .len() + .saturating_add(user.as_bytes().len()) + .saturating_add(1) + .saturating_add(pushkey.len()) + .saturating_add(1); - let mut p = Vec::with_capacity(len); - p.extend_from_slice(sigil); - p.extend_from_slice(user.as_bytes()); - p.push(0xFF); - p.extend_from_slice(pushkey.as_bytes()); - p.push(0xFF); - p - }, + let mut p = Vec::with_capacity(len); + p.extend_from_slice(sigil); + p.extend_from_slice(user.as_bytes()); + p.push(0xFF); + p.extend_from_slice(pushkey.as_bytes()); + p.push(0xFF); + p + }, + } } }