Skip to content

g1t/crates/runner/src/docker/http.rs

342 lines12,953 bytesCodeBlameRaw
1//! Just enough HTTP/1.1 to stand between the Docker CLI and the Engine:
2//! message heads, and bodies framed by `Content-Length`, chunked, or the
3//! end of the connection. Bodies pass through as they come (a chunk at a
4//! time, so `docker logs -f` and a pull's progress keep streaming), unless
5//! the proxy reads one whole to change it.
6
7use std::io::{self, BufRead, Read, Write};
8
9/// The most a message head may be. The Engine's and the CLI's are a few
10/// hundred bytes.
11const MAX_HEAD: usize = 64 * 1024;
12/// The most a body read whole may be: a container's config or inspection
13/// is a few kilobytes.
14pub(crate) const MAX_BODY: u64 = 16 * 1024 * 1024;
15
16/// A request's or a response's start line and headers.
17#[derive(Clone, Debug, PartialEq)]
18pub(crate) struct Head {
19 pub(crate) start: String,
20 pub(crate) headers: Vec<(String, String)>,
21}
22
23impl Head {
24 pub(crate) fn header(&self, name: &str) -> Option<&str> {
25 self.headers.iter().find(|(key, _)| key.eq_ignore_ascii_case(name)).map(|(_, value)| value.as_str())
26 }
27
28 pub(crate) fn remove_header(&mut self, name: &str) {
29 self.headers.retain(|(key, _)| !key.eq_ignore_ascii_case(name));
30 }
31
32 pub(crate) fn set_header(&mut self, name: &str, value: &str) {
33 self.remove_header(name);
34 self.headers.push((name.to_owned(), value.to_owned()));
35 }
36
37 /// A request's method.
38 pub(crate) fn method(&self) -> &str {
39 self.start.split(' ').next().unwrap_or_default()
40 }
41
42 /// A request's target: its path and query.
43 pub(crate) fn target(&self) -> &str {
44 self.start.split(' ').nth(1).unwrap_or_default()
45 }
46
47 pub(crate) fn set_target(&mut self, target: &str) {
48 let mut parts: Vec<&str> = self.start.splitn(3, ' ').collect();
49 if parts.len() == 3 {
50 parts[1] = target;
51 self.start = parts.join(" ");
52 }
53 }
54
55 /// A response's status code.
56 pub(crate) fn status(&self) -> u16 {
57 self.start.split(' ').nth(1).and_then(|code| code.parse().ok()).unwrap_or(0)
58 }
59
60 /// Whether the request asks to leave HTTP (`docker attach`, `exec`,
61 /// BuildKit's `/grpc` and `/session`).
62 pub(crate) fn upgrades(&self) -> bool {
63 self.header("upgrade").is_some() || self.header("connection").is_some_and(|value| value.to_ascii_lowercase().contains("upgrade"))
64 }
65
66 pub(crate) fn to_bytes(&self) -> Vec<u8> {
67 let mut out = String::with_capacity(256);
68 out.push_str(&self.start);
69 out.push_str("\r\n");
70 for (name, value) in &self.headers {
71 out.push_str(name);
72 out.push_str(": ");
73 out.push_str(value);
74 out.push_str("\r\n");
75 }
76 out.push_str("\r\n");
77 out.into_bytes()
78 }
79}
80
81/// Reads a message head. `None` when the connection ends before one starts.
82pub(crate) fn read_head<R: BufRead>(reader: &mut R) -> io::Result<Option<Head>> {
83 let mut lines: Vec<String> = Vec::new();
84 let mut size = 0;
85 loop {
86 let mut line = Vec::new();
87 let read = reader.read_until(b'\n', &mut line)?;
88 if read == 0 {
89 if lines.is_empty() {
90 return Ok(None);
91 }
92 return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "the connection ended inside a message head"));
93 }
94 size += read;
95 if size > MAX_HEAD {
96 return Err(io::Error::new(io::ErrorKind::InvalidData, "a message head too large"));
97 }
98 let text = String::from_utf8_lossy(&line).trim_end_matches(['\r', '\n']).to_owned();
99 if text.is_empty() {
100 // Blank lines before a request are allowed, and skipped.
101 if lines.is_empty() {
102 continue;
103 }
104 break;
105 }
106 lines.push(text);
107 }
108 let start = lines.remove(0);
109 let headers = lines
110 .into_iter()
111 .filter_map(|line| line.split_once(':').map(|(name, value)| (name.trim().to_owned(), value.trim().to_owned())))
112 .collect();
113 Ok(Some(Head { start, headers }))
114}
115
116/// How a message's body ends.
117#[derive(Clone, Copy, Debug, PartialEq)]
118pub(crate) enum Body {
119 None,
120 Length(u64),
121 Chunked,
122 /// Until the connection closes: a response with no length.
123 UntilClose,
124}
125
126fn chunked(head: &Head) -> bool {
127 head.header("transfer-encoding").is_some_and(|value| value.to_ascii_lowercase().contains("chunked"))
128}
129
130fn length(head: &Head) -> Option<u64> {
131 head.header("content-length").and_then(|value| value.trim().parse().ok())
132}
133
134/// A request's body: chunked, a length, or none.
135pub(crate) fn request_body(head: &Head) -> Body {
136 if chunked(head) {
137 Body::Chunked
138 } else {
139 match length(head) {
140 Some(0) | None => Body::None,
141 Some(n) => Body::Length(n),
142 }
143 }
144}
145
146/// A response's body, which also depends on what was asked.
147pub(crate) fn response_body(head: &Head, method: &str) -> Body {
148 let status = head.status();
149 if method.eq_ignore_ascii_case("HEAD") || (100..200).contains(&status) || status == 204 || status == 304 {
150 return Body::None;
151 }
152 if chunked(head) {
153 return Body::Chunked;
154 }
155 match length(head) {
156 Some(0) => Body::None,
157 Some(n) => Body::Length(n),
158 None => Body::UntilClose,
159 }
160}
161
162/// Copies a body as it is framed, flushing as each piece arrives.
163pub(crate) fn copy_body<R: BufRead, W: Write>(reader: &mut R, writer: &mut W, body: Body) -> io::Result<()> {
164 match body {
165 Body::None => Ok(()),
166 Body::Length(n) => {
167 let copied = io::copy(&mut reader.take(n), writer)?;
168 writer.flush()?;
169 if copied < n {
170 return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "the body ended early"));
171 }
172 Ok(())
173 }
174 Body::UntilClose => {
175 io::copy(reader, writer)?;
176 writer.flush()
177 }
178 Body::Chunked => loop {
179 let mut line = Vec::new();
180 if reader.read_until(b'\n', &mut line)? == 0 {
181 return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "a chunked body ended early"));
182 }
183 writer.write_all(&line)?;
184 let size = chunk_size(&line)?;
185 if size == 0 {
186 // Trailers, then a blank line.
187 loop {
188 let mut trailer = Vec::new();
189 if reader.read_until(b'\n', &mut trailer)? == 0 {
190 break;
191 }
192 writer.write_all(&trailer)?;
193 if trailer == b"\r\n" || trailer == b"\n" {
194 break;
195 }
196 }
197 writer.flush()?;
198 return Ok(());
199 }
200 // The chunk and its CRLF.
201 let copied = io::copy(&mut reader.take(size + 2), writer)?;
202 writer.flush()?;
203 if copied < size + 2 {
204 return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "a chunk ended early"));
205 }
206 },
207 }
208}
209
210fn chunk_size(line: &[u8]) -> io::Result<u64> {
211 let text = String::from_utf8_lossy(line);
212 let hex = text.trim().split(';').next().unwrap_or_default().trim();
213 u64::from_str_radix(hex, 16).map_err(|_| io::Error::new(io::ErrorKind::InvalidData, format!("not a chunk size: {hex:?}")))
214}
215
216/// Reads a whole body, unframed.
217pub(crate) fn read_body<R: BufRead>(reader: &mut R, body: Body) -> io::Result<Vec<u8>> {
218 let mut out = Vec::new();
219 match body {
220 Body::None => {}
221 Body::Length(n) => {
222 if n > MAX_BODY {
223 return Err(io::Error::new(io::ErrorKind::InvalidData, "a body too large to read"));
224 }
225 reader.take(n).read_to_end(&mut out)?;
226 if (out.len() as u64) < n {
227 return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "the body ended early"));
228 }
229 }
230 Body::UntilClose => {
231 reader.take(MAX_BODY).read_to_end(&mut out)?;
232 }
233 Body::Chunked => loop {
234 let mut line = Vec::new();
235 if reader.read_until(b'\n', &mut line)? == 0 {
236 return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "a chunked body ended early"));
237 }
238 let size = chunk_size(&line)?;
239 if size == 0 {
240 loop {
241 let mut trailer = Vec::new();
242 if reader.read_until(b'\n', &mut trailer)? == 0 || trailer == b"\r\n" || trailer == b"\n" {
243 break;
244 }
245 }
246 break;
247 }
248 if out.len() as u64 + size > MAX_BODY {
249 return Err(io::Error::new(io::ErrorKind::InvalidData, "a body too large to read"));
250 }
251 let mut chunk = Vec::new();
252 reader.take(size).read_to_end(&mut chunk)?;
253 out.extend_from_slice(&chunk);
254 let mut crlf = Vec::new();
255 reader.read_until(b'\n', &mut crlf)?;
256 },
257 }
258 Ok(out)
259}
260
261/// A head with its body framed by length, for a body the proxy rewrote.
262pub(crate) fn with_length(mut head: Head, body: &[u8]) -> Vec<u8> {
263 head.remove_header("transfer-encoding");
264 head.set_header("Content-Length", &body.len().to_string());
265 let mut out = head.to_bytes();
266 out.extend_from_slice(body);
267 out
268}
269
270/// A whole JSON response, as the Engine words its errors.
271pub(crate) fn json_response(status: u16, reason: &str, body: &serde_json::Value) -> Vec<u8> {
272 let text = body.to_string();
273 let head = Head {
274 start: format!("HTTP/1.1 {status} {reason}"),
275 headers: vec![("Content-Type".into(), "application/json".into())],
276 };
277 with_length(head, text.as_bytes())
278}
279
280/// An empty `200 OK`.
281pub(crate) fn empty_ok() -> Vec<u8> {
282 with_length(Head { start: "HTTP/1.1 200 OK".into(), headers: Vec::new() }, b"")
283}
284
285#[cfg(test)]
286mod tests {
287 use super::*;
288 use std::io::BufReader;
289
290 #[test]
291 fn heads_are_read_and_written_back() {
292 let raw = b"POST /v1.47/containers/create?name=db HTTP/1.1\r\nHost: api.moby.localhost\r\nContent-Type: application/json\r\nContent-Length: 2\r\n\r\n{}";
293 let mut reader = BufReader::new(&raw[..]);
294 let head = read_head(&mut reader).unwrap().unwrap();
295 assert_eq!(head.method(), "POST");
296 assert_eq!(head.target(), "/v1.47/containers/create?name=db");
297 assert_eq!(head.header("content-length"), Some("2"));
298 assert_eq!(request_body(&head), Body::Length(2));
299 assert_eq!(read_body(&mut reader, Body::Length(2)).unwrap(), b"{}");
300 assert!(read_head(&mut reader).unwrap().is_none());
301 let mut again = head.clone();
302 again.set_target("/v1.47/containers/create");
303 assert!(String::from_utf8(again.to_bytes()).unwrap().starts_with("POST /v1.47/containers/create HTTP/1.1\r\n"));
304 }
305
306 #[test]
307 fn chunked_bodies_copy_verbatim_and_read_unframed() {
308 let raw = b"5\r\nhello\r\n6;ext=1\r\n world\r\n0\r\n\r\nNEXT";
309 let mut copied = Vec::new();
310 let mut reader = BufReader::new(&raw[..]);
311 copy_body(&mut reader, &mut copied, Body::Chunked).unwrap();
312 assert_eq!(copied, &raw[..raw.len() - 4]);
313 let mut rest = String::new();
314 reader.read_to_string(&mut rest).unwrap();
315 assert_eq!(rest, "NEXT");
316 let mut reader = BufReader::new(&raw[..]);
317 assert_eq!(read_body(&mut reader, Body::Chunked).unwrap(), b"hello world");
318 }
319
320 #[test]
321 fn response_bodies_follow_the_request_and_the_status() {
322 let head = |start: &str, headers: &[(&str, &str)]| Head {
323 start: start.into(),
324 headers: headers.iter().map(|(k, v)| (k.to_string(), v.to_string())).collect(),
325 };
326 assert_eq!(response_body(&head("HTTP/1.1 200 OK", &[("Content-Length", "10")]), "HEAD"), Body::None);
327 assert_eq!(response_body(&head("HTTP/1.1 204 No Content", &[]), "POST"), Body::None);
328 assert_eq!(response_body(&head("HTTP/1.1 101 UPGRADED", &[]), "POST"), Body::None);
329 assert_eq!(response_body(&head("HTTP/1.1 200 OK", &[("Transfer-Encoding", "chunked")]), "GET"), Body::Chunked);
330 assert_eq!(response_body(&head("HTTP/1.1 200 OK", &[]), "GET"), Body::UntilClose);
331 assert!(head("POST /grpc HTTP/1.1", &[("Connection", "Upgrade"), ("Upgrade", "h2c")]).upgrades());
332 }
333
334 #[test]
335 fn a_rewritten_body_gets_its_length() {
336 let head = Head { start: "HTTP/1.1 200 OK".into(), headers: vec![("Transfer-Encoding".into(), "chunked".into())] };
337 let out = String::from_utf8(with_length(head, b"{\"a\":1}")).unwrap();
338 assert!(out.contains("Content-Length: 7\r\n"));
339 assert!(!out.to_ascii_lowercase().contains("chunked"));
340 assert!(out.ends_with("\r\n\r\n{\"a\":1}"));
341 }
342}