File
Blob: src/workerd/server/container-client.h
| 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 | #pragma once |
| 6 | |
| 7 | #include <workerd/io/container.capnp.h> |
| 8 | #include <workerd/io/io-channels.h> |
| 9 | #include <workerd/server/channel-token.h> |
| 10 | #include <workerd/server/docker-api.capnp.h> |
| 11 | |
| 12 | #include <capnp/compat/byte-stream.h> |
| 13 | #include <capnp/compat/json.h> |
| 14 | #include <capnp/list.h> |
| 15 | #include <capnp/message.h> |
| 16 | #include <kj/async-io.h> |
| 17 | #include <kj/async.h> |
| 18 | #include <kj/compat/http.h> |
| 19 | #include <kj/filesystem.h> |
| 20 | #include <kj/map.h> |
| 21 | #include <kj/refcount.h> |
| 22 | #include <kj/string.h> |
| 23 | |
| 24 | #include <atomic> |
| 25 | |
| 26 | namespace workerd::server { |
| 27 | |
| 28 | // Distinguishes how an egress mapping should proxy matched connections. |
| 29 | enum class EgressProtocol : uint8_t { |
| 30 | HTTP, // Parse HTTP inside the CONNECT tunnel and forward via worker request(). |
| 31 | HTTPS, // Same as HTTP but with TLS interception (CA-cert injected). |
| 32 | TCP, // Forward raw bytes via worker connect(). |
| 33 | }; |
| 34 | |
| 35 | // Decode a JSON string into a Cap'n Proto message of type T. The MallocMessageBuilder is |
| 36 | // heap-allocated and returned as an owned pointer so that the decoded data outlives this |
| 37 | // call. Callers must keep the returned message alive while accessing the root via |
| 38 | // message->getRoot<T>(). A previous version allocated the builder on the stack and returned |
| 39 | // a Builder (which is just a pointer into the message's arena); that caused every caller to |
| 40 | // dereference freed memory after the function returned. |
| 41 | template <typename T> |
| 42 | kj::Own<capnp::MallocMessageBuilder> decodeJsonResponse(kj::StringPtr response) { |
| 43 | auto message = kj::heap<capnp::MallocMessageBuilder>(); |
| 44 | capnp::JsonCodec codec; |
| 45 | codec.handleByAnnotation<T>(); |
| 46 | auto jsonRoot = message->initRoot<T>(); |
| 47 | codec.decode(response, jsonRoot); |
| 48 | return message; |
| 49 | } |
| 50 | |
| 51 | // Docker-based implementation that implements the rpc::Container::Server interface |
| 52 | // so it can be used as a rpc::Container::Client via kj::heap<ContainerClient>(). |
| 53 | // This allows the Container JSG class to use Docker directly without knowing |
| 54 | // it's talking to Docker instead of a real RPC service. |
| 55 | // |
| 56 | // ContainerClient is reference-counted to support actor reconnection with inactivity timeouts. |
| 57 | // When setInactivityTimeout() is called, a timer holds a reference to prevent premature |
| 58 | // destruction. The ContainerClient can be shared across multiple actor lifetimes |
| 59 | class ContainerClient final: public rpc::Container::Server, public kj::Refcounted { |
| 60 | public: |
| 61 | ContainerClient(capnp::ByteStreamFactory& byteStreamFactory, |
| 62 | kj::Timer& timer, |
| 63 | kj::Network& network, |
| 64 | kj::String dockerPath, |
| 65 | kj::String containerName, |
| 66 | kj::String imageName, |
| 67 | kj::String containerEgressInterceptorImage, |
| 68 | kj::TaskSet& waitUntilTasks, |
| 69 | kj::Promise<void> pendingCleanup, |
| 70 | kj::Function<void(kj::Promise<void>)> cleanupCallback, |
| 71 | ChannelTokenHandler& channelTokenHandler); |
| 72 | |
| 73 | ~ContainerClient() noexcept(false); |
| 74 | |
| 75 | // Implement rpc::Container::Server interface |
| 76 | kj::Promise<void> status(StatusContext context) override; |
| 77 | kj::Promise<void> start(StartContext context) override; |
| 78 | kj::Promise<void> monitor(MonitorContext context) override; |
| 79 | kj::Promise<void> destroy(DestroyContext context) override; |
| 80 | kj::Promise<void> signal(SignalContext context) override; |
| 81 | kj::Promise<void> exec(ExecContext context) override; |
| 82 | kj::Promise<void> getTcpPort(GetTcpPortContext context) override; |
| 83 | kj::Promise<void> listenTcp(ListenTcpContext context) override; |
| 84 | kj::Promise<void> setInactivityTimeout(SetInactivityTimeoutContext context) override; |
| 85 | kj::Promise<void> setEgressHttp(SetEgressHttpContext context) override; |
| 86 | kj::Promise<void> setEgressHttps(SetEgressHttpsContext context) override; |
| 87 | kj::Promise<void> setEgressTcp(SetEgressTcpContext context) override; |
| 88 | kj::Promise<void> snapshotDirectory(SnapshotDirectoryContext context) override; |
| 89 | kj::Promise<void> snapshotContainer(SnapshotContainerContext context) override; |
| 90 | kj::Promise<void> inspect(InspectContext context) override; |
| 91 | |
| 92 | kj::Own<ContainerClient> addRef(); |
| 93 | |
| 94 | private: |
| 95 | capnp::ByteStreamFactory& byteStreamFactory; |
| 96 | kj::HttpHeaderTable headerTable; |
| 97 | kj::Timer& timer; |
| 98 | kj::Network& network; |
| 99 | kj::String dockerPath; |
| 100 | kj::String containerName; |
| 101 | kj::String sidecarContainerName; |
| 102 | kj::String imageName; |
| 103 | |
| 104 | // Container egress interceptor image name (sidecar for egress proxy) |
| 105 | kj::String containerEgressInterceptorImage; |
| 106 | |
| 107 | kj::TaskSet& waitUntilTasks; |
| 108 | |
| 109 | // Forked promise representing pending cleanup from a previous ContainerClient for the same |
| 110 | // container ID. status() co_awaits a branch so that Docker inspect only runs after any |
| 111 | // in-flight DELETE from the previous client has settled (either completed or been cancelled |
| 112 | // via containerCleanupCanceler, in which case the .catch_() resolves it immediately). |
| 113 | kj::ForkedPromise<void> pendingCleanup; |
| 114 | |
| 115 | static constexpr kj::StringPtr defaultEnv[] = {"CLOUDFLARE_COUNTRY_A2=XX"_kj, |
| 116 | "CLOUDFLARE_DEPLOYMENT_ID=xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx"_kj, |
| 117 | "CLOUDFLARE_LOCATION=loc01"_kj, "CLOUDFLARE_REGION=REGN"_kj, |
| 118 | "CLOUDFLARE_APPLICATION_ID=xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx"_kj, |
| 119 | "CLOUDFLARE_DURABLE_OBJECT_ID=xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"_kj}; |
| 120 | |
| 121 | // Docker-specific Port implementation |
| 122 | class DockerPort; |
| 123 | class DockerProcessHandle; |
| 124 | |
| 125 | // EgressHttpService handles CONNECT requests from proxy-anything sidecar |
| 126 | friend class EgressHttpService; |
| 127 | |
| 128 | struct Label { |
| 129 | kj::String name; |
| 130 | kj::String value; |
| 131 | }; |
| 132 | |
| 133 | struct InspectResponse { |
| 134 | bool isRunning; |
| 135 | kj::Array<Label> labels; |
| 136 | }; |
| 137 | |
| 138 | struct IPAMConfigResult { |
| 139 | kj::String gateway; |
| 140 | kj::String subnet; |
| 141 | }; |
| 142 | |
| 143 | struct SidecarInspectResponse { |
| 144 | uint16_t ingressHostPort; |
| 145 | }; |
| 146 | |
| 147 | struct SnapshotRestoreMount { |
| 148 | kj::Path restorePath; |
| 149 | kj::String sourceVolume; |
| 150 | kj::String cloneVolume; |
| 151 | }; |
| 152 | |
| 153 | struct ImageInspectResponse { |
| 154 | kj::String id; |
| 155 | uint64_t size; |
| 156 | }; |
| 157 | |
| 158 | struct ExecInspectResponse { |
| 159 | int32_t exitCode; |
| 160 | bool running; |
| 161 | uint32_t pid; |
| 162 | }; |
| 163 | |
| 164 | kj::Promise<kj::Maybe<InspectResponse>> inspectContainer(); |
| 165 | |
| 166 | kj::Promise<void> updateSidecarEgressPort(uint16_t ingressHostPort, uint16_t egressPort); |
| 167 | kj::Promise<void> updateSidecarEgressConfig(uint16_t ingressHostPort, uint16_t egressPort); |
| 168 | kj::Promise<void> createContainer(kj::StringPtr effectiveImage, |
| 169 | kj::Maybe<capnp::List<capnp::Text>::Reader> entrypoint, |
| 170 | kj::Maybe<capnp::List<capnp::Text>::Reader> environment, |
| 171 | kj::ArrayPtr<const SnapshotRestoreMount> restoreMounts, |
| 172 | rpc::Container::StartParams::Reader params); |
| 173 | kj::Promise<kj::String> createExec(capnp::List<capnp::Text>::Reader cmd, |
| 174 | rpc::Container::ExecOptions::Reader params, |
| 175 | bool attachStdout, |
| 176 | bool attachStderr); |
| 177 | kj::Promise<kj::Own<kj::AsyncIoStream>> startExec(kj::String execId); |
| 178 | kj::Promise<ExecInspectResponse> inspectExec(kj::StringPtr execId); |
| 179 | kj::Promise<void> runSimpleExec(kj::ArrayPtr<const kj::String> cmd); |
| 180 | kj::Promise<void> startContainer(); |
| 181 | kj::Promise<void> stopContainer(); |
| 182 | kj::Promise<void> killContainer(uint32_t signal); |
| 183 | kj::Promise<void> destroyContainer(); |
| 184 | |
| 185 | // Docker volume management for snapshots |
| 186 | kj::Promise<void> createVolume(kj::StringPtr volumeName); |
| 187 | kj::Promise<void> deleteVolume(kj::String volumeName); |
| 188 | kj::Promise<void> commitContainer(kj::StringPtr imageRef); |
| 189 | kj::Promise<ImageInspectResponse> inspectImage(kj::StringPtr imageRef); |
| 190 | kj::Promise<void> deleteImage(kj::String imageRef); |
| 191 | kj::Promise<kj::String> createTempContainerWithVolume( |
| 192 | kj::StringPtr volumeName, kj::StringPtr mountPath); |
| 193 | // Creates a writable clone volume by copying an existing snapshot volume through a |
| 194 | // short-lived helper container. The caller mounts the returned clone into the app |
| 195 | // container with NoCopy=true so the restored path masks any image contents there. |
| 196 | kj::Promise<void> cloneSnapshot(SnapshotRestoreMount& snapshot); |
| 197 | kj::Promise<void> deleteTempContainer(kj::String tempContainerId); |
| 198 | |
| 199 | // Sidecar container management (for egress proxy) |
| 200 | // Inspect the sidecar container to retrieve the port to ingress to |
| 201 | kj::Promise<kj::Maybe<SidecarInspectResponse>> inspectSidecar(); |
| 202 | kj::Promise<void> createSidecarContainer(uint16_t egressPort, kj::String networkCidr); |
| 203 | kj::Promise<void> startSidecarContainer(); |
| 204 | kj::Promise<void> destroySidecarContainer(); |
| 205 | kj::Promise<void> monitorSidecarContainer(); |
| 206 | |
| 207 | // Cleanup callback invoked from the destructor. Receives the joined cleanup promise so |
| 208 | // ActorNamespace can wrap it with the canceler, store it for the next ContainerClient |
| 209 | // to await, and add a branch to waitUntilTasks to keep the cleanup tasks alive. |
| 210 | kj::Function<void(kj::Promise<void>)> cleanupCallback; |
| 211 | |
| 212 | // For redeeming channel tokens received via setEgressHttp / setEgressHttps. |
| 213 | ChannelTokenHandler& channelTokenHandler; |
| 214 | |
| 215 | // Opaque implementation struct holding egress mappings. Defined in container-client.c++ to |
| 216 | // avoid pulling heavy types (kj::OneOf, kj::CidrRange, kj::Vector) into server.c++ which |
| 217 | // includes this header. |
| 218 | struct EgressState; |
| 219 | kj::Own<EgressState> egressState; |
| 220 | |
| 221 | // Insert or replace an egress mapping. |
| 222 | struct EgressMapping; |
| 223 | void upsertEgressMapping(EgressMapping mapping); |
| 224 | kj::Vector<kj::String> getDnsAllowHostnames() const; |
| 225 | |
| 226 | // Find a matching egress mapping for the given destination address (host:port format). |
| 227 | // Returns an addRef'd Own so the channel stays alive even if the mapping is later replaced. |
| 228 | kj::Maybe<kj::Own<workerd::IoChannelFactory::SubrequestChannel>> findEgressMapping( |
| 229 | kj::StringPtr destAddr, |
| 230 | uint16_t defaultPort, |
| 231 | kj::Maybe<kj::StringPtr> hostname, |
| 232 | EgressProtocol protocol); |
| 233 | |
| 234 | kj::Promise<void> writeFileToContainer(kj::StringPtr container, |
| 235 | kj::StringPtr dir, |
| 236 | kj::StringPtr filename, |
| 237 | kj::ArrayPtr<const kj::byte> content); |
| 238 | kj::Promise<void> readCACert(); |
| 239 | kj::Promise<void> injectCACert(); |
| 240 | |
| 241 | // Whether general internet access is enabled for this container, when known. |
| 242 | kj::Maybe<bool> internetEnabled = kj::none; |
| 243 | |
| 244 | std::atomic_bool containerStarted = false; |
| 245 | std::atomic_bool containerSidecarStarted = false; |
| 246 | std::atomic_bool egressListenerStarted = false; |
| 247 | std::atomic_bool caCertInjected = false; |
| 248 | |
| 249 | // Writable clone volumes currently owned by the app container, or by an in-flight start() |
| 250 | // that still needs failure cleanup. |
| 251 | kj::Vector<kj::String> snapshotClones; |
| 252 | |
| 253 | // CA cert read from the sidecar after it starts. |
| 254 | kj::Maybe<kj::String> caCert; |
| 255 | |
| 256 | kj::Maybe<kj::Own<kj::HttpServer>> egressHttpServer; |
| 257 | kj::Maybe<kj::Promise<void>> egressListenerTask; |
| 258 | |
| 259 | uint16_t egressListenerPort = 0; |
| 260 | kj::Maybe<uint16_t> sidecarIngressHostPort; |
| 261 | |
| 262 | // All mutating RPCs need to ask and wait on an RpcTurn before doing any mutations. |
| 263 | // monitor() is an exception. It waits for all pending mutating RPCs without joining |
| 264 | // the queue itself. |
| 265 | kj::ForkedPromise<void> mutationQueue = kj::Promise<void>(kj::READY_NOW).fork(); |
| 266 | |
| 267 | struct RpcTurn { |
| 268 | kj::Promise<void> ready; |
| 269 | kj::Own<kj::PromiseFulfiller<void>> done; |
| 270 | }; |
| 271 | // Get a turn to run mutating RPC. |
| 272 | // Callers will receive a RpcTurn where they can wait and then resolve |
| 273 | // when they finish through a KJ defer. |
| 274 | RpcTurn getRpcTurn(); |
| 275 | |
| 276 | // Get the Docker bridge network gateway IP and subnet. |
| 277 | kj::Promise<IPAMConfigResult> getDockerBridgeIPAMConfig(); |
| 278 | // Check if the Docker daemon has IPv6 enabled by inspecting the default bridge network's |
| 279 | // IPAM config for IPv6 subnets. |
| 280 | kj::Promise<bool> isDaemonIpv6Enabled(); |
| 281 | // Start the egress listener on the given address. If port is 0, an OS-chosen port is used. |
| 282 | kj::Promise<uint16_t> startEgressListener(kj::String listenAddress, uint16_t port = 0); |
| 283 | void stopEgressListener(); |
| 284 | // Ensure the egress listener is started exactly once. |
| 285 | // Uses egressListenerStarted as a guard. Called from setEgressHttp() and status(). |
| 286 | // If port is non-zero, binds to that specific port (for reconnecting to an existing sidecar). |
| 287 | kj::Promise<void> ensureEgressListenerStarted(uint16_t port = 0); |
| 288 | // Ensure the egress listener and sidecar container are started exactly once. |
| 289 | // Uses containerSidecarStarted as a guard. Called from both start() and setEgressHttp(). |
| 290 | kj::Promise<void> ensureSidecarStarted(); |
| 291 | }; |
| 292 | |
| 293 | } // namespace workerd::server |