1#![allow(clippy::disallowed_types)]
12
13use std::time::Duration;
14
15use http::{HeaderName, HeaderValue, Method};
16
17use crate::{
18 ffi::{handles::BodyReader, pumps::pump_hyper_body_to_channel_limited},
19 http::client::{fetch_request, FetchError},
20 parse_node_addr, Body, CoreError, FfiResponse, IrohEndpoint,
21};
22
23#[allow(clippy::too_many_arguments)]
28pub async fn fetch(
29 endpoint: &IrohEndpoint,
30 remote_node_id: &str,
31 url: &str,
32 method: &str,
33 headers: &[(String, String)],
34 req_body_reader: Option<BodyReader>,
35 fetch_token: Option<u64>,
36 direct_addrs: Option<&[std::net::SocketAddr]>,
37 timeout: Option<Duration>,
38 decompress: bool,
39 max_response_body_bytes: Option<usize>,
40) -> Result<FfiResponse, CoreError> {
41 {
43 let lower = url.to_ascii_lowercase();
44 if lower.starts_with("https://") || lower.starts_with("http://") {
45 let scheme_end = lower
46 .find("://")
47 .map(|i| i.saturating_add(3))
48 .unwrap_or(lower.len());
49 return Err(CoreError::invalid_input(format!(
50 "iroh-http URLs must use the \"httpi://\" scheme, not \"{}\". \
51 Example: httpi://nodeId/path",
52 &url[..scheme_end]
53 )));
54 }
55 }
56
57 let http_method = Method::from_bytes(method.as_bytes())
59 .map_err(|_| CoreError::invalid_input(format!("invalid HTTP method {:?}", method)))?;
60 for (name, value) in headers {
61 HeaderName::from_bytes(name.as_bytes())
62 .map_err(|_| CoreError::invalid_input(format!("invalid header name {:?}", name)))?;
63 HeaderValue::from_str(value).map_err(|_| {
64 CoreError::invalid_input(format!("invalid header value for {:?}", name))
65 })?;
66 }
67
68 let parsed = parse_node_addr(remote_node_id)?;
71 let mut addr = iroh::EndpointAddr::new(parsed.node_id);
72 for relay in parsed.relay_urls {
73 addr = addr.with_relay_url(relay);
74 }
75 for a in &parsed.direct_addrs {
76 addr = addr.with_ip_addr(*a);
77 }
78 if let Some(addrs) = direct_addrs {
79 for a in addrs {
80 addr = addr.with_ip_addr(*a);
81 }
82 }
83 let remote_str = crate::base32_encode(parsed.node_id.as_bytes());
84 let path = extract_path(url);
85
86 let mut req_builder = hyper::Request::builder()
88 .method(http_method)
89 .uri(&path)
90 .header(hyper::header::HOST, &remote_str);
91
92 {
97 let has_accept_encoding = headers
98 .iter()
99 .any(|(k, _)| k.eq_ignore_ascii_case("accept-encoding"));
100 if !has_accept_encoding {
101 req_builder = req_builder.header("accept-encoding", "zstd");
102 }
103 }
104 for (k, v) in headers {
105 req_builder = req_builder.header(k.as_str(), v.as_str());
106 }
107
108 let req_body: Body = if let Some(reader) = req_body_reader {
109 Body::new(reader)
110 } else {
111 Body::empty()
112 };
113 let req = req_builder
114 .body(req_body)
115 .map_err(|e| CoreError::internal(format!("build request: {e}")))?;
116
117 if remote_str == endpoint.node_id() {
122 return self_fetch(
125 endpoint,
126 req,
127 &remote_str,
128 &path,
129 fetch_token,
130 timeout,
131 max_response_body_bytes,
132 )
133 .await;
134 }
135
136 let cfg = crate::http::server::stack::StackConfig {
143 timeout,
144 decompression: decompress,
145 ..crate::http::server::stack::StackConfig::default()
146 };
147
148 let fetch_fut = async {
149 fetch_request(endpoint, &addr, req, &cfg)
150 .await
151 .map_err(fetch_error_to_core)
152 };
153 let resp = run_with_cancel_and_timeout(endpoint, fetch_token, None, fetch_fut).await?;
154
155 package_response(endpoint, resp, &remote_str, &path, max_response_body_bytes).await
156}
157
158#[allow(clippy::too_many_arguments)]
173async fn self_fetch(
174 endpoint: &IrohEndpoint,
175 mut req: hyper::Request<Body>,
176 remote_str: &str,
177 path: &str,
178 fetch_token: Option<u64>,
179 timeout: Option<Duration>,
180 max_response_body_bytes: Option<usize>,
181) -> Result<FfiResponse, CoreError> {
182 use tower::ServiceExt;
183
184 let svc = endpoint.local_service().ok_or_else(|| {
185 CoreError::connection_failed(
186 "self-request: this node has no active server to handle a request to its \
187 own node id. Call serve() before fetching httpi://<your-own-node-id>/…",
188 )
189 })?;
190
191 req.extensions_mut()
195 .insert(crate::http::server::RemoteNodeId(std::sync::Arc::new(
196 remote_str.to_string(),
197 )));
198
199 let dispatch = async move {
201 match svc.oneshot(req).await {
202 Ok(resp) => resp,
203 Err(never) => match never {},
204 }
205 };
206
207 let dispatch_fut = async { Ok::<_, CoreError>(dispatch.await) };
208 let resp = run_with_cancel_and_timeout(endpoint, fetch_token, timeout, dispatch_fut).await?;
209
210 package_response(endpoint, resp, remote_str, path, max_response_body_bytes).await
211}
212
213async fn run_with_cancel_and_timeout<F, T>(
218 endpoint: &IrohEndpoint,
219 fetch_token: Option<u64>,
220 timeout: Option<Duration>,
221 fut: F,
222) -> Result<T, CoreError>
223where
224 F: std::future::Future<Output = Result<T, CoreError>>,
225{
226 let cancel_notify = fetch_token.and_then(|t| endpoint.handles().get_fetch_cancel_notify(t));
227 let timed = async {
228 match timeout {
229 Some(t) => tokio::time::timeout(t, fut)
230 .await
231 .map_err(|_| CoreError::timeout("request timed out"))?,
232 None => fut.await,
233 }
234 };
235 let result = match cancel_notify {
236 Some(notify) => {
237 tokio::select! {
238 _ = notify.notified() => Err(CoreError::cancelled()),
239 r = timed => r,
240 }
241 }
242 None => timed.await,
243 };
244 if let Some(t) = fetch_token {
245 endpoint.handles().remove_fetch_token(t);
246 }
247 result
248}
249
250fn fetch_error_to_core(e: FetchError) -> CoreError {
253 match e {
254 FetchError::ConnectionFailed { detail, .. } => CoreError::connection_failed(detail),
255 FetchError::RequestBodyFailed { detail, .. } => {
256 CoreError::internal(format!("request body failed: {detail}"))
257 }
258 FetchError::HeaderTooLarge { detail } => CoreError::header_too_large(detail),
259 FetchError::BodyTooLarge => CoreError::body_too_large("response body too large"),
260 FetchError::Timeout => CoreError::timeout("request timed out"),
261 FetchError::Cancelled => CoreError::cancelled(),
262 FetchError::Internal(msg) => CoreError::internal(msg),
263 }
264}
265
266pub(crate) fn extract_path(url: &str) -> String {
271 let url = url.split('#').next().unwrap_or(url);
275 if let Some(rest) = url.strip_prefix("httpi://") {
276 if let Some(slash) = rest.find('/') {
277 return rest[slash..].to_string();
278 }
279 return "/".to_string();
280 }
281 if url.starts_with('/') {
282 return url.to_string();
283 }
284 format!("/{url}")
285}
286
287async fn package_response(
295 endpoint: &IrohEndpoint,
296 resp: hyper::Response<Body>,
297 remote_str: &str,
298 path: &str,
299 max_response_body_bytes: Option<usize>,
300) -> Result<FfiResponse, CoreError> {
301 let max_header_size = endpoint.max_header_size();
302 let max_response_body_bytes =
304 max_response_body_bytes.unwrap_or_else(|| endpoint.max_response_body_bytes());
305 let handles = endpoint.handles();
306
307 let status = resp.status().as_u16();
308 let header_bytes: usize = resp
311 .headers()
312 .iter()
313 .map(|(k, v)| {
314 k.as_str()
315 .len()
316 .saturating_add(v.as_bytes().len())
317 .saturating_add(4) })
319 .fold(16usize, |acc, x| acc.saturating_add(x)); if header_bytes > max_header_size {
321 return Err(CoreError::header_too_large(format!(
322 "response header size {header_bytes} exceeds limit {max_header_size}"
323 )));
324 }
325
326 let mut resp_headers: Vec<(String, String)> = Vec::new();
327 for (k, v) in resp.headers().iter() {
328 match v.to_str() {
329 Ok(s) => resp_headers.push((k.as_str().to_string(), s.to_string())),
330 Err(_) => {
331 return Err(CoreError::invalid_input(format!(
332 "non-UTF8 response header value for '{}'",
333 k.as_str()
334 )));
335 }
336 }
337 }
338
339 let response_url = format!("httpi://{remote_str}{path}");
340
341 if matches!(status, 204 | 205 | 304) {
347 drop(resp.into_body());
351 return Ok(FfiResponse {
352 status,
353 headers: resp_headers,
354 body_handle: 0,
355 url: response_url,
356 });
357 }
358
359 let mut guard = handles.insert_guard();
361 let (res_writer, res_reader) = handles.make_body_channel();
362 let body = resp.into_body();
363 let frame_timeout = res_writer.drain_timeout;
364 tokio::spawn(pump_hyper_body_to_channel_limited(
365 body,
366 res_writer,
367 Some(max_response_body_bytes),
368 frame_timeout,
369 None,
370 ));
371
372 let body_handle = guard.insert_reader(res_reader)?;
373 guard.commit();
374 Ok(FfiResponse {
375 status,
376 headers: resp_headers,
377 body_handle,
378 url: response_url,
379 })
380}
381
382#[cfg(test)]
383mod tests {
384 use super::extract_path;
385
386 #[test]
387 fn extracts_path_and_query() {
388 assert_eq!(extract_path("httpi://node/a/b?c=1"), "/a/b?c=1");
389 assert_eq!(extract_path("httpi://node/"), "/");
390 assert_eq!(extract_path("httpi://node"), "/");
391 }
392
393 #[test]
394 fn strips_fragment_from_request_target() {
395 assert_eq!(extract_path("httpi://node/a?b=1#secret"), "/a?b=1");
398 assert_eq!(extract_path("httpi://node/a#frag"), "/a");
399 assert_eq!(extract_path("httpi://node/#frag"), "/");
400 assert_eq!(extract_path("httpi://node#frag"), "/");
401 assert_eq!(extract_path("/path?q=1#frag"), "/path?q=1");
402 }
403}