// Copyright (c) 2017-2022 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 "channel-token.h" #include #include #include #include #include #include #include #include #include namespace kj { class TlsContext; } namespace workerd::jsg { class V8System; } namespace workerd::server { using api::pyodide::PythonConfig; // Implements the single-tenant Workers Runtime server / CLI. // // The purpose of this class is to implement the core logic independently of the CLI itself, // in such a way that it can be unit-tested. workerd.c++ implements the CLI wrapper around this. class Server final: private kj::TaskSet::ErrorHandler, private ChannelTokenHandler::Resolver { public: Server(kj::Filesystem& fs, kj::Timer& timer, const kj::MonotonicClock& monotonicClock, kj::Network& network, kj::EntropySource& entropySource, Worker::LoggingOptions loggingOptions, kj::Function reportConfigError); ~Server() noexcept; // Permit experimental features to be used. These features may break backwards compatibility // in the future. void allowExperimental() { experimental = true; } void overrideSocket(kj::String name, kj::Own port) { socketOverrides.upsert(kj::mv(name), kj::mv(port)); } void overrideSocket(kj::String name, kj::String addr) { socketOverrides.upsert(kj::mv(name), kj::mv(addr)); } void overrideDirectory(kj::String name, kj::String path) { directoryOverrides.upsert(kj::mv(name), kj::mv(path)); } void overrideExternal(kj::String name, kj::String addr) { externalOverrides.upsert(kj::mv(name), kj::mv(addr)); } void enableInspector(kj::String addr) { inspectorOverride = kj::mv(addr); } void enableControl(uint fd) { controlOverride = kj::heap(fd); } void enableDebugPort(kj::String addr) { debugPortOverride = kj::mv(addr); } void setPackageDiskCacheRoot(kj::Maybe>&& dir) { pythonConfig.packageDiskCacheRoot = kj::mv(dir); } void setPyodideDiskCacheRoot(kj::Maybe>&& dir) { pythonConfig.pyodideDiskCacheRoot = kj::mv(dir); } void setPythonCreateSnapshot() { pythonConfig.createSnapshot = true; } void setPythonCreateBaselineSnapshot() { pythonConfig.createBaselineSnapshot = true; } void setPythonLoadSnapshot(kj::String snapshot) { pythonConfig.loadSnapshotFromDisk = kj::mv(snapshot); } void setPythonSnapshotDirectory(kj::Maybe>&& dir) { pythonConfig.snapshotDirectory = kj::mv(dir); } // Set the compatibility date to use for all workers. When set, workers in the config must NOT // specify compatibilityDate (an error is reported if they do). This is used for testing to // ensure tests run with both old and new compat dates. void setTestCompatibilityDateOverride(kj::String date) { testCompatibilityDateOverride = kj::mv(date); } // Runs the server using the given config. kj::Promise run(jsg::V8System& v8System, config::Config::Reader conf, kj::Promise drainWhen = kj::NEVER_DONE); // Executes one or more tests. By default, all exported test handlers from all entrypoints to // all services in the config are executed. Glob patterns can be specified to match specific // service and entrypoint names. // // The returned promise resolves true if at least one test ran and no tests failed. kj::Promise test(jsg::V8System& v8System, config::Config::Reader conf, kj::StringPtr servicePattern = "*"_kj, kj::StringPtr entrypointPattern = "*"_kj); struct Durable { kj::String uniqueKey; bool isEvictable; bool enableSql; kj::Maybe containerOptions; }; struct Ephemeral { bool isEvictable; bool enableSql; }; using ActorConfig = kj::OneOf; class InspectorService; class InspectorServiceIsolateRegistrar; void handleReportConfigError(kj::String error) { reportConfigError(kj::mv(error)); } private: kj::Filesystem& fs; kj::Timer& timer; // monotonicClock must produce time values consistent with those produced by timer whenever // timer updates, but monotonicClock updates continuously (not just when system I/O is polled). const kj::MonotonicClock& monotonicClock; kj::Network& network; kj::EntropySource& entropySource; kj::Function reportConfigError; PythonConfig pythonConfig = PythonConfig{.packageDiskCacheRoot = kj::none, .pyodideDiskCacheRoot = kj::none, .createSnapshot = false, .createBaselineSnapshot = false, .loadSnapshotFromDisk = kj::none}; bool experimental = false; // When set, overrides compatibilityDate for all workers and enforces that workers don't // specify their own compatibilityDate. kj::Maybe testCompatibilityDateOverride; Worker::LoggingOptions loggingOptions; kj::Own memoryCacheProvider; ChannelTokenHandler channelTokenHandler; kj::HashMap>> socketOverrides; kj::HashMap directoryOverrides; // Overrides from the command line. // // String overrides are left as strings rather than parsed by the caller in order to reuse the // code that parses strings from the config file. kj::HashMap externalOverrides; kj::Maybe inspectorOverride; kj::Maybe> inspectorIsolateRegistrar; kj::Maybe> controlOverride; kj::Maybe debugPortOverride; struct GlobalContext; // General context needed to construct workers. Initialized early in run(). kj::Own globalContext; class Service; kj::Own invalidConfigServiceSingleton; class ActorClass; kj::Own invalidConfigActorClassSingleton; // Information about all known actor namespaces. Maps serviceName -> className -> config. // This needs to be populated in advance of constructing any services, in order to be able to // correctly construct dependent services. kj::HashMap> actorConfigs; kj::HashMap> services; class WorkerLoaderNamespace; kj::HashMap> workerLoaderNamespaces; kj::Vector> anonymousWorkerLoaderNamespaces; kj::Own> fatalFulfiller; // An HttpServer object maintained in a linked list. struct ListedHttpServer { Server& owner; kj::HttpServer httpServer; kj::ListLink link; template ListedHttpServer(Server& owner, Params&&... params) : owner(owner), httpServer(kj::fwd(params)...) { owner.httpServers.add(*this); }; ~ListedHttpServer() noexcept(false) { owner.httpServers.remove(*this); } }; // All active HttpServer objects -- used to implement drain(). kj::List httpServers; // Especially includes server loop tasks to listen on sockets. Any error is considered fatal. kj::TaskSet tasks; // Reports an exception thrown by a task in `tasks`. void taskFailed(kj::Exception&& exception) override; // Tell all HttpServers to drain once the drainWhen promise resolves. // This causes them to disconnect any connections that do not have a // request in flight. kj::Promise handleDrain(kj::Promise drainWhen); kj::Own makeTlsContext(config::TlsOptions::Reader conf); kj::Promise> makeTlsNetworkAddress(config::TlsOptions::Reader conf, kj::StringPtr addrStr, kj::Maybe certificateHost, uint defaultPort = 0); class HttpRewriter; kj::Own makeInvalidConfigService(); kj::Own makeExternalService(kj::StringPtr name, config::ExternalServer::Reader conf, kj::HttpHeaderTable::Builder& headerTableBuilder); kj::Own makeNetworkService(config::Network::Reader conf); kj::Own makeDiskDirectoryService(kj::StringPtr name, config::DiskDirectory::Reader conf, kj::HttpHeaderTable::Builder& headerTableBuilder); kj::Promise> makeWorker(kj::StringPtr name, config::Worker::Reader conf, capnp::List::Reader extensions); kj::Promise> makeService(config::Service::Reader conf, kj::HttpHeaderTable::Builder& headerTableBuilder, capnp::List::Reader extensions); // Aborts all actors in this server except those in namespaces marked with `preventEviction`. void abortAllActors(kj::Maybe reason); // Aborts all actors, cancels all alarms, and deletes all underlying storage for evictable // namespaces. After this, DOs can be recreated with clean state. Useful for test isolation. void deleteAllActors(kj::Maybe reason); // Can only be called in the link stage. // // May return a new object or may return a fake-own around a long-lived object. kj::Own lookupService( config::ServiceDesignator::Reader designator, kj::String errorContext); // Like lookupService() but looks up an actor class (especially for use as a facet class). // Returns none on a config error. kj::Own lookupActorClass( config::ServiceDesignator::Reader designator, kj::String errorContext); // Pretty similar to lookupService() and lookupActorClass(), but these callbacks are called by // the `ChannelTokenHandler` when decoding tokens. kj::Own resolveEntrypoint( kj::StringPtr serviceName, kj::Maybe entrypoint, Frankenvalue props) override; kj::Own resolveActorClass( kj::StringPtr serviceName, kj::Maybe entrypoint, Frankenvalue props) override; kj::Array encodeChannelToken(IoChannelFactory::ChannelTokenUsage usage, kj::StringPtr serviceName, kj::Maybe entrypoint, Frankenvalue& props); void decodeChannelToken(IoChannelFactory::ChannelTokenUsage usage, kj::ArrayPtr token, kj::FunctionParam entrypoint, Frankenvalue props)> callback); kj::Promise listenHttp(kj::Own listener, kj::Own service, kj::StringPtr physicalProtocol, kj::Own rewriter); kj::Promise listenTcp( kj::Own listener, kj::Own service, kj::StringPtr addrStr); kj::Promise listenDebugPort(kj::Own listener); class InvalidConfigService; class InvalidConfigActorClass; class ExternalHttpService; class ExternalTcpService; class NetworkService; class DiskDirectoryService; class WorkerService; class WorkerEntrypointService; class WorkerdBootstrapImpl; class HttpListener; class TcpListener; class DebugPortListener; struct ErrorReporter; struct ConfigErrorReporter; struct DynamicErrorReporter; struct WorkerDef; kj::Promise> makeWorkerImpl(kj::StringPtr name, WorkerDef def, capnp::List::Reader extensions, ErrorReporter& errorReporter); kj::Promise startServices(jsg::V8System& v8System, config::Config::Reader config, kj::HttpHeaderTable::Builder& headerTableBuilder, kj::ForkedPromise& forkedDrainWhen); kj::Promise listenOnSockets(config::Config::Reader config, kj::HttpHeaderTable::Builder& headerTableBuilder, kj::ForkedPromise& forkedDrainWhen, bool forTest = false); // Parsed socket protocol/TLS config. Extracted from the switch in listenOnSockets() to avoid // goto-over-initialization inside a coroutine, which triggers a clang optimizer crash. struct SocketTypeConfig { uint defaultPort = 0; config::HttpOptions::Reader httpOptions; kj::Maybe> tls; kj::StringPtr physicalProtocol; }; kj::Maybe parseSocketType(config::Socket::Reader sock, kj::StringPtr name); void unlinkWorkerLoaders(); kj::Promise preloadPython( kj::StringPtr workerName, const WorkerDef& workerDef, ErrorReporter& errorReporter); friend struct FutureSubrequestChannel; friend struct FutureActorClassChannel; }; } // namespace workerd::server