// Copyright (c) 2025 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #include "container-client.h" #include "ada.h" #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include namespace workerd::server { namespace { constexpr uint16_t SIDECAR_INGRESS_PORT = 39001; constexpr kj::StringPtr SIDECAR_DNS_SERVERS[] = { "1.1.1.1"_kj, "8.8.8.8"_kj, }; // Default limit for JSON API responses (16 MiB — Docker JSON responses are small). constexpr uint64_t MAX_JSON_RESPONSE_SIZE = 16ULL * 1024 * 1024; constexpr kj::StringPtr SNAPSHOT_VOLUME_PREFIX = "workerd-snap-"_kj; constexpr kj::StringPtr SNAPSHOT_CLONE_VOLUME_PREFIX = "workerd-snap-clone-"_kj; constexpr kj::StringPtr CONTAINER_SNAPSHOT_IMAGE_PREFIX = "workerd-container-snap-"_kj; constexpr kj::StringPtr SNAPSHOT_VOLUME_CREATED_AT_LABEL = "dev.workerd.snapshot-created-at"_kj; // Prefix applied to user-supplied labels when writing them to the Docker container, and // stripped back out when reading them via inspect(). Lets us distinguish labels the worker // set via start() from labels that came from the image (via Dockerfile LABEL) or engine. constexpr kj::StringPtr WORKERD_LABEL_PREFIX = "workerd-"_kj; constexpr auto SNAPSHOT_STALE_AGE = 30 * kj::DAYS; // Maximum size of a snapshot tar archive held in memory during snapshot create/restore. constexpr size_t MAX_SNAPSHOT_TAR_SIZE = 1ULL * 1024 * 1024 * 1024; // 1 GiB static_assert(static_cast(MAX_SNAPSHOT_TAR_SIZE) == MAX_SNAPSHOT_TAR_SIZE, "MAX_SNAPSHOT_TAR_SIZE must be exactly representable as double"); // POSIX tar stores file size in an 11-digit octal header field. constexpr size_t MAX_TAR_CONTENT_SIZE = 8ull * 1024 * 1024 * 1024; // Ensures the stale-volume check runs at most once per process. std::atomic_bool staleSnapshotVolumeCheckScheduled = false; struct ParsedAddress { kj::OneOf destination; kj::Maybe port; }; struct HostAndPort { kj::String host; kj::Maybe port; }; struct DockerResponse { kj::uint statusCode; kj::String body; }; struct DockerBinaryResponse { kj::uint statusCode; kj::Array body; }; struct DockerStreamedResponse { kj::uint statusCode; kj::String statusText; kj::Own connection; }; // Validates an absolute path for snapshot use and returns the parsed component path. // Rejects relative paths, embedded null bytes, and path traversal components (".."). kj::Path parseAbsolutePath(kj::StringPtr path) { JSG_REQUIRE( path.size() > 0 && path[0] == '/', Error, "Snapshot path must be absolute, got: ", path); JSG_REQUIRE(path.findFirst('\0') == kj::none, Error, "Snapshot path must not contain null bytes"); try { return kj::Path::parse(path.slice(1)); } catch (kj::Exception& e) { JSG_FAIL_REQUIRE( Error, "Snapshot path contains invalid components: ", path, "; ", e.getDescription()); } } // Parse and validate a snapshot ID. Throws an error if the snapshot ID is invalid. kj::String parseSnapshotId(kj::StringPtr snapshotId) { KJ_IF_SOME(uuid, UUID::fromString(snapshotId)) { auto s = uuid.toString(); JSG_REQUIRE(s == snapshotId, Error, "Invalid snapshot ID", snapshotId); return s; } else { JSG_FAIL_REQUIRE(Error, "Invalid snapshot ID", snapshotId); } } // Really similar to BufferedInputStreamWrapper, but Async... // We need this because of Docker's exec keeping a bidirectional connection // needing to own the IoStream after writing and reading headers, as it does // "Upgrade: tcp". class BufferedAsyncIoStream final: public kj::AsyncIoStream { public: BufferedAsyncIoStream(kj::Own inner, kj::Array buffered) : inner(kj::mv(inner)), buffered(kj::mv(buffered)) {} kj::Promise tryRead(void* dst, size_t minBytes, size_t maxBytes) override { KJ_REQUIRE(minBytes <= maxBytes, minBytes, maxBytes); auto out = kj::arrayPtr(reinterpret_cast(dst), maxBytes); size_t copied = 0; auto bufferedRemaining = buffered.size() - bufferedOffset; if (bufferedRemaining > 0) { auto toCopy = kj::min(maxBytes, bufferedRemaining); out.first(toCopy).copyFrom(buffered.asPtr().slice(bufferedOffset, bufferedOffset + toCopy)); bufferedOffset += toCopy; copied = toCopy; if (copied >= minBytes || copied == maxBytes) { co_return copied; } } auto read = co_await inner->tryRead(out.begin() + copied, minBytes - copied, maxBytes - copied); co_return copied + read; } kj::Maybe tryGetLength() override { KJ_IF_SOME(innerLength, inner->tryGetLength()) { return innerLength + (buffered.size() - bufferedOffset); } return kj::none; } kj::Promise pumpTo(kj::AsyncOutputStream& output, uint64_t amount) override { uint64_t pumped = 0; auto bufferedRemaining = buffered.size() - bufferedOffset; if (bufferedRemaining > 0) { auto toWrite = static_cast(kj::min(amount, static_cast(bufferedRemaining))); co_await output.write(buffered.asPtr().slice(bufferedOffset, bufferedOffset + toWrite)); bufferedOffset += toWrite; pumped += toWrite; if (pumped == amount) { co_return pumped; } } co_return pumped + co_await inner->pumpTo(output, amount - pumped); } kj::Promise write(kj::ArrayPtr buffer) override { return inner->write(buffer); } kj::Promise write(kj::ArrayPtr> pieces) override { return inner->write(pieces); } kj::Maybe> tryPumpFrom( kj::AsyncInputStream& input, uint64_t amount = kj::maxValue) override { return inner->tryPumpFrom(input, amount); } kj::Promise whenWriteDisconnected() override { return inner->whenWriteDisconnected(); } void abortWrite(kj::Exception&& exception) override { inner->abortWrite(kj::mv(exception)); } void shutdownWrite() override { inner->shutdownWrite(); } void abortRead() override { inner->abortRead(); } void getsockopt(int level, int option, void* value, kj::uint* length) override { inner->getsockopt(level, option, value, length); } void setsockopt(int level, int option, const void* value, kj::uint length) override { inner->setsockopt(level, option, value, length); } void getsockname(struct sockaddr* addr, kj::uint* length) override { inner->getsockname(addr, length); } void getpeername(struct sockaddr* addr, kj::uint* length) override { inner->getpeername(addr, length); } kj::Maybe getFd() const override { return inner->getFd(); } private: kj::Own inner; kj::Array buffered; size_t bufferedOffset = 0; }; // Docker exec uses a single hijacked stream for stdin and stdout/stderr. Keep that stream in a // small refcounted holder so the returned stdin ByteStream and the output demux task can share it. class SharedExecConnection final: public kj::Refcounted { public: explicit SharedExecConnection(kj::Own connection) : connection(kj::mv(connection)) {} kj::Own connection; bool stdinOpened = false; bool stdinClosed = false; }; class DockerExecStdinStream final: public capnp::ExplicitEndOutputStream { public: explicit DockerExecStdinStream(kj::Own sharedConnection) : sharedConnection(kj::mv(sharedConnection)) {} kj::Promise write(kj::ArrayPtr buffer) override { return sharedConnection->connection->write(buffer); } kj::Promise write(kj::ArrayPtr> pieces) override { return sharedConnection->connection->write(pieces); } kj::Promise whenWriteDisconnected() override { return sharedConnection->connection->whenWriteDisconnected(); } kj::Promise end() override { if (!sharedConnection->stdinClosed) { sharedConnection->connection->shutdownWrite(); sharedConnection->stdinClosed = true; } return kj::READY_NOW; } private: kj::Own sharedConnection; }; // Strips a port suffix from a string, returning the host and port separately. // For IPv6, expects brackets: "[::1]:8080" -> ("::1", 8080) // For IPv4: "10.0.0.1:8080" -> ("10.0.0.1", 8080) // If no port, returns the host as-is with no port. HostAndPort stripPort(kj::StringPtr str) { if (str.startsWith("[")) { // Bracketed IPv6: "[ipv6]" or "[ipv6]:port" size_t closeBracket = KJ_REQUIRE_NONNULL(str.findLast(']'), "Unclosed '[' in address string.", str); auto host = str.slice(1, closeBracket); if (str.size() > closeBracket + 1) { KJ_REQUIRE( str.slice(closeBracket + 1).startsWith(":"), "Expected port suffix after ']'.", str); auto port = KJ_REQUIRE_NONNULL( str.slice(closeBracket + 2).tryParseAs(), "Invalid port number.", str); return {kj::str(host), port}; } return {kj::str(host), kj::none}; } // No brackets - check if there's exactly one colon (IPv4 with port) // IPv6 without brackets has 2+ colons and no port suffix supported KJ_IF_SOME(colonPos, str.findLast(':')) { auto afterColon = str.slice(colonPos + 1); KJ_IF_SOME(port, afterColon.tryParseAs()) { // Valid port - but only treat as port for IPv4 (check no other colons before) auto beforeColon = str.first(colonPos); if (beforeColon.findFirst(':') == kj::none) { return {kj::str(beforeColon), port}; } } } return {kj::str(str), kj::none}; } // Build a CidrRange from a host string, adding /32 or /128 prefix if not present. kj::CidrRange makeCidr(kj::StringPtr host) { if (host.findFirst('/') != kj::none) { return kj::CidrRange(host); } // No CIDR prefix - add /32 for IPv4, /128 for IPv6 bool isIpv6 = host.findFirst(':') != kj::none; return kj::CidrRange(kj::str(host, isIpv6 ? "/128" : "/32")); } kj::Maybe tryMakeCidr(kj::StringPtr host) { kj::Maybe cidr; KJ_IF_SOME(_, kj::runCatchingExceptions([&]() { cidr = makeCidr(host); })) { return kj::none; } return kj::mv(cidr); } // normalizeHostname normalizes the hostname. It's designed to receive the hostname when // proxy-everything sends the HTTP CONNECT with the X-Hostname hint. kj::String normalizeHostname(kj::StringPtr hostname) { auto url = kj::str("http://", hostname); auto parsed = ada::parse({url.begin(), url.size()}, nullptr); KJ_REQUIRE(parsed.has_value(), "Invalid X-Hostname URL hint.", hostname); auto normalizedHostname = parsed->get_hostname(); return kj::heapString(normalizedHostname.data(), normalizedHostname.size()); } // hostnameGlobMatches should match patterns like: // cloudflare.*.com // cloudflare.com // cloudflare // * // // hostname must be normalized beforehand bool hostnameGlobMatches(kj::StringPtr pattern, kj::StringPtr hostname) { size_t patternIndex = 0; size_t hostnameIndex = 0; size_t restartHostnameIndex = 0; kj::Maybe starPatternIndex; while (hostnameIndex < hostname.size()) { if (patternIndex < pattern.size() && pattern[patternIndex] == '*') { starPatternIndex = patternIndex++; restartHostnameIndex = hostnameIndex; continue; } if (patternIndex < pattern.size() && pattern[patternIndex] == hostname[hostnameIndex]) { ++patternIndex; ++hostnameIndex; continue; } KJ_IF_SOME(starIndex, starPatternIndex) { patternIndex = starIndex + 1; hostnameIndex = ++restartHostnameIndex; continue; } return false; } while (patternIndex < pattern.size() && pattern[patternIndex] == '*') { ++patternIndex; } return patternIndex == pattern.size(); } kj::Maybe getHeader(const kj::HttpHeaders& headers, kj::StringPtr name) { kj::Maybe result; headers.forEach([&](kj::StringPtr headerName, kj::StringPtr value) { if (result == kj::none && workerd::strcaseeq(headerName, name)) { result = value; } }); return result; } // Parses "host[:port]" strings. Handles: // - IPv4: "10.0.0.1", "10.0.0.1:8080", "10.0.0.0/8", "10.0.0.0/8:8080" // - IPv6 with brackets: "[::1]", "[::1]:8080", "[fe80::1]", "[fe80::/10]:8080" // - IPv6 without brackets: "::1", "fe80::1", "fe80::/10" ParsedAddress parseHostPort(kj::StringPtr str) { auto hostAndPort = stripPort(str); KJ_REQUIRE(hostAndPort.host.size() > 0, "Host must not be empty.", str); KJ_IF_SOME(cidr, tryMakeCidr(hostAndPort.host)) { return { .destination = kj::mv(cidr), .port = hostAndPort.port, }; } return { .destination = workerd::toLower(hostAndPort.host), .port = hostAndPort.port, }; } kj::StringPtr signalToString(uint32_t signal) { switch (signal) { case 1: return "SIGHUP"_kj; // Hangup case 2: return "SIGINT"_kj; // Interrupt case 3: return "SIGQUIT"_kj; // Quit case 4: return "SIGILL"_kj; // Illegal instruction case 5: return "SIGTRAP"_kj; // Trace trap case 6: return "SIGABRT"_kj; // Abort case 7: return "SIGBUS"_kj; // Bus error case 8: return "SIGFPE"_kj; // Floating point exception case 9: return "SIGKILL"_kj; // Kill case 10: return "SIGUSR1"_kj; // User signal 1 case 11: return "SIGSEGV"_kj; // Segmentation violation case 12: return "SIGUSR2"_kj; // User signal 2 case 13: return "SIGPIPE"_kj; // Broken pipe case 14: return "SIGALRM"_kj; // Alarm clock case 15: return "SIGTERM"_kj; // Termination case 16: return "SIGSTKFLT"_kj; // Stack fault (Linux) case 17: return "SIGCHLD"_kj; // Child status changed case 18: return "SIGCONT"_kj; // Continue case 19: return "SIGSTOP"_kj; // Stop case 20: return "SIGTSTP"_kj; // Terminal stop case 21: return "SIGTTIN"_kj; // Background read from tty case 22: return "SIGTTOU"_kj; // Background write to tty case 23: return "SIGURG"_kj; // Urgent condition on socket case 24: return "SIGXCPU"_kj; // CPU limit exceeded case 25: return "SIGXFSZ"_kj; // File size limit exceeded case 26: return "SIGVTALRM"_kj; // Virtual alarm clock case 27: return "SIGPROF"_kj; // Profiling alarm clock case 28: return "SIGWINCH"_kj; // Window size change case 29: return "SIGIO"_kj; // I/O now possible case 30: return "SIGPWR"_kj; // Power failure restart (Linux) case 31: return "SIGSYS"_kj; // Bad system call default: return "SIGKILL"_kj; } } void writeTarField(kj::ArrayPtr field, kj::StringPtr value) { auto len = kj::min(value.size(), field.size()); field.first(len).copyFrom(value.asBytes().first(len)); } // createTarWithFile creates simple tar files without importing a full blown TAR library. // It's a pretty limited method that creates a single tar file with a single file on it, // as the Docker API only accepts tars. kj::Array createTarWithFile( kj::StringPtr filename, kj::ArrayPtr content) { KJ_REQUIRE(filename.size() < 100, "tar filename must be < 100 bytes"); KJ_REQUIRE(content.size() < MAX_TAR_CONTENT_SIZE, "tar content too large for 11-digit octal"); size_t paddedSize = (content.size() + 511) & ~static_cast(511); size_t totalSize = 512 + paddedSize + 1024; auto tar = kj::heapArray(totalSize); tar.asPtr().fill(0); auto header = tar.first(512); writeTarField(header.slice(0, 100), filename); writeTarField(header.slice(100, 108), "0000644"_kj); writeTarField(header.slice(108, 116), "0000000"_kj); writeTarField(header.slice(116, 124), "0000000"_kj); { char sizeBuf[12]; snprintf(sizeBuf, sizeof(sizeBuf), "%011" PRIo64, static_cast(content.size())); writeTarField(header.slice(124, 136), kj::StringPtr(sizeBuf)); } writeTarField(header.slice(136, 148), "00000000000"_kj); header[156] = '0'; writeTarField(header.slice(257, 263), "ustar"_kj); writeTarField(header.slice(263, 265), "00"_kj); header.slice(148, 156).fill(' '); uint32_t checksum = 0; for (auto byte: header) checksum += byte; { char checksumBuf[8]; snprintf(checksumBuf, sizeof(checksumBuf), "%06o ", checksum); writeTarField(header.slice(148, 155), kj::StringPtr(checksumBuf)); } tar.slice(512, 512 + content.size()).copyFrom(content); return tar; } // Shared Docker API HTTP helper. Connects to the Docker socket, sends a request with an // optional body, and reads the response as raw bytes. kj::Promise dockerApiRequestRaw(kj::Network& network, kj::String dockerPath, kj::HttpMethod method, kj::String endpoint, kj::Maybe> body, kj::StringPtr contentType, uint64_t maxResponseSize) { kj::HttpHeaderTable headerTable; auto address = co_await network.parseAddress(dockerPath); auto connection = co_await address->connect(); auto httpClient = kj::newHttpClient(headerTable, *connection).attach(kj::mv(connection)); kj::HttpHeaders headers(headerTable); headers.setPtr(kj::HttpHeaderId::HOST, "localhost"); KJ_IF_SOME(requestBody, body) { headers.setPtr(kj::HttpHeaderId::CONTENT_TYPE, contentType); headers.set(kj::HttpHeaderId::CONTENT_LENGTH, kj::str(requestBody.size())); auto req = httpClient->request(method, endpoint, headers, requestBody.size()); { auto stream = kj::mv(req.body); co_await stream->write(requestBody); } auto response = co_await req.response; auto result = co_await response.body->readAllBytes(maxResponseSize); co_return DockerBinaryResponse{.statusCode = response.statusCode, .body = kj::mv(result)}; } else { auto req = httpClient->request(method, endpoint, headers); { auto stream = kj::mv(req.body); } auto response = co_await req.response; auto result = co_await response.body->readAllBytes(maxResponseSize); co_return DockerBinaryResponse{.statusCode = response.statusCode, .body = kj::mv(result)}; } } kj::Promise dockerApiRequest(kj::Network& network, kj::String dockerPath, kj::HttpMethod method, kj::String endpoint, kj::Maybe body = kj::none) { kj::Maybe> bodyBytes; KJ_IF_SOME(b, body) { bodyBytes = b.asBytes(); } auto raw = co_await dockerApiRequestRaw(network, kj::mv(dockerPath), method, kj::mv(endpoint), bodyBytes, "application/json"_kj, MAX_JSON_RESPONSE_SIZE); co_return DockerResponse{.statusCode = raw.statusCode, .body = kj::str(raw.body.asChars())}; } kj::Promise dockerApiBinaryRequest(kj::Network& network, kj::String dockerPath, kj::HttpMethod method, kj::String endpoint, kj::Maybe> body, uint64_t maxResponseSize) { kj::Maybe> bodyBytes; KJ_IF_SOME(b, body) { bodyBytes = b.asPtr(); } co_return co_await dockerApiRequestRaw(network, kj::mv(dockerPath), method, kj::mv(endpoint), bodyBytes, "application/x-tar"_kj, maxResponseSize); } kj::Promise deleteVolume(kj::Network& network, kj::String dockerPath, kj::String volumeName) { auto response = co_await dockerApiRequest( network, kj::mv(dockerPath), kj::HttpMethod::DELETE, kj::str("/volumes/", volumeName)); if (response.statusCode != 204 && response.statusCode != 404) { KJ_LOG(WARNING, "failed to delete volume", volumeName, response.statusCode, response.body); } } kj::Promise deleteVolumes( kj::Network& network, kj::String dockerPath, kj::Array snapshotCloneVolumes) { kj::Vector> volumeDeletes; volumeDeletes.reserve(snapshotCloneVolumes.size()); for (auto& volumeName: snapshotCloneVolumes) { auto logName = kj::str(volumeName); volumeDeletes.add(deleteVolume(network, kj::str(dockerPath), kj::mv(volumeName)) .catch_([logName = kj::mv(logName)](kj::Exception&& e) { KJ_LOG(WARNING, "failed to delete volume", logName, e); })); } co_await kj::joinPromises(volumeDeletes.releaseAsArray()); } kj::Promise removeContainer( kj::Network& network, kj::String dockerPath, kj::String containerName, bool wait = true) { auto endpoint = kj::str("/containers/", containerName, "?force=true"); auto response = co_await dockerApiRequest( network, kj::str(dockerPath), kj::HttpMethod::DELETE, kj::mv(endpoint)); // 204 means the container was removed. // 404 means it was already gone. // 409 means removal is already in progress, which is fine for our teardown paths. KJ_REQUIRE(response.statusCode == 204 || response.statusCode == 404 || response.statusCode == 409, "Removing a container failed with: ", response.body); // If removal succeeded or is already in progress, wait for Docker to report the container as // fully removed before proceeding with any follow-up cleanup like deleting mounted volumes. if (wait && (response.statusCode == 204 || response.statusCode == 409)) { response = co_await dockerApiRequest(network, kj::mv(dockerPath), kj::HttpMethod::POST, kj::str("/containers/", containerName, "/wait?condition=removed")); // 200 means Docker observed the removal. 404 means the container disappeared before the wait // request was processed, which is also fine. KJ_REQUIRE(response.statusCode == 200 || response.statusCode == 404, "Waiting for container removal failed with: ", response.statusCode, response.body); } } kj::Maybe tryFindHttpHeaderEnd(kj::ArrayPtr bytes) { for (auto i: kj::zeroTo(bytes.size())) { if (i + 4 > bytes.size()) { return kj::none; } if (bytes[i] == '\r' && bytes[i + 1] == '\n' && bytes[i + 2] == '\r' && bytes[i + 3] == '\n') { return i; } } return kj::none; } // readDockerStreamedResponse is necessary because Docker streamed responses // require an open bidirectional stream after the response headers have been read. kj::Promise readDockerStreamedResponse( kj::Own connection) { kj::Vector buffer; auto& input = *connection; while (true) { KJ_IF_SOME(headerEnd, tryFindHttpHeaderEnd(buffer.asPtr())) { auto parsedHeaders = kj::heapArray(headerEnd + 2); for (auto i: kj::zeroTo(parsedHeaders.size())) { parsedHeaders[i] = static_cast(buffer[i]); } kj::HttpHeaderTable headerTable; kj::HttpHeaders headers(headerTable); auto parsedResponse = headers.tryParseResponse(parsedHeaders.asPtr()); headers.takeOwnership(kj::mv(parsedHeaders)); kj::uint statusCode = 0; kj::String statusText; KJ_SWITCH_ONEOF(parsedResponse) { KJ_CASE_ONEOF(response, kj::HttpHeaders::Response) { statusCode = response.statusCode; statusText = kj::str(response.statusText); } KJ_CASE_ONEOF(protocolError, kj::HttpHeaders::ProtocolError) { KJ_FAIL_REQUIRE("Docker streamed response returned malformed HTTP headers: ", protocolError.statusMessage, ": ", protocolError.description); } } auto bodyOffset = headerEnd + 4; auto prefetchedBytes = kj::heapArray(buffer.asPtr().slice(bodyOffset)); kj::Own prefixedConnection = kj::mv(connection); if (prefetchedBytes.size() > 0) { prefixedConnection = kj::heap(kj::mv(prefixedConnection), kj::mv(prefetchedBytes)); } co_return DockerStreamedResponse{ .statusCode = statusCode, .statusText = kj::mv(statusText), .connection = kj::mv(prefixedConnection), }; } auto scratch = kj::heapArray(4096); auto amount = co_await input.tryRead(scratch.begin(), 1, scratch.size()); KJ_REQUIRE(amount > 0, "EOF while waiting for Docker streamed response headers"); buffer.addAll(scratch.first(amount)); KJ_REQUIRE(buffer.size() <= 65536, "Docker streamed response headers exceeded 64KiB"); } } kj::Promise dockerApiStreamedRequest(kj::Network& network, kj::String dockerPath, kj::HttpMethod method, kj::String endpoint, const kj::HttpHeaders& headers, kj::Maybe> body = kj::none) { auto address = co_await network.parseAddress(dockerPath); auto connection = co_await address->connect(); auto requestHeaders = headers.serializeRequest(method, endpoint); KJ_IF_SOME(requestBody, body) { kj::ArrayPtr pieces[] = {requestHeaders.asBytes(), requestBody}; co_await connection->write(kj::arrayPtr(pieces)); } else { co_await connection->write(requestHeaders.asBytes()); } co_return co_await readDockerStreamedResponse(kj::mv(connection)); } // Docker multiplexed stream frames: 1 byte stream ID + 3 reserved + 4 bytes big-endian length. constexpr size_t DOCKER_FRAME_HEADER_SIZE = 8; uint32_t parseDockerFrameLength(kj::ArrayPtr frameHeader) { KJ_REQUIRE(frameHeader.size() >= DOCKER_FRAME_HEADER_SIZE, "Docker raw stream header too short"); return (static_cast(frameHeader[4]) << 24) | (static_cast(frameHeader[5]) << 16) | (static_cast(frameHeader[6]) << 8) | static_cast(frameHeader[7]); } void detachEnd(kj::Maybe> stream) { KJ_IF_SOME(s, stream) { s->end().attach(kj::mv(s)).detach([](kj::Exception&&) {}); } } // demuxDockerExecOutput demuxes the input from Docker to passed stdout/stderr. kj::Promise demuxDockerExecOutput(kj::AsyncInputStream& input, kj::Maybe> stdoutWriter, kj::Maybe> stderrWriter, bool combinedOutput) { kj::Vector buffer; size_t offset = 0; auto compactBuffer = [&]() { if (offset == 0) { return; } kj::Vector compacted; compacted.addAll(buffer.asPtr().slice(offset)); buffer = kj::mv(compacted); offset = 0; }; auto ensureBytes = [&](size_t count) -> kj::Promise { while (buffer.size() - offset < count) { compactBuffer(); auto scratch = kj::heapArray(4096); auto amount = co_await input.tryRead(scratch.begin(), 1, scratch.size()); if (amount == 0) { co_return false; } buffer.addAll(scratch.first(amount)); } co_return true; }; try { while (co_await ensureBytes(DOCKER_FRAME_HEADER_SIZE)) { auto frameHeader = buffer.asPtr().slice(offset, offset + DOCKER_FRAME_HEADER_SIZE); auto streamId = frameHeader[0]; auto frameLength = parseDockerFrameLength(frameHeader); KJ_REQUIRE(co_await ensureBytes(DOCKER_FRAME_HEADER_SIZE + frameLength), "Docker exec raw stream ended in the middle of a frame"); auto payload = buffer.asPtr().slice( offset + DOCKER_FRAME_HEADER_SIZE, offset + DOCKER_FRAME_HEADER_SIZE + frameLength); if (streamId == 1) { KJ_IF_SOME(out, stdoutWriter) { co_await out->write(payload); } } else { if (streamId == 2) { if (combinedOutput) { KJ_IF_SOME(out, stdoutWriter) { co_await out->write(payload); } } else { KJ_IF_SOME(err, stderrWriter) { co_await err->write(payload); } } } } offset += DOCKER_FRAME_HEADER_SIZE + frameLength; } if (buffer.size() != offset) { KJ_FAIL_REQUIRE("Docker exec raw stream ended with a truncated frame header"); } // We need to detach ourselves from the end() as the user might've // decided to not read them altogether. detachEnd(kj::mv(stdoutWriter)); detachEnd(kj::mv(stderrWriter)); } catch (...) { auto exception = kj::getCaughtExceptionAsKj(); KJ_IF_SOME(out, stdoutWriter) { out->abortWrite(exception.clone()); } KJ_IF_SOME(err, stderrWriter) { err->abortWrite(exception.clone()); } kj::throwFatalException(kj::mv(exception)); } } kj::String currentSnapshotVolumeTimestamp() { return kj::str((kj::systemPreciseCalendarClock().now() - kj::UNIX_EPOCH) / kj::SECONDS); } kj::Maybe tryGetSnapshotCreatedAt(capnp::JsonValue::Reader labels) { if (!labels.isObject()) { return kj::none; } for (auto field: labels.getObject()) { if (field.getName() != SNAPSHOT_VOLUME_CREATED_AT_LABEL) { continue; } auto value = field.getValue(); if (!value.isString()) { return kj::none; } return value.getString().tryParseAs(); } return kj::none; } kj::Promise warnAboutStaleSnapshotVolumes(kj::Network& network, kj::String dockerPath) { capnp::JsonCodec codec; codec.handleByAnnotation(); capnp::MallocMessageBuilder filterMessage; auto filters = filterMessage.initRoot(); auto names = filters.initName(1); names.set(0, SNAPSHOT_VOLUME_PREFIX); auto response = co_await dockerApiRequest(network, kj::mv(dockerPath), kj::HttpMethod::GET, kj::str("/volumes?filters=", kj::encodeUriComponent(codec.encode(filters)))); if (response.statusCode != 200) { co_return; } auto message = decodeJsonResponse(response.body); auto root = message->getRoot(); auto now = kj::systemPreciseCalendarClock().now(); kj::Vector staleVolumes; for (auto volume: root.getVolumes()) { KJ_IF_SOME(createdAtSeconds, tryGetSnapshotCreatedAt(volume.getLabels())) { auto createdAt = kj::UNIX_EPOCH + createdAtSeconds * kj::SECONDS; if (now - createdAt >= SNAPSHOT_STALE_AGE) { staleVolumes.add(kj::str(volume.getName())); } } } if (!staleVolumes.empty()) { KJ_LOG(WARNING, "the following snapshot volumes were created 30+ days ago and may be stale", kj::strArray(staleVolumes, ", ")); } } // Returns the gateway IP on Linux for direct container access. // Returns kj::none on macOS where Docker Desktop routes host-gateway to host loopback. kj::Maybe gatewayForPlatform(kj::String gateway) { #ifdef __APPLE__ return kj::none; #else return kj::mv(gateway); #endif } kj::Maybe tryParsePublishedHostPort(capnp::json::Value::Reader portMappingValue) { if (portMappingValue.isNull()) { return kj::none; } JSG_REQUIRE( portMappingValue.isArray(), Error, "Malformed ContainerInspect port mapping response"); auto bindings = portMappingValue.getArray(); if (bindings.size() == 0) { return kj::none; } auto binding = bindings[0]; JSG_REQUIRE(binding.isObject(), Error, "Malformed ContainerInspect port binding response"); for (auto field: binding.getObject()) { if (field.getName() == "HostPort") { auto value = field.getValue(); JSG_REQUIRE(value.isString(), Error, "Malformed ContainerInspect port binding response"); kj::StringPtr hostPort = value.getString(); return KJ_REQUIRE_NONNULL( hostPort.tryParseAs(), "Malformed ContainerInspect host port"); } } KJ_FAIL_REQUIRE("Malformed ContainerInspect port binding response: missing HostPort"); } } // namespace // Represents a parsed egress mapping. IP/CIDR mappings match destination IPs, // while hostnameGlob mappings match either HTTP hostnames or TLS SNI depending on protocol. // Defined here (not in the header) to avoid pulling kj::OneOf, kj::CidrRange, and // kj::Vector into server.c++ which includes container-client.h. struct ContainerClient::EgressMapping { kj::OneOf destination; uint16_t port; // 0 means match all ports EgressProtocol protocol; kj::Own channel; }; // Holds all egress mapping state. Stored via kj::Own in ContainerClient // so that the EgressMapping type is not visible in container-client.h. struct ContainerClient::EgressState { kj::Vector mappings; }; ContainerClient::ContainerClient(capnp::ByteStreamFactory& byteStreamFactory, kj::Timer& timer, kj::Network& network, kj::String dockerPath, kj::String containerName, kj::String imageName, kj::String containerEgressInterceptorImage, kj::TaskSet& waitUntilTasks, kj::Promise pendingCleanup, kj::Function)> cleanupCallback, ChannelTokenHandler& channelTokenHandler) : byteStreamFactory(byteStreamFactory), timer(timer), network(network), dockerPath(kj::mv(dockerPath)), containerName(kj::encodeUriComponent(kj::str(containerName))), sidecarContainerName(kj::encodeUriComponent(kj::str(containerName, "-proxy"))), imageName(kj::mv(imageName)), containerEgressInterceptorImage(kj::mv(containerEgressInterceptorImage)), waitUntilTasks(waitUntilTasks), pendingCleanup(kj::mv(pendingCleanup).fork()), cleanupCallback(kj::mv(cleanupCallback)), channelTokenHandler(channelTokenHandler), egressState(kj::heap()) { if (!staleSnapshotVolumeCheckScheduled.exchange(true)) { waitUntilTasks.add(warnAboutStaleSnapshotVolumes(network, kj::str(this->dockerPath)) .catch_([](kj::Exception&& e) { KJ_LOG(WARNING, "failed to inspect snapshot volumes for staleness", e); })); } } ContainerClient::~ContainerClient() noexcept(false) { stopEgressListener(); // Best-effort cleanup for both containers. auto sidecarCleanup = removeContainer(network, kj::str(dockerPath), kj::str(sidecarContainerName), false) .catch_([](kj::Exception&&) {}); // Also try to delete any cloned snapshot volumes. auto volumes = snapshotClones.releaseAsArray(); auto mainCleanup = removeContainer(network, kj::str(dockerPath), kj::str(containerName)) .catch_([](kj::Exception&&) {}) .then([&network = network, dockerPath = kj::str(dockerPath), volumes = kj::mv(volumes)]() mutable { return deleteVolumes(network, kj::mv(dockerPath), kj::mv(volumes)); }).catch_([](kj::Exception&&) {}); // Pass the joined cleanup promise to the callback. The callback wraps it with the // canceler (so a future client creation can cancel it), stores it so the next // ContainerClient can await it, and adds a branch to waitUntilTasks to keep the // underlying I/O alive. cleanupCallback(kj::joinPromises(kj::arr(kj::mv(sidecarCleanup), kj::mv(mainCleanup)))); } // Docker-specific Port implementation that implements rpc::Container::Port::Server // It does a HTTP CONNECT to the proxy-everything sidecar port. class ContainerClient::DockerPort final: public rpc::Container::Port::Server { public: DockerPort(ContainerClient& containerClient, kj::String containerHost, uint16_t containerPort) : containerClient(containerClient), containerHost(kj::mv(containerHost)), containerPort(containerPort) {} kj::Promise connect(ConnectContext context) override { auto mappedPort = JSG_REQUIRE_NONNULL(containerClient.sidecarIngressHostPort, Error, "connect(): Container ingress proxy is not running."); auto dstAddr = kj::str(containerHost, ":", containerPort); auto address = co_await containerClient.network.parseAddress(kj::str("127.0.0.1:", mappedPort)); kj::HttpHeaderTable::Builder headerTableBuilder; auto xDstAddrHeader = headerTableBuilder.add("X-Dst-Addr"); auto headerTable = headerTableBuilder.build(); kj::HttpHeaders headers(*headerTable); headers.set(xDstAddrHeader, kj::str(dstAddr)); auto proxyConnection = co_await address->connect(); auto httpClient = kj::newHttpClient(*headerTable, *proxyConnection) .attach(kj::mv(proxyConnection), kj::mv(headerTable)); auto connectRequest = httpClient->connect(dstAddr, headers, {}); auto status = co_await kj::mv(connectRequest.status); if (status.statusCode == 400) { throw JSG_KJ_EXCEPTION( DISCONNECTED, Error, "Container is not listening to port ", containerPort); } if (status.statusCode < 200 || status.statusCode >= 300) { KJ_IF_SOME(errorBody, status.errorBody) { auto errorBodyText = co_await errorBody->readAllText(); JSG_FAIL_REQUIRE(Error, "Connecting to container port through proxy-everything failed: [", status.statusCode, "] ", status.statusText, " ", errorBodyText); } JSG_FAIL_REQUIRE(Error, "Connecting to container port through proxy-everything failed: [", status.statusCode, "] ", status.statusText); } auto connection = kj::mv(connectRequest.connection); auto upPipe = kj::newOneWayPipe(); auto upEnd = kj::mv(upPipe.in); auto results = context.getResults(); results.setUp(containerClient.byteStreamFactory.kjToCapnp(kj::mv(upPipe.out))); auto downEnd = containerClient.byteStreamFactory.capnpToKj(context.getParams().getDown()); pumpTask = kj::joinPromisesFailFast(kj::arr(upEnd->pumpTo(*connection), connection->pumpTo(*downEnd))) .ignoreResult() .attach(kj::mv(httpClient), kj::mv(upEnd), kj::mv(connection), kj::mv(downEnd)); co_return; } private: // ContainerClient is owned by the Worker::Actor and keeps it alive. ContainerClient& containerClient; kj::String containerHost; uint16_t containerPort; kj::Maybe> pumpTask; }; class ContainerClient::DockerProcessHandle final: public rpc::Container::ProcessHandle::Server { public: DockerProcessHandle(ContainerClient& containerClient, kj::String execId, kj::Own connection, kj::Maybe stdoutWriter, kj::Maybe stderrWriter, bool combinedOutput) : containerClient(containerClient.addRef()), execId(kj::mv(execId)), sharedConnection(kj::refcounted(kj::mv(connection))) { kj::Maybe> stdoutStream = kj::none; KJ_IF_SOME(out, stdoutWriter) { stdoutStream = this->containerClient->byteStreamFactory.capnpToKjExplicitEnd(out); } else { stdoutStream = capnp::ExplicitEndOutputStream::wrap(newNullOutputStream(), []() {}); } kj::Maybe> stderrStream = kj::none; KJ_IF_SOME(err, stderrWriter) { stderrStream = this->containerClient->byteStreamFactory.capnpToKjExplicitEnd(err); } else if (!combinedOutput) { stderrStream = capnp::ExplicitEndOutputStream::wrap(newNullOutputStream(), []() {}); } // Always drain the Docker exec stream. This lets wait() use stream closure as the primary // process-completion signal, even when stdout/stderr are ignored. auto task = demuxDockerExecOutput( *sharedConnection->connection, kj::mv(stdoutStream), kj::mv(stderrStream), combinedOutput) .attach(this->containerClient->addRef(), kj::addRef(*sharedConnection)); streamClosedTask = kj::mv(task).fork(); } kj::Promise wait(WaitContext context) override { waitStarted = true; if (!sharedConnection->stdinOpened && !sharedConnection->stdinClosed) { sharedConnection->connection->shutdownWrite(); sharedConnection->stdinClosed = true; } co_await KJ_ASSERT_NONNULL(streamClosedTask).addBranch(); // Docker's exec-inspect state can lag slightly behind the hijacked stream closing, so after // we observe EOF we allow a short bounded retry window to obtain the final exit code. for (auto attempt: kj::zeroTo(20)) { auto inspect = co_await containerClient->inspectExec(execId); if (!inspect.running) { context.getResults().setExitCode(inspect.exitCode); co_return; } if (attempt + 1 < 20) { co_await containerClient->timer.afterDelay(50 * kj::MILLISECONDS); } } JSG_FAIL_REQUIRE(Error, "Docker exec stream closed before exit status became available."); } kj::Promise stdinWriter(StdinWriterContext context) override { JSG_REQUIRE(!waitStarted, Error, "Process stdinWriter() cannot be called after wait()."); JSG_REQUIRE( !sharedConnection->stdinOpened, Error, "Process stdinWriter() can only be called once."); sharedConnection->stdinOpened = true; context.getResults().setWriter(containerClient->byteStreamFactory.kjToCapnp( kj::heap(kj::addRef(*sharedConnection)))); co_return; } kj::Promise kill(KillContext context) override { auto inspect = co_await containerClient->inspectExec(execId); JSG_REQUIRE(inspect.pid > 0, Error, "Exec process does not have a visible pid to signal."); auto signal = kj::str("-", signalToString(context.getParams().getSigno())); auto pid = kj::str(inspect.pid); auto cmd = kj::arr(kj::str("kill"), kj::mv(signal), kj::mv(pid)); co_await containerClient->runSimpleExec(cmd.asPtr()); } private: kj::Own containerClient; kj::String execId; kj::Own sharedConnection; bool waitStarted = false; kj::Maybe> streamClosedTask; }; // ConnectResponse adapter for TCP egress. Since we've already accepted the sidecar's HTTP // CONNECT before calling worker->connect(), this adapter simply records the worker's // accept/reject decision without sending anything on the wire. class TcpEgressConnectResponse final: public kj::HttpService::ConnectResponse { public: void accept(uint statusCode, kj::StringPtr statusText, const kj::HttpHeaders& headers) override { // Worker accepted the connection. Nothing additional to do since we already // accepted the sidecar CONNECT. } kj::Own reject(uint statusCode, kj::StringPtr statusText, const kj::HttpHeaders& headers, kj::Maybe expectedBodySize = kj::none) override { KJ_FAIL_REQUIRE("TCP egress worker rejected the connection: ", statusCode, " ", statusText); } }; // HTTP service that handles HTTP CONNECT requests from the container sidecar (proxy-everything). // When the sidecar intercepts container egress traffic, it sends HTTP CONNECT to this service. // After accepting the CONNECT, the tunnel carries the actual HTTP request from the container, // which we parse and forward to the appropriate SubrequestChannel based on egressMappings. // Inner HTTP service that handles requests inside the CONNECT tunnel. // Forwards requests to the worker binding via SubrequestChannel. class InnerEgressService final: public kj::HttpService { public: using ChannelLookup = kj::Function>()>; InnerEgressService(ChannelLookup lookupChannel, kj::StringPtr destAddr, bool isTls = false) : lookupChannel(kj::mv(lookupChannel)), destAddr(kj::str(destAddr)), isTls(isTls) {} kj::Promise request(kj::HttpMethod method, kj::StringPtr requestUri, const kj::HttpHeaders& headers, kj::AsyncInputStream& requestBody, Response& response) override { // Look up the channel on each request so we always use the latest mapping, // even if it was replaced via interceptOutboundHttp while the tunnel is open. auto channel = KJ_REQUIRE_NONNULL(lookupChannel(), "egress mapping disappeared during active tunnel"); IoChannelFactory::SubrequestMetadata metadata; auto worker = channel->startRequest(kj::mv(metadata)); auto urlForWorker = kj::str(requestUri); // Probably only a path, try to get it from Host: if (requestUri.startsWith("/")) { auto scheme = isTls ? "https://"_kj : "http://"_kj; auto baseUrl = kj::str(scheme, destAddr); // Use Host: when possible KJ_IF_SOME(host, headers.get(kj::HttpHeaderId::HOST)) { baseUrl = kj::str(scheme, host); } // Parse url, if invalid, try to use the original requestUri (http:/// KJ_IF_SOME(parsedUrl, jsg::Url::tryParse(requestUri, baseUrl.asPtr())) { urlForWorker = kj::str(parsedUrl.getHref()); } else { urlForWorker = kj::str(baseUrl, requestUri); } } co_await worker->request(method, urlForWorker, headers, requestBody, response); } private: ChannelLookup lookupChannel; kj::String destAddr; bool isTls; }; kj::Promise pumpBidirectional(kj::AsyncIoStream& a, kj::AsyncIoStream& b) { auto aToB = a.pumpTo(b).then([&b](uint64_t) { b.shutdownWrite(); }); auto bToA = b.pumpTo(a).then([&a](uint64_t) { a.shutdownWrite(); }); co_await kj::joinPromisesFailFast(kj::arr(kj::mv(aToB), kj::mv(bToA))); } // Outer HTTP service that handles CONNECT requests from the sidecar. class EgressHttpService final: public kj::HttpService { public: EgressHttpService(ContainerClient& containerClient, kj::HttpHeaderTable& headerTable) : containerClient(containerClient), headerTable(headerTable) {} kj::Promise request(kj::HttpMethod method, kj::StringPtr url, const kj::HttpHeaders& headers, kj::AsyncInputStream& requestBody, Response& response) override { // Regular HTTP requests are not expected - we only handle CONNECT co_return co_await response.sendError(405, "Method Not Allowed", headerTable); } kj::Promise connect(kj::StringPtr host, const kj::HttpHeaders& headers, kj::AsyncIoStream& connection, ConnectResponse& response, kj::HttpConnectSettings settings) override { auto destAddr = kj::str(host); if (co_await handleConnectMode(destAddr, headers, connection, response, "X-Tls-Sni", /*defaultPort=*/443, EgressProtocol::HTTPS)) { co_return; } if (co_await handleConnectMode(destAddr, headers, connection, response, "X-Hostname", /*defaultPort=*/80, EgressProtocol::HTTP)) { co_return; } // Try raw TCP mapping before falling through to passthrough. TCP mappings match // on IP/CIDR + port without any application-layer hostname information. if (co_await handleTcpConnect(destAddr, connection, response)) { co_return; } kj::HttpHeaders responseHeaders(headerTable); // 202 is interpreted by proxy-everything as "just send bytes as-is". // If the connection was TLS, it's useful so we just proxy transparently // to the internet. response.accept(202, "Accepted", responseHeaders); co_await passThroughConnection(destAddr, connection); } private: kj::Promise handleConnectMode(kj::StringPtr destAddr, const kj::HttpHeaders& headers, kj::AsyncIoStream& connection, ConnectResponse& response, kj::StringPtr hostnameHeader, uint16_t defaultPort, EgressProtocol protocol) { kj::Maybe requestHostname; KJ_IF_SOME(value, getHeader(headers, hostnameHeader)) { requestHostname = kj::str(value); } auto mapping = containerClient.findEgressMapping(destAddr, defaultPort, requestHostname.map([](auto& hostname) { return kj::Maybe(hostname); }).orDefault(kj::none), protocol); if (requestHostname == kj::none && mapping == kj::none) { co_return false; } if (mapping != kj::none) { kj::HttpHeaders responseHeaders(headerTable); response.accept(200, "OK", responseHeaders); bool isTls = (protocol == EgressProtocol::HTTPS); auto innerService = kj::heap( [&client = containerClient, addr = kj::str(destAddr), hostname = requestHostname.map([](auto& value) { return kj::str(value); }), defaultPort, protocol]() mutable -> kj::Maybe> { return client.findEgressMapping(addr, defaultPort, hostname.map([](auto& value) { return kj::Maybe(value); }).orDefault(kj::none), protocol); }, destAddr, isTls); auto innerServer = kj::heap(containerClient.timer, headerTable, *innerService); co_await innerServer->listenHttpCleanDrain(connection); co_return true; } kj::HttpHeaders responseHeaders(headerTable); response.accept(202, "Accepted", responseHeaders); co_await passThroughConnection(destAddr, connection); co_return true; } // Handles raw TCP egress by forwarding the sidecar tunnel to the worker's connect() handler. // Returns true if a TCP mapping matched and the connection was handled. kj::Promise handleTcpConnect( kj::StringPtr destAddr, kj::AsyncIoStream& connection, ConnectResponse& response) { // For TCP, we match on IP:port only — no hostname matching since raw TCP // doesn't carry application-layer hostname information. auto mapping = containerClient.findEgressMapping( destAddr, /*defaultPort=*/0, /*hostname=*/kj::none, EgressProtocol::TCP); if (mapping == kj::none) { co_return false; } auto& channel = KJ_ASSERT_NONNULL(mapping); kj::HttpHeaders responseHeaders(headerTable); // 202 tells proxy-everything to send bytes as-is without attempting to // interpret the stream (e.g. if the underlying TCP carries TLS). response.accept(202, "Accepted", responseHeaders); IoChannelFactory::SubrequestMetadata metadata; auto worker = channel->startRequest(kj::mv(metadata)); // Bridge the sidecar tunnel to the worker's connect() handler. The worker entrypoint // is expected to implement connect() (e.g., a WorkerEntrypoint that proxies TCP). // We provide a simple ConnectResponse adapter since we've already accepted the // sidecar's CONNECT above. TcpEgressConnectResponse tcpResponse; kj::HttpHeaders connectHeaders(headerTable); co_await worker->connect(destAddr, connectHeaders, connection, tcpResponse, {}); co_return true; } kj::Promise passThroughConnection(kj::StringPtr destAddr, kj::AsyncIoStream& connection) { if (!containerClient.internetEnabled.orDefault(false)) { connection.shutdownWrite(); co_return; } auto addr = co_await containerClient.network.parseAddress(destAddr); auto destConn = co_await addr->connect(); co_await pumpBidirectional(connection, *destConn); co_return; } ContainerClient& containerClient; kj::HttpHeaderTable& headerTable; }; kj::Promise ContainerClient::getDockerBridgeIPAMConfig() { auto response = co_await dockerApiRequest( network, kj::str(dockerPath), kj::HttpMethod::GET, kj::str("/networks/bridge")); if (response.statusCode == 200) { auto message = decodeJsonResponse(response.body); auto jsonRoot = message->getRoot(); auto ipamConfig = jsonRoot.getIpam().getConfig(); if (ipamConfig.size() > 0) { auto config = ipamConfig[0]; co_return IPAMConfigResult{ .gateway = kj::str(config.getGateway()), .subnet = kj::str(config.getSubnet()), }; } } JSG_FAIL_REQUIRE(Error, "Failed to get bridge. " "Status: ", response.statusCode, ", Body: ", response.body); } kj::Promise ContainerClient::isDaemonIpv6Enabled() { // Inspect the default bridge network. When the Docker daemon has "ipv6": true in // daemon.json, the default bridge gets an IPv6 IPAM subnet entry (e.g. "fd00::/80"). auto response = co_await dockerApiRequest( network, kj::str(dockerPath), kj::HttpMethod::GET, kj::str("/networks/bridge")); if (response.statusCode != 200) { co_return false; } auto message = decodeJsonResponse(response.body); auto jsonRoot = message->getRoot(); for (auto config: jsonRoot.getIpam().getConfig()) { // IPv6 subnets contain ':' (e.g. "fd00::/80", "2001:db8::/64") if (kj::StringPtr(config.getSubnet()).findFirst(':') != kj::none) { co_return true; } } co_return false; } kj::Promise ContainerClient::startEgressListener( kj::String listenAddress, uint16_t port) { auto service = kj::heap(*this, headerTable); auto httpServer = kj::heap(timer, headerTable, *service); auto& httpServerRef = *httpServer; egressHttpServer = httpServer.attach(kj::mv(service)); auto addr = co_await network.parseAddress(kj::str(listenAddress, ":", port)); // The gateway IP from Docker's bridge network is not always bindable on the host. // On WSL with Docker Desktop, 172.17.0.1 lives inside the Docker VM, not on the WSL host's // interfaces. In that case, fall back to loopback — the sidecar reaches the host via // host-gateway (which Docker Desktop maps to the host loopback) so 127.0.0.1 works correctly. kj::Own listener; KJ_IF_SOME(e, kj::runCatchingExceptions([&]() { listener = addr->listen(); })) { KJ_LOG(WARNING, "Could not bind egress listener to gateway address, falling back to loopback", listenAddress, e); auto fallbackAddr = co_await network.parseAddress(kj::str("127.0.0.1:", port)); listener = fallbackAddr->listen(); } uint16_t chosenPort = listener->getPort(); egressListenerTask = httpServerRef.listenHttp(*listener) .attach(kj::mv(listener)) .eagerlyEvaluate([](kj::Exception&& e) { LOG_EXCEPTION( "Workerd could not listen in the TCP port to proxy traffic off the docker container", e); }); co_return chosenPort; } void ContainerClient::stopEgressListener() { egressListenerTask = kj::none; egressHttpServer = kj::none; egressListenerStarted.store(false, std::memory_order_release); } kj::Promise ContainerClient::writeFileToContainer(kj::StringPtr container, kj::StringPtr dir, kj::StringPtr filename, kj::ArrayPtr content) { kj::HttpHeaderTable table; auto address = co_await network.parseAddress(kj::str(dockerPath)); auto connection = co_await address->connect(); auto httpClient = kj::newHttpClient(table, *connection).attach(kj::mv(connection)); auto tar = createTarWithFile(filename, content); kj::HttpHeaders headers(table); headers.setPtr(kj::HttpHeaderId::HOST, "localhost"); headers.setPtr(kj::HttpHeaderId::CONTENT_TYPE, "application/x-tar"); headers.set(kj::HttpHeaderId::CONTENT_LENGTH, kj::str(tar.size())); auto endpoint = kj::str("/containers/", container, "/archive?path=", kj::encodeUriComponent(dir)); auto req = httpClient->request(kj::HttpMethod::PUT, endpoint, headers, tar.size()); { auto body = kj::mv(req.body); co_await body->write(tar.asBytes()); } auto response = co_await req.response; auto result = co_await response.body->readAllText(); JSG_REQUIRE(response.statusCode == 200, Error, "Failed to write file ", dir, "/", filename, " to container [", response.statusCode, "] ", result); } static constexpr kj::StringPtr cloudflareCaDir = "/etc"_kj; static constexpr kj::StringPtr cloudflareCaFilename = "cloudflare/certs/cloudflare-containers-ca.crt"_kj; kj::Promise ContainerClient::readCACert() { auto ingressPort = KJ_REQUIRE_NONNULL( sidecarIngressHostPort, "Cannot read CA cert: sidecar ingress port not known"); auto response = co_await dockerApiRequest( network, kj::str("127.0.0.1:", ingressPort), kj::HttpMethod::GET, kj::str("/ca")); JSG_REQUIRE(response.statusCode == 200, Error, "Failed to read CA cert from sidecar: ", response.statusCode, " ", response.body); caCert = kj::mv(response.body); } kj::Promise ContainerClient::injectCACert() { if (caCertInjected.exchange(true, std::memory_order_acquire)) { co_return; } bool succeeded = false; KJ_DEFER(if (!succeeded) caCertInjected.store(false, std::memory_order_release)); if (caCert == kj::none) { co_await readCACert(); } auto& cert = KJ_REQUIRE_NONNULL(caCert, "CA cert not read from sidecar yet"); co_await writeFileToContainer( containerName, cloudflareCaDir, cloudflareCaFilename, cert.asBytes()); succeeded = true; } kj::Promise> ContainerClient::inspectContainer() { auto endpoint = kj::str("/containers/", containerName, "/json"); auto response = co_await dockerApiRequest( network, kj::str(dockerPath), kj::HttpMethod::GET, kj::mv(endpoint)); // We check if the container with the given name exist, and if it's not, // we simply return false while avoiding an unnecessary error. if (response.statusCode == 404) { co_return kj::none; } JSG_REQUIRE(response.statusCode == 200, Error, "Container inspect failed"); // Parse JSON response auto message = decodeJsonResponse(response.body); auto jsonRoot = message->getRoot(); // Look for Status field in the JSON object JSG_REQUIRE(jsonRoot.hasState(), Error, "Malformed ContainerInspect response"); auto state = jsonRoot.getState(); JSG_REQUIRE(state.hasStatus(), Error, "Malformed ContainerInspect response"); auto status = state.getStatus(); // Treat both "running" and "restarting" as running. The "restarting" state occurs when // Docker is automatically restarting a container (due to restart policy). From the user's // perspective, a restarting container is still "alive" and should be treated as running // so that start() correctly refuses to start a duplicate and destroy() can clean it up. bool running = status == "running" || status == "restarting"; kj::Vector