Skip to content

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

727 lines35,025 bytesCodeBlame
1//! The Engine API as a job sees it: `/var/run/docker.sock` is this proxy,
2//! and the job's own Docker Engine is behind it. It passes everything
3//! through, and changes the few requests that would not work in a
4//! sandbox as they are:
5//!
6//! - **A container's network.** A sandbox cannot route a container's
7//! bridge network out (Cloudflare Containers allow no iptables and no IP
8//! forwarding), so containers join the job's own network (`host`), the
9//! one its guardrails apply to. The names a container would have had on
10//! its network (its name, its aliases, its Compose service) resolve to
11//! 127.0.0.1 in later containers and in the job's own steps, and a port
12//! published under another number (`-p 8080:80`) is forwarded to it.
13//! `docker inspect` reports the ports as published, for tools that look
14//! a container's port up.
15//! - **Networks joined later** (`docker network connect`): the container
16//! already has the job's network, so the request only adds its aliases.
17//! - **The classic builder** (`DOCKER_BUILDKIT=0`): its `RUN` steps use the
18//! job's network too. BuildKit's are handled where they start (oci.rs).
19//! - **Miners**: a container whose image or command names one is refused,
20//! as a step's script would be.
21
22use std::collections::{BTreeMap, BTreeSet};
23use std::io::{self, BufReader, Read, Write};
24use std::sync::mpsc;
25use std::sync::{Arc, Mutex};
26
27use serde_json::{Map, Value, json};
28
29use super::http;
30
31/// What the proxy remembers across connections.
32#[derive(Default)]
33pub(crate) struct State {
34 /// Every name a container has been given, which all now mean 127.0.0.1.
35 pub(crate) aliases: BTreeSet<String>,
36 /// Containers moved to the job's network: by id and by name, the ports
37 /// they publish (`80/tcp` → the host port, as text).
38 pub(crate) published: BTreeMap<String, BTreeMap<String, String>>,
39}
40
41/// What the proxy asks of the sandbox around it.
42pub(crate) trait Host: Send + Sync {
43 /// Makes names resolve to 127.0.0.1 in the job's own steps.
44 fn add_hosts(&self, names: &[String]);
45 /// Forwards a host port to a container's port on the job's network.
46 fn forward(&self, host_port: u16, container_port: u16);
47 /// Makes sure the Engine is running; why not, if it cannot be.
48 fn ensure_engine(&self) -> Result<(), String>;
49 /// Connects to the Engine itself.
50 fn connect(&self) -> io::Result<Box<dyn Duplex>>;
51}
52
53/// A connection, readable and writable from two threads.
54pub(crate) trait Duplex: Read + Write + Send {
55 fn try_clone_box(&self) -> io::Result<Box<dyn Duplex>>;
56 fn shutdown_both(&self);
57 /// Ends what this side sends, and goes on reading.
58 fn shutdown_write(&self);
59}
60
61impl Duplex for std::net::TcpStream {
62 fn try_clone_box(&self) -> io::Result<Box<dyn Duplex>> {
63 Ok(Box::new(self.try_clone()?))
64 }
65 fn shutdown_both(&self) {
66 let _ = self.shutdown(std::net::Shutdown::Both);
67 }
68 fn shutdown_write(&self) {
69 let _ = self.shutdown(std::net::Shutdown::Write);
70 }
71}
72
73#[cfg(unix)]
74impl Duplex for std::os::unix::net::UnixStream {
75 fn try_clone_box(&self) -> io::Result<Box<dyn Duplex>> {
76 Ok(Box::new(self.try_clone()?))
77 }
78 fn shutdown_both(&self) {
79 let _ = self.shutdown(std::net::Shutdown::Both);
80 }
81 fn shutdown_write(&self) {
82 let _ = self.shutdown(std::net::Shutdown::Write);
83 }
84}
85
86/// The Engine's routes start with an optional `/v1.NN`.
87fn route(target: &str) -> (&str, &str) {
88 let (path, query) = target.split_once('?').unwrap_or((target, ""));
89 let path = match path.strip_prefix("/v") {
90 Some(rest) if rest.starts_with(|c: char| c.is_ascii_digit()) => rest.find('/').map_or(path, |slash| &rest[slash..]),
91 _ => path,
92 };
93 (path, query)
94}
95
96fn query_value<'a>(query: &'a str, name: &str) -> Option<&'a str> {
97 query.split('&').find_map(|pair| pair.split_once('=').filter(|(key, _)| *key == name).map(|(_, value)| value))
98}
99
100fn decode(text: &str) -> String {
101 let bytes = text.as_bytes();
102 let mut out = Vec::with_capacity(bytes.len());
103 let mut i = 0;
104 while i < bytes.len() {
105 let hex = |b: u8| (b as char).to_digit(16);
106 match bytes[i] {
107 b'%' if i + 2 < bytes.len() && hex(bytes[i + 1]).is_some() && hex(bytes[i + 2]).is_some() => {
108 out.push((hex(bytes[i + 1]).unwrap_or(0) * 16 + hex(bytes[i + 2]).unwrap_or(0)) as u8);
109 i += 3;
110 continue;
111 }
112 b'+' => out.push(b' '),
113 byte => out.push(byte),
114 }
115 i += 1;
116 }
117 String::from_utf8_lossy(&out).into_owned()
118}
119
120/// What a request is, to the proxy.
121#[derive(Debug, PartialEq)]
122pub(crate) enum Kind {
123 CreateContainer { name: Option<String> },
124 InspectContainer { id: String },
125 ConnectNetwork,
126 DisconnectNetwork,
127 ClassicBuild,
128 Other,
129}
130
131pub(crate) fn classify(method: &str, target: &str) -> Kind {
132 let (path, query) = route(target);
133 let parts: Vec<&str> = path.trim_matches('/').split('/').collect();
134 match (method, parts.as_slice()) {
135 ("POST", ["containers", "create"]) => Kind::CreateContainer { name: query_value(query, "name").map(decode).filter(|n| !n.is_empty()) },
136 ("GET", ["containers", id, "json"]) => Kind::InspectContainer { id: decode(id) },
137 ("POST", ["networks", _, "connect"]) => Kind::ConnectNetwork,
138 ("POST", ["networks", _, "disconnect"]) => Kind::DisconnectNetwork,
139 ("POST", ["build"]) => Kind::ClassicBuild,
140 _ => Kind::Other,
141 }
142}
143
144/// Whether a network mode is one a sandbox can run as it is.
145fn keeps_network(mode: &str) -> bool {
146 matches!(mode, "host" | "none") || mode.starts_with("container:")
147}
148
149/// Whether a name can stand in `/etc/hosts`.
150fn host_name(name: &str) -> bool {
151 !name.is_empty()
152 && name.len() <= 253
153 && name != "localhost"
154 && name.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '.' | '_'))
155 && !name.starts_with(['-', '.'])
156}
157
158/// What changing a container's create request came to.
159#[derive(Debug, Default, PartialEq)]
160pub(crate) struct Created {
161 /// The network it asked for, now the job's.
162 pub(crate) moved_from: Option<String>,
163 /// Its own names, now meaning 127.0.0.1.
164 pub(crate) aliases: Vec<String>,
165 /// `(host port, container port)` pairs to forward.
166 pub(crate) forwards: Vec<(u16, u16)>,
167 /// What `docker inspect` should say it publishes.
168 pub(crate) published: BTreeMap<String, String>,
169}
170
171fn strings(value: Option<&Value>) -> Vec<String> {
172 value.and_then(Value::as_array).map(|items| items.iter().filter_map(|v| v.as_str().map(str::to_owned)).collect()).unwrap_or_default()
173}
174
175/// Moves a container to the job's network, unless it asked for `host`,
176/// `none` or another container's. `known` is every name given so far.
177pub(crate) fn rewrite_create(body: &mut Value, name: Option<&str>, known: &BTreeSet<String>) -> Created {
178 let mut created = Created::default();
179 let Some(config) = body.as_object_mut() else { return created };
180 let host_config = config.entry("HostConfig").or_insert_with(|| json!({}));
181 if !host_config.is_object() {
182 *host_config = json!({});
183 }
184 let mode = host_config.get("NetworkMode").and_then(Value::as_str).unwrap_or_default().to_owned();
185 if keeps_network(&mode) {
186 return created;
187 }
188 created.moved_from = Some(if mode.is_empty() || mode == "default" { "bridge".to_owned() } else { mode.clone() });
189
190 // Its names: its own, its aliases on each network, its links' aliases
191 // and its Compose service.
192 let mut aliases: Vec<String> = Vec::new();
193 if let Some(name) = name {
194 aliases.push(name.trim_start_matches('/').to_owned());
195 }
196 if let Some(Value::Object(endpoints)) = config.get("NetworkingConfig").and_then(|n| n.get("EndpointsConfig")) {
197 for endpoint in endpoints.values() {
198 aliases.extend(strings(endpoint.get("Aliases")));
199 aliases.extend(strings(endpoint.get("DNSNames")));
200 }
201 }
202 if let Some(service) = config.get("Labels").and_then(|l| l.get("com.docker.compose.service")).and_then(Value::as_str) {
203 aliases.push(service.to_owned());
204 }
205 let host_config = config.get_mut("HostConfig").and_then(Value::as_object_mut).expect("made above");
206 for link in strings(host_config.get("Links")) {
207 // `name:alias`, or `/name:/container/alias` as the Engine stores them.
208 if let Some((_, alias)) = link.split_once(':') {
209 aliases.push(alias.rsplit('/').next().unwrap_or(alias).to_owned());
210 }
211 }
212 let mut seen = BTreeSet::new();
213 aliases.retain(|alias| host_name(alias) && seen.insert(alias.clone()));
214 created.aliases = aliases.clone();
215
216 // Ports: published under the same number, they need nothing; under
217 // another, a forward. A port left to the Engine to choose is the
218 // container's own.
219 if let Some(Value::Object(bindings)) = host_config.get("PortBindings") {
220 for (port, hosts) in bindings {
221 let (number, protocol) = port.split_once('/').unwrap_or((port, "tcp"));
222 let Ok(container_port) = number.parse::<u16>() else { continue };
223 let host_port = hosts
224 .as_array()
225 .and_then(|list| list.iter().find_map(|h| h.get("HostPort").and_then(Value::as_str).filter(|p| !p.is_empty())))
226 .and_then(|p| p.parse::<u16>().ok())
227 .unwrap_or(container_port);
228 created.published.insert(format!("{container_port}/{protocol}"), host_port.to_string());
229 if host_port != container_port && protocol == "tcp" {
230 created.forwards.push((host_port, container_port));
231 }
232 }
233 }
234
235 host_config.insert("NetworkMode".into(), json!("host"));
236 host_config.remove("Links");
237 host_config.insert("PublishAllPorts".into(), json!(false));
238 let mut extra = strings(host_config.get("ExtraHosts"));
239 for alias in known.iter().chain(aliases.iter()) {
240 let entry = format!("{alias}:127.0.0.1");
241 if !extra.iter().any(|e| e.split(':').next() == Some(alias.as_str())) {
242 extra.push(entry);
243 }
244 }
245 host_config.insert("ExtraHosts".into(), json!(extra));
246 config.remove("NetworkingConfig");
247 config.remove("MacAddress");
248 created
249}
250
251/// What `docker network connect` adds: the container's aliases there.
252pub(crate) fn connect_aliases(body: &Value) -> Vec<String> {
253 let mut aliases = strings(body.get("EndpointConfig").and_then(|e| e.get("Aliases")));
254 aliases.extend(strings(body.get("EndpointConfig").and_then(|e| e.get("DNSNames"))));
255 aliases.retain(|alias| host_name(alias));
256 aliases
257}
258
259/// A container's inspection, with the ports it publishes filled in.
260pub(crate) fn rewrite_inspect(body: &mut Value, published: &BTreeMap<String, String>) {
261 let ports: Map<String, Value> = published
262 .iter()
263 .map(|(port, host)| (port.clone(), json!([{ "HostIp": "0.0.0.0", "HostPort": host }])))
264 .collect();
265 if let Some(settings) = body.get_mut("NetworkSettings").and_then(Value::as_object_mut) {
266 settings.insert("Ports".into(), Value::Object(ports));
267 }
268}
269
270/// The classic builder's `RUN` steps on the job's network.
271pub(crate) fn rewrite_build(target: &str) -> String {
272 let (path, query) = target.split_once('?').unwrap_or((target, ""));
273 let mode = query_value(query, "networkmode").unwrap_or_default();
274 if keeps_network(mode) {
275 return target.to_owned();
276 }
277 let mut pairs: Vec<&str> = query.split('&').filter(|pair| !pair.is_empty() && !pair.starts_with("networkmode=")).collect();
278 pairs.push("networkmode=host");
279 format!("{path}?{}", pairs.join("&"))
280}
281
282/// What the response thread is to do with the next response.
283enum Pending {
284 /// Pass it through.
285 Plain { method: String },
286 /// A container was created: note its id against its ports.
287 Created { name: Option<String>, published: BTreeMap<String, String> },
288 /// Fill in an inspection's ports.
289 Inspect { published: BTreeMap<String, String> },
290 /// Answer it here: the Engine was never asked.
291 Answer(Vec<u8>),
292 /// The connection leaves HTTP after this response.
293 Upgrade,
294}
295
296/// Serves one connection from a client, until either side closes it.
297pub(crate) fn serve(client: Box<dyn Duplex>, host: Arc<dyn Host>, state: Arc<Mutex<State>>) {
298 if let Err(problem) = host.ensure_engine() {
299 let mut client = client;
300 let mut reader = BufReader::new(client.try_clone_box().expect("a clone"));
301 // Answer whatever was asked, so the CLI says why.
302 if let Ok(Some(head)) = http::read_head(&mut reader) {
303 let _ = http::read_body(&mut reader, http::request_body(&head));
304 let _ = client.write_all(&http::json_response(503, "Service Unavailable", &json!({ "message": problem })));
305 }
306 client.shutdown_both();
307 return;
308 }
309 let upstream = match host.connect() {
310 Ok(upstream) => upstream,
311 Err(_) => {
312 client.shutdown_both();
313 return;
314 }
315 };
316 let (Ok(client_writer), Ok(upstream_reader)) = (client.try_clone_box(), upstream.try_clone_box()) else { return };
317 let (sender, pending) = mpsc::channel::<Pending>();
318 let responses = {
319 let state = state.clone();
320 std::thread::spawn(move || answer(upstream_reader, client_writer, pending, state))
321 };
322 let _ = forward(client, upstream, sender, &host, &state);
323 let _ = responses.join();
324}
325
326/// Requests, client to Engine.
327fn forward(client: Box<dyn Duplex>, mut upstream: Box<dyn Duplex>, pending: mpsc::Sender<Pending>, host: &Arc<dyn Host>, state: &Arc<Mutex<State>>) -> io::Result<()> {
328 let closer = client.try_clone_box()?;
329 let upstream_closer = upstream.try_clone_box()?;
330 let mut reader = BufReader::new(client);
331 let result = (|| -> io::Result<()> {
332 loop {
333 let Some(mut head) = http::read_head(&mut reader)? else { return Ok(()) };
334 let body = http::request_body(&head);
335 let method = head.method().to_owned();
336 match classify(&method, head.target()) {
337 Kind::CreateContainer { name } => {
338 let raw = http::read_body(&mut reader, body)?;
339 let Ok(mut config) = serde_json::from_slice::<Value>(&raw) else {
340 upstream.write_all(&http::with_length(head, &raw))?;
341 let _ = pending.send(Pending::Plain { method });
342 continue;
343 };
344 if let Some(miner) = miner_in_config(&config) {
345 let refusal = json!({ "message": format!("g1t does not run cryptocurrency miners ({miner}). This container was not created.") });
346 let _ = pending.send(Pending::Answer(http::json_response(403, "Forbidden", &refusal)));
347 continue;
348 }
349 let known = state.lock().map(|s| s.aliases.clone()).unwrap_or_default();
350 let created = rewrite_create(&mut config, name.as_deref(), &known);
351 if created.moved_from.is_some() {
352 let fresh: Vec<String> = created.aliases.iter().filter(|a| !known.contains(*a)).cloned().collect();
353 if let Ok(mut state) = state.lock() {
354 state.aliases.extend(created.aliases.iter().cloned());
355 // By name now; by id once the Engine says it.
356 if let Some(name) = &name {
357 state.published.insert(name.trim_start_matches('/').to_owned(), created.published.clone());
358 }
359 }
360 if !fresh.is_empty() {
361 host.add_hosts(&fresh);
362 }
363 for (host_port, container_port) in &created.forwards {
364 host.forward(*host_port, *container_port);
365 }
366 }
367 let text = serde_json::to_vec(&config).unwrap_or(raw);
368 upstream.write_all(&http::with_length(head, &text))?;
369 let _ = pending.send(if created.moved_from.is_some() {
370 Pending::Created { name, published: created.published }
371 } else {
372 Pending::Plain { method }
373 });
374 }
375 Kind::InspectContainer { id } => {
376 let published = state.lock().ok().and_then(|s| published_for(&s, &id));
377 upstream.write_all(&head.to_bytes())?;
378 http::copy_body(&mut reader, &mut upstream, body)?;
379 let _ = pending.send(match published {
380 Some(published) => Pending::Inspect { published },
381 None => Pending::Plain { method },
382 });
383 }
384 kind @ (Kind::ConnectNetwork | Kind::DisconnectNetwork) => {
385 let raw = http::read_body(&mut reader, body)?;
386 let value: Value = serde_json::from_slice(&raw).unwrap_or(Value::Null);
387 let container = value.get("Container").and_then(Value::as_str).unwrap_or_default().to_owned();
388 let moved = state.lock().ok().is_some_and(|s| published_for(&s, &container).is_some());
389 if moved {
390 // The container is on the job's network already.
391 if kind == Kind::ConnectNetwork {
392 let aliases = connect_aliases(&value);
393 let fresh: Vec<String> = match state.lock() {
394 Ok(mut state) => aliases.into_iter().filter(|a| state.aliases.insert(a.clone())).collect(),
395 Err(_) => Vec::new(),
396 };
397 if !fresh.is_empty() {
398 host.add_hosts(&fresh);
399 }
400 }
401 let _ = pending.send(Pending::Answer(http::empty_ok()));
402 } else {
403 upstream.write_all(&http::with_length(head, &raw))?;
404 let _ = pending.send(Pending::Plain { method });
405 }
406 }
407 Kind::ClassicBuild => {
408 let target = rewrite_build(head.target());
409 head.set_target(&target);
410 upstream.write_all(&head.to_bytes())?;
411 http::copy_body(&mut reader, &mut upstream, body)?;
412 let _ = pending.send(Pending::Plain { method });
413 }
414 Kind::Other => {
415 let upgrade = head.upgrades();
416 upstream.write_all(&head.to_bytes())?;
417 if upgrade {
418 let _ = pending.send(Pending::Upgrade);
419 // From here the two ends speak to each other.
420 io::copy(&mut reader, &mut upstream)?;
421 upstream_closer.shutdown_write();
422 return Ok(());
423 }
424 http::copy_body(&mut reader, &mut upstream, body)?;
425 let _ = pending.send(Pending::Plain { method });
426 }
427 }
428 }
429 })();
430 // The client is done asking; the Engine's answers may still be coming.
431 drop(pending);
432 if result.is_err() {
433 closer.shutdown_both();
434 upstream_closer.shutdown_both();
435 }
436 result
437}
438
439/// Responses, Engine to client, in the order they were asked for.
440fn answer(upstream: Box<dyn Duplex>, mut client: Box<dyn Duplex>, pending: mpsc::Receiver<Pending>, state: Arc<Mutex<State>>) {
441 let closer = upstream.try_clone_box().ok();
442 let mut reader = BufReader::new(upstream);
443 let result = (|| -> io::Result<()> {
444 while let Ok(next) = pending.recv() {
445 if let Pending::Answer(bytes) = next {
446 client.write_all(&bytes)?;
447 client.flush()?;
448 continue;
449 }
450 let Some(head) = http::read_head(&mut reader)? else { return Ok(()) };
451 match next {
452 Pending::Upgrade => {
453 client.write_all(&head.to_bytes())?;
454 client.flush()?;
455 io::copy(&mut reader, &mut client)?;
456 client.shutdown_write();
457 return Err(io::Error::other("upgraded"));
458 }
459 Pending::Plain { method } => {
460 let body = http::response_body(&head, &method);
461 client.write_all(&head.to_bytes())?;
462 http::copy_body(&mut reader, &mut client, body)?;
463 }
464 Pending::Created { name, published } => {
465 let body = http::read_body(&mut reader, http::response_body(&head, "POST"))?;
466 if (200..300).contains(&head.status())
467 && let Some(id) = serde_json::from_slice::<Value>(&body).ok().and_then(|v| v.get("Id").and_then(Value::as_str).map(str::to_owned))
468 && let Ok(mut state) = state.lock()
469 {
470 state.published.insert(id, published.clone());
471 if let Some(name) = name {
472 state.published.insert(name.trim_start_matches('/').to_owned(), published);
473 }
474 }
475 client.write_all(&http::with_length(head, &body))?;
476 }
477 Pending::Inspect { published } => {
478 let body = http::read_body(&mut reader, http::response_body(&head, "GET"))?;
479 let rewritten = match serde_json::from_slice::<Value>(&body) {
480 Ok(mut value) if head.status() == 200 => {
481 rewrite_inspect(&mut value, &published);
482 serde_json::to_vec(&value).unwrap_or(body)
483 }
484 _ => body,
485 };
486 client.write_all(&http::with_length(head, &rewritten))?;
487 }
488 Pending::Answer(_) => unreachable!("answered above"),
489 }
490 client.flush()?;
491 }
492 Ok(())
493 })();
494 if matches!(&result, Err(error) if error.to_string() == "upgraded") {
495 // A hijacked connection ends when both ends have said so.
496 return;
497 }
498 client.shutdown_both();
499 if let Some(closer) = closer {
500 closer.shutdown_both();
501 }
502}
503
504/// The ports of a container moved to the job's network, by its id, the
505/// start of its id, or its name.
506fn published_for(state: &State, id: &str) -> Option<BTreeMap<String, String>> {
507 let id = id.trim_start_matches('/');
508 if id.is_empty() {
509 return None;
510 }
511 if let Some(found) = state.published.get(id) {
512 return Some(found.clone());
513 }
514 if id.len() >= 6 && id.chars().all(|c| c.is_ascii_hexdigit()) {
515 let mut matches = state.published.iter().filter(|(key, _)| key.len() == 64 && key.starts_with(id));
516 if let (Some((_, found)), None) = (matches.next(), matches.next()) {
517 return Some(found.clone());
518 }
519 }
520 None
521}
522
523/// The miner a container's image or command names, if any.
524fn miner_in_config(config: &Value) -> Option<&'static str> {
525 let mut words = vec![config.get("Image").and_then(Value::as_str).unwrap_or_default().to_owned()];
526 for key in ["Entrypoint", "Cmd"] {
527 match config.get(key) {
528 Some(Value::String(text)) => words.push(text.clone()),
529 other => words.extend(strings(other)),
530 }
531 }
532 crate::abuse::miner_in(&words.join(" "))
533}
534
535#[cfg(test)]
536mod tests {
537 use super::*;
538
539 #[test]
540 fn routes_are_read_with_or_without_a_version() {
541 assert_eq!(classify("POST", "/v1.47/containers/create?name=db"), Kind::CreateContainer { name: Some("db".into()) });
542 assert_eq!(classify("POST", "/containers/create"), Kind::CreateContainer { name: None });
543 assert_eq!(classify("GET", "/v1.51/containers/abc123/json?size=false"), Kind::InspectContainer { id: "abc123".into() });
544 assert_eq!(classify("POST", "/v1.47/networks/app_default/connect"), Kind::ConnectNetwork);
545 assert_eq!(classify("POST", "/v1.47/build?t=x"), Kind::ClassicBuild);
546 assert_eq!(classify("GET", "/v1.47/containers/json"), Kind::Other);
547 assert_eq!(classify("POST", "/v1.47/containers/abc/start"), Kind::Other);
548 }
549
550 #[test]
551 fn a_compose_service_moves_to_the_jobs_network_and_keeps_its_names() {
552 let mut body = json!({
553 "Image": "postgres:17",
554 "Labels": { "com.docker.compose.service": "db", "com.docker.compose.project": "app" },
555 "HostConfig": {
556 "NetworkMode": "app_default",
557 "PortBindings": { "5432/tcp": [{ "HostIp": "", "HostPort": "15432" }], "8080/tcp": [{ "HostPort": "" }] },
558 "Links": ["/cache:/app-db-1/redis"],
559 "ExtraHosts": ["host.docker.internal:host-gateway"],
560 },
561 "NetworkingConfig": { "EndpointsConfig": { "app_default": { "Aliases": ["db", "app-db-1"], "MacAddress": "x" } } },
562 });
563 let known: BTreeSet<String> = ["cache".to_owned()].into();
564 let created = rewrite_create(&mut body, Some("app-db-1"), &known);
565 assert_eq!(created.moved_from.as_deref(), Some("app_default"));
566 assert_eq!(created.aliases, vec!["app-db-1", "db", "redis"]);
567 assert_eq!(created.forwards, vec![(15432, 5432)]);
568 assert_eq!(created.published["5432/tcp"], "15432");
569 assert_eq!(created.published["8080/tcp"], "8080");
570 assert_eq!(body["HostConfig"]["NetworkMode"], "host");
571 assert!(body.get("NetworkingConfig").is_none());
572 assert!(body["HostConfig"].get("Links").is_none());
573 let extra = strings(body["HostConfig"].get("ExtraHosts"));
574 assert!(extra.contains(&"host.docker.internal:host-gateway".to_owned()));
575 for name in ["cache", "db", "app-db-1", "redis"] {
576 assert!(extra.contains(&format!("{name}:127.0.0.1")), "{name} in {extra:?}");
577 }
578 }
579
580 #[test]
581 fn host_none_and_shared_networks_are_left_alone() {
582 for mode in ["host", "none", "container:abc"] {
583 let mut body = json!({ "Image": "alpine", "HostConfig": { "NetworkMode": mode } });
584 let before = body.clone();
585 assert_eq!(rewrite_create(&mut body, Some("x"), &BTreeSet::new()), Created::default());
586 assert_eq!(body, before);
587 }
588 // The default network, named or not.
589 let mut body = json!({ "Image": "alpine" });
590 assert_eq!(rewrite_create(&mut body, None, &BTreeSet::new()).moved_from.as_deref(), Some("bridge"));
591 assert_eq!(body["HostConfig"]["NetworkMode"], "host");
592 }
593
594 #[test]
595 fn names_that_cannot_be_hosts_are_skipped() {
596 let mut body = json!({ "HostConfig": {}, "NetworkingConfig": { "EndpointsConfig": { "n": { "Aliases": ["ok-name", "bad name", "localhost", ""] } } } });
597 assert_eq!(rewrite_create(&mut body, None, &BTreeSet::new()).aliases, vec!["ok-name"]);
598 }
599
600 #[test]
601 fn inspections_report_the_ports_published() {
602 let mut body = json!({ "Id": "abc", "NetworkSettings": { "Ports": {}, "Networks": { "host": {} } } });
603 rewrite_inspect(&mut body, &[("5432/tcp".to_owned(), "15432".to_owned())].into());
604 assert_eq!(body["NetworkSettings"]["Ports"]["5432/tcp"][0]["HostPort"], "15432");
605 }
606
607 #[test]
608 fn the_classic_builder_runs_on_the_jobs_network() {
609 assert_eq!(rewrite_build("/v1.47/build?t=app&networkmode=default"), "/v1.47/build?t=app&networkmode=host");
610 assert_eq!(rewrite_build("/build"), "/build?networkmode=host");
611 assert_eq!(rewrite_build("/build?networkmode=none"), "/build?networkmode=none");
612 }
613
614 #[test]
615 fn ids_are_found_by_their_start_or_a_name() {
616 let mut state = State::default();
617 let id = "a".repeat(64);
618 state.published.insert(id.clone(), [("80/tcp".to_owned(), "8080".to_owned())].into());
619 state.published.insert("web".into(), [("80/tcp".to_owned(), "8080".to_owned())].into());
620 assert!(published_for(&state, "aaaaaaaaaaaa").is_some());
621 assert!(published_for(&state, "/web").is_some());
622 assert!(published_for(&state, "aaa").is_none());
623 assert!(published_for(&state, "other").is_none());
624 }
625
626 #[test]
627 fn miners_are_refused_by_image_or_command() {
628 assert!(miner_in_config(&json!({ "Image": "metal3d/xmrig" })).is_some());
629 assert!(miner_in_config(&json!({ "Image": "alpine", "Cmd": ["sh", "-c", "./xmrig -o stratum+tcp://pool"] })).is_some());
630 assert!(miner_in_config(&json!({ "Image": "postgres:17", "Cmd": ["postgres"] })).is_none());
631 }
632
633 /// A pretend Engine on TCP, and the proxy in front of it, end to end.
634 #[test]
635 fn requests_and_responses_pass_through_in_order() {
636 use std::net::{TcpListener, TcpStream};
637 let engine = TcpListener::bind("127.0.0.1:0").unwrap();
638 let engine_addr = engine.local_addr().unwrap();
639 // The Engine: a create, an inspect, a chunked stream, then a hijack.
640 let fake = std::thread::spawn(move || {
641 let (stream, _) = engine.accept().unwrap();
642 let mut reader = BufReader::new(stream.try_clone().unwrap());
643 let mut writer = stream;
644 let mut seen = Vec::new();
645 while let Some(head) = http::read_head(&mut reader).unwrap() {
646 let body = http::read_body(&mut reader, http::request_body(&head)).unwrap();
647 seen.push((head.start.clone(), String::from_utf8_lossy(&body).into_owned()));
648 if head.target().contains("/containers/create") {
649 writer.write_all(&http::json_response(201, "Created", &json!({ "Id": "c".repeat(64), "Warnings": [] }))).unwrap();
650 } else if head.target().contains("/json") {
651 writer
652 .write_all(b"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nTransfer-Encoding: chunked\r\n\r\n")
653 .unwrap();
654 let text = json!({ "Id": "c".repeat(64), "NetworkSettings": { "Ports": {} } }).to_string();
655 writer.write_all(format!("{:x}\r\n{text}\r\n0\r\n\r\n", text.len()).as_bytes()).unwrap();
656 } else if head.target().contains("/logs") {
657 writer.write_all(b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n3\r\none\r\n3\r\ntwo\r\n0\r\n\r\n").unwrap();
658 } else if head.upgrades() {
659 writer.write_all(b"HTTP/1.1 101 UPGRADED\r\nConnection: Upgrade\r\nUpgrade: tcp\r\n\r\n").unwrap();
660 let mut buffer = [0u8; 5];
661 reader.read_exact(&mut buffer).unwrap();
662 writer.write_all(&buffer).unwrap();
663 break;
664 }
665 }
666 seen
667 });
668
669 struct Fake(std::net::SocketAddr, Mutex<Vec<(u16, u16)>>, Mutex<Vec<String>>);
670 impl Host for Fake {
671 fn add_hosts(&self, names: &[String]) {
672 self.2.lock().unwrap().extend(names.iter().cloned());
673 }
674 fn forward(&self, host_port: u16, container_port: u16) {
675 self.1.lock().unwrap().push((host_port, container_port));
676 }
677 fn ensure_engine(&self) -> Result<(), String> {
678 Ok(())
679 }
680 fn connect(&self) -> io::Result<Box<dyn Duplex>> {
681 Ok(Box::new(TcpStream::connect(self.0)?))
682 }
683 }
684 let host = Arc::new(Fake(engine_addr, Mutex::new(Vec::new()), Mutex::new(Vec::new())));
685 let state = Arc::new(Mutex::new(State::default()));
686 let proxy = TcpListener::bind("127.0.0.1:0").unwrap();
687 let proxy_addr = proxy.local_addr().unwrap();
688 let (host2, state2) = (host.clone() as Arc<dyn Host>, state.clone());
689 std::thread::spawn(move || {
690 let (stream, _) = proxy.accept().unwrap();
691 serve(Box::new(stream), host2, state2);
692 });
693
694 let client = TcpStream::connect(proxy_addr).unwrap();
695 let mut reader = BufReader::new(client.try_clone().unwrap());
696 let mut writer = client;
697 let create = json!({ "Image": "redis", "HostConfig": { "PortBindings": { "6379/tcp": [{ "HostPort": "16379" }] } }, "NetworkingConfig": { "EndpointsConfig": { "job": { "Aliases": ["redis"] } } } }).to_string();
698 // Pipelined: all asked before any answer is read.
699 writer
700 .write_all(format!("POST /v1.47/containers/create?name=cache HTTP/1.1\r\nHost: docker\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{create}", create.len()).as_bytes())
701 .unwrap();
702 writer.write_all(b"GET /v1.47/containers/cache/json HTTP/1.1\r\nHost: docker\r\n\r\n").unwrap();
703 writer.write_all(b"GET /v1.47/containers/cache/logs?follow=1 HTTP/1.1\r\nHost: docker\r\n\r\n").unwrap();
704 writer.write_all(b"POST /v1.47/containers/cache/attach?stream=1 HTTP/1.1\r\nHost: docker\r\nConnection: Upgrade\r\nUpgrade: tcp\r\n\r\nhello").unwrap();
705
706 let created = http::read_head(&mut reader).unwrap().unwrap();
707 assert_eq!(created.status(), 201);
708 let body: Value = serde_json::from_slice(&http::read_body(&mut reader, http::response_body(&created, "POST")).unwrap()).unwrap();
709 assert_eq!(body["Id"].as_str().unwrap().len(), 64);
710 let inspected = http::read_head(&mut reader).unwrap().unwrap();
711 let body: Value = serde_json::from_slice(&http::read_body(&mut reader, http::response_body(&inspected, "GET")).unwrap()).unwrap();
712 assert_eq!(body["NetworkSettings"]["Ports"]["6379/tcp"][0]["HostPort"], "16379");
713 let logs = http::read_head(&mut reader).unwrap().unwrap();
714 assert_eq!(http::read_body(&mut reader, http::response_body(&logs, "GET")).unwrap(), b"onetwo");
715 let upgraded = http::read_head(&mut reader).unwrap().unwrap();
716 assert_eq!(upgraded.status(), 101);
717 let mut echoed = [0u8; 5];
718 reader.read_exact(&mut echoed).unwrap();
719 assert_eq!(&echoed, b"hello");
720
721 let seen = fake.join().unwrap();
722 assert!(seen[0].1.contains("\"NetworkMode\":\"host\""), "{}", seen[0].1);
723 assert!(seen[0].1.contains("redis:127.0.0.1"));
724 assert_eq!(*host.1.lock().unwrap(), vec![(16379, 6379)]);
725 assert_eq!(*host.2.lock().unwrap(), vec!["cache".to_owned(), "redis".to_owned()]);
726 }
727}