// Copyright (c) 2025 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #pragma once #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include namespace workerd::server { // Distinguishes how an egress mapping should proxy matched connections. enum class EgressProtocol : uint8_t { HTTP, // Parse HTTP inside the CONNECT tunnel and forward via worker request(). HTTPS, // Same as HTTP but with TLS interception (CA-cert injected). TCP, // Forward raw bytes via worker connect(). }; // Decode a JSON string into a Cap'n Proto message of type T. The MallocMessageBuilder is // heap-allocated and returned as an owned pointer so that the decoded data outlives this // call. Callers must keep the returned message alive while accessing the root via // message->getRoot(). A previous version allocated the builder on the stack and returned // a Builder (which is just a pointer into the message's arena); that caused every caller to // dereference freed memory after the function returned. template kj::Own decodeJsonResponse(kj::StringPtr response) { auto message = kj::heap(); capnp::JsonCodec codec; codec.handleByAnnotation(); auto jsonRoot = message->initRoot(); codec.decode(response, jsonRoot); return message; } // Docker-based implementation that implements the rpc::Container::Server interface // so it can be used as a rpc::Container::Client via kj::heap(). // This allows the Container JSG class to use Docker directly without knowing // it's talking to Docker instead of a real RPC service. // // ContainerClient is reference-counted to support actor reconnection with inactivity timeouts. // When setInactivityTimeout() is called, a timer holds a reference to prevent premature // destruction. The ContainerClient can be shared across multiple actor lifetimes class ContainerClient final: public rpc::Container::Server, public kj::Refcounted { public: ContainerClient(capnp::ByteStreamFactory& byteStreamFactory, kj::Timer& timer, kj::Network& network, kj::String dockerPath, kj::String containerName, kj::String imageName, kj::String containerEgressInterceptorImage, kj::TaskSet& waitUntilTasks, kj::Promise pendingCleanup, kj::Function)> cleanupCallback, ChannelTokenHandler& channelTokenHandler); ~ContainerClient() noexcept(false); // Implement rpc::Container::Server interface kj::Promise status(StatusContext context) override; kj::Promise start(StartContext context) override; kj::Promise monitor(MonitorContext context) override; kj::Promise destroy(DestroyContext context) override; kj::Promise signal(SignalContext context) override; kj::Promise exec(ExecContext context) override; kj::Promise getTcpPort(GetTcpPortContext context) override; kj::Promise listenTcp(ListenTcpContext context) override; kj::Promise setInactivityTimeout(SetInactivityTimeoutContext context) override; kj::Promise setEgressHttp(SetEgressHttpContext context) override; kj::Promise setEgressHttps(SetEgressHttpsContext context) override; kj::Promise setEgressTcp(SetEgressTcpContext context) override; kj::Promise snapshotDirectory(SnapshotDirectoryContext context) override; kj::Promise snapshotContainer(SnapshotContainerContext context) override; kj::Promise inspect(InspectContext context) override; kj::Own addRef(); private: capnp::ByteStreamFactory& byteStreamFactory; kj::HttpHeaderTable headerTable; kj::Timer& timer; kj::Network& network; kj::String dockerPath; kj::String containerName; kj::String sidecarContainerName; kj::String imageName; // Container egress interceptor image name (sidecar for egress proxy) kj::String containerEgressInterceptorImage; kj::TaskSet& waitUntilTasks; // Forked promise representing pending cleanup from a previous ContainerClient for the same // container ID. status() co_awaits a branch so that Docker inspect only runs after any // in-flight DELETE from the previous client has settled (either completed or been cancelled // via containerCleanupCanceler, in which case the .catch_() resolves it immediately). kj::ForkedPromise pendingCleanup; static constexpr kj::StringPtr defaultEnv[] = {"CLOUDFLARE_COUNTRY_A2=XX"_kj, "CLOUDFLARE_DEPLOYMENT_ID=xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx"_kj, "CLOUDFLARE_LOCATION=loc01"_kj, "CLOUDFLARE_REGION=REGN"_kj, "CLOUDFLARE_APPLICATION_ID=xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx"_kj, "CLOUDFLARE_DURABLE_OBJECT_ID=xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"_kj}; // Docker-specific Port implementation class DockerPort; class DockerProcessHandle; // EgressHttpService handles CONNECT requests from proxy-anything sidecar friend class EgressHttpService; struct Label { kj::String name; kj::String value; }; struct InspectResponse { bool isRunning; kj::Array