Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_net/src/http/client.rs

20.4 KiB, 62 runs

created by r1870400018:9529, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Minimal async HTTPS client built on `tokio` + `tokio_rustls`.
2//!
3//! The pattern this module implements -- `TcpStream::connect` → TLS wrap via
4//! `TlsConnector::from(Arc<ClientConfig>)` → write a `HttpMessage`-shaped
5//! request → read a `HttpMessage` response -- already existed inside
6//! `fe2o3_steel/tests/client.rs` for test harness purposes. This module hoists
7//! it into `fe2o3_net` as a reusable primitive that any crate in the
8//! workspace can call without reinventing it.
9//!
10//! Design choices kept deliberately small:
11//!
12//! - One request per connection, closed via `Connection: close`. No keep-alive,
13//! no pipelining, no HTTP/2. Sufficient for RFC 8555 ACME traffic and for
14//! the outbound HTTPS needs of SMTP webhooks, WebSocket handshakes to
15//! remote servers and similar short-lived call patterns.
16//! - No trust store is bundled. The caller supplies an
17//! `Arc<rustls::ClientConfig>` that already carries whatever root anchors
18//! they want to trust, and `fe2o3_net` stays free of `webpki-roots` or
19//! `rustls-native-certs`. The ACME client under `fe2o3_net/src/acme/`
20//! compiles in its own pinned Let's Encrypt root anchors rather than
21//! pulling a generic trust store.
22//! - Responses are read with `HttpMessage::read` using the existing default
23//! chunk sizes from `fe2o3_net::constant`. Chunked transfer encoding is
24//! not supported: ACME API responses always carry a `Content-Length`
25//! header, and that is the only production caller for now.
26//!
27//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
28//! Anthropic Claude
29
30use crate::{
31 constant,
32 http::{
33 header::HttpMethod,
34 msg::{
35 HttpMessage,
36 ReadLimits,
37 },
38 },
39};
40
41use oxedyne_fe2o3_core::prelude::*;
42
43use std::{
44 net::SocketAddr,
45 pin::Pin,
46 sync::Arc,
47};
48
49use tokio::{
50 io::{
51 AsyncRead,
52 AsyncWrite,
53 AsyncWriteExt,
54 },
55 net::TcpStream,
56};
57use tokio_rustls::{
58 rustls::{
59 pki_types::ServerName,
60 ClientConfig,
61 },
62 TlsConnector,
63};
64
65
66// ┌───────────────────────────────────────────────────────────────────────────┐
67// │ REQUEST FORMATTING │
68// └───────────────────────────────────────────────────────────────────────────┘
69
70/// `host` goes into the `Host:` header and `path` is the request target, query
71/// string included. `Host`, `Content-Length` and `Connection: close` are always
72/// written here, so a caller must not repeat them in `headers`.
73///
74/// Factored out so the byte layout can be tested without a TLS socket.
75pub fn format_request(
76 method: HttpMethod,
77 host: &str,
78 path: &str,
79 headers: &[(&str, &str)],
80 body: &[u8],
81)
82 -> Vec<u8>
83{
84 let mut out = String::with_capacity(256 + body.len());
85 out.push_str(&fmt!("{} {} HTTP/1.1\r\n", method, path));
86 out.push_str(&fmt!("Host: {}\r\n", host));
87 out.push_str("Connection: close\r\n");
88 for (name, value) in headers {
89 out.push_str(&fmt!("{}: {}\r\n", name, value));
90 }
91 out.push_str(&fmt!("Content-Length: {}\r\n", body.len()));
92 out.push_str("\r\n");
93 let mut bytes = out.into_bytes();
94 bytes.extend_from_slice(body);
95 bytes
96}
97
98
99// ┌───────────────────────────────────────────────────────────────────────────┐
100// │ THE EXCHANGE │
101// └───────────────────────────────────────────────────────────────────────────┘
102
103/// The half of a request that does not care whether the stream beneath it is
104/// TLS-wrapped, and so is shared by all four entry points below.
105async fn exchange<S>(
106 stream: &mut S,
107 request_bytes: &[u8],
108 peer: &str,
109 limits: Option<&ReadLimits>,
110)
111 -> Outcome<HttpMessage>
112where
113 S: AsyncRead + AsyncWrite + Unpin,
114{
115 match stream.write_all(request_bytes).await {
116 Ok(()) => (),
117 Err(e) => return Err(err!(e,
118 "Failed to write HTTP request body to {}.", peer;
119 IO, Network, Wire, Write)),
120 }
121 match stream.flush().await {
122 Ok(()) => (),
123 Err(e) => return Err(err!(e,
124 "Failed to flush HTTP request to {}.", peer;
125 IO, Network, Wire, Write)),
126 }
127
128 let result = HttpMessage::read::<
129 { constant::HTTP_DEFAULT_HEADER_CHUNK_SIZE },
130 { constant::HTTP_DEFAULT_BODY_CHUNK_SIZE },
131 _,
132 >(
133 Pin::new(stream),
134 &Vec::new(),
135 Some(false),
136 limits,
137 ).await;
138
139 match result {
140 Ok((Some(msg), _remnant)) => Ok(msg),
141 Ok((None, _)) => Err(err!(
142 "Server at {} closed the connection before sending a \
143 complete HTTP response.",
144 peer;
145 IO, Network, Wire, Read, Missing)),
146 // Kept distinct, so a caller that set `limits` can tell an answer larger than it
147 // allows from one that broke off, without reading the text.
148 Err(e) if e.tags().contains(&ErrTag::TooBig) => Err(err!(e,
149 "The HTTP response from {} is larger than this caller reads.", peer;
150 IO, Network, Wire, Read, TooBig)),
151 Err(e) => Err(err!(e,
152 "Failed to read or parse the HTTP response from {}.", peer;
153 IO, Network, Wire, Read)),
154 }
155}
156
157
158// ┌───────────────────────────────────────────────────────────────────────────┐
159// │ HTTP REQUEST (PLAIN) │
160// └───────────────────────────────────────────────────────────────────────────┘
161
162/// The sibling of [`https_request`] for an upstream over loopback or another
163/// trusted segment, where TLS buys nothing: a proxied app binding
164/// `127.0.0.1:<port>` need not present a certificate for traffic that never
165/// leaves the host, and insisting on one would mean an internal CA to rotate and
166/// a handshake on every hit.
167pub async fn http_request(
168 host: &str,
169 port: u16,
170 method: HttpMethod,
171 path: &str,
172 headers: &[(&str, &str)],
173 body: &[u8],
174)
175 -> Outcome<HttpMessage>
176{
177 http_request_limited(host, port, method, path, headers, body, None).await
178}
179
180/// [`http_request`] with the response bounded by `limits`, for a caller that
181/// does not trust the peer to answer in proportion: a health probe, say, whose
182/// every reply is kept.
183pub async fn http_request_limited(
184 host: &str,
185 port: u16,
186 method: HttpMethod,
187 path: &str,
188 headers: &[(&str, &str)],
189 body: &[u8],
190 limits: Option<&ReadLimits>,
191)
192 -> Outcome<HttpMessage>
193{
194 let request_bytes = format_request(method, host, path, headers, body);
195 let peer = fmt!("{}:{}", host, port);
196
197 let mut stream = match TcpStream::connect((host, port)).await {
198 Ok(s) => s,
199 Err(e) => return Err(err!(e,
200 "Failed to open a TCP connection to {}.", peer;
201 IO, Network, Init)),
202 };
203
204 exchange(&mut stream, &request_bytes, &peer, limits).await
205}
206
207/// Dials an address the caller has already vetted, rather than a host name this
208/// would resolve for itself.
209///
210/// The distinction is the whole point. A server that connects somewhere its
211/// user named must check the address first (see
212/// [`crate::addr::resolve_public`]), and a check is worthless if the name is
213/// then resolved a second time to dial it: the answer can change in between,
214/// and DNS rebinding is precisely that trick. So the caller resolves once,
215/// vets what came back, and hands the surviving address here. `host` is still
216/// needed, but only for the `Host` header the origin server reads.
217///
218/// `limits` bounds the response, so a caller fetching a page on a user's
219/// behalf can cap what it is willing to read.
220pub async fn http_request_at(
221 addr: SocketAddr,
222 host: &str,
223 method: HttpMethod,
224 path: &str,
225 headers: &[(&str, &str)],
226 body: &[u8],
227 limits: Option<&ReadLimits>,
228)
229 -> Outcome<HttpMessage>
230{
231 let request_bytes = format_request(method, host, path, headers, body);
232 let peer = fmt!("{} ({})", host, addr);
233
234 let mut stream = match TcpStream::connect(addr).await {
235 Ok(s) => s,
236 Err(e) => return Err(err!(e,
237 "Failed to open a TCP connection to {}.", peer;
238 IO, Network, Init)),
239 };
240
241 exchange(&mut stream, &request_bytes, &peer, limits).await
242}
243
244
245// ┌───────────────────────────────────────────────────────────────────────────┐
246// │ HTTPS REQUEST │
247// └───────────────────────────────────────────────────────────────────────────┘
248
249/// The request is formed as [`format_request`] describes, and `tls_config` is
250/// the rustls configuration the caller built, normally a trust store of root CAs
251/// and no client auth.
252///
253/// Each step's error is tagged `IO`, `Network` and, where it applies, `Wire`, so
254/// a caller can tell a connect failure from a handshake failure from a
255/// response-parse failure without reading the text.
256pub async fn https_request(
257 host: &str,
258 port: u16,
259 method: HttpMethod,
260 path: &str,
261 headers: &[(&str, &str)],
262 body: &[u8],
263 tls_config: Arc<ClientConfig>,
264)
265 -> Outcome<HttpMessage>
266{
267 https_request_limited(host, port, method, path, headers, body, tls_config, None).await
268}
269
270/// [`https_request`] with the response bounded by `limits`, as
271/// [`http_request_limited`] is.
272pub async fn https_request_limited(
273 host: &str,
274 port: u16,
275 method: HttpMethod,
276 path: &str,
277 headers: &[(&str, &str)],
278 body: &[u8],
279 tls_config: Arc<ClientConfig>,
280 limits: Option<&ReadLimits>,
281)
282 -> Outcome<HttpMessage>
283{
284 // Format the request bytes up front so any failure from this point on is
285 // a real network or TLS fault, not a local formatting bug.
286 let request_bytes = format_request(method, host, path, headers, body);
287 let peer = fmt!("{}:{}", host, port);
288
289 // TCP connect to the remote server.
290 let tcp = match TcpStream::connect((host, port)).await {
291 Ok(s) => s,
292 Err(e) => return Err(err!(e,
293 "Failed to open a TCP connection to {}.", peer;
294 IO, Network, Init)),
295 };
296
297 let mut stream = res!(tls_wrap(tcp, host, &peer, tls_config).await);
298 exchange(&mut stream, &request_bytes, &peer, limits).await
299}
300
301/// The TLS sibling of [`http_request_at`], and vetted for the same reason: the
302/// address is dialled as given, while `host` names the certificate that must
303/// validate and fills the `Host` header. Pinning the address does not weaken
304/// the TLS check -- the server still has to present a certificate for the name
305/// the caller asked for.
306pub async fn https_request_at(
307 addr: SocketAddr,
308 host: &str,
309 method: HttpMethod,
310 path: &str,
311 headers: &[(&str, &str)],
312 body: &[u8],
313 tls_config: Arc<ClientConfig>,
314 limits: Option<&ReadLimits>,
315)
316 -> Outcome<HttpMessage>
317{
318 let request_bytes = format_request(method, host, path, headers, body);
319 let peer = fmt!("{} ({})", host, addr);
320
321 let tcp = match TcpStream::connect(addr).await {
322 Ok(s) => s,
323 Err(e) => return Err(err!(e,
324 "Failed to open a TCP connection to {}.", peer;
325 IO, Network, Init)),
326 };
327
328 let mut stream = res!(tls_wrap(tcp, host, &peer, tls_config).await);
329 exchange(&mut stream, &request_bytes, &peer, limits).await
330}
331
332/// rustls needs the host name as a validated `ServerName`, so that it can send
333/// the right SNI and check the server certificate's SANs against it.
334async fn tls_wrap(
335 tcp: TcpStream,
336 host: &str,
337 peer: &str,
338 tls_config: Arc<ClientConfig>,
339)
340 -> Outcome<tokio_rustls::client::TlsStream<TcpStream>>
341{
342 let server_name = match ServerName::try_from(host.to_string()) {
343 Ok(n) => n,
344 Err(e) => return Err(err!(e,
345 "Host {:?} is not a valid DNS name for TLS SNI.", host;
346 IO, Network, Invalid, Input)),
347 };
348 let connector = TlsConnector::from(tls_config);
349 match connector.connect(server_name, tcp).await {
350 Ok(s) => Ok(s),
351 Err(e) => Err(err!(e,
352 "TLS handshake with {} failed.", peer;
353 IO, Network, Init)),
354 }
355}
356
357
358// ┌───────────────────────────────────────────────────────────────────────────┐
359// │ TESTS │
360// └───────────────────────────────────────────────────────────────────────────┘
361
362#[cfg(test)]
363mod tests {
364 use super::*;
365
366 /// Split the wire bytes at the `\r\n\r\n` boundary between header block
367 /// and body so assertions can inspect them separately. The header block
368 /// keeps the `\r\n` that terminates its last header line so that every
369 /// header line in the returned string ends with `\r\n` consistently
370 /// (the empty-line half of the separator is dropped). Returns
371 /// `(header_block, body)`.
372 fn split_wire(bytes: &[u8]) -> (String, Vec<u8>) {
373 let sep = b"\r\n\r\n";
374 let pos = bytes.windows(sep.len())
375 .position(|w| w == sep)
376 .expect("wire bytes did not contain an HTTP header/body separator");
377 let header = String::from_utf8(bytes[..pos + 2].to_vec())
378 .expect("header block was not valid UTF-8");
379 let body = bytes[pos + sep.len()..].to_vec();
380 (header, body)
381 }
382
383 fn count_occurrences(haystack: &str, needle: &str) -> usize {
384 haystack.matches(needle).count()
385 }
386
387 /// A GET request with no body must emit the correct request line, a
388 /// `Host` header, `Connection: close`, and `Content-Length: 0`, with an
389 /// empty body.
390 #[test]
391 fn test_format_request_get_no_body() -> Outcome<()> {
392 let bytes = format_request(
393 HttpMethod::GET,
394 "acme-v02.api.letsencrypt.org",
395 "/directory",
396 &[],
397 &[],
398 );
399 let (header, body) = split_wire(&bytes);
400
401 if !header.starts_with("GET /directory HTTP/1.1\r\n") {
402 return Err(err!(
403 "Expected request line 'GET /directory HTTP/1.1', got first \
404 line: {:?}.",
405 header.lines().next().unwrap_or("");
406 Test, Mismatch));
407 }
408 if !header.contains("Host: acme-v02.api.letsencrypt.org\r\n") {
409 return Err(err!(
410 "Missing or wrong Host header in:\n{}", header;
411 Test, Missing));
412 }
413 if !header.contains("Connection: close\r\n") {
414 return Err(err!(
415 "Missing Connection: close header in:\n{}", header;
416 Test, Missing));
417 }
418 if !header.contains("Content-Length: 0\r\n") {
419 return Err(err!(
420 "Missing Content-Length: 0 header in:\n{}", header;
421 Test, Missing));
422 }
423 if !body.is_empty() {
424 return Err(err!(
425 "Expected empty body for a GET request, got {} bytes.",
426 body.len();
427 Test, Mismatch));
428 }
429 Ok(())
430 }
431
432 /// A POST with a body must emit the correct Content-Length and place the
433 /// body bytes verbatim after the header terminator.
434 #[test]
435 fn test_format_request_post_with_body() -> Outcome<()> {
436 let payload = br#"{"protected":"...","payload":"...","signature":"..."}"#;
437 let bytes = format_request(
438 HttpMethod::POST,
439 "acme-v02.api.letsencrypt.org",
440 "/acme/new-order",
441 &[("Content-Type", "application/jose+json")],
442 payload,
443 );
444 let (header, body) = split_wire(&bytes);
445
446 if !header.starts_with("POST /acme/new-order HTTP/1.1\r\n") {
447 return Err(err!(
448 "Expected request line 'POST /acme/new-order HTTP/1.1', got \
449 first line: {:?}.",
450 header.lines().next().unwrap_or("");
451 Test, Mismatch));
452 }
453 if !header.contains("Content-Type: application/jose+json\r\n") {
454 return Err(err!(
455 "Missing or wrong Content-Type header in:\n{}", header;
456 Test, Missing));
457 }
458 let expected_len_line = fmt!("Content-Length: {}\r\n", payload.len());
459 if !header.contains(&expected_len_line) {
460 return Err(err!(
461 "Missing or wrong {:?} header in:\n{}",
462 expected_len_line, header;
463 Test, Mismatch));
464 }
465 if body != payload {
466 return Err(err!(
467 "Body bytes did not round-trip: expected {} bytes, got {}.",
468 payload.len(), body.len();
469 Test, Mismatch));
470 }
471 Ok(())
472 }
473
474 /// Custom headers supplied by the caller must appear in the header block,
475 /// without duplicating `Host`, `Connection` or `Content-Length`.
476 #[test]
477 fn test_format_request_custom_headers() -> Outcome<()> {
478 let bytes = format_request(
479 HttpMethod::POST,
480 "example.test",
481 "/acme/order/1",
482 &[
483 ("Content-Type", "application/jose+json"),
484 ("User-Agent", "hematite-acme/0.5"),
485 ("Accept", "application/json"),
486 ],
487 b"{}",
488 );
489 let (header, _body) = split_wire(&bytes);
490
491 // Our three custom headers must each appear exactly once.
492 for name in ["Content-Type", "User-Agent", "Accept"] {
493 let line_prefix = fmt!("{}: ", name);
494 if count_occurrences(&header, &line_prefix) != 1 {
495 return Err(err!(
496 "Expected exactly one {:?} header in:\n{}",
497 line_prefix, header;
498 Test, Mismatch));
499 }
500 }
501
502 // Managed headers must still appear exactly once.
503 for needle in [
504 "Host: example.test\r\n",
505 "Connection: close\r\n",
506 "Content-Length: 2\r\n",
507 ] {
508 if count_occurrences(&header, needle) != 1 {
509 return Err(err!(
510 "Expected exactly one occurrence of {:?} in:\n{}",
511 needle, header;
512 Test, Mismatch));
513 }
514 }
515 Ok(())
516 }
517
518 /// Header block must always end with an empty line (`\r\n\r\n`), even
519 /// when no custom headers are supplied.
520 #[test]
521 fn test_format_request_terminator() -> Outcome<()> {
522 let bytes = format_request(
523 HttpMethod::GET,
524 "example.test",
525 "/",
526 &[],
527 &[],
528 );
529 let sep = b"\r\n\r\n";
530 if !bytes.windows(sep.len()).any(|w| w == sep) {
531 return Err(err!(
532 "Formatted request does not contain the CRLFCRLF header \
533 terminator required by RFC 7230 §3.";
534 Test, Missing));
535 }
536 Ok(())
537 }
538}