Skip to content
File

Blob: src/workerd/server/container-client.c++

107.8 KB
1// Copyright (c) 2025 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "container-client.h"
6 
7#include "ada.h"
8 
9#include <workerd/io/container.capnp.h>
10#include <workerd/io/worker-interface.h>
11#include <workerd/jsg/jsg.h>
12#include <workerd/jsg/url.h>
13#include <workerd/server/docker-api.capnp.h>
14#include <workerd/util/stream-utils.h>
15#include <workerd/util/strings.h>
16#include <workerd/util/uuid.h>
17 
18#include <stdio.h>
19 
20#include <capnp/compat/json.h>
21#include <capnp/message.h>
22#include <kj/async-io.h>
23#include <kj/async.h>
24#include <kj/cidr.h>
25#include <kj/compat/http.h>
26#include <kj/debug.h>
27#include <kj/encoding.h>
28#include <kj/exception.h>
29#include <kj/string.h>
30 
31#include <limits>
32 
33namespace workerd::server {
34 
35namespace {
36 
37constexpr uint16_t SIDECAR_INGRESS_PORT = 39001;
38 
39constexpr kj::StringPtr SIDECAR_DNS_SERVERS[] = {
40 "1.1.1.1"_kj,
41 "8.8.8.8"_kj,
42};
43 
44// Default limit for JSON API responses (16 MiB — Docker JSON responses are small).
45constexpr uint64_t MAX_JSON_RESPONSE_SIZE = 16ULL * 1024 * 1024;
46 
47constexpr kj::StringPtr SNAPSHOT_VOLUME_PREFIX = "workerd-snap-"_kj;
48constexpr kj::StringPtr SNAPSHOT_CLONE_VOLUME_PREFIX = "workerd-snap-clone-"_kj;
49constexpr kj::StringPtr CONTAINER_SNAPSHOT_IMAGE_PREFIX = "workerd-container-snap-"_kj;
50constexpr kj::StringPtr SNAPSHOT_VOLUME_CREATED_AT_LABEL = "dev.workerd.snapshot-created-at"_kj;
51 
52// Prefix applied to user-supplied labels when writing them to the Docker container, and
53// stripped back out when reading them via inspect(). Lets us distinguish labels the worker
54// set via start() from labels that came from the image (via Dockerfile LABEL) or engine.
55constexpr kj::StringPtr WORKERD_LABEL_PREFIX = "workerd-"_kj;
56constexpr auto SNAPSHOT_STALE_AGE = 30 * kj::DAYS;
57 
58// Maximum size of a snapshot tar archive held in memory during snapshot create/restore.
59constexpr size_t MAX_SNAPSHOT_TAR_SIZE = 1ULL * 1024 * 1024 * 1024; // 1 GiB
60static_assert(static_cast<double>(MAX_SNAPSHOT_TAR_SIZE) == MAX_SNAPSHOT_TAR_SIZE,
61 "MAX_SNAPSHOT_TAR_SIZE must be exactly representable as double");
62 
63// POSIX tar stores file size in an 11-digit octal header field.
64constexpr size_t MAX_TAR_CONTENT_SIZE = 8ull * 1024 * 1024 * 1024;
65 
66// Ensures the stale-volume check runs at most once per process.
67std::atomic_bool staleSnapshotVolumeCheckScheduled = false;
68 
69struct ParsedAddress {
70 kj::OneOf<kj::CidrRange, kj::String> destination;
71 kj::Maybe<uint16_t> port;
72};
73 
74struct HostAndPort {
75 kj::String host;
76 kj::Maybe<uint16_t> port;
77};
78 
79struct DockerResponse {
80 kj::uint statusCode;
81 kj::String body;
82};
83 
84struct DockerBinaryResponse {
85 kj::uint statusCode;
86 kj::Array<kj::byte> body;
87};
88 
89struct DockerStreamedResponse {
90 kj::uint statusCode;
91 kj::String statusText;
92 kj::Own<kj::AsyncIoStream> connection;
93};
94 
95// Validates an absolute path for snapshot use and returns the parsed component path.
96// Rejects relative paths, embedded null bytes, and path traversal components ("..").
97kj::Path parseAbsolutePath(kj::StringPtr path) {
98 JSG_REQUIRE(
99 path.size() > 0 && path[0] == '/', Error, "Snapshot path must be absolute, got: ", path);
100 
101 JSG_REQUIRE(path.findFirst('\0') == kj::none, Error, "Snapshot path must not contain null bytes");
102 
103 try {
104 return kj::Path::parse(path.slice(1));
105 } catch (kj::Exception& e) {
106 JSG_FAIL_REQUIRE(
107 Error, "Snapshot path contains invalid components: ", path, "; ", e.getDescription());
108 }
109}
110 
111// Parse and validate a snapshot ID. Throws an error if the snapshot ID is invalid.
112kj::String parseSnapshotId(kj::StringPtr snapshotId) {
113 KJ_IF_SOME(uuid, UUID::fromString(snapshotId)) {
114 auto s = uuid.toString();
115 JSG_REQUIRE(s == snapshotId, Error, "Invalid snapshot ID", snapshotId);
116 return s;
117 } else {
118 JSG_FAIL_REQUIRE(Error, "Invalid snapshot ID", snapshotId);
119 }
120}
121 
122// Really similar to BufferedInputStreamWrapper, but Async...
123// We need this because of Docker's exec keeping a bidirectional connection
124// needing to own the IoStream after writing and reading headers, as it does
125// "Upgrade: tcp".
126class BufferedAsyncIoStream final: public kj::AsyncIoStream {
127 public:
128 BufferedAsyncIoStream(kj::Own<kj::AsyncIoStream> inner, kj::Array<kj::byte> buffered)
129 : inner(kj::mv(inner)),
130 buffered(kj::mv(buffered)) {}
131 
132 kj::Promise<size_t> tryRead(void* dst, size_t minBytes, size_t maxBytes) override {
133 KJ_REQUIRE(minBytes <= maxBytes, minBytes, maxBytes);
134 
135 auto out = kj::arrayPtr(reinterpret_cast<kj::byte*>(dst), maxBytes);
136 size_t copied = 0;
137 
138 auto bufferedRemaining = buffered.size() - bufferedOffset;
139 if (bufferedRemaining > 0) {
140 auto toCopy = kj::min(maxBytes, bufferedRemaining);
141 out.first(toCopy).copyFrom(buffered.asPtr().slice(bufferedOffset, bufferedOffset + toCopy));
142 bufferedOffset += toCopy;
143 copied = toCopy;
144 
145 if (copied >= minBytes || copied == maxBytes) {
146 co_return copied;
147 }
148 }
149 
150 auto read = co_await inner->tryRead(out.begin() + copied, minBytes - copied, maxBytes - copied);
151 co_return copied + read;
152 }
153 
154 kj::Maybe<uint64_t> tryGetLength() override {
155 KJ_IF_SOME(innerLength, inner->tryGetLength()) {
156 return innerLength + (buffered.size() - bufferedOffset);
157 }
158 return kj::none;
159 }
160 
161 kj::Promise<uint64_t> pumpTo(kj::AsyncOutputStream& output, uint64_t amount) override {
162 uint64_t pumped = 0;
163 auto bufferedRemaining = buffered.size() - bufferedOffset;
164 if (bufferedRemaining > 0) {
165 auto toWrite = static_cast<size_t>(kj::min(amount, static_cast<uint64_t>(bufferedRemaining)));
166 co_await output.write(buffered.asPtr().slice(bufferedOffset, bufferedOffset + toWrite));
167 bufferedOffset += toWrite;
168 pumped += toWrite;
169 
170 if (pumped == amount) {
171 co_return pumped;
172 }
173 }
174 
175 co_return pumped + co_await inner->pumpTo(output, amount - pumped);
176 }
177 
178 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
179 return inner->write(buffer);
180 }
181 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
182 return inner->write(pieces);
183 }
184 kj::Maybe<kj::Promise<uint64_t>> tryPumpFrom(
185 kj::AsyncInputStream& input, uint64_t amount = kj::maxValue) override {
186 return inner->tryPumpFrom(input, amount);
187 }
188 kj::Promise<void> whenWriteDisconnected() override {
189 return inner->whenWriteDisconnected();
190 }
191 void abortWrite(kj::Exception&& exception) override {
192 inner->abortWrite(kj::mv(exception));
193 }
194 
195 void shutdownWrite() override {
196 inner->shutdownWrite();
197 }
198 void abortRead() override {
199 inner->abortRead();
200 }
201 void getsockopt(int level, int option, void* value, kj::uint* length) override {
202 inner->getsockopt(level, option, value, length);
203 }
204 void setsockopt(int level, int option, const void* value, kj::uint length) override {
205 inner->setsockopt(level, option, value, length);
206 }
207 void getsockname(struct sockaddr* addr, kj::uint* length) override {
208 inner->getsockname(addr, length);
209 }
210 void getpeername(struct sockaddr* addr, kj::uint* length) override {
211 inner->getpeername(addr, length);
212 }
213 kj::Maybe<int> getFd() const override {
214 return inner->getFd();
215 }
216 
217 private:
218 kj::Own<kj::AsyncIoStream> inner;
219 kj::Array<kj::byte> buffered;
220 size_t bufferedOffset = 0;
221};
222 
223// Docker exec uses a single hijacked stream for stdin and stdout/stderr. Keep that stream in a
224// small refcounted holder so the returned stdin ByteStream and the output demux task can share it.
225class SharedExecConnection final: public kj::Refcounted {
226 public:
227 explicit SharedExecConnection(kj::Own<kj::AsyncIoStream> connection)
228 : connection(kj::mv(connection)) {}
229 
230 kj::Own<kj::AsyncIoStream> connection;
231 bool stdinOpened = false;
232 bool stdinClosed = false;
233};
234 
235class DockerExecStdinStream final: public capnp::ExplicitEndOutputStream {
236 public:
237 explicit DockerExecStdinStream(kj::Own<SharedExecConnection> sharedConnection)
238 : sharedConnection(kj::mv(sharedConnection)) {}
239 
240 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
241 return sharedConnection->connection->write(buffer);
242 }
243 
244 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
245 return sharedConnection->connection->write(pieces);
246 }
247 
248 kj::Promise<void> whenWriteDisconnected() override {
249 return sharedConnection->connection->whenWriteDisconnected();
250 }
251 
252 kj::Promise<void> end() override {
253 if (!sharedConnection->stdinClosed) {
254 sharedConnection->connection->shutdownWrite();
255 sharedConnection->stdinClosed = true;
256 }
257 return kj::READY_NOW;
258 }
259 
260 private:
261 kj::Own<SharedExecConnection> sharedConnection;
262};
263 
264// Strips a port suffix from a string, returning the host and port separately.
265// For IPv6, expects brackets: "[::1]:8080" -> ("::1", 8080)
266// For IPv4: "10.0.0.1:8080" -> ("10.0.0.1", 8080)
267// If no port, returns the host as-is with no port.
268HostAndPort stripPort(kj::StringPtr str) {
269 if (str.startsWith("[")) {
270 // Bracketed IPv6: "[ipv6]" or "[ipv6]:port"
271 size_t closeBracket =
272 KJ_REQUIRE_NONNULL(str.findLast(']'), "Unclosed '[' in address string.", str);
273 
274 auto host = str.slice(1, closeBracket);
275 
276 if (str.size() > closeBracket + 1) {
277 KJ_REQUIRE(
278 str.slice(closeBracket + 1).startsWith(":"), "Expected port suffix after ']'.", str);
279 auto port = KJ_REQUIRE_NONNULL(
280 str.slice(closeBracket + 2).tryParseAs<uint16_t>(), "Invalid port number.", str);
281 return {kj::str(host), port};
282 }
283 return {kj::str(host), kj::none};
284 }
285 
286 // No brackets - check if there's exactly one colon (IPv4 with port)
287 // IPv6 without brackets has 2+ colons and no port suffix supported
288 KJ_IF_SOME(colonPos, str.findLast(':')) {
289 auto afterColon = str.slice(colonPos + 1);
290 KJ_IF_SOME(port, afterColon.tryParseAs<uint16_t>()) {
291 // Valid port - but only treat as port for IPv4 (check no other colons before)
292 auto beforeColon = str.first(colonPos);
293 if (beforeColon.findFirst(':') == kj::none) {
294 return {kj::str(beforeColon), port};
295 }
296 }
297 }
298 
299 return {kj::str(str), kj::none};
300}
301 
302// Build a CidrRange from a host string, adding /32 or /128 prefix if not present.
303kj::CidrRange makeCidr(kj::StringPtr host) {
304 if (host.findFirst('/') != kj::none) {
305 return kj::CidrRange(host);
306 }
307 // No CIDR prefix - add /32 for IPv4, /128 for IPv6
308 bool isIpv6 = host.findFirst(':') != kj::none;
309 return kj::CidrRange(kj::str(host, isIpv6 ? "/128" : "/32"));
310}
311 
312kj::Maybe<kj::CidrRange> tryMakeCidr(kj::StringPtr host) {
313 kj::Maybe<kj::CidrRange> cidr;
314 KJ_IF_SOME(_, kj::runCatchingExceptions([&]() { cidr = makeCidr(host); })) {
315 return kj::none;
316 }
317 
318 return kj::mv(cidr);
319}
320 
321// normalizeHostname normalizes the hostname. It's designed to receive the hostname when
322// proxy-everything sends the HTTP CONNECT with the X-Hostname hint.
323kj::String normalizeHostname(kj::StringPtr hostname) {
324 auto url = kj::str("http://", hostname);
325 auto parsed = ada::parse<ada::url_aggregator>({url.begin(), url.size()}, nullptr);
326 KJ_REQUIRE(parsed.has_value(), "Invalid X-Hostname URL hint.", hostname);
327 auto normalizedHostname = parsed->get_hostname();
328 return kj::heapString(normalizedHostname.data(), normalizedHostname.size());
329}
330 
331// hostnameGlobMatches should match patterns like:
332// cloudflare.*.com
333// cloudflare.com
334// cloudflare
335// *
336//
337// hostname must be normalized beforehand
338bool hostnameGlobMatches(kj::StringPtr pattern, kj::StringPtr hostname) {
339 size_t patternIndex = 0;
340 size_t hostnameIndex = 0;
341 size_t restartHostnameIndex = 0;
342 kj::Maybe<size_t> starPatternIndex;
343 
344 while (hostnameIndex < hostname.size()) {
345 if (patternIndex < pattern.size() && pattern[patternIndex] == '*') {
346 starPatternIndex = patternIndex++;
347 restartHostnameIndex = hostnameIndex;
348 continue;
349 }
350 
351 if (patternIndex < pattern.size() && pattern[patternIndex] == hostname[hostnameIndex]) {
352 ++patternIndex;
353 ++hostnameIndex;
354 continue;
355 }
356 
357 KJ_IF_SOME(starIndex, starPatternIndex) {
358 patternIndex = starIndex + 1;
359 hostnameIndex = ++restartHostnameIndex;
360 continue;
361 }
362 
363 return false;
364 }
365 
366 while (patternIndex < pattern.size() && pattern[patternIndex] == '*') {
367 ++patternIndex;
368 }
369 
370 return patternIndex == pattern.size();
371}
372 
373kj::Maybe<kj::StringPtr> getHeader(const kj::HttpHeaders& headers, kj::StringPtr name) {
374 kj::Maybe<kj::StringPtr> result;
375 headers.forEach([&](kj::StringPtr headerName, kj::StringPtr value) {
376 if (result == kj::none && workerd::strcaseeq(headerName, name)) {
377 result = value;
378 }
379 });
380 return result;
381}
382 
383// Parses "host[:port]" strings. Handles:
384// - IPv4: "10.0.0.1", "10.0.0.1:8080", "10.0.0.0/8", "10.0.0.0/8:8080"
385// - IPv6 with brackets: "[::1]", "[::1]:8080", "[fe80::1]", "[fe80::/10]:8080"
386// - IPv6 without brackets: "::1", "fe80::1", "fe80::/10"
387ParsedAddress parseHostPort(kj::StringPtr str) {
388 auto hostAndPort = stripPort(str);
389 KJ_REQUIRE(hostAndPort.host.size() > 0, "Host must not be empty.", str);
390 
391 KJ_IF_SOME(cidr, tryMakeCidr(hostAndPort.host)) {
392 return {
393 .destination = kj::mv(cidr),
394 .port = hostAndPort.port,
395 };
396 }
397 
398 return {
399 .destination = workerd::toLower(hostAndPort.host),
400 .port = hostAndPort.port,
401 };
402}
403 
404kj::StringPtr signalToString(uint32_t signal) {
405 switch (signal) {
406 case 1:
407 return "SIGHUP"_kj; // Hangup
408 case 2:
409 return "SIGINT"_kj; // Interrupt
410 case 3:
411 return "SIGQUIT"_kj; // Quit
412 case 4:
413 return "SIGILL"_kj; // Illegal instruction
414 case 5:
415 return "SIGTRAP"_kj; // Trace trap
416 case 6:
417 return "SIGABRT"_kj; // Abort
418 case 7:
419 return "SIGBUS"_kj; // Bus error
420 case 8:
421 return "SIGFPE"_kj; // Floating point exception
422 case 9:
423 return "SIGKILL"_kj; // Kill
424 case 10:
425 return "SIGUSR1"_kj; // User signal 1
426 case 11:
427 return "SIGSEGV"_kj; // Segmentation violation
428 case 12:
429 return "SIGUSR2"_kj; // User signal 2
430 case 13:
431 return "SIGPIPE"_kj; // Broken pipe
432 case 14:
433 return "SIGALRM"_kj; // Alarm clock
434 case 15:
435 return "SIGTERM"_kj; // Termination
436 case 16:
437 return "SIGSTKFLT"_kj; // Stack fault (Linux)
438 case 17:
439 return "SIGCHLD"_kj; // Child status changed
440 case 18:
441 return "SIGCONT"_kj; // Continue
442 case 19:
443 return "SIGSTOP"_kj; // Stop
444 case 20:
445 return "SIGTSTP"_kj; // Terminal stop
446 case 21:
447 return "SIGTTIN"_kj; // Background read from tty
448 case 22:
449 return "SIGTTOU"_kj; // Background write to tty
450 case 23:
451 return "SIGURG"_kj; // Urgent condition on socket
452 case 24:
453 return "SIGXCPU"_kj; // CPU limit exceeded
454 case 25:
455 return "SIGXFSZ"_kj; // File size limit exceeded
456 case 26:
457 return "SIGVTALRM"_kj; // Virtual alarm clock
458 case 27:
459 return "SIGPROF"_kj; // Profiling alarm clock
460 case 28:
461 return "SIGWINCH"_kj; // Window size change
462 case 29:
463 return "SIGIO"_kj; // I/O now possible
464 case 30:
465 return "SIGPWR"_kj; // Power failure restart (Linux)
466 case 31:
467 return "SIGSYS"_kj; // Bad system call
468 default:
469 return "SIGKILL"_kj;
470 }
471}
472 
473void writeTarField(kj::ArrayPtr<kj::byte> field, kj::StringPtr value) {
474 auto len = kj::min(value.size(), field.size());
475 field.first(len).copyFrom(value.asBytes().first(len));
476}
477 
478// createTarWithFile creates simple tar files without importing a full blown TAR library.
479// It's a pretty limited method that creates a single tar file with a single file on it,
480// as the Docker API only accepts tars.
481kj::Array<kj::byte> createTarWithFile(
482 kj::StringPtr filename, kj::ArrayPtr<const kj::byte> content) {
483 KJ_REQUIRE(filename.size() < 100, "tar filename must be < 100 bytes");
484 KJ_REQUIRE(content.size() < MAX_TAR_CONTENT_SIZE, "tar content too large for 11-digit octal");
485 
486 size_t paddedSize = (content.size() + 511) & ~static_cast<size_t>(511);
487 size_t totalSize = 512 + paddedSize + 1024;
488 auto tar = kj::heapArray<kj::byte>(totalSize);
489 tar.asPtr().fill(0);
490 
491 auto header = tar.first(512);
492 writeTarField(header.slice(0, 100), filename);
493 writeTarField(header.slice(100, 108), "0000644"_kj);
494 writeTarField(header.slice(108, 116), "0000000"_kj);
495 writeTarField(header.slice(116, 124), "0000000"_kj);
496 
497 {
498 char sizeBuf[12];
499 snprintf(sizeBuf, sizeof(sizeBuf), "%011" PRIo64, static_cast<uint64_t>(content.size()));
500 writeTarField(header.slice(124, 136), kj::StringPtr(sizeBuf));
501 }
502 
503 writeTarField(header.slice(136, 148), "00000000000"_kj);
504 header[156] = '0';
505 writeTarField(header.slice(257, 263), "ustar"_kj);
506 writeTarField(header.slice(263, 265), "00"_kj);
507 
508 header.slice(148, 156).fill(' ');
509 uint32_t checksum = 0;
510 for (auto byte: header) checksum += byte;
511 
512 {
513 char checksumBuf[8];
514 snprintf(checksumBuf, sizeof(checksumBuf), "%06o ", checksum);
515 writeTarField(header.slice(148, 155), kj::StringPtr(checksumBuf));
516 }
517 
518 tar.slice(512, 512 + content.size()).copyFrom(content);
519 return tar;
520}
521 
522// Shared Docker API HTTP helper. Connects to the Docker socket, sends a request with an
523// optional body, and reads the response as raw bytes.
524kj::Promise<DockerBinaryResponse> dockerApiRequestRaw(kj::Network& network,
525 kj::String dockerPath,
526 kj::HttpMethod method,
527 kj::String endpoint,
528 kj::Maybe<kj::ArrayPtr<const kj::byte>> body,
529 kj::StringPtr contentType,
530 uint64_t maxResponseSize) {
531 kj::HttpHeaderTable headerTable;
532 auto address = co_await network.parseAddress(dockerPath);
533 auto connection = co_await address->connect();
534 auto httpClient = kj::newHttpClient(headerTable, *connection).attach(kj::mv(connection));
535 kj::HttpHeaders headers(headerTable);
536 headers.setPtr(kj::HttpHeaderId::HOST, "localhost");
537 
538 KJ_IF_SOME(requestBody, body) {
539 headers.setPtr(kj::HttpHeaderId::CONTENT_TYPE, contentType);
540 headers.set(kj::HttpHeaderId::CONTENT_LENGTH, kj::str(requestBody.size()));
541 
542 auto req = httpClient->request(method, endpoint, headers, requestBody.size());
543 {
544 auto stream = kj::mv(req.body);
545 co_await stream->write(requestBody);
546 }
547 auto response = co_await req.response;
548 auto result = co_await response.body->readAllBytes(maxResponseSize);
549 co_return DockerBinaryResponse{.statusCode = response.statusCode, .body = kj::mv(result)};
550 } else {
551 auto req = httpClient->request(method, endpoint, headers);
552 { auto stream = kj::mv(req.body); }
553 auto response = co_await req.response;
554 auto result = co_await response.body->readAllBytes(maxResponseSize);
555 co_return DockerBinaryResponse{.statusCode = response.statusCode, .body = kj::mv(result)};
556 }
557}
558 
559kj::Promise<DockerResponse> dockerApiRequest(kj::Network& network,
560 kj::String dockerPath,
561 kj::HttpMethod method,
562 kj::String endpoint,
563 kj::Maybe<kj::String> body = kj::none) {
564 kj::Maybe<kj::ArrayPtr<const kj::byte>> bodyBytes;
565 KJ_IF_SOME(b, body) {
566 bodyBytes = b.asBytes();
567 }
568 auto raw = co_await dockerApiRequestRaw(network, kj::mv(dockerPath), method, kj::mv(endpoint),
569 bodyBytes, "application/json"_kj, MAX_JSON_RESPONSE_SIZE);
570 co_return DockerResponse{.statusCode = raw.statusCode, .body = kj::str(raw.body.asChars())};
571}
572 
573kj::Promise<DockerBinaryResponse> dockerApiBinaryRequest(kj::Network& network,
574 kj::String dockerPath,
575 kj::HttpMethod method,
576 kj::String endpoint,
577 kj::Maybe<kj::Array<kj::byte>> body,
578 uint64_t maxResponseSize) {
579 kj::Maybe<kj::ArrayPtr<const kj::byte>> bodyBytes;
580 KJ_IF_SOME(b, body) {
581 bodyBytes = b.asPtr();
582 }
583 co_return co_await dockerApiRequestRaw(network, kj::mv(dockerPath), method, kj::mv(endpoint),
584 bodyBytes, "application/x-tar"_kj, maxResponseSize);
585}
586 
587kj::Promise<void> deleteVolume(kj::Network& network, kj::String dockerPath, kj::String volumeName) {
588 auto response = co_await dockerApiRequest(
589 network, kj::mv(dockerPath), kj::HttpMethod::DELETE, kj::str("/volumes/", volumeName));
590 if (response.statusCode != 204 && response.statusCode != 404) {
591 KJ_LOG(WARNING, "failed to delete volume", volumeName, response.statusCode, response.body);
592 }
593}
594 
595kj::Promise<void> deleteVolumes(
596 kj::Network& network, kj::String dockerPath, kj::Array<kj::String> snapshotCloneVolumes) {
597 kj::Vector<kj::Promise<void>> volumeDeletes;
598 volumeDeletes.reserve(snapshotCloneVolumes.size());
599 for (auto& volumeName: snapshotCloneVolumes) {
600 auto logName = kj::str(volumeName);
601 volumeDeletes.add(deleteVolume(network, kj::str(dockerPath), kj::mv(volumeName))
602 .catch_([logName = kj::mv(logName)](kj::Exception&& e) {
603 KJ_LOG(WARNING, "failed to delete volume", logName, e);
604 }));
605 }
606 co_await kj::joinPromises(volumeDeletes.releaseAsArray());
607}
608 
609kj::Promise<void> removeContainer(
610 kj::Network& network, kj::String dockerPath, kj::String containerName, bool wait = true) {
611 auto endpoint = kj::str("/containers/", containerName, "?force=true");
612 auto response = co_await dockerApiRequest(
613 network, kj::str(dockerPath), kj::HttpMethod::DELETE, kj::mv(endpoint));
614 // 204 means the container was removed.
615 // 404 means it was already gone.
616 // 409 means removal is already in progress, which is fine for our teardown paths.
617 KJ_REQUIRE(response.statusCode == 204 || response.statusCode == 404 || response.statusCode == 409,
618 "Removing a container failed with: ", response.body);
619 
620 // If removal succeeded or is already in progress, wait for Docker to report the container as
621 // fully removed before proceeding with any follow-up cleanup like deleting mounted volumes.
622 if (wait && (response.statusCode == 204 || response.statusCode == 409)) {
623 response = co_await dockerApiRequest(network, kj::mv(dockerPath), kj::HttpMethod::POST,
624 kj::str("/containers/", containerName, "/wait?condition=removed"));
625 // 200 means Docker observed the removal. 404 means the container disappeared before the wait
626 // request was processed, which is also fine.
627 KJ_REQUIRE(response.statusCode == 200 || response.statusCode == 404,
628 "Waiting for container removal failed with: ", response.statusCode, response.body);
629 }
630}
631 
632kj::Maybe<size_t> tryFindHttpHeaderEnd(kj::ArrayPtr<const kj::byte> bytes) {
633 for (auto i: kj::zeroTo(bytes.size())) {
634 if (i + 4 > bytes.size()) {
635 return kj::none;
636 }
637 if (bytes[i] == '\r' && bytes[i + 1] == '\n' && bytes[i + 2] == '\r' && bytes[i + 3] == '\n') {
638 return i;
639 }
640 }
641 return kj::none;
642}
643 
644// readDockerStreamedResponse is necessary because Docker streamed responses
645// require an open bidirectional stream after the response headers have been read.
646kj::Promise<DockerStreamedResponse> readDockerStreamedResponse(
647 kj::Own<kj::AsyncIoStream> connection) {
648 kj::Vector<kj::byte> buffer;
649 auto& input = *connection;
650 
651 while (true) {
652 KJ_IF_SOME(headerEnd, tryFindHttpHeaderEnd(buffer.asPtr())) {
653 auto parsedHeaders = kj::heapArray<char>(headerEnd + 2);
654 for (auto i: kj::zeroTo(parsedHeaders.size())) {
655 parsedHeaders[i] = static_cast<char>(buffer[i]);
656 }
657 
658 kj::HttpHeaderTable headerTable;
659 kj::HttpHeaders headers(headerTable);
660 auto parsedResponse = headers.tryParseResponse(parsedHeaders.asPtr());
661 headers.takeOwnership(kj::mv(parsedHeaders));
662 
663 kj::uint statusCode = 0;
664 kj::String statusText;
665 KJ_SWITCH_ONEOF(parsedResponse) {
666 KJ_CASE_ONEOF(response, kj::HttpHeaders::Response) {
667 statusCode = response.statusCode;
668 statusText = kj::str(response.statusText);
669 }
670 KJ_CASE_ONEOF(protocolError, kj::HttpHeaders::ProtocolError) {
671 KJ_FAIL_REQUIRE("Docker streamed response returned malformed HTTP headers: ",
672 protocolError.statusMessage, ": ", protocolError.description);
673 }
674 }
675 
676 auto bodyOffset = headerEnd + 4;
677 auto prefetchedBytes = kj::heapArray(buffer.asPtr().slice(bodyOffset));
678 kj::Own<kj::AsyncIoStream> prefixedConnection = kj::mv(connection);
679 if (prefetchedBytes.size() > 0) {
680 prefixedConnection =
681 kj::heap<BufferedAsyncIoStream>(kj::mv(prefixedConnection), kj::mv(prefetchedBytes));
682 }
683 co_return DockerStreamedResponse{
684 .statusCode = statusCode,
685 .statusText = kj::mv(statusText),
686 .connection = kj::mv(prefixedConnection),
687 };
688 }
689 
690 auto scratch = kj::heapArray<kj::byte>(4096);
691 auto amount = co_await input.tryRead(scratch.begin(), 1, scratch.size());
692 KJ_REQUIRE(amount > 0, "EOF while waiting for Docker streamed response headers");
693 buffer.addAll(scratch.first(amount));
694 KJ_REQUIRE(buffer.size() <= 65536, "Docker streamed response headers exceeded 64KiB");
695 }
696}
697 
698kj::Promise<DockerStreamedResponse> dockerApiStreamedRequest(kj::Network& network,
699 kj::String dockerPath,
700 kj::HttpMethod method,
701 kj::String endpoint,
702 const kj::HttpHeaders& headers,
703 kj::Maybe<kj::ArrayPtr<const kj::byte>> body = kj::none) {
704 auto address = co_await network.parseAddress(dockerPath);
705 auto connection = co_await address->connect();
706 
707 auto requestHeaders = headers.serializeRequest(method, endpoint);
708 KJ_IF_SOME(requestBody, body) {
709 kj::ArrayPtr<const kj::byte> pieces[] = {requestHeaders.asBytes(), requestBody};
710 co_await connection->write(kj::arrayPtr(pieces));
711 } else {
712 co_await connection->write(requestHeaders.asBytes());
713 }
714 
715 co_return co_await readDockerStreamedResponse(kj::mv(connection));
716}
717 
718// Docker multiplexed stream frames: 1 byte stream ID + 3 reserved + 4 bytes big-endian length.
719constexpr size_t DOCKER_FRAME_HEADER_SIZE = 8;
720 
721uint32_t parseDockerFrameLength(kj::ArrayPtr<const kj::byte> frameHeader) {
722 KJ_REQUIRE(frameHeader.size() >= DOCKER_FRAME_HEADER_SIZE, "Docker raw stream header too short");
723 return (static_cast<uint32_t>(frameHeader[4]) << 24) |
724 (static_cast<uint32_t>(frameHeader[5]) << 16) | (static_cast<uint32_t>(frameHeader[6]) << 8) |
725 static_cast<uint32_t>(frameHeader[7]);
726}
727 
728void detachEnd(kj::Maybe<kj::Own<capnp::ExplicitEndOutputStream>> stream) {
729 KJ_IF_SOME(s, stream) {
730 s->end().attach(kj::mv(s)).detach([](kj::Exception&&) {});
731 }
732}
733 
734// demuxDockerExecOutput demuxes the input from Docker to passed stdout/stderr.
735kj::Promise<void> demuxDockerExecOutput(kj::AsyncInputStream& input,
736 kj::Maybe<kj::Own<capnp::ExplicitEndOutputStream>> stdoutWriter,
737 kj::Maybe<kj::Own<capnp::ExplicitEndOutputStream>> stderrWriter,
738 bool combinedOutput) {
739 kj::Vector<kj::byte> buffer;
740 size_t offset = 0;
741 
742 auto compactBuffer = [&]() {
743 if (offset == 0) {
744 return;
745 }
746 
747 kj::Vector<kj::byte> compacted;
748 compacted.addAll(buffer.asPtr().slice(offset));
749 buffer = kj::mv(compacted);
750 offset = 0;
751 };
752 
753 auto ensureBytes = [&](size_t count) -> kj::Promise<bool> {
754 while (buffer.size() - offset < count) {
755 compactBuffer();
756 auto scratch = kj::heapArray<kj::byte>(4096);
757 auto amount = co_await input.tryRead(scratch.begin(), 1, scratch.size());
758 if (amount == 0) {
759 co_return false;
760 }
761 buffer.addAll(scratch.first(amount));
762 }
763 co_return true;
764 };
765 
766 try {
767 while (co_await ensureBytes(DOCKER_FRAME_HEADER_SIZE)) {
768 auto frameHeader = buffer.asPtr().slice(offset, offset + DOCKER_FRAME_HEADER_SIZE);
769 auto streamId = frameHeader[0];
770 auto frameLength = parseDockerFrameLength(frameHeader);
771 KJ_REQUIRE(co_await ensureBytes(DOCKER_FRAME_HEADER_SIZE + frameLength),
772 "Docker exec raw stream ended in the middle of a frame");
773 
774 auto payload = buffer.asPtr().slice(
775 offset + DOCKER_FRAME_HEADER_SIZE, offset + DOCKER_FRAME_HEADER_SIZE + frameLength);
776 if (streamId == 1) {
777 KJ_IF_SOME(out, stdoutWriter) {
778 co_await out->write(payload);
779 }
780 } else {
781 if (streamId == 2) {
782 if (combinedOutput) {
783 KJ_IF_SOME(out, stdoutWriter) {
784 co_await out->write(payload);
785 }
786 } else {
787 KJ_IF_SOME(err, stderrWriter) {
788 co_await err->write(payload);
789 }
790 }
791 }
792 }
793 
794 offset += DOCKER_FRAME_HEADER_SIZE + frameLength;
795 }
796 
797 if (buffer.size() != offset) {
798 KJ_FAIL_REQUIRE("Docker exec raw stream ended with a truncated frame header");
799 }
800 
801 // We need to detach ourselves from the end() as the user might've
802 // decided to not read them altogether.
803 detachEnd(kj::mv(stdoutWriter));
804 detachEnd(kj::mv(stderrWriter));
805 } catch (...) {
806 auto exception = kj::getCaughtExceptionAsKj();
807 
808 KJ_IF_SOME(out, stdoutWriter) {
809 out->abortWrite(exception.clone());
810 }
811 KJ_IF_SOME(err, stderrWriter) {
812 err->abortWrite(exception.clone());
813 }
814 kj::throwFatalException(kj::mv(exception));
815 }
816}
817 
818kj::String currentSnapshotVolumeTimestamp() {
819 return kj::str((kj::systemPreciseCalendarClock().now() - kj::UNIX_EPOCH) / kj::SECONDS);
820}
821 
822kj::Maybe<int64_t> tryGetSnapshotCreatedAt(capnp::JsonValue::Reader labels) {
823 if (!labels.isObject()) {
824 return kj::none;
825 }
826 
827 for (auto field: labels.getObject()) {
828 if (field.getName() != SNAPSHOT_VOLUME_CREATED_AT_LABEL) {
829 continue;
830 }
831 
832 auto value = field.getValue();
833 if (!value.isString()) {
834 return kj::none;
835 }
836 return value.getString().tryParseAs<int64_t>();
837 }
838 
839 return kj::none;
840}
841 
842kj::Promise<void> warnAboutStaleSnapshotVolumes(kj::Network& network, kj::String dockerPath) {
843 capnp::JsonCodec codec;
844 codec.handleByAnnotation<docker_api::Docker::VolumeListFilters>();
845 capnp::MallocMessageBuilder filterMessage;
846 auto filters = filterMessage.initRoot<docker_api::Docker::VolumeListFilters>();
847 auto names = filters.initName(1);
848 names.set(0, SNAPSHOT_VOLUME_PREFIX);
849 
850 auto response = co_await dockerApiRequest(network, kj::mv(dockerPath), kj::HttpMethod::GET,
851 kj::str("/volumes?filters=", kj::encodeUriComponent(codec.encode(filters))));
852 if (response.statusCode != 200) {
853 co_return;
854 }
855 
856 auto message = decodeJsonResponse<docker_api::Docker::VolumeListResponse>(response.body);
857 auto root = message->getRoot<docker_api::Docker::VolumeListResponse>();
858 auto now = kj::systemPreciseCalendarClock().now();
859 kj::Vector<kj::String> staleVolumes;
860 
861 for (auto volume: root.getVolumes()) {
862 KJ_IF_SOME(createdAtSeconds, tryGetSnapshotCreatedAt(volume.getLabels())) {
863 auto createdAt = kj::UNIX_EPOCH + createdAtSeconds * kj::SECONDS;
864 if (now - createdAt >= SNAPSHOT_STALE_AGE) {
865 staleVolumes.add(kj::str(volume.getName()));
866 }
867 }
868 }
869 
870 if (!staleVolumes.empty()) {
871 KJ_LOG(WARNING, "the following snapshot volumes were created 30+ days ago and may be stale",
872 kj::strArray(staleVolumes, ", "));
873 }
874}
875 
876// Returns the gateway IP on Linux for direct container access.
877// Returns kj::none on macOS where Docker Desktop routes host-gateway to host loopback.
878kj::Maybe<kj::String> gatewayForPlatform(kj::String gateway) {
879#ifdef __APPLE__
880 return kj::none;
881#else
882 return kj::mv(gateway);
883#endif
884}
885 
886kj::Maybe<uint16_t> tryParsePublishedHostPort(capnp::json::Value::Reader portMappingValue) {
887 if (portMappingValue.isNull()) {
888 return kj::none;
889 }
890 
891 JSG_REQUIRE(
892 portMappingValue.isArray(), Error, "Malformed ContainerInspect port mapping response");
893 auto bindings = portMappingValue.getArray();
894 if (bindings.size() == 0) {
895 return kj::none;
896 }
897 
898 auto binding = bindings[0];
899 JSG_REQUIRE(binding.isObject(), Error, "Malformed ContainerInspect port binding response");
900 for (auto field: binding.getObject()) {
901 if (field.getName() == "HostPort") {
902 auto value = field.getValue();
903 JSG_REQUIRE(value.isString(), Error, "Malformed ContainerInspect port binding response");
904 kj::StringPtr hostPort = value.getString();
905 return KJ_REQUIRE_NONNULL(
906 hostPort.tryParseAs<uint16_t>(), "Malformed ContainerInspect host port");
907 }
908 }
909 
910 KJ_FAIL_REQUIRE("Malformed ContainerInspect port binding response: missing HostPort");
911}
912 
913} // namespace
914 
915// Represents a parsed egress mapping. IP/CIDR mappings match destination IPs,
916// while hostnameGlob mappings match either HTTP hostnames or TLS SNI depending on protocol.
917// Defined here (not in the header) to avoid pulling kj::OneOf, kj::CidrRange, and
918// kj::Vector into server.c++ which includes container-client.h.
919struct ContainerClient::EgressMapping {
920 kj::OneOf<kj::CidrRange, kj::String> destination;
921 uint16_t port; // 0 means match all ports
922 EgressProtocol protocol;
923 kj::Own<workerd::IoChannelFactory::SubrequestChannel> channel;
924};
925 
926// Holds all egress mapping state. Stored via kj::Own<EgressState> in ContainerClient
927// so that the EgressMapping type is not visible in container-client.h.
928struct ContainerClient::EgressState {
929 kj::Vector<EgressMapping> mappings;
930};
931 
932ContainerClient::ContainerClient(capnp::ByteStreamFactory& byteStreamFactory,
933 kj::Timer& timer,
934 kj::Network& network,
935 kj::String dockerPath,
936 kj::String containerName,
937 kj::String imageName,
938 kj::String containerEgressInterceptorImage,
939 kj::TaskSet& waitUntilTasks,
940 kj::Promise<void> pendingCleanup,
941 kj::Function<void(kj::Promise<void>)> cleanupCallback,
942 ChannelTokenHandler& channelTokenHandler)
943 : byteStreamFactory(byteStreamFactory),
944 timer(timer),
945 network(network),
946 dockerPath(kj::mv(dockerPath)),
947 containerName(kj::encodeUriComponent(kj::str(containerName))),
948 sidecarContainerName(kj::encodeUriComponent(kj::str(containerName, "-proxy"))),
949 imageName(kj::mv(imageName)),
950 containerEgressInterceptorImage(kj::mv(containerEgressInterceptorImage)),
951 waitUntilTasks(waitUntilTasks),
952 pendingCleanup(kj::mv(pendingCleanup).fork()),
953 cleanupCallback(kj::mv(cleanupCallback)),
954 channelTokenHandler(channelTokenHandler),
955 egressState(kj::heap<EgressState>()) {
956 if (!staleSnapshotVolumeCheckScheduled.exchange(true)) {
957 waitUntilTasks.add(warnAboutStaleSnapshotVolumes(network, kj::str(this->dockerPath))
958 .catch_([](kj::Exception&& e) {
959 KJ_LOG(WARNING, "failed to inspect snapshot volumes for staleness", e);
960 }));
961 }
962}
963 
964ContainerClient::~ContainerClient() noexcept(false) {
965 stopEgressListener();
966 
967 // Best-effort cleanup for both containers.
968 auto sidecarCleanup =
969 removeContainer(network, kj::str(dockerPath), kj::str(sidecarContainerName), false)
970 .catch_([](kj::Exception&&) {});
971 
972 // Also try to delete any cloned snapshot volumes.
973 auto volumes = snapshotClones.releaseAsArray();
974 auto mainCleanup = removeContainer(network, kj::str(dockerPath), kj::str(containerName))
975 .catch_([](kj::Exception&&) {})
976 .then([&network = network, dockerPath = kj::str(dockerPath),
977 volumes = kj::mv(volumes)]() mutable {
978 return deleteVolumes(network, kj::mv(dockerPath), kj::mv(volumes));
979 }).catch_([](kj::Exception&&) {});
980 
981 // Pass the joined cleanup promise to the callback. The callback wraps it with the
982 // canceler (so a future client creation can cancel it), stores it so the next
983 // ContainerClient can await it, and adds a branch to waitUntilTasks to keep the
984 // underlying I/O alive.
985 cleanupCallback(kj::joinPromises(kj::arr(kj::mv(sidecarCleanup), kj::mv(mainCleanup))));
986}
987 
988// Docker-specific Port implementation that implements rpc::Container::Port::Server
989// It does a HTTP CONNECT to the proxy-everything sidecar port.
990class ContainerClient::DockerPort final: public rpc::Container::Port::Server {
991 public:
992 DockerPort(ContainerClient& containerClient, kj::String containerHost, uint16_t containerPort)
993 : containerClient(containerClient),
994 containerHost(kj::mv(containerHost)),
995 containerPort(containerPort) {}
996 
997 kj::Promise<void> connect(ConnectContext context) override {
998 auto mappedPort = JSG_REQUIRE_NONNULL(containerClient.sidecarIngressHostPort, Error,
999 "connect(): Container ingress proxy is not running.");
1000 
1001 auto dstAddr = kj::str(containerHost, ":", containerPort);
1002 
1003 auto address = co_await containerClient.network.parseAddress(kj::str("127.0.0.1:", mappedPort));
1004 
1005 kj::HttpHeaderTable::Builder headerTableBuilder;
1006 auto xDstAddrHeader = headerTableBuilder.add("X-Dst-Addr");
1007 auto headerTable = headerTableBuilder.build();
1008 kj::HttpHeaders headers(*headerTable);
1009 headers.set(xDstAddrHeader, kj::str(dstAddr));
1010 
1011 auto proxyConnection = co_await address->connect();
1012 auto httpClient = kj::newHttpClient(*headerTable, *proxyConnection)
1013 .attach(kj::mv(proxyConnection), kj::mv(headerTable));
1014 auto connectRequest = httpClient->connect(dstAddr, headers, {});
1015 auto status = co_await kj::mv(connectRequest.status);
1016 
1017 if (status.statusCode == 400) {
1018 throw JSG_KJ_EXCEPTION(
1019 DISCONNECTED, Error, "Container is not listening to port ", containerPort);
1020 }
1021 
1022 if (status.statusCode < 200 || status.statusCode >= 300) {
1023 KJ_IF_SOME(errorBody, status.errorBody) {
1024 auto errorBodyText = co_await errorBody->readAllText();
1025 JSG_FAIL_REQUIRE(Error, "Connecting to container port through proxy-everything failed: [",
1026 status.statusCode, "] ", status.statusText, " ", errorBodyText);
1027 }
1028 
1029 JSG_FAIL_REQUIRE(Error, "Connecting to container port through proxy-everything failed: [",
1030 status.statusCode, "] ", status.statusText);
1031 }
1032 
1033 auto connection = kj::mv(connectRequest.connection);
1034 auto upPipe = kj::newOneWayPipe();
1035 auto upEnd = kj::mv(upPipe.in);
1036 auto results = context.getResults();
1037 results.setUp(containerClient.byteStreamFactory.kjToCapnp(kj::mv(upPipe.out)));
1038 auto downEnd = containerClient.byteStreamFactory.capnpToKj(context.getParams().getDown());
1039 
1040 pumpTask =
1041 kj::joinPromisesFailFast(kj::arr(upEnd->pumpTo(*connection), connection->pumpTo(*downEnd)))
1042 .ignoreResult()
1043 .attach(kj::mv(httpClient), kj::mv(upEnd), kj::mv(connection), kj::mv(downEnd));
1044 co_return;
1045 }
1046 
1047 private:
1048 // ContainerClient is owned by the Worker::Actor and keeps it alive.
1049 ContainerClient& containerClient;
1050 kj::String containerHost;
1051 uint16_t containerPort;
1052 kj::Maybe<kj::Promise<void>> pumpTask;
1053};
1054 
1055class ContainerClient::DockerProcessHandle final: public rpc::Container::ProcessHandle::Server {
1056 public:
1057 DockerProcessHandle(ContainerClient& containerClient,
1058 kj::String execId,
1059 kj::Own<kj::AsyncIoStream> connection,
1060 kj::Maybe<capnp::ByteStream::Client> stdoutWriter,
1061 kj::Maybe<capnp::ByteStream::Client> stderrWriter,
1062 bool combinedOutput)
1063 : containerClient(containerClient.addRef()),
1064 execId(kj::mv(execId)),
1065 sharedConnection(kj::refcounted<SharedExecConnection>(kj::mv(connection))) {
1066 kj::Maybe<kj::Own<capnp::ExplicitEndOutputStream>> stdoutStream = kj::none;
1067 KJ_IF_SOME(out, stdoutWriter) {
1068 stdoutStream = this->containerClient->byteStreamFactory.capnpToKjExplicitEnd(out);
1069 } else {
1070 stdoutStream = capnp::ExplicitEndOutputStream::wrap(newNullOutputStream(), []() {});
1071 }
1072 
1073 kj::Maybe<kj::Own<capnp::ExplicitEndOutputStream>> stderrStream = kj::none;
1074 KJ_IF_SOME(err, stderrWriter) {
1075 stderrStream = this->containerClient->byteStreamFactory.capnpToKjExplicitEnd(err);
1076 } else if (!combinedOutput) {
1077 stderrStream = capnp::ExplicitEndOutputStream::wrap(newNullOutputStream(), []() {});
1078 }
1079 
1080 // Always drain the Docker exec stream. This lets wait() use stream closure as the primary
1081 // process-completion signal, even when stdout/stderr are ignored.
1082 auto task = demuxDockerExecOutput(
1083 *sharedConnection->connection, kj::mv(stdoutStream), kj::mv(stderrStream), combinedOutput)
1084 .attach(this->containerClient->addRef(), kj::addRef(*sharedConnection));
1085 streamClosedTask = kj::mv(task).fork();
1086 }
1087 
1088 kj::Promise<void> wait(WaitContext context) override {
1089 waitStarted = true;
1090 if (!sharedConnection->stdinOpened && !sharedConnection->stdinClosed) {
1091 sharedConnection->connection->shutdownWrite();
1092 sharedConnection->stdinClosed = true;
1093 }
1094 
1095 co_await KJ_ASSERT_NONNULL(streamClosedTask).addBranch();
1096 
1097 // Docker's exec-inspect state can lag slightly behind the hijacked stream closing, so after
1098 // we observe EOF we allow a short bounded retry window to obtain the final exit code.
1099 for (auto attempt: kj::zeroTo(20)) {
1100 auto inspect = co_await containerClient->inspectExec(execId);
1101 if (!inspect.running) {
1102 context.getResults().setExitCode(inspect.exitCode);
1103 co_return;
1104 }
1105 
1106 if (attempt + 1 < 20) {
1107 co_await containerClient->timer.afterDelay(50 * kj::MILLISECONDS);
1108 }
1109 }
1110 
1111 JSG_FAIL_REQUIRE(Error, "Docker exec stream closed before exit status became available.");
1112 }
1113 
1114 kj::Promise<void> stdinWriter(StdinWriterContext context) override {
1115 JSG_REQUIRE(!waitStarted, Error, "Process stdinWriter() cannot be called after wait().");
1116 JSG_REQUIRE(
1117 !sharedConnection->stdinOpened, Error, "Process stdinWriter() can only be called once.");
1118 
1119 sharedConnection->stdinOpened = true;
1120 context.getResults().setWriter(containerClient->byteStreamFactory.kjToCapnp(
1121 kj::heap<DockerExecStdinStream>(kj::addRef(*sharedConnection))));
1122 co_return;
1123 }
1124 
1125 kj::Promise<void> kill(KillContext context) override {
1126 auto inspect = co_await containerClient->inspectExec(execId);
1127 JSG_REQUIRE(inspect.pid > 0, Error, "Exec process does not have a visible pid to signal.");
1128 
1129 auto signal = kj::str("-", signalToString(context.getParams().getSigno()));
1130 auto pid = kj::str(inspect.pid);
1131 auto cmd = kj::arr(kj::str("kill"), kj::mv(signal), kj::mv(pid));
1132 co_await containerClient->runSimpleExec(cmd.asPtr());
1133 }
1134 
1135 private:
1136 kj::Own<ContainerClient> containerClient;
1137 kj::String execId;
1138 kj::Own<SharedExecConnection> sharedConnection;
1139 bool waitStarted = false;
1140 kj::Maybe<kj::ForkedPromise<void>> streamClosedTask;
1141};
1142 
1143// ConnectResponse adapter for TCP egress. Since we've already accepted the sidecar's HTTP
1144// CONNECT before calling worker->connect(), this adapter simply records the worker's
1145// accept/reject decision without sending anything on the wire.
1146class TcpEgressConnectResponse final: public kj::HttpService::ConnectResponse {
1147 public:
1148 void accept(uint statusCode, kj::StringPtr statusText, const kj::HttpHeaders& headers) override {
1149 // Worker accepted the connection. Nothing additional to do since we already
1150 // accepted the sidecar CONNECT.
1151 }
1152 
1153 kj::Own<kj::AsyncOutputStream> reject(uint statusCode,
1154 kj::StringPtr statusText,
1155 const kj::HttpHeaders& headers,
1156 kj::Maybe<uint64_t> expectedBodySize = kj::none) override {
1157 KJ_FAIL_REQUIRE("TCP egress worker rejected the connection: ", statusCode, " ", statusText);
1158 }
1159};
1160 
1161// HTTP service that handles HTTP CONNECT requests from the container sidecar (proxy-everything).
1162// When the sidecar intercepts container egress traffic, it sends HTTP CONNECT to this service.
1163// After accepting the CONNECT, the tunnel carries the actual HTTP request from the container,
1164// which we parse and forward to the appropriate SubrequestChannel based on egressMappings.
1165// Inner HTTP service that handles requests inside the CONNECT tunnel.
1166// Forwards requests to the worker binding via SubrequestChannel.
1167class InnerEgressService final: public kj::HttpService {
1168 public:
1169 using ChannelLookup = kj::Function<kj::Maybe<kj::Own<IoChannelFactory::SubrequestChannel>>()>;
1170 
1171 InnerEgressService(ChannelLookup lookupChannel, kj::StringPtr destAddr, bool isTls = false)
1172 : lookupChannel(kj::mv(lookupChannel)),
1173 destAddr(kj::str(destAddr)),
1174 isTls(isTls) {}
1175 
1176 kj::Promise<void> request(kj::HttpMethod method,
1177 kj::StringPtr requestUri,
1178 const kj::HttpHeaders& headers,
1179 kj::AsyncInputStream& requestBody,
1180 Response& response) override {
1181 // Look up the channel on each request so we always use the latest mapping,
1182 // even if it was replaced via interceptOutboundHttp while the tunnel is open.
1183 auto channel =
1184 KJ_REQUIRE_NONNULL(lookupChannel(), "egress mapping disappeared during active tunnel");
1185 
1186 IoChannelFactory::SubrequestMetadata metadata;
1187 auto worker = channel->startRequest(kj::mv(metadata));
1188 auto urlForWorker = kj::str(requestUri);
1189 // Probably only a path, try to get it from Host:
1190 if (requestUri.startsWith("/")) {
1191 auto scheme = isTls ? "https://"_kj : "http://"_kj;
1192 auto baseUrl = kj::str(scheme, destAddr);
1193 // Use Host: when possible
1194 KJ_IF_SOME(host, headers.get(kj::HttpHeaderId::HOST)) {
1195 baseUrl = kj::str(scheme, host);
1196 }
1197 
1198 // Parse url, if invalid, try to use the original requestUri (http://<ip>/<path>
1199 KJ_IF_SOME(parsedUrl, jsg::Url::tryParse(requestUri, baseUrl.asPtr())) {
1200 urlForWorker = kj::str(parsedUrl.getHref());
1201 } else {
1202 urlForWorker = kj::str(baseUrl, requestUri);
1203 }
1204 }
1205 
1206 co_await worker->request(method, urlForWorker, headers, requestBody, response);
1207 }
1208 
1209 private:
1210 ChannelLookup lookupChannel;
1211 kj::String destAddr;
1212 bool isTls;
1213};
1214 
1215kj::Promise<void> pumpBidirectional(kj::AsyncIoStream& a, kj::AsyncIoStream& b) {
1216 auto aToB = a.pumpTo(b).then([&b](uint64_t) { b.shutdownWrite(); });
1217 auto bToA = b.pumpTo(a).then([&a](uint64_t) { a.shutdownWrite(); });
1218 co_await kj::joinPromisesFailFast(kj::arr(kj::mv(aToB), kj::mv(bToA)));
1219}
1220 
1221// Outer HTTP service that handles CONNECT requests from the sidecar.
1222class EgressHttpService final: public kj::HttpService {
1223 public:
1224 EgressHttpService(ContainerClient& containerClient, kj::HttpHeaderTable& headerTable)
1225 : containerClient(containerClient),
1226 headerTable(headerTable) {}
1227 
1228 kj::Promise<void> request(kj::HttpMethod method,
1229 kj::StringPtr url,
1230 const kj::HttpHeaders& headers,
1231 kj::AsyncInputStream& requestBody,
1232 Response& response) override {
1233 // Regular HTTP requests are not expected - we only handle CONNECT
1234 co_return co_await response.sendError(405, "Method Not Allowed", headerTable);
1235 }
1236 
1237 kj::Promise<void> connect(kj::StringPtr host,
1238 const kj::HttpHeaders& headers,
1239 kj::AsyncIoStream& connection,
1240 ConnectResponse& response,
1241 kj::HttpConnectSettings settings) override {
1242 auto destAddr = kj::str(host);
1243 if (co_await handleConnectMode(destAddr, headers, connection, response, "X-Tls-Sni",
1244 /*defaultPort=*/443, EgressProtocol::HTTPS)) {
1245 co_return;
1246 }
1247 
1248 if (co_await handleConnectMode(destAddr, headers, connection, response, "X-Hostname",
1249 /*defaultPort=*/80, EgressProtocol::HTTP)) {
1250 co_return;
1251 }
1252 
1253 // Try raw TCP mapping before falling through to passthrough. TCP mappings match
1254 // on IP/CIDR + port without any application-layer hostname information.
1255 if (co_await handleTcpConnect(destAddr, connection, response)) {
1256 co_return;
1257 }
1258 
1259 kj::HttpHeaders responseHeaders(headerTable);
1260 // 202 is interpreted by proxy-everything as "just send bytes as-is".
1261 // If the connection was TLS, it's useful so we just proxy transparently
1262 // to the internet.
1263 response.accept(202, "Accepted", responseHeaders);
1264 
1265 co_await passThroughConnection(destAddr, connection);
1266 }
1267 
1268 private:
1269 kj::Promise<bool> handleConnectMode(kj::StringPtr destAddr,
1270 const kj::HttpHeaders& headers,
1271 kj::AsyncIoStream& connection,
1272 ConnectResponse& response,
1273 kj::StringPtr hostnameHeader,
1274 uint16_t defaultPort,
1275 EgressProtocol protocol) {
1276 kj::Maybe<kj::String> requestHostname;
1277 KJ_IF_SOME(value, getHeader(headers, hostnameHeader)) {
1278 requestHostname = kj::str(value);
1279 }
1280 
1281 auto mapping = containerClient.findEgressMapping(destAddr, defaultPort,
1282 requestHostname.map([](auto& hostname) {
1283 return kj::Maybe<kj::StringPtr>(hostname);
1284 }).orDefault(kj::none),
1285 protocol);
1286 
1287 if (requestHostname == kj::none && mapping == kj::none) {
1288 co_return false;
1289 }
1290 
1291 if (mapping != kj::none) {
1292 kj::HttpHeaders responseHeaders(headerTable);
1293 response.accept(200, "OK", responseHeaders);
1294 
1295 bool isTls = (protocol == EgressProtocol::HTTPS);
1296 auto innerService = kj::heap<InnerEgressService>(
1297 [&client = containerClient, addr = kj::str(destAddr),
1298 hostname = requestHostname.map([](auto& value) { return kj::str(value); }),
1299 defaultPort,
1300 protocol]() mutable -> kj::Maybe<kj::Own<IoChannelFactory::SubrequestChannel>> {
1301 return client.findEgressMapping(addr, defaultPort,
1302 hostname.map([](auto& value) {
1303 return kj::Maybe<kj::StringPtr>(value);
1304 }).orDefault(kj::none),
1305 protocol);
1306 },
1307 destAddr, isTls);
1308 auto innerServer =
1309 kj::heap<kj::HttpServer>(containerClient.timer, headerTable, *innerService);
1310 
1311 co_await innerServer->listenHttpCleanDrain(connection);
1312 co_return true;
1313 }
1314 
1315 kj::HttpHeaders responseHeaders(headerTable);
1316 response.accept(202, "Accepted", responseHeaders);
1317 
1318 co_await passThroughConnection(destAddr, connection);
1319 co_return true;
1320 }
1321 
1322 // Handles raw TCP egress by forwarding the sidecar tunnel to the worker's connect() handler.
1323 // Returns true if a TCP mapping matched and the connection was handled.
1324 kj::Promise<bool> handleTcpConnect(
1325 kj::StringPtr destAddr, kj::AsyncIoStream& connection, ConnectResponse& response) {
1326 // For TCP, we match on IP:port only — no hostname matching since raw TCP
1327 // doesn't carry application-layer hostname information.
1328 auto mapping = containerClient.findEgressMapping(
1329 destAddr, /*defaultPort=*/0, /*hostname=*/kj::none, EgressProtocol::TCP);
1330 
1331 if (mapping == kj::none) {
1332 co_return false;
1333 }
1334 
1335 auto& channel = KJ_ASSERT_NONNULL(mapping);
1336 
1337 kj::HttpHeaders responseHeaders(headerTable);
1338 // 202 tells proxy-everything to send bytes as-is without attempting to
1339 // interpret the stream (e.g. if the underlying TCP carries TLS).
1340 response.accept(202, "Accepted", responseHeaders);
1341 
1342 IoChannelFactory::SubrequestMetadata metadata;
1343 auto worker = channel->startRequest(kj::mv(metadata));
1344 
1345 // Bridge the sidecar tunnel to the worker's connect() handler. The worker entrypoint
1346 // is expected to implement connect() (e.g., a WorkerEntrypoint that proxies TCP).
1347 // We provide a simple ConnectResponse adapter since we've already accepted the
1348 // sidecar's CONNECT above.
1349 TcpEgressConnectResponse tcpResponse;
1350 kj::HttpHeaders connectHeaders(headerTable);
1351 co_await worker->connect(destAddr, connectHeaders, connection, tcpResponse, {});
1352 
1353 co_return true;
1354 }
1355 
1356 kj::Promise<void> passThroughConnection(kj::StringPtr destAddr, kj::AsyncIoStream& connection) {
1357 if (!containerClient.internetEnabled.orDefault(false)) {
1358 connection.shutdownWrite();
1359 co_return;
1360 }
1361 
1362 auto addr = co_await containerClient.network.parseAddress(destAddr);
1363 auto destConn = co_await addr->connect();
1364 co_await pumpBidirectional(connection, *destConn);
1365 co_return;
1366 }
1367 
1368 ContainerClient& containerClient;
1369 kj::HttpHeaderTable& headerTable;
1370};
1371 
1372kj::Promise<ContainerClient::IPAMConfigResult> ContainerClient::getDockerBridgeIPAMConfig() {
1373 auto response = co_await dockerApiRequest(
1374 network, kj::str(dockerPath), kj::HttpMethod::GET, kj::str("/networks/bridge"));
1375 if (response.statusCode == 200) {
1376 auto message = decodeJsonResponse<docker_api::Docker::NetworkInspectResponse>(response.body);
1377 auto jsonRoot = message->getRoot<docker_api::Docker::NetworkInspectResponse>();
1378 auto ipamConfig = jsonRoot.getIpam().getConfig();
1379 if (ipamConfig.size() > 0) {
1380 auto config = ipamConfig[0];
1381 co_return IPAMConfigResult{
1382 .gateway = kj::str(config.getGateway()),
1383 .subnet = kj::str(config.getSubnet()),
1384 };
1385 }
1386 }
1387 
1388 JSG_FAIL_REQUIRE(Error,
1389 "Failed to get bridge. "
1390 "Status: ",
1391 response.statusCode, ", Body: ", response.body);
1392}
1393 
1394kj::Promise<bool> ContainerClient::isDaemonIpv6Enabled() {
1395 // Inspect the default bridge network. When the Docker daemon has "ipv6": true in
1396 // daemon.json, the default bridge gets an IPv6 IPAM subnet entry (e.g. "fd00::/80").
1397 auto response = co_await dockerApiRequest(
1398 network, kj::str(dockerPath), kj::HttpMethod::GET, kj::str("/networks/bridge"));
1399 
1400 if (response.statusCode != 200) {
1401 co_return false;
1402 }
1403 
1404 auto message = decodeJsonResponse<docker_api::Docker::NetworkInspectResponse>(response.body);
1405 auto jsonRoot = message->getRoot<docker_api::Docker::NetworkInspectResponse>();
1406 for (auto config: jsonRoot.getIpam().getConfig()) {
1407 // IPv6 subnets contain ':' (e.g. "fd00::/80", "2001:db8::/64")
1408 if (kj::StringPtr(config.getSubnet()).findFirst(':') != kj::none) {
1409 co_return true;
1410 }
1411 }
1412 
1413 co_return false;
1414}
1415 
1416kj::Promise<uint16_t> ContainerClient::startEgressListener(
1417 kj::String listenAddress, uint16_t port) {
1418 auto service = kj::heap<EgressHttpService>(*this, headerTable);
1419 auto httpServer = kj::heap<kj::HttpServer>(timer, headerTable, *service);
1420 auto& httpServerRef = *httpServer;
1421 
1422 egressHttpServer = httpServer.attach(kj::mv(service));
1423 
1424 auto addr = co_await network.parseAddress(kj::str(listenAddress, ":", port));
1425 // The gateway IP from Docker's bridge network is not always bindable on the host.
1426 // On WSL with Docker Desktop, 172.17.0.1 lives inside the Docker VM, not on the WSL host's
1427 // interfaces. In that case, fall back to loopback — the sidecar reaches the host via
1428 // host-gateway (which Docker Desktop maps to the host loopback) so 127.0.0.1 works correctly.
1429 kj::Own<kj::ConnectionReceiver> listener;
1430 KJ_IF_SOME(e, kj::runCatchingExceptions([&]() { listener = addr->listen(); })) {
1431 KJ_LOG(WARNING, "Could not bind egress listener to gateway address, falling back to loopback",
1432 listenAddress, e);
1433 auto fallbackAddr = co_await network.parseAddress(kj::str("127.0.0.1:", port));
1434 listener = fallbackAddr->listen();
1435 }
1436 
1437 uint16_t chosenPort = listener->getPort();
1438 
1439 egressListenerTask = httpServerRef.listenHttp(*listener)
1440 .attach(kj::mv(listener))
1441 .eagerlyEvaluate([](kj::Exception&& e) {
1442 LOG_EXCEPTION(
1443 "Workerd could not listen in the TCP port to proxy traffic off the docker container", e);
1444 });
1445 
1446 co_return chosenPort;
1447}
1448 
1449void ContainerClient::stopEgressListener() {
1450 egressListenerTask = kj::none;
1451 egressHttpServer = kj::none;
1452 egressListenerStarted.store(false, std::memory_order_release);
1453}
1454 
1455kj::Promise<void> ContainerClient::writeFileToContainer(kj::StringPtr container,
1456 kj::StringPtr dir,
1457 kj::StringPtr filename,
1458 kj::ArrayPtr<const kj::byte> content) {
1459 kj::HttpHeaderTable table;
1460 auto address = co_await network.parseAddress(kj::str(dockerPath));
1461 auto connection = co_await address->connect();
1462 auto httpClient = kj::newHttpClient(table, *connection).attach(kj::mv(connection));
1463 
1464 auto tar = createTarWithFile(filename, content);
1465 
1466 kj::HttpHeaders headers(table);
1467 headers.setPtr(kj::HttpHeaderId::HOST, "localhost");
1468 headers.setPtr(kj::HttpHeaderId::CONTENT_TYPE, "application/x-tar");
1469 headers.set(kj::HttpHeaderId::CONTENT_LENGTH, kj::str(tar.size()));
1470 
1471 auto endpoint = kj::str("/containers/", container, "/archive?path=", kj::encodeUriComponent(dir));
1472 auto req = httpClient->request(kj::HttpMethod::PUT, endpoint, headers, tar.size());
1473 {
1474 auto body = kj::mv(req.body);
1475 co_await body->write(tar.asBytes());
1476 }
1477 
1478 auto response = co_await req.response;
1479 auto result = co_await response.body->readAllText();
1480 JSG_REQUIRE(response.statusCode == 200, Error, "Failed to write file ", dir, "/", filename,
1481 " to container [", response.statusCode, "] ", result);
1482}
1483 
1484static constexpr kj::StringPtr cloudflareCaDir = "/etc"_kj;
1485static constexpr kj::StringPtr cloudflareCaFilename =
1486 "cloudflare/certs/cloudflare-containers-ca.crt"_kj;
1487 
1488kj::Promise<void> ContainerClient::readCACert() {
1489 auto ingressPort = KJ_REQUIRE_NONNULL(
1490 sidecarIngressHostPort, "Cannot read CA cert: sidecar ingress port not known");
1491 
1492 auto response = co_await dockerApiRequest(
1493 network, kj::str("127.0.0.1:", ingressPort), kj::HttpMethod::GET, kj::str("/ca"));
1494 
1495 JSG_REQUIRE(response.statusCode == 200, Error,
1496 "Failed to read CA cert from sidecar: ", response.statusCode, " ", response.body);
1497 
1498 caCert = kj::mv(response.body);
1499}
1500 
1501kj::Promise<void> ContainerClient::injectCACert() {
1502 if (caCertInjected.exchange(true, std::memory_order_acquire)) {
1503 co_return;
1504 }
1505 
1506 bool succeeded = false;
1507 KJ_DEFER(if (!succeeded) caCertInjected.store(false, std::memory_order_release));
1508 
1509 if (caCert == kj::none) {
1510 co_await readCACert();
1511 }
1512 
1513 auto& cert = KJ_REQUIRE_NONNULL(caCert, "CA cert not read from sidecar yet");
1514 co_await writeFileToContainer(
1515 containerName, cloudflareCaDir, cloudflareCaFilename, cert.asBytes());
1516 
1517 succeeded = true;
1518}
1519 
1520kj::Promise<kj::Maybe<ContainerClient::InspectResponse>> ContainerClient::inspectContainer() {
1521 auto endpoint = kj::str("/containers/", containerName, "/json");
1522 
1523 auto response = co_await dockerApiRequest(
1524 network, kj::str(dockerPath), kj::HttpMethod::GET, kj::mv(endpoint));
1525 // We check if the container with the given name exist, and if it's not,
1526 // we simply return false while avoiding an unnecessary error.
1527 if (response.statusCode == 404) {
1528 co_return kj::none;
1529 }
1530 
1531 JSG_REQUIRE(response.statusCode == 200, Error, "Container inspect failed");
1532 // Parse JSON response
1533 auto message = decodeJsonResponse<docker_api::Docker::ContainerInspectResponse>(response.body);
1534 auto jsonRoot = message->getRoot<docker_api::Docker::ContainerInspectResponse>();
1535 
1536 // Look for Status field in the JSON object
1537 JSG_REQUIRE(jsonRoot.hasState(), Error, "Malformed ContainerInspect response");
1538 auto state = jsonRoot.getState();
1539 JSG_REQUIRE(state.hasStatus(), Error, "Malformed ContainerInspect response");
1540 auto status = state.getStatus();
1541 // Treat both "running" and "restarting" as running. The "restarting" state occurs when
1542 // Docker is automatically restarting a container (due to restart policy). From the user's
1543 // perspective, a restarting container is still "alive" and should be treated as running
1544 // so that start() correctly refuses to start a duplicate and destroy() can clean it up.
1545 bool running = status == "running" || status == "restarting";
1546 
1547 kj::Vector<Label> labels;
1548 if (jsonRoot.hasConfig() && jsonRoot.getConfig().hasLabels()) {
1549 auto labelsJson = jsonRoot.getConfig().getLabels();
1550 if (labelsJson.isObject()) {
1551 for (auto field: labelsJson.getObject()) {
1552 kj::StringPtr name = field.getName();
1553 if (!name.startsWith(WORKERD_LABEL_PREFIX)) continue;
1554 auto value = field.getValue();
1555 JSG_REQUIRE(value.isString(), Error, "Malformed ContainerInspect label value");
1556 labels.add(Label{
1557 .name = kj::str(name.slice(WORKERD_LABEL_PREFIX.size())),
1558 .value = kj::str(value.getString()),
1559 });
1560 }
1561 }
1562 }
1563 
1564 co_return InspectResponse{.isRunning = running, .labels = labels.releaseAsArray()};
1565}
1566 
1567kj::Promise<kj::Maybe<ContainerClient::SidecarInspectResponse>> ContainerClient::inspectSidecar() {
1568 auto endpoint = kj::str("/containers/", sidecarContainerName, "/json");
1569 auto response = co_await dockerApiRequest(
1570 network, kj::str(dockerPath), kj::HttpMethod::GET, kj::mv(endpoint));
1571 
1572 if (response.statusCode == 404) {
1573 co_return kj::none;
1574 }
1575 
1576 JSG_REQUIRE(response.statusCode == 200, Error, "Sidecar container inspect failed");
1577 
1578 auto message = decodeJsonResponse<docker_api::Docker::ContainerInspectResponse>(response.body);
1579 auto jsonRoot = message->getRoot<docker_api::Docker::ContainerInspectResponse>();
1580 
1581 // Check if sidecar is actually running
1582 bool running = false;
1583 if (jsonRoot.hasState()) {
1584 auto state = jsonRoot.getState();
1585 if (state.hasStatus()) {
1586 auto status = state.getStatus();
1587 running = status == "running" || status == "restarting";
1588 }
1589 }
1590 
1591 if (!running) {
1592 co_return kj::none;
1593 }
1594 
1595 kj::Maybe<uint16_t> ingressHostPort;
1596 
1597 auto ingressPortKey = kj::str(SIDECAR_INGRESS_PORT, "/tcp");
1598 for (auto portMapping: jsonRoot.getNetworkSettings().getPorts().getObject()) {
1599 if (portMapping.getName() != ingressPortKey) {
1600 continue;
1601 }
1602 
1603 ingressHostPort = tryParsePublishedHostPort(portMapping.getValue());
1604 break;
1605 }
1606 
1607 auto requiredIngressHostPort =
1608 KJ_REQUIRE_NONNULL(ingressHostPort, "running sidecar missing ingress host port");
1609 
1610 co_return SidecarInspectResponse{
1611 .ingressHostPort = requiredIngressHostPort,
1612 };
1613}
1614 
1615kj::Promise<void> ContainerClient::updateSidecarEgressPort(
1616 uint16_t ingressHostPort, uint16_t egressPort) {
1617 capnp::JsonCodec codec;
1618 codec.handleByAnnotation<docker_api::ProxyEverything::Port>();
1619 capnp::MallocMessageBuilder message;
1620 auto jsonRoot = message.initRoot<docker_api::ProxyEverything::Port>();
1621 jsonRoot.setPort(egressPort);
1622 
1623 auto body = codec.encode(jsonRoot);
1624 auto response = co_await dockerApiRequest(network, kj::str("127.0.0.1:", ingressHostPort),
1625 kj::HttpMethod::PUT, kj::str("/egress"), kj::mv(body));
1626 
1627 JSG_REQUIRE(response.statusCode >= 200 && response.statusCode < 300, Error,
1628 "Updating sidecar egress port failed with: ", response.statusCode, " ", response.body);
1629}
1630 
1631kj::Promise<void> ContainerClient::updateSidecarEgressConfig(
1632 uint16_t ingressHostPort, uint16_t egressPort) {
1633 capnp::JsonCodec codec;
1634 codec.handleByAnnotation<docker_api::ProxyEverything::Port>();
1635 capnp::MallocMessageBuilder message;
1636 auto jsonRoot = message.initRoot<docker_api::ProxyEverything::Port>();
1637 jsonRoot.setPort(egressPort);
1638 
1639 auto allowHostnames = getDnsAllowHostnames();
1640 auto dns = jsonRoot.initDns();
1641 auto allowHostnamesList = dns.initAllowHostnames(allowHostnames.size());
1642 for (auto i: kj::indices(allowHostnames)) {
1643 allowHostnamesList.set(i, allowHostnames[i]);
1644 }
1645 
1646 KJ_IF_SOME(enabled, internetEnabled) {
1647 jsonRoot.initInternet().setEnabled(enabled);
1648 }
1649 
1650 auto body = codec.encode(jsonRoot);
1651 auto response = co_await dockerApiRequest(network, kj::str("127.0.0.1:", ingressHostPort),
1652 kj::HttpMethod::PUT, kj::str("/egress"), kj::mv(body));
1653 
1654 JSG_REQUIRE(response.statusCode >= 200 && response.statusCode < 300, Error,
1655 "Updating sidecar egress config failed with: ", response.statusCode, " ", response.body);
1656}
1657 
1658kj::Promise<void> ContainerClient::createContainer(kj::StringPtr effectiveImage,
1659 kj::Maybe<capnp::List<capnp::Text>::Reader> entrypoint,
1660 kj::Maybe<capnp::List<capnp::Text>::Reader> environment,
1661 kj::ArrayPtr<const SnapshotRestoreMount> restoreMounts,
1662 rpc::Container::StartParams::Reader params) {
1663 capnp::JsonCodec codec;
1664 codec.handleByAnnotation<docker_api::Docker::ContainerCreateRequest>();
1665 capnp::MallocMessageBuilder message;
1666 auto jsonRoot = message.initRoot<docker_api::Docker::ContainerCreateRequest>();
1667 jsonRoot.setImage(effectiveImage);
1668 // Add entrypoint if provided
1669 KJ_IF_SOME(ep, entrypoint) {
1670 auto jsonCmd = jsonRoot.initCmd(ep.size());
1671 for (uint32_t i: kj::zeroTo(ep.size())) {
1672 jsonCmd.set(i, ep[i]);
1673 }
1674 }
1675 
1676 auto envSize = environment.map([](auto& env) { return env.size(); }).orDefault(0);
1677 auto jsonEnv = jsonRoot.initEnv(envSize + kj::size(defaultEnv));
1678 
1679 KJ_IF_SOME(env, environment) {
1680 for (uint32_t i: kj::zeroTo(env.size())) {
1681 jsonEnv.set(i, env[i]);
1682 }
1683 }
1684 
1685 for (uint32_t i: kj::zeroTo(kj::size(defaultEnv))) {
1686 jsonEnv.set(envSize + i, defaultEnv[i]);
1687 }
1688 
1689 // Pass user-supplied labels as Docker object labels, prefixed so we can distinguish
1690 // them from image/engine labels when reading back via inspect().
1691 if (params.hasLabels()) {
1692 auto lbls = params.getLabels();
1693 auto labelsObj = jsonRoot.initLabels().initObject(lbls.size());
1694 for (auto i: kj::zeroTo(lbls.size())) {
1695 labelsObj[i].setName(kj::str(WORKERD_LABEL_PREFIX, lbls[i].getName()));
1696 labelsObj[i].initValue().setString(lbls[i].getValue());
1697 }
1698 }
1699 
1700 auto hostConfig = jsonRoot.initHostConfig();
1701 // We need to set a restart policy to avoid having ambiguous states
1702 // where the container we're managing is stuck at "exited" state.
1703 hostConfig.initRestartPolicy().setName("on-failure");
1704 
1705 hostConfig.setNetworkMode(kj::str("container:", sidecarContainerName));
1706 
1707 // When containersPidNamespace is NOT enabled, use host PID namespace for backwards compatibility.
1708 // This allows the container to see processes on the host.
1709 if (!params.getCompatibilityFlags().getContainersPidNamespace()) {
1710 hostConfig.setPidMode("host");
1711 }
1712 
1713 if (restoreMounts.size() > 0) {
1714 auto mounts = hostConfig.initMounts(restoreMounts.size());
1715 for (auto i: kj::indices(restoreMounts)) {
1716 auto mount = mounts[i];
1717 auto& restoreMount = restoreMounts[i];
1718 mount.setType("volume");
1719 mount.setSource(restoreMount.cloneVolume);
1720 mount.setTarget(restoreMount.restorePath.toString(true));
1721 mount.initVolumeOptions().setNoCopy(true);
1722 }
1723 }
1724 
1725 auto response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
1726 kj::str("/containers/create?name=", containerName), codec.encode(jsonRoot));
1727 
1728 // statusCode 409 refers to "conflict". Occurs when a container with the given name exists.
1729 // In that case we destroy and re-create the container. We retry a few times with delays
1730 // because Docker may take a moment to fully release the container name after removal.
1731 constexpr int MAX_RETRIES = 3;
1732 constexpr auto RETRY_DELAY = 100 * kj::MILLISECONDS;
1733 
1734 for (int attempt = 0; response.statusCode == 409 && attempt < MAX_RETRIES; ++attempt) {
1735 co_await removeContainer(network, kj::str(dockerPath), kj::str(containerName));
1736 co_await timer.afterDelay(RETRY_DELAY);
1737 response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
1738 kj::str("/containers/create?name=", containerName), codec.encode(jsonRoot));
1739 }
1740 
1741 // statusCode 201 refers to "container created successfully"
1742 if (response.statusCode != 201) {
1743 JSG_REQUIRE(
1744 response.statusCode != 404, Error, "No such image available named ", effectiveImage);
1745 JSG_REQUIRE(response.statusCode != 409, Error, "Container already exists");
1746 JSG_FAIL_REQUIRE(
1747 Error, "Create container failed with [", response.statusCode, "] ", response.body);
1748 }
1749}
1750 
1751kj::Promise<kj::String> ContainerClient::createExec(capnp::List<capnp::Text>::Reader cmd,
1752 rpc::Container::ExecOptions::Reader params,
1753 bool attachStdout,
1754 bool attachStderr) {
1755 capnp::JsonCodec codec;
1756 codec.handleByAnnotation<docker_api::Docker::ExecCreateRequest>();
1757 
1758 capnp::MallocMessageBuilder message;
1759 auto request = message.initRoot<docker_api::Docker::ExecCreateRequest>();
1760 request.setAttachStdin(true);
1761 request.setAttachStdout(attachStdout);
1762 request.setAttachStderr(attachStderr);
1763 request.setTty(false);
1764 
1765 auto jsonCmd = request.initCmd(cmd.size());
1766 for (auto i: kj::zeroTo(cmd.size())) {
1767 jsonCmd.set(i, cmd[i]);
1768 }
1769 
1770 if (params.hasEnv()) {
1771 auto env = params.getEnv();
1772 auto jsonEnv = request.initEnv(env.size());
1773 for (auto i: kj::zeroTo(env.size())) {
1774 jsonEnv.set(i, env[i]);
1775 }
1776 }
1777 
1778 if (params.hasWorkingDirectory()) {
1779 request.setWorkingDir(params.getWorkingDirectory());
1780 }
1781 
1782 if (params.hasUser()) {
1783 request.setUser(params.getUser());
1784 }
1785 
1786 auto response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
1787 kj::str("/containers/", containerName, "/exec"), codec.encode(request));
1788 JSG_REQUIRE(response.statusCode == 201, Error, "Creating Docker exec failed with [",
1789 response.statusCode, "] ", response.body);
1790 
1791 auto parsed = decodeJsonResponse<docker_api::Docker::ExecCreateResponse>(response.body);
1792 co_return kj::str(parsed->getRoot<docker_api::Docker::ExecCreateResponse>().getId());
1793}
1794 
1795kj::Promise<kj::Own<kj::AsyncIoStream>> ContainerClient::startExec(kj::String execId) {
1796 capnp::JsonCodec codec;
1797 codec.handleByAnnotation<docker_api::Docker::ExecStartRequest>();
1798 
1799 capnp::MallocMessageBuilder message;
1800 auto requestBody = message.initRoot<docker_api::Docker::ExecStartRequest>();
1801 requestBody.setDetach(false);
1802 requestBody.setTty(false);
1803 auto encodedBody = codec.encode(requestBody);
1804 
1805 // Exec attach uses HTTP connection hijacking. A plain POST can succeed with 200 OK but then not
1806 // behave like the raw stream Docker's CLI expects, so we must request the upgrade explicitly.
1807 kj::HttpHeaderTable headerTable;
1808 kj::HttpHeaders headers(headerTable);
1809 headers.setPtr(kj::HttpHeaderId::HOST, "localhost");
1810 headers.setPtr(kj::HttpHeaderId::CONNECTION, "Upgrade");
1811 // ... Why not CONNECT or WebSockets, Docker?
1812 headers.setPtr(kj::HttpHeaderId::UPGRADE, "tcp");
1813 headers.setPtr(kj::HttpHeaderId::CONTENT_TYPE, "application/json");
1814 headers.set(kj::HttpHeaderId::CONTENT_LENGTH, kj::str(encodedBody.size()));
1815 kj::ArrayPtr<const kj::byte> encodedBodyBytes = encodedBody.asBytes();
1816 
1817 auto response = co_await dockerApiStreamedRequest(network, kj::str(dockerPath),
1818 kj::HttpMethod::POST, kj::str("/exec/", execId, "/start"), headers, encodedBodyBytes);
1819 if (response.statusCode != 101) {
1820 auto errorBodyBytes = co_await response.connection->readAllBytes(MAX_JSON_RESPONSE_SIZE);
1821 auto errorBody = kj::str(errorBodyBytes.asChars());
1822 JSG_FAIL_REQUIRE(Error, "Starting Docker exec failed with [", response.statusCode, "] ",
1823 response.statusText, " ", errorBody);
1824 }
1825 
1826 co_return kj::mv(response.connection);
1827}
1828 
1829kj::Promise<ContainerClient::ExecInspectResponse> ContainerClient::inspectExec(
1830 kj::StringPtr execId) {
1831 auto response = co_await dockerApiRequest(
1832 network, kj::str(dockerPath), kj::HttpMethod::GET, kj::str("/exec/", execId, "/json"));
1833 JSG_REQUIRE(response.statusCode == 200, Error, "Inspecting Docker exec failed with [",
1834 response.statusCode, "] ", response.body);
1835 
1836 auto parsed = decodeJsonResponse<docker_api::Docker::ExecInspectResponse>(response.body);
1837 auto root = parsed->getRoot<docker_api::Docker::ExecInspectResponse>();
1838 auto exitCodeValue = root.getExitCode();
1839 auto exitCode = exitCodeValue.isNumber() ? static_cast<int32_t>(exitCodeValue.getNumber()) : 0;
1840 co_return ExecInspectResponse{
1841 .exitCode = exitCode,
1842 .running = root.getRunning(),
1843 .pid = root.getPid(),
1844 };
1845}
1846 
1847kj::Promise<void> ContainerClient::runSimpleExec(kj::ArrayPtr<const kj::String> cmd) {
1848 capnp::JsonCodec codec;
1849 codec.handleByAnnotation<docker_api::Docker::ExecCreateRequest>();
1850 
1851 capnp::MallocMessageBuilder createMessage;
1852 auto createRequest = createMessage.initRoot<docker_api::Docker::ExecCreateRequest>();
1853 createRequest.setAttachStdin(false);
1854 createRequest.setAttachStdout(false);
1855 createRequest.setAttachStderr(false);
1856 createRequest.setTty(false);
1857 
1858 auto jsonCmd = createRequest.initCmd(cmd.size());
1859 for (auto i: kj::indices(cmd)) {
1860 jsonCmd.set(i, cmd[i]);
1861 }
1862 
1863 auto createResponse =
1864 co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
1865 kj::str("/containers/", containerName, "/exec"), codec.encode(createRequest));
1866 JSG_REQUIRE(createResponse.statusCode == 201, Error, "Creating helper Docker exec failed with [",
1867 createResponse.statusCode, "] ", createResponse.body);
1868 
1869 auto parsedCreate =
1870 decodeJsonResponse<docker_api::Docker::ExecCreateResponse>(createResponse.body);
1871 auto execId = kj::str(parsedCreate->getRoot<docker_api::Docker::ExecCreateResponse>().getId());
1872 
1873 capnp::JsonCodec startCodec;
1874 startCodec.handleByAnnotation<docker_api::Docker::ExecStartRequest>();
1875 
1876 capnp::MallocMessageBuilder startMessage;
1877 auto startRequest = startMessage.initRoot<docker_api::Docker::ExecStartRequest>();
1878 startRequest.setDetach(true);
1879 startRequest.setTty(false);
1880 
1881 auto startResponse = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
1882 kj::str("/exec/", execId, "/start"), startCodec.encode(startRequest));
1883 JSG_REQUIRE(startResponse.statusCode == 200, Error, "Starting helper Docker exec failed with [",
1884 startResponse.statusCode, "] ", startResponse.body);
1885 
1886 while (true) {
1887 auto inspect = co_await inspectExec(execId);
1888 if (!inspect.running) {
1889 JSG_REQUIRE(inspect.exitCode == 0, Error, "Helper Docker exec failed with exit code ",
1890 inspect.exitCode);
1891 co_return;
1892 }
1893 co_await timer.afterDelay(50 * kj::MILLISECONDS);
1894 }
1895}
1896 
1897kj::Promise<void> ContainerClient::startContainer() {
1898 auto endpoint = kj::str("/containers/", containerName, "/start");
1899 // We have to send an empty body since docker API will throw an error if we don't.
1900 auto response = co_await dockerApiRequest(
1901 network, kj::str(dockerPath), kj::HttpMethod::POST, kj::mv(endpoint), kj::str(""));
1902 // statusCode 304 refers to "container already started"
1903 JSG_REQUIRE(response.statusCode != 304, Error, "Container already started");
1904 // statusCode 204 refers to "no error"
1905 JSG_REQUIRE(response.statusCode == 204, Error, "Starting container failed with: ", response.body);
1906}
1907 
1908kj::Promise<void> ContainerClient::stopContainer() {
1909 auto endpoint = kj::str("/containers/", containerName, "/stop");
1910 auto response = co_await dockerApiRequest(
1911 network, kj::str(dockerPath), kj::HttpMethod::POST, kj::mv(endpoint));
1912 // statusCode 204 refers to "no error"
1913 // statusCode 304 refers to "container already stopped"
1914 // Both are fine to avoid when stop container is called.
1915 JSG_REQUIRE(response.statusCode == 204 || response.statusCode == 304, Error,
1916 "Stopping container failed with: ", response.body);
1917}
1918 
1919kj::Promise<void> ContainerClient::killContainer(uint32_t signal) {
1920 auto endpoint = kj::str("/containers/", containerName, "/kill?signal=", signalToString(signal));
1921 auto response = co_await dockerApiRequest(
1922 network, kj::str(dockerPath), kj::HttpMethod::POST, kj::mv(endpoint));
1923 // statusCode 409 refers to "container is not running"
1924 // We should not throw an error when the container is already not running.
1925 JSG_REQUIRE(response.statusCode == 204 || response.statusCode == 409, Error,
1926 "Stopping container failed with: ", response.body);
1927}
1928 
1929// Destroys the container.
1930// No-op when the container does not exist.
1931// Wait for the container to actually be stopped and removed when it exists.
1932kj::Promise<void> ContainerClient::destroyContainer() {
1933 co_await removeContainer(network, kj::str(dockerPath), kj::str(containerName));
1934 co_await deleteVolumes(network, kj::str(dockerPath), snapshotClones.releaseAsArray());
1935}
1936 
1937// Creates the sidecar container that owns the shared network namespace.
1938// The application container joins this namespace and all ingress/egress goes through it.
1939kj::Promise<void> ContainerClient::createSidecarContainer(
1940 uint16_t egressPort, kj::String networkCidr) {
1941 // Equivalent to: docker run --cap-add=NET_ADMIN -p <random-host>:39001 ...
1942 capnp::JsonCodec codec;
1943 codec.handleByAnnotation<docker_api::Docker::ContainerCreateRequest>();
1944 capnp::MallocMessageBuilder message;
1945 auto jsonRoot = message.initRoot<docker_api::Docker::ContainerCreateRequest>();
1946 jsonRoot.setImage(containerEgressInterceptorImage);
1947 
1948 auto ipv6Enabled = co_await isDaemonIpv6Enabled();
1949 
1950 // determined by the number of flags we need to pass to proxy-everything
1951 uint32_t cmdSize =
1952 8; // --http-egress-port <port> --http-ingress-address 0.0.0.0:<port> --docker-gateway-cidr <cidr> --dns-enabled --tls-intercept
1953 if (!ipv6Enabled) cmdSize += 1; // --disable-ipv6
1954 
1955 auto cmd = jsonRoot.initCmd(cmdSize);
1956 uint32_t idx = 0;
1957 cmd.set(idx++, "--http-egress-port");
1958 cmd.set(idx++, kj::str(egressPort));
1959 cmd.set(idx++, "--http-ingress-address");
1960 cmd.set(idx++, kj::str("0.0.0.0:", SIDECAR_INGRESS_PORT));
1961 cmd.set(idx++, "--docker-gateway-cidr");
1962 cmd.set(idx++, networkCidr);
1963 cmd.set(idx++, "--dns-enabled");
1964 cmd.set(idx++, "--tls-intercept");
1965 if (!ipv6Enabled) {
1966 cmd.set(idx++, "--disable-ipv6");
1967 }
1968 
1969 jsonRoot.initExposedPorts().setRaw(kj::str("{\"", SIDECAR_INGRESS_PORT, "/tcp\":{}}"));
1970 
1971 auto hostConfig = jsonRoot.initHostConfig();
1972 hostConfig.setPublishAllPorts(true);
1973 hostConfig.setNetworkMode("bridge");
1974 auto dns = hostConfig.initDns(kj::size(SIDECAR_DNS_SERVERS));
1975 for (auto i: kj::indices(SIDECAR_DNS_SERVERS)) {
1976 dns.set(i, SIDECAR_DNS_SERVERS[i]);
1977 }
1978 
1979 auto extraHosts = hostConfig.initExtraHosts(1);
1980 extraHosts.set(0, "host.docker.internal:host-gateway"_kj);
1981 
1982 // Sidecar needs NET_ADMIN capability for iptables/TPROXY
1983 auto capAdd = hostConfig.initCapAdd(1);
1984 capAdd.set(0, "NET_ADMIN");
1985 
1986 auto response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
1987 kj::str("/containers/create?name=", sidecarContainerName), codec.encode(jsonRoot));
1988 
1989 if (response.statusCode == 409) {
1990 // Already created, nothing to do
1991 co_return;
1992 }
1993 
1994 if (response.statusCode != 201) {
1995 JSG_REQUIRE(response.statusCode != 404, Error, "No such image available named ",
1996 containerEgressInterceptorImage,
1997 ". Please ensure the container egress interceptor image is built and available.");
1998 JSG_FAIL_REQUIRE(Error, "Failed to create the networking sidecar [", response.statusCode, "] ",
1999 response.body);
2000 }
2001}
2002 
2003kj::Promise<void> ContainerClient::startSidecarContainer() {
2004 auto endpoint = kj::str("/containers/", sidecarContainerName, "/start");
2005 auto response = co_await dockerApiRequest(
2006 network, kj::str(dockerPath), kj::HttpMethod::POST, kj::mv(endpoint), kj::str(""));
2007 // statusCode 304 refers to "container already started"
2008 // statusCode 204 refers to "request succeeded"
2009 JSG_REQUIRE(response.statusCode == 204 || response.statusCode == 304, Error,
2010 "Starting network sidecar container failed with: ", response.statusCode, response.body);
2011}
2012 
2013kj::Promise<void> ContainerClient::destroySidecarContainer() {
2014 co_await removeContainer(network, kj::str(dockerPath), kj::str(sidecarContainerName));
2015}
2016 
2017kj::Promise<void> ContainerClient::createVolume(kj::StringPtr volumeName) {
2018 capnp::JsonCodec codec;
2019 codec.handleByAnnotation<docker_api::Docker::VolumeCreateRequest>();
2020 capnp::MallocMessageBuilder message;
2021 auto req = message.initRoot<docker_api::Docker::VolumeCreateRequest>();
2022 req.setName(volumeName);
2023 auto labels = req.initLabels().initObject(1);
2024 labels[0].setName(SNAPSHOT_VOLUME_CREATED_AT_LABEL);
2025 labels[0].initValue().setString(currentSnapshotVolumeTimestamp());
2026 
2027 auto response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
2028 kj::str("/volumes/create"), codec.encode(req));
2029 // Docker returns 201 for new volumes and 200 for existing ones.
2030 JSG_REQUIRE(response.statusCode == 201 || response.statusCode == 200, Error,
2031 "Failed to create Docker volume '", volumeName, "': ", response.statusCode, " ",
2032 response.body);
2033}
2034 
2035kj::Promise<void> ContainerClient::deleteVolume(kj::String volumeName) {
2036 auto response = co_await dockerApiRequest(
2037 network, kj::str(dockerPath), kj::HttpMethod::DELETE, kj::str("/volumes/", volumeName));
2038 // 204 = deleted, 404 = not found (both are fine)
2039 JSG_REQUIRE(response.statusCode == 204 || response.statusCode == 404, Error,
2040 "Failed to delete Docker volume '", volumeName, "': ", response.statusCode, " ",
2041 response.body);
2042}
2043 
2044kj::Promise<void> ContainerClient::commitContainer(kj::StringPtr imageRef) {
2045 auto response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
2046 kj::str("/commit?container=", containerName,
2047 "&pause=true&repo=", kj::encodeUriComponent(imageRef)),
2048 kj::str(""));
2049 JSG_REQUIRE(response.statusCode == 201, Error, "Failed to commit container to image '", imageRef,
2050 "': ", response.statusCode, " ", response.body);
2051}
2052 
2053kj::Promise<ContainerClient::ImageInspectResponse> ContainerClient::inspectImage(
2054 kj::StringPtr imageRef) {
2055 auto response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::GET,
2056 kj::str("/images/", kj::encodeUriComponent(imageRef), "/json"));
2057 JSG_REQUIRE(response.statusCode == 200, Error, "Failed to inspect Docker image '", imageRef,
2058 "': ", response.statusCode, " ", response.body);
2059 
2060 auto message = decodeJsonResponse<docker_api::Docker::ImageInspectResponse>(response.body);
2061 auto root = message->getRoot<docker_api::Docker::ImageInspectResponse>();
2062 co_return ImageInspectResponse{kj::str(root.getId()), root.getSize()};
2063}
2064 
2065kj::Promise<void> ContainerClient::deleteImage(kj::String imageRef) {
2066 auto response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::DELETE,
2067 kj::str("/images/", kj::encodeUriComponent(imageRef), "?noprune=true"));
2068 JSG_REQUIRE(response.statusCode == 200 || response.statusCode == 404, Error,
2069 "Failed to delete Docker image '", imageRef, "': ", response.statusCode, " ", response.body);
2070}
2071 
2072kj::Promise<kj::String> ContainerClient::createTempContainerWithVolume(
2073 kj::StringPtr volumeName, kj::StringPtr mountPath) {
2074 capnp::JsonCodec codec;
2075 codec.handleByAnnotation<docker_api::Docker::ContainerCreateRequest>();
2076 capnp::MallocMessageBuilder message;
2077 auto jsonRoot = message.initRoot<docker_api::Docker::ContainerCreateRequest>();
2078 jsonRoot.setImage(imageName);
2079 
2080 auto hostConfig = jsonRoot.initHostConfig();
2081 auto binds = hostConfig.initBinds(1);
2082 binds.set(0, kj::str(volumeName, ":", mountPath));
2083 
2084 auto response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
2085 kj::str("/containers/create"), codec.encode(jsonRoot));
2086 JSG_REQUIRE(response.statusCode == 201, Error, "Failed to create temp container for volume '",
2087 volumeName, "': ", response.statusCode, " ", response.body);
2088 
2089 auto respMessage = decodeJsonResponse<docker_api::Docker::ContainerCreateResponse>(response.body);
2090 auto respRoot = respMessage->getRoot<docker_api::Docker::ContainerCreateResponse>();
2091 co_return kj::str(respRoot.getId());
2092}
2093 
2094kj::Promise<void> ContainerClient::cloneSnapshot(SnapshotRestoreMount& snapshot) {
2095 co_await createVolume(snapshot.cloneVolume);
2096 
2097 bool cloneCommitted = false;
2098 KJ_DEFER(if (!cloneCommitted) {
2099 waitUntilTasks.add(deleteVolume(kj::str(snapshot.cloneVolume)).catch_([](kj::Exception&&) {
2100 }).attach(addRef()));
2101 });
2102 
2103 capnp::JsonCodec codec;
2104 codec.handleByAnnotation<docker_api::Docker::ContainerCreateRequest>();
2105 capnp::MallocMessageBuilder message;
2106 auto jsonRoot = message.initRoot<docker_api::Docker::ContainerCreateRequest>();
2107 jsonRoot.setImage(containerEgressInterceptorImage);
2108 jsonRoot.setEntrypoint("/bin/cp");
2109 
2110 // Run `/bin/cp -a /src/. /dst/` so the clone volume gets the snapshot contents directly.
2111 auto cmd = jsonRoot.initCmd(3);
2112 cmd.set(0, "-a");
2113 cmd.set(1, "/src/.");
2114 cmd.set(2, "/dst/");
2115 
2116 auto hostConfig = jsonRoot.initHostConfig();
2117 auto binds = hostConfig.initBinds(2);
2118 binds.set(0, kj::str(snapshot.sourceVolume, ":/src:ro"));
2119 binds.set(1, kj::str(snapshot.cloneVolume, ":/dst"));
2120 
2121 auto createResponse = co_await dockerApiRequest(network, kj::str(dockerPath),
2122 kj::HttpMethod::POST, kj::str("/containers/create"), codec.encode(jsonRoot));
2123 JSG_REQUIRE(createResponse.statusCode == 201, Error,
2124 "Failed to create snapshot clone helper container for volume '", snapshot.sourceVolume,
2125 "': ", createResponse.statusCode, " ", createResponse.body);
2126 
2127 auto createMessage =
2128 decodeJsonResponse<docker_api::Docker::ContainerCreateResponse>(createResponse.body);
2129 auto createRoot = createMessage->getRoot<docker_api::Docker::ContainerCreateResponse>();
2130 auto helperContainerId = kj::str(createRoot.getId());
2131 bool helperDeleted = false;
2132 KJ_DEFER(if (!helperDeleted) {
2133 waitUntilTasks.add(
2134 deleteTempContainer(kj::str(helperContainerId)).catch_([](kj::Exception&&) {
2135 }).attach(addRef()));
2136 });
2137 
2138 auto startResponse = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
2139 kj::str("/containers/", helperContainerId, "/start"), kj::str(""));
2140 JSG_REQUIRE(startResponse.statusCode == 204, Error,
2141 "Failed to start snapshot clone helper container '", helperContainerId,
2142 "': ", startResponse.statusCode, " ", startResponse.body);
2143 
2144 auto waitResponse = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::POST,
2145 kj::str("/containers/", helperContainerId, "/wait?condition=not-running"));
2146 JSG_REQUIRE(waitResponse.statusCode == 200, Error,
2147 "Failed waiting for snapshot clone helper container '", helperContainerId,
2148 "': ", waitResponse.statusCode, " ", waitResponse.body);
2149 
2150 auto waitMessage =
2151 decodeJsonResponse<docker_api::Docker::ContainerMonitorResponse>(waitResponse.body);
2152 auto waitRoot = waitMessage->getRoot<docker_api::Docker::ContainerMonitorResponse>();
2153 // A non-zero exit means the copy failed and the clone volume contents are incomplete.
2154 JSG_REQUIRE(waitRoot.getStatusCode() == 0, Error, "Snapshot clone helper container '",
2155 helperContainerId, "' exited with status ", waitRoot.getStatusCode());
2156 
2157 co_await deleteTempContainer(kj::str(helperContainerId));
2158 helperDeleted = true;
2159 cloneCommitted = true;
2160 snapshotClones.add(kj::str(snapshot.cloneVolume));
2161}
2162 
2163kj::Promise<void> ContainerClient::deleteTempContainer(kj::String tempContainerId) {
2164 auto response = co_await dockerApiRequest(network, kj::str(dockerPath), kj::HttpMethod::DELETE,
2165 kj::str("/containers/", tempContainerId, "?force=true"));
2166 // 204 = deleted, 404 = not found (both are fine).
2167 KJ_REQUIRE(response.statusCode == 204 || response.statusCode == 404,
2168 "Failed to delete temp container", tempContainerId, response.statusCode, response.body);
2169}
2170 
2171ContainerClient::RpcTurn ContainerClient::getRpcTurn() {
2172 auto paf = kj::newPromiseAndFulfiller<void>();
2173 auto prev = mutationQueue.addBranch();
2174 mutationQueue = paf.promise.fork();
2175 return {kj::mv(prev), kj::mv(paf.fulfiller)};
2176}
2177 
2178kj::Promise<void> ContainerClient::status(StatusContext context) {
2179 // Wait for any pending cleanup from a previous ContainerClient (Docker DELETE).
2180 // If the cleanup was already cancelled via containerCleanupCanceler the .catch_()
2181 // in the destructor resolves it immediately, so this is a no-op in that case.
2182 co_await pendingCleanup.addBranch();
2183 
2184 auto [ready, done] = getRpcTurn();
2185 co_await ready;
2186 KJ_DEFER(done->fulfill());
2187 
2188 bool isRunning = false;
2189 KJ_IF_SOME(info, co_await inspectContainer()) {
2190 isRunning = info.isRunning;
2191 }
2192 containerStarted.store(isRunning, std::memory_order_release);
2193 containerSidecarStarted.store(false, std::memory_order_release);
2194 this->sidecarIngressHostPort = kj::none;
2195 
2196 if (isRunning) {
2197 // If the sidecar container is already running (e.g. workerd restarted while
2198 // containers stayed up), recover its published ingress port, then configure
2199 // it to use our current egress listener port.
2200 auto sidecar = KJ_REQUIRE_NONNULL(co_await inspectSidecar(),
2201 "Recovered running container without a running networking sidecar");
2202 containerSidecarStarted.store(true, std::memory_order_release);
2203 this->sidecarIngressHostPort = sidecar.ingressHostPort;
2204 co_await ensureEgressListenerStarted();
2205 co_await updateSidecarEgressPort(sidecar.ingressHostPort, egressListenerPort);
2206 co_await readCACert();
2207 }
2208 
2209 context.getResults().setRunning(isRunning);
2210}
2211 
2212kj::Promise<void> ContainerClient::inspect(InspectContext context) {
2213 auto maybeResp = co_await inspectContainer();
2214 auto info = context.getResults().initInfo();
2215 KJ_IF_SOME(resp, maybeResp) {
2216 if (resp.isRunning) {
2217 auto started = info.initStarted();
2218 auto list = started.initLabels(resp.labels.size());
2219 for (auto i: kj::indices(resp.labels)) {
2220 list[i].setName(resp.labels[i].name);
2221 list[i].setValue(resp.labels[i].value);
2222 }
2223 co_return;
2224 }
2225 }
2226 info.setNone();
2227}
2228 
2229kj::Promise<void> ContainerClient::start(StartContext context) {
2230 auto [ready, done] = getRpcTurn();
2231 co_await ready;
2232 KJ_DEFER(done->fulfill());
2233 
2234 const auto params = context.getParams();
2235 
2236 // Get the lists directly from Cap'n Proto
2237 kj::Maybe<capnp::List<capnp::Text>::Reader> entrypoint = kj::none;
2238 kj::Maybe<capnp::List<capnp::Text>::Reader> environment = kj::none;
2239 
2240 if (params.hasEntrypoint()) {
2241 entrypoint = params.getEntrypoint();
2242 }
2243 
2244 if (params.hasEnvironmentVariables()) {
2245 environment = params.getEnvironmentVariables();
2246 }
2247 
2248 internetEnabled = params.getEnableInternet();
2249 
2250 kj::String effectiveImage = kj::str(imageName);
2251 if (params.hasContainerSnapshotId()) {
2252 auto snapshotId = parseSnapshotId(params.getContainerSnapshotId());
2253 effectiveImage = kj::str(CONTAINER_SNAPSHOT_IMAGE_PREFIX, snapshotId);
2254 co_await inspectImage(effectiveImage);
2255 }
2256 
2257 // If startup fails after we clone any snapshot volumes, tear down the app container first and
2258 // then delete those clone volumes so we don't leave mounted Docker volumes behind.
2259 KJ_DEFER(if (!containerStarted.load(std::memory_order_acquire)) {
2260 waitUntilTasks.add(destroyContainer().attach(addRef()));
2261 });
2262 
2263 kj::Vector<SnapshotRestoreMount> restoreMounts;
2264 if (params.hasDirectorySnapshots()) {
2265 auto snapshotList = params.getDirectorySnapshots();
2266 restoreMounts.reserve(snapshotList.size());
2267 for (auto i: kj::zeroTo(snapshotList.size())) {
2268 auto entry = snapshotList[i];
2269 auto snapshotId = parseSnapshotId(entry.getSnapshotId());
2270 
2271 auto restorePath = parseAbsolutePath(entry.getRestorePath());
2272 JSG_REQUIRE(restorePath.toString(true) != "/", Error,
2273 "Directory snapshot cannot be restored to root directory.");
2274 
2275 auto sourceVolume = kj::str(SNAPSHOT_VOLUME_PREFIX, snapshotId);
2276 
2277 auto inspectResp = co_await dockerApiRequest(
2278 network, kj::str(dockerPath), kj::HttpMethod::GET, kj::str("/volumes/", sourceVolume));
2279 JSG_REQUIRE(inspectResp.statusCode == 200, Error, "Snapshot '", snapshotId,
2280 "' not found (volume '", sourceVolume, "' does not exist)");
2281 
2282 restoreMounts.add(SnapshotRestoreMount{kj::mv(restorePath), kj::mv(sourceVolume),
2283 kj::str(SNAPSHOT_CLONE_VOLUME_PREFIX, randomUUID(kj::none))});
2284 }
2285 
2286 for (auto& restoreMount: restoreMounts) {
2287 co_await cloneSnapshot(restoreMount);
2288 }
2289 }
2290 
2291 co_await ensureEgressListenerStarted();
2292 co_await ensureSidecarStarted();
2293 
2294 // Refresh the sidecar's egress configuration on every start(). This is required because:
2295 // - `internetEnabled` may have changed since the sidecar was originally created.
2296 // - The DNS allow-list may have changed (egress mappings added/removed between starts).
2297 // The sidecar applies this synchronously and atomically (proxy-everything's PUT /egress
2298 // updates an atomic.Pointer before the response). On the cold path, ensureSidecarStarted()
2299 // has already pushed an initial config; this second push is a single fast PUT and is
2300 // idempotent.
2301 KJ_IF_SOME(ingressHostPort, sidecarIngressHostPort) {
2302 co_await updateSidecarEgressConfig(ingressHostPort, egressListenerPort);
2303 }
2304 
2305 caCertInjected.store(false, std::memory_order_release);
2306 co_await createContainer(effectiveImage, entrypoint, environment, restoreMounts.asPtr(), params);
2307 
2308 for (auto& mapping: egressState->mappings) {
2309 if (mapping.protocol == EgressProtocol::HTTPS) {
2310 co_await injectCACert();
2311 break;
2312 }
2313 }
2314 
2315 co_await startContainer();
2316 
2317 containerStarted.store(true, std::memory_order_release);
2318}
2319 
2320kj::Promise<void> ContainerClient::monitor(MonitorContext context) {
2321 // Wait for any in-progress mutating RPCs (e.g. start()) to complete
2322 // before issuing the Docker wait request.
2323 co_await mutationQueue.addBranch();
2324 
2325 // If start() ran but failed (e.g. snapshot restore error), containerStarted
2326 // remains false. Reject immediately rather than hanging on Docker /wait for a
2327 // container that was never started.
2328 JSG_REQUIRE(containerStarted.load(std::memory_order_acquire), Error, "Container failed to start");
2329 
2330 auto results = context.getResults();
2331 KJ_DEFER(containerStarted.store(false, std::memory_order_release));
2332 
2333 auto endpoint = kj::str("/containers/", containerName, "/wait");
2334 auto response = co_await dockerApiRequest(
2335 network, kj::str(dockerPath), kj::HttpMethod::POST, kj::mv(endpoint));
2336 
2337 JSG_REQUIRE(response.statusCode == 200, Error,
2338 "Monitoring container failed with: ", response.statusCode, " ", response.body);
2339 
2340 auto message = decodeJsonResponse<docker_api::Docker::ContainerMonitorResponse>(response.body);
2341 auto jsonRoot = message->getRoot<docker_api::Docker::ContainerMonitorResponse>();
2342 results.setExitCode(jsonRoot.getStatusCode());
2343}
2344 
2345kj::Promise<void> ContainerClient::destroy(DestroyContext context) {
2346 auto [ready, done] = getRpcTurn();
2347 co_await ready;
2348 KJ_DEFER(done->fulfill());
2349 
2350 // Tear down the app container; the sidecar is kept warm so a subsequent start() can
2351 // reuse it (creating a fresh sidecar — and waiting for its iptables / network namespace
2352 // setup — costs 4–8 seconds under Docker daemon contention; reusing the warm sidecar
2353 // skips that cost entirely). The sidecar's egress configuration is refreshed on the
2354 // next start() via updateSidecarEgressConfig.
2355 //
2356 // The sidecar will be cleaned up when the ContainerClient is destroyed (cleanupCallback
2357 // in the destructor), or on the next start() if state is inconsistent (e.g. workerd
2358 // restart left an orphaned sidecar; status() recovery handles that case).
2359 co_await destroyContainer();
2360}
2361 
2362kj::Promise<void> ContainerClient::signal(SignalContext context) {
2363 auto [ready, done] = getRpcTurn();
2364 co_await ready;
2365 KJ_DEFER(done->fulfill());
2366 
2367 const auto params = context.getParams();
2368 co_await killContainer(params.getSigno());
2369}
2370 
2371kj::Promise<void> ContainerClient::exec(ExecContext context) {
2372 auto [ready, done] = getRpcTurn();
2373 co_await ready;
2374 KJ_DEFER(done->fulfill());
2375 
2376 JSG_REQUIRE(containerStarted.load(std::memory_order_acquire), Error,
2377 "exec() requires a running container.");
2378 
2379 auto request = context.getParams();
2380 auto execParams = request.getParams();
2381 // Always attach stdout/stderr to Docker so the hijacked stream lifetime continues to track the
2382 // process even when the JS API requested "ignore". We discard ignored output locally.
2383 bool attachStdout = true;
2384 bool attachStderr = true;
2385 
2386 auto execId = co_await createExec(request.getCmd(), execParams, attachStdout, attachStderr);
2387 kj::Own<kj::AsyncIoStream> execConnection = co_await startExec(kj::str(execId));
2388 kj::Maybe<capnp::ByteStream::Client> stdoutWriter = kj::none;
2389 if (request.hasStdoutWriter()) {
2390 stdoutWriter = request.getStdoutWriter();
2391 }
2392 
2393 kj::Maybe<capnp::ByteStream::Client> stderrWriter = kj::none;
2394 if (request.hasStderrWriter()) {
2395 stderrWriter = request.getStderrWriter();
2396 }
2397 
2398 // Retrying is not great, however Docker's inspectExec might return running = false
2399 // before it has fully spawned the process (as startExec() returns before
2400 // even docker has spawned the process...)
2401 ExecInspectResponse inspect{.exitCode = 0, .running = false, .pid = 0};
2402 for (auto attempt: kj::zeroTo(20)) {
2403 inspect = co_await inspectExec(execId);
2404 if (inspect.pid != 0 || !inspect.running || attempt + 1 == 20) {
2405 break;
2406 }
2407 
2408 co_await timer.afterDelay(50 * kj::MILLISECONDS);
2409 }
2410 
2411 auto process = context.getResults().initProcess();
2412 process.setPid(static_cast<int32_t>(inspect.pid));
2413 process.setHandle(kj::heap<DockerProcessHandle>(*this, kj::mv(execId), kj::mv(execConnection),
2414 kj::mv(stdoutWriter), kj::mv(stderrWriter), execParams.getCombinedOutput()));
2415}
2416 
2417kj::Promise<void> ContainerClient::setInactivityTimeout(SetInactivityTimeoutContext context) {
2418 auto [ready, done] = getRpcTurn();
2419 co_await ready;
2420 KJ_DEFER(done->fulfill());
2421 
2422 auto params = context.getParams();
2423 auto durationMs = params.getDurationMs();
2424 
2425 JSG_REQUIRE(
2426 durationMs > 0, Error, "setInactivityTimeout() requires durationMs > 0, got ", durationMs);
2427 
2428 auto timeout = durationMs * kj::MILLISECONDS;
2429 
2430 // Add a timer task that holds a reference to this ContainerClient.
2431 waitUntilTasks.add(timer.afterDelay(timeout).then([self = kj::addRef(*this)]() {
2432 // This callback does nothing but drop the reference
2433 }));
2434 
2435 co_return;
2436}
2437 
2438kj::Promise<void> ContainerClient::snapshotDirectory(SnapshotDirectoryContext context) {
2439 auto [ready, done] = getRpcTurn();
2440 co_await ready;
2441 KJ_DEFER(done->fulfill());
2442 
2443 const auto params = context.getParams();
2444 
2445 const auto dir = parseAbsolutePath(params.getDir()).toString(true);
2446 
2447 auto name = params.hasName() && params.getName().size() > 0
2448 ? kj::Maybe<kj::String>(kj::str(params.getName()))
2449 : kj::Maybe<kj::String>(kj::none);
2450 
2451 JSG_REQUIRE(containerStarted.load(std::memory_order_acquire), Error,
2452 "snapshotDirectory() requires a running container.");
2453 
2454 auto snapshotId = randomUUID(kj::none);
2455 
2456 // GET tar archive of the directory CONTENTS from the running container.
2457 // The trailing "/." tells Docker to return contents without the directory wrapper,
2458 // so the tar entries are relative to the directory (e.g., "hello/aaa.txt" instead
2459 // of "data/hello/aaa.txt"). This decouples storage from the directory name,
2460 // allowing restore to a different mount point.
2461 // Append "/." to the path to get directory contents without the directory wrapper.
2462 // For dir == "/", this is just "/."; for others, e.g. "/app/data" → "/app/data/.".
2463 auto archivePath = dir == "/" ? kj::str("/.") : kj::str(dir, "/.");
2464 auto tarResponse = co_await dockerApiBinaryRequest(network, kj::str(dockerPath),
2465 kj::HttpMethod::GET,
2466 kj::str("/containers/", containerName, "/archive?path=", kj::encodeUriComponent(archivePath)),
2467 kj::none, MAX_SNAPSHOT_TAR_SIZE);
2468 
2469 if (tarResponse.statusCode == 404) {
2470 JSG_FAIL_REQUIRE(Error, "snapshotDirectory(): directory not found in container: ", dir);
2471 }
2472 JSG_REQUIRE(tarResponse.statusCode == 200, Error,
2473 "snapshotDirectory(): failed to read directory '", dir,
2474 "' from container: ", tarResponse.statusCode);
2475 
2476 auto tarSize = static_cast<uint64_t>(tarResponse.body.size());
2477 
2478 // Create a Docker volume to store the snapshot contents. If anything after this
2479 // fails, clean up the volume so we don't leak Docker resources on retries.
2480 auto volumeName = kj::str(SNAPSHOT_VOLUME_PREFIX, snapshotId);
2481 co_await createVolume(volumeName);
2482 bool volumeCommitted = false;
2483 KJ_DEFER(if (!volumeCommitted) {
2484 waitUntilTasks.add(
2485 deleteVolume(kj::str(volumeName)).catch_([](kj::Exception&&) {}).attach(addRef()));
2486 });
2487 
2488 // Store the contents tar in the volume via a temp container mounted at /mnt.
2489 auto tempId = co_await createTempContainerWithVolume(volumeName, "/mnt");
2490 KJ_DEFER(waitUntilTasks.add(deleteTempContainer(kj::str(tempId)).attach(addRef())));
2491 
2492 auto putResponse = co_await dockerApiBinaryRequest(network, kj::str(dockerPath),
2493 kj::HttpMethod::PUT, kj::str("/containers/", tempId, "/archive?path=/mnt"),
2494 kj::mv(tarResponse.body), MAX_JSON_RESPONSE_SIZE);
2495 JSG_REQUIRE(putResponse.statusCode == 200, Error,
2496 "snapshotDirectory(): failed to store snapshot in volume '", volumeName,
2497 "': ", putResponse.statusCode);
2498 
2499 volumeCommitted = true;
2500 KJ_LOG(INFO, "created snapshot volume", volumeName, dir, tarSize);
2501 
2502 auto result = context.getResults().initSnapshot();
2503 result.setId(snapshotId);
2504 result.setSize(tarSize);
2505 result.setDir(dir);
2506 KJ_IF_SOME(n, name) {
2507 result.setName(n);
2508 }
2509}
2510 
2511kj::Promise<void> ContainerClient::snapshotContainer(SnapshotContainerContext context) {
2512 auto [ready, done] = getRpcTurn();
2513 co_await ready;
2514 KJ_DEFER(done->fulfill());
2515 
2516 const auto params = context.getParams();
2517 
2518 JSG_REQUIRE(containerStarted.load(std::memory_order_acquire), Error,
2519 "snapshotContainer() requires a running container.");
2520 
2521 auto snapshotId = randomUUID(kj::none);
2522 auto imageRef = kj::str(CONTAINER_SNAPSHOT_IMAGE_PREFIX, snapshotId);
2523 bool imageCommitted = false;
2524 KJ_DEFER(if (imageCommitted) {
2525 waitUntilTasks.add(
2526 deleteImage(kj::str(imageRef)).catch_([](kj::Exception&&) {}).attach(addRef()));
2527 });
2528 
2529 co_await commitContainer(imageRef);
2530 imageCommitted = true;
2531 
2532 auto image = co_await inspectImage(imageRef);
2533 
2534 auto result = context.getResults().initSnapshot();
2535 result.setId(snapshotId);
2536 result.setSize(image.size);
2537 if (params.hasName() && params.getName().size() > 0) {
2538 result.setName(params.getName());
2539 }
2540 
2541 imageCommitted = false;
2542}
2543 
2544kj::Promise<void> ContainerClient::getTcpPort(GetTcpPortContext context) {
2545 co_await mutationQueue.addBranch();
2546 
2547 const auto params = context.getParams();
2548 uint16_t port = params.getPort();
2549 auto results = context.getResults();
2550 auto dockerPort = kj::heap<DockerPort>(*this, kj::str("127.0.0.1"), port);
2551 results.setPort(kj::mv(dockerPort));
2552 co_return;
2553}
2554 
2555kj::Promise<void> ContainerClient::listenTcp(ListenTcpContext context) {
2556 KJ_UNIMPLEMENTED("listenTcp not implemented for Docker containers - use port mapping instead");
2557}
2558 
2559void ContainerClient::upsertEgressMapping(EgressMapping mapping) {
2560 for (auto& m: egressState->mappings) {
2561 // If the mapping differs in port or protocol, we skip it as it's
2562 // not the same.
2563 if (m.port != mapping.port || m.protocol != mapping.protocol) {
2564 continue;
2565 }
2566 
2567 bool matches = false;
2568 KJ_SWITCH_ONEOF(m.destination) {
2569 KJ_CASE_ONEOF(existingCidr, kj::CidrRange) {
2570 KJ_IF_SOME(newCidr, mapping.destination.tryGet<kj::CidrRange>()) {
2571 matches = existingCidr.toString() == newCidr.toString();
2572 }
2573 }
2574 KJ_CASE_ONEOF(existingHostnameGlob, kj::String) {
2575 KJ_IF_SOME(newHostnameGlob, mapping.destination.tryGet<kj::String>()) {
2576 matches = existingHostnameGlob == newHostnameGlob;
2577 }
2578 }
2579 }
2580 
2581 if (matches) {
2582 m.channel = kj::mv(mapping.channel);
2583 return;
2584 }
2585 }
2586 
2587 egressState->mappings.add(kj::mv(mapping));
2588}
2589 
2590kj::Vector<kj::String> ContainerClient::getDnsAllowHostnames() const {
2591 // result N can be at most size of egressState->mappings.
2592 kj::Vector<kj::String> result;
2593 
2594 for (auto& mapping: egressState->mappings) {
2595 KJ_SWITCH_ONEOF(mapping.destination) {
2596 KJ_CASE_ONEOF(_, kj::CidrRange) {
2597 result.add(kj::str("*"));
2598 return result;
2599 }
2600 KJ_CASE_ONEOF(hostnameGlob, kj::String) {
2601 bool alreadyPresent = false;
2602 // Check if we have the hostnameGlob already present in the DNS allow
2603 // list.
2604 for (auto& existing: result) {
2605 if (existing == hostnameGlob) {
2606 alreadyPresent = true;
2607 break;
2608 }
2609 }
2610 
2611 if (!alreadyPresent) {
2612 result.add(kj::str(hostnameGlob));
2613 }
2614 }
2615 }
2616 }
2617 
2618 return result;
2619}
2620 
2621kj::Maybe<kj::Own<workerd::IoChannelFactory::SubrequestChannel>> ContainerClient::findEgressMapping(
2622 kj::StringPtr destAddr,
2623 uint16_t defaultPort,
2624 kj::Maybe<kj::StringPtr> hostname,
2625 EgressProtocol protocol) {
2626 auto hostAndPort = stripPort(destAddr);
2627 uint16_t port = hostAndPort.port.orDefault(defaultPort);
2628 kj::Maybe<kj::String> normalizedHostname;
2629 KJ_IF_SOME(hostnameValue, hostname) {
2630 normalizedHostname = normalizeHostname(hostnameValue);
2631 }
2632 
2633 for (auto& mapping: egressState->mappings) {
2634 // Mappings can differ in port, protocol and the cidr/hostname.
2635 // Users can specify things like google.com:7070, or 0.0.0.0:7070. On top of that,
2636 // they might want TLS interception (HTTPS) or raw TCP forwarding.
2637 if (mapping.protocol != protocol) {
2638 continue;
2639 }
2640 
2641 if (mapping.port != 0 && mapping.port != port) {
2642 continue;
2643 }
2644 
2645 KJ_SWITCH_ONEOF(mapping.destination) {
2646 KJ_CASE_ONEOF(cidr, kj::CidrRange) {
2647 if (cidr.matches(hostAndPort.host)) {
2648 return kj::addRef(*mapping.channel);
2649 }
2650 }
2651 KJ_CASE_ONEOF(hostnameGlob, kj::String) {
2652 KJ_IF_SOME(hostnameValue, normalizedHostname) {
2653 if (hostnameGlobMatches(hostnameGlob, hostnameValue)) {
2654 return kj::addRef(*mapping.channel);
2655 }
2656 }
2657 }
2658 }
2659 }
2660 
2661 return kj::none;
2662}
2663 
2664kj::Promise<void> ContainerClient::ensureSidecarStarted() {
2665 if (containerSidecarStarted.exchange(true, std::memory_order_acquire)) {
2666 co_return;
2667 }
2668 
2669 // We need to call destroy here, it's mandatory that this is a fresh sidecar
2670 // start. Maybe we lost track of it on a previous workerd restart.
2671 co_await destroySidecarContainer();
2672 
2673 KJ_ON_SCOPE_FAILURE(containerSidecarStarted.store(false, std::memory_order_release));
2674 
2675 auto ipamConfig = co_await getDockerBridgeIPAMConfig();
2676 co_await createSidecarContainer(egressListenerPort, kj::mv(ipamConfig.subnet));
2677 co_await startSidecarContainer();
2678 
2679 auto sidecar = KJ_REQUIRE_NONNULL(co_await inspectSidecar(), "started sidecar not running");
2680 this->sidecarIngressHostPort = sidecar.ingressHostPort;
2681 
2682 // Wait for the sidecar's HTTP server to be ready by calling updateSidecarEgressConfig
2683 // in a retry loop with a per-attempt timeout.
2684 constexpr int MAX_READY_RETRIES = 10;
2685 constexpr auto READY_RETRY_DELAY = 200 * kj::MILLISECONDS;
2686 constexpr auto READY_ATTEMPT_TIMEOUT = 2 * kj::SECONDS;
2687 for (int attempt = 0;; ++attempt) {
2688 kj::Maybe<kj::Exception> maybeError;
2689 try {
2690 co_await timer.timeoutAfter(READY_ATTEMPT_TIMEOUT,
2691 updateSidecarEgressConfig(sidecar.ingressHostPort, egressListenerPort));
2692 } catch (...) {
2693 maybeError = kj::getCaughtExceptionAsKj();
2694 }
2695 
2696 if (maybeError == kj::none) break;
2697 if (attempt >= MAX_READY_RETRIES - 1)
2698 kj::throwFatalException(kj::mv(KJ_REQUIRE_NONNULL(maybeError)));
2699 co_await timer.afterDelay(READY_RETRY_DELAY);
2700 }
2701 
2702 co_await readCACert();
2703}
2704 
2705kj::Promise<void> ContainerClient::ensureEgressListenerStarted(uint16_t port) {
2706 if (egressListenerStarted.exchange(true, std::memory_order_acquire)) {
2707 co_return;
2708 }
2709 
2710 KJ_ON_SCOPE_FAILURE(egressListenerStarted.store(false, std::memory_order_release));
2711 
2712 // Determine the listen address: on Linux, use the Docker bridge gateway IP
2713 // and fall back to loopback (Docker Desktop
2714 // routes host-gateway to host loopback through the VM).
2715 auto ipamConfig = co_await getDockerBridgeIPAMConfig();
2716 egressListenerPort = co_await startEgressListener(
2717 gatewayForPlatform(kj::mv(ipamConfig.gateway)).orDefault(kj::str("127.0.0.1")), port);
2718}
2719 
2720kj::Promise<void> ContainerClient::setEgressHttp(SetEgressHttpContext context) {
2721 auto [ready, done] = getRpcTurn();
2722 co_await ready;
2723 KJ_DEFER(done->fulfill());
2724 
2725 auto params = context.getParams();
2726 auto hostPortStr = kj::str(params.getHostPort());
2727 auto tokenBytes = params.getChannelToken();
2728 
2729 auto parsed = parseHostPort(hostPortStr);
2730 uint16_t port = parsed.port.orDefault(80);
2731 
2732 co_await ensureEgressListenerStarted();
2733 
2734 if (containerStarted.load(std::memory_order_acquire)) {
2735 // Only try to create and start a sidecar container
2736 // if the user container is running.
2737 co_await ensureSidecarStarted();
2738 }
2739 
2740 auto subrequestChannel = channelTokenHandler.decodeSubrequestChannelToken(
2741 workerd::IoChannelFactory::ChannelTokenUsage::RPC, tokenBytes);
2742 
2743 upsertEgressMapping(EgressMapping{
2744 .destination = kj::mv(parsed.destination),
2745 .port = port,
2746 .protocol = EgressProtocol::HTTP,
2747 .channel = kj::mv(subrequestChannel),
2748 });
2749 
2750 KJ_IF_SOME(ingressHostPort, sidecarIngressHostPort) {
2751 co_await updateSidecarEgressConfig(ingressHostPort, egressListenerPort);
2752 }
2753 
2754 co_return;
2755}
2756 
2757kj::Promise<void> ContainerClient::setEgressHttps(SetEgressHttpsContext context) {
2758 auto [ready, done] = getRpcTurn();
2759 co_await ready;
2760 KJ_DEFER(done->fulfill());
2761 
2762 auto params = context.getParams();
2763 auto hostPortStr = kj::str(params.getHostPort());
2764 auto tokenBytes = params.getChannelToken();
2765 
2766 auto parsed = parseHostPort(hostPortStr);
2767 uint16_t port = parsed.port.orDefault(443);
2768 
2769 co_await ensureEgressListenerStarted();
2770 
2771 if (containerStarted.load(std::memory_order_acquire)) {
2772 co_await injectCACert();
2773 }
2774 
2775 auto subrequestChannel = channelTokenHandler.decodeSubrequestChannelToken(
2776 workerd::IoChannelFactory::ChannelTokenUsage::RPC, tokenBytes);
2777 
2778 upsertEgressMapping(EgressMapping{
2779 .destination = kj::mv(parsed.destination),
2780 .port = port,
2781 .protocol = EgressProtocol::HTTPS,
2782 .channel = kj::mv(subrequestChannel),
2783 });
2784 
2785 KJ_IF_SOME(ingressHostPort, sidecarIngressHostPort) {
2786 co_await updateSidecarEgressConfig(ingressHostPort, egressListenerPort);
2787 }
2788 
2789 co_return;
2790}
2791 
2792kj::Promise<void> ContainerClient::setEgressTcp(SetEgressTcpContext context) {
2793 auto [ready, done] = getRpcTurn();
2794 co_await ready;
2795 KJ_DEFER(done->fulfill());
2796 
2797 auto params = context.getParams();
2798 auto hostPortStr = kj::str(params.getHostPort());
2799 auto tokenBytes = params.getChannelToken();
2800 
2801 auto parsed = parseHostPort(hostPortStr);
2802 // For TCP, default to port 0 (match all ports) when no port is specified.
2803 uint16_t port = parsed.port.orDefault(0);
2804 
2805 co_await ensureEgressListenerStarted();
2806 
2807 if (containerStarted.load(std::memory_order_acquire)) {
2808 co_await ensureSidecarStarted();
2809 }
2810 
2811 auto subrequestChannel = channelTokenHandler.decodeSubrequestChannelToken(
2812 workerd::IoChannelFactory::ChannelTokenUsage::RPC, tokenBytes);
2813 
2814 upsertEgressMapping(EgressMapping{
2815 .destination = kj::mv(parsed.destination),
2816 .port = port,
2817 .protocol = EgressProtocol::TCP,
2818 .channel = kj::mv(subrequestChannel),
2819 });
2820 
2821 KJ_IF_SOME(ingressHostPort, sidecarIngressHostPort) {
2822 co_await updateSidecarEgressConfig(ingressHostPort, egressListenerPort);
2823 }
2824 
2825 co_return;
2826}
2827 
2828kj::Own<ContainerClient> ContainerClient::addRef() {
2829 return kj::addRef(*this);
2830}
2831 
2832} // namespace workerd::server