File
Blob: src/workerd/server/container-client.c++
| 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 | |
| 33 | namespace workerd::server { |
| 34 | |
| 35 | namespace { |
| 36 | |
| 37 | constexpr uint16_t SIDECAR_INGRESS_PORT = 39001; |
| 38 | |
| 39 | constexpr 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). |
| 45 | constexpr uint64_t MAX_JSON_RESPONSE_SIZE = 16ULL * 1024 * 1024; |
| 46 | |
| 47 | constexpr kj::StringPtr SNAPSHOT_VOLUME_PREFIX = "workerd-snap-"_kj; |
| 48 | constexpr kj::StringPtr SNAPSHOT_CLONE_VOLUME_PREFIX = "workerd-snap-clone-"_kj; |
| 49 | constexpr kj::StringPtr CONTAINER_SNAPSHOT_IMAGE_PREFIX = "workerd-container-snap-"_kj; |
| 50 | constexpr 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. |
| 55 | constexpr kj::StringPtr WORKERD_LABEL_PREFIX = "workerd-"_kj; |
| 56 | constexpr auto SNAPSHOT_STALE_AGE = 30 * kj::DAYS; |
| 57 | |
| 58 | // Maximum size of a snapshot tar archive held in memory during snapshot create/restore. |
| 59 | constexpr size_t MAX_SNAPSHOT_TAR_SIZE = 1ULL * 1024 * 1024 * 1024; // 1 GiB |
| 60 | static_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. |
| 64 | constexpr size_t MAX_TAR_CONTENT_SIZE = 8ull * 1024 * 1024 * 1024; |
| 65 | |
| 66 | // Ensures the stale-volume check runs at most once per process. |
| 67 | std::atomic_bool staleSnapshotVolumeCheckScheduled = false; |
| 68 | |
| 69 | struct ParsedAddress { |
| 70 | kj::OneOf<kj::CidrRange, kj::String> destination; |
| 71 | kj::Maybe<uint16_t> port; |
| 72 | }; |
| 73 | |
| 74 | struct HostAndPort { |
| 75 | kj::String host; |
| 76 | kj::Maybe<uint16_t> port; |
| 77 | }; |
| 78 | |
| 79 | struct DockerResponse { |
| 80 | kj::uint statusCode; |
| 81 | kj::String body; |
| 82 | }; |
| 83 | |
| 84 | struct DockerBinaryResponse { |
| 85 | kj::uint statusCode; |
| 86 | kj::Array<kj::byte> body; |
| 87 | }; |
| 88 | |
| 89 | struct 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 (".."). |
| 97 | kj::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. |
| 112 | kj::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". |
| 126 | class 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. |
| 225 | class 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 | |
| 235 | class 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. |
| 268 | HostAndPort 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. |
| 303 | kj::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 | |
| 312 | kj::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. |
| 323 | kj::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 |
| 338 | bool 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 | |
| 373 | kj::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" |
| 387 | ParsedAddress 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 | |
| 404 | kj::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 | |
| 473 | void 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. |
| 481 | kj::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. |
| 524 | kj::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 | |
| 559 | kj::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 | |
| 573 | kj::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 | |
| 587 | kj::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 | |
| 595 | kj::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 | |
| 609 | kj::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 | |
| 632 | kj::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. |
| 646 | kj::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 | |
| 698 | kj::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. |
| 719 | constexpr size_t DOCKER_FRAME_HEADER_SIZE = 8; |
| 720 | |
| 721 | uint32_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 | |
| 728 | void 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. |
| 735 | kj::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 | |
| 818 | kj::String currentSnapshotVolumeTimestamp() { |
| 819 | return kj::str((kj::systemPreciseCalendarClock().now() - kj::UNIX_EPOCH) / kj::SECONDS); |
| 820 | } |
| 821 | |
| 822 | kj::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 | |
| 842 | kj::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. |
| 878 | kj::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 | |
| 886 | kj::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. |
| 919 | struct 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. |
| 928 | struct ContainerClient::EgressState { |
| 929 | kj::Vector<EgressMapping> mappings; |
| 930 | }; |
| 931 | |
| 932 | ContainerClient::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 | |
| 964 | ContainerClient::~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. |
| 990 | class 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 | |
| 1055 | class 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. |
| 1146 | class 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. |
| 1167 | class 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 | |
| 1215 | kj::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. |
| 1222 | class 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 | |
| 1372 | kj::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 | |
| 1394 | kj::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 | |
| 1416 | kj::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 | |
| 1449 | void ContainerClient::stopEgressListener() { |
| 1450 | egressListenerTask = kj::none; |
| 1451 | egressHttpServer = kj::none; |
| 1452 | egressListenerStarted.store(false, std::memory_order_release); |
| 1453 | } |
| 1454 | |
| 1455 | kj::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 | |
| 1484 | static constexpr kj::StringPtr cloudflareCaDir = "/etc"_kj; |
| 1485 | static constexpr kj::StringPtr cloudflareCaFilename = |
| 1486 | "cloudflare/certs/cloudflare-containers-ca.crt"_kj; |
| 1487 | |
| 1488 | kj::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 | |
| 1501 | kj::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 | |
| 1520 | kj::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 | |
| 1567 | kj::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 | |
| 1615 | kj::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 | |
| 1631 | kj::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 | |
| 1658 | kj::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 | |
| 1751 | kj::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 | |
| 1795 | kj::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 | |
| 1829 | kj::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 | |
| 1847 | kj::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 | |
| 1897 | kj::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 | |
| 1908 | kj::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 | |
| 1919 | kj::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. |
| 1932 | kj::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. |
| 1939 | kj::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 | |
| 2003 | kj::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 | |
| 2013 | kj::Promise<void> ContainerClient::destroySidecarContainer() { |
| 2014 | co_await removeContainer(network, kj::str(dockerPath), kj::str(sidecarContainerName)); |
| 2015 | } |
| 2016 | |
| 2017 | kj::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 | |
| 2035 | kj::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 | |
| 2044 | kj::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 | |
| 2053 | kj::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 | |
| 2065 | kj::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 | |
| 2072 | kj::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 | |
| 2094 | kj::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 | |
| 2163 | kj::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 | |
| 2171 | ContainerClient::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 | |
| 2178 | kj::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 | |
| 2212 | kj::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 | |
| 2229 | kj::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 | |
| 2320 | kj::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 | |
| 2345 | kj::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 | |
| 2362 | kj::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 | |
| 2371 | kj::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 | |
| 2417 | kj::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 | |
| 2438 | kj::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 | |
| 2511 | kj::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 | |
| 2544 | kj::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 | |
| 2555 | kj::Promise<void> ContainerClient::listenTcp(ListenTcpContext context) { |
| 2556 | KJ_UNIMPLEMENTED("listenTcp not implemented for Docker containers - use port mapping instead"); |
| 2557 | } |
| 2558 | |
| 2559 | void 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 | |
| 2590 | kj::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 | |
| 2621 | kj::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 | |
| 2664 | kj::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 | |
| 2705 | kj::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 | |
| 2720 | kj::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 | |
| 2757 | kj::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 | |
| 2792 | kj::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 | |
| 2828 | kj::Own<ContainerClient> ContainerClient::addRef() { |
| 2829 | return kj::addRef(*this); |
| 2830 | } |
| 2831 | |
| 2832 | } // namespace workerd::server |