| 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 | |
| 22 | use std::collections::{BTreeMap, BTreeSet}; |
| 23 | use std::io::{self, BufReader, Read, Write}; |
| 24 | use std::sync::mpsc; |
| 25 | use std::sync::{Arc, Mutex}; |
| 26 | |
| 27 | use serde_json::{Map, Value, json}; |
| 28 | |
| 29 | use super::http; |
| 30 | |
| 31 | /// What the proxy remembers across connections. |
| 32 | #[derive(Default)] |
| 33 | pub(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. |
| 42 | pub(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. |
| 54 | pub(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 | |
| 61 | impl 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)] |
| 74 | impl 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`. |
| 87 | fn 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 | |
| 96 | fn 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 | |
| 100 | fn 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)] |
| 122 | pub(crate) enum Kind { |
| 123 | CreateContainer { name: Option<String> }, |
| 124 | InspectContainer { id: String }, |
| 125 | ConnectNetwork, |
| 126 | DisconnectNetwork, |
| 127 | ClassicBuild, |
| 128 | Other, |
| 129 | } |
| 130 | |
| 131 | pub(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. |
| 145 | fn 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`. |
| 150 | fn 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)] |
| 160 | pub(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 | |
| 171 | fn 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. |
| 177 | pub(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. |
| 252 | pub(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. |
| 260 | pub(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. |
| 271 | pub(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. |
| 283 | enum 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. |
| 297 | pub(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. |
| 327 | fn 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. |
| 440 | fn 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. |
| 506 | fn 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. |
| 524 | fn 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)] |
| 536 | mod 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 | } |