File
Blob: src/workerd/server/server.h
| 1 | // Copyright (c) 2017-2022 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 "channel-token.h" |
| 8 | |
| 9 | #include <workerd/api/memory-cache.h> |
| 10 | #include <workerd/api/pyodide/pyodide.h> |
| 11 | #include <workerd/io/worker.h> |
| 12 | #include <workerd/server/workerd.capnp.h> |
| 13 | |
| 14 | #include <kj/async-io.h> |
| 15 | #include <kj/compat/http.h> |
| 16 | #include <kj/filesystem.h> |
| 17 | #include <kj/map.h> |
| 18 | #include <kj/one-of.h> |
| 19 | |
| 20 | namespace kj { |
| 21 | class TlsContext; |
| 22 | } |
| 23 | |
| 24 | namespace workerd::jsg { |
| 25 | class V8System; |
| 26 | } |
| 27 | |
| 28 | namespace workerd::server { |
| 29 | |
| 30 | using api::pyodide::PythonConfig; |
| 31 | |
| 32 | // Implements the single-tenant Workers Runtime server / CLI. |
| 33 | // |
| 34 | // The purpose of this class is to implement the core logic independently of the CLI itself, |
| 35 | // in such a way that it can be unit-tested. workerd.c++ implements the CLI wrapper around this. |
| 36 | class Server final: private kj::TaskSet::ErrorHandler, private ChannelTokenHandler::Resolver { |
| 37 | public: |
| 38 | Server(kj::Filesystem& fs, |
| 39 | kj::Timer& timer, |
| 40 | const kj::MonotonicClock& monotonicClock, |
| 41 | kj::Network& network, |
| 42 | kj::EntropySource& entropySource, |
| 43 | Worker::LoggingOptions loggingOptions, |
| 44 | kj::Function<void(kj::String)> reportConfigError); |
| 45 | ~Server() noexcept; |
| 46 | |
| 47 | // Permit experimental features to be used. These features may break backwards compatibility |
| 48 | // in the future. |
| 49 | void allowExperimental() { |
| 50 | experimental = true; |
| 51 | } |
| 52 | |
| 53 | void overrideSocket(kj::String name, kj::Own<kj::ConnectionReceiver> port) { |
| 54 | socketOverrides.upsert(kj::mv(name), kj::mv(port)); |
| 55 | } |
| 56 | void overrideSocket(kj::String name, kj::String addr) { |
| 57 | socketOverrides.upsert(kj::mv(name), kj::mv(addr)); |
| 58 | } |
| 59 | void overrideDirectory(kj::String name, kj::String path) { |
| 60 | directoryOverrides.upsert(kj::mv(name), kj::mv(path)); |
| 61 | } |
| 62 | void overrideExternal(kj::String name, kj::String addr) { |
| 63 | externalOverrides.upsert(kj::mv(name), kj::mv(addr)); |
| 64 | } |
| 65 | void enableInspector(kj::String addr) { |
| 66 | inspectorOverride = kj::mv(addr); |
| 67 | } |
| 68 | void enableControl(uint fd) { |
| 69 | controlOverride = kj::heap<kj::FdOutputStream>(fd); |
| 70 | } |
| 71 | void enableDebugPort(kj::String addr) { |
| 72 | debugPortOverride = kj::mv(addr); |
| 73 | } |
| 74 | void setPackageDiskCacheRoot(kj::Maybe<kj::Own<const kj::Directory>>&& dir) { |
| 75 | pythonConfig.packageDiskCacheRoot = kj::mv(dir); |
| 76 | } |
| 77 | void setPyodideDiskCacheRoot(kj::Maybe<kj::Own<const kj::Directory>>&& dir) { |
| 78 | pythonConfig.pyodideDiskCacheRoot = kj::mv(dir); |
| 79 | } |
| 80 | void setPythonCreateSnapshot() { |
| 81 | pythonConfig.createSnapshot = true; |
| 82 | } |
| 83 | void setPythonCreateBaselineSnapshot() { |
| 84 | pythonConfig.createBaselineSnapshot = true; |
| 85 | } |
| 86 | void setPythonLoadSnapshot(kj::String snapshot) { |
| 87 | pythonConfig.loadSnapshotFromDisk = kj::mv(snapshot); |
| 88 | } |
| 89 | void setPythonSnapshotDirectory(kj::Maybe<kj::Own<const kj::Directory>>&& dir) { |
| 90 | pythonConfig.snapshotDirectory = kj::mv(dir); |
| 91 | } |
| 92 | |
| 93 | // Set the compatibility date to use for all workers. When set, workers in the config must NOT |
| 94 | // specify compatibilityDate (an error is reported if they do). This is used for testing to |
| 95 | // ensure tests run with both old and new compat dates. |
| 96 | void setTestCompatibilityDateOverride(kj::String date) { |
| 97 | testCompatibilityDateOverride = kj::mv(date); |
| 98 | } |
| 99 | |
| 100 | // Runs the server using the given config. |
| 101 | kj::Promise<void> run(jsg::V8System& v8System, |
| 102 | config::Config::Reader conf, |
| 103 | kj::Promise<void> drainWhen = kj::NEVER_DONE); |
| 104 | |
| 105 | // Executes one or more tests. By default, all exported test handlers from all entrypoints to |
| 106 | // all services in the config are executed. Glob patterns can be specified to match specific |
| 107 | // service and entrypoint names. |
| 108 | // |
| 109 | // The returned promise resolves true if at least one test ran and no tests failed. |
| 110 | kj::Promise<bool> test(jsg::V8System& v8System, |
| 111 | config::Config::Reader conf, |
| 112 | kj::StringPtr servicePattern = "*"_kj, |
| 113 | kj::StringPtr entrypointPattern = "*"_kj); |
| 114 | |
| 115 | struct Durable { |
| 116 | kj::String uniqueKey; |
| 117 | bool isEvictable; |
| 118 | bool enableSql; |
| 119 | kj::Maybe<config::Worker::DurableObjectNamespace::ContainerOptions::Reader> containerOptions; |
| 120 | }; |
| 121 | struct Ephemeral { |
| 122 | bool isEvictable; |
| 123 | bool enableSql; |
| 124 | }; |
| 125 | using ActorConfig = kj::OneOf<Durable, Ephemeral>; |
| 126 | |
| 127 | class InspectorService; |
| 128 | class InspectorServiceIsolateRegistrar; |
| 129 | |
| 130 | void handleReportConfigError(kj::String error) { |
| 131 | reportConfigError(kj::mv(error)); |
| 132 | } |
| 133 | |
| 134 | private: |
| 135 | kj::Filesystem& fs; |
| 136 | kj::Timer& timer; |
| 137 | // monotonicClock must produce time values consistent with those produced by timer whenever |
| 138 | // timer updates, but monotonicClock updates continuously (not just when system I/O is polled). |
| 139 | const kj::MonotonicClock& monotonicClock; |
| 140 | kj::Network& network; |
| 141 | kj::EntropySource& entropySource; |
| 142 | kj::Function<void(kj::String)> reportConfigError; |
| 143 | PythonConfig pythonConfig = PythonConfig{.packageDiskCacheRoot = kj::none, |
| 144 | .pyodideDiskCacheRoot = kj::none, |
| 145 | .createSnapshot = false, |
| 146 | .createBaselineSnapshot = false, |
| 147 | .loadSnapshotFromDisk = kj::none}; |
| 148 | |
| 149 | bool experimental = false; |
| 150 | |
| 151 | // When set, overrides compatibilityDate for all workers and enforces that workers don't |
| 152 | // specify their own compatibilityDate. |
| 153 | kj::Maybe<kj::String> testCompatibilityDateOverride; |
| 154 | |
| 155 | Worker::LoggingOptions loggingOptions; |
| 156 | |
| 157 | kj::Own<api::MemoryCacheProvider> memoryCacheProvider; |
| 158 | |
| 159 | ChannelTokenHandler channelTokenHandler; |
| 160 | |
| 161 | kj::HashMap<kj::String, kj::OneOf<kj::String, kj::Own<kj::ConnectionReceiver>>> socketOverrides; |
| 162 | kj::HashMap<kj::String, kj::String> directoryOverrides; |
| 163 | |
| 164 | // Overrides from the command line. |
| 165 | // |
| 166 | // String overrides are left as strings rather than parsed by the caller in order to reuse the |
| 167 | // code that parses strings from the config file. |
| 168 | kj::HashMap<kj::String, kj::String> externalOverrides; |
| 169 | |
| 170 | kj::Maybe<kj::String> inspectorOverride; |
| 171 | kj::Maybe<kj::Own<InspectorServiceIsolateRegistrar>> inspectorIsolateRegistrar; |
| 172 | kj::Maybe<kj::Own<kj::FdOutputStream>> controlOverride; |
| 173 | kj::Maybe<kj::String> debugPortOverride; |
| 174 | |
| 175 | struct GlobalContext; |
| 176 | // General context needed to construct workers. Initialized early in run(). |
| 177 | kj::Own<GlobalContext> globalContext; |
| 178 | |
| 179 | class Service; |
| 180 | kj::Own<Service> invalidConfigServiceSingleton; |
| 181 | |
| 182 | class ActorClass; |
| 183 | kj::Own<ActorClass> invalidConfigActorClassSingleton; |
| 184 | |
| 185 | // Information about all known actor namespaces. Maps serviceName -> className -> config. |
| 186 | // This needs to be populated in advance of constructing any services, in order to be able to |
| 187 | // correctly construct dependent services. |
| 188 | kj::HashMap<kj::String, kj::HashMap<kj::String, ActorConfig>> actorConfigs; |
| 189 | |
| 190 | kj::HashMap<kj::String, kj::Own<Service>> services; |
| 191 | |
| 192 | class WorkerLoaderNamespace; |
| 193 | kj::HashMap<kj::String, kj::Rc<WorkerLoaderNamespace>> workerLoaderNamespaces; |
| 194 | kj::Vector<kj::Rc<WorkerLoaderNamespace>> anonymousWorkerLoaderNamespaces; |
| 195 | |
| 196 | kj::Own<kj::PromiseFulfiller<void>> fatalFulfiller; |
| 197 | |
| 198 | // An HttpServer object maintained in a linked list. |
| 199 | struct ListedHttpServer { |
| 200 | Server& owner; |
| 201 | kj::HttpServer httpServer; |
| 202 | kj::ListLink<ListedHttpServer> link; |
| 203 | |
| 204 | template <typename... Params> |
| 205 | ListedHttpServer(Server& owner, Params&&... params) |
| 206 | : owner(owner), |
| 207 | httpServer(kj::fwd<Params>(params)...) { |
| 208 | owner.httpServers.add(*this); |
| 209 | }; |
| 210 | ~ListedHttpServer() noexcept(false) { |
| 211 | owner.httpServers.remove(*this); |
| 212 | } |
| 213 | }; |
| 214 | |
| 215 | // All active HttpServer objects -- used to implement drain(). |
| 216 | kj::List<ListedHttpServer, &ListedHttpServer::link> httpServers; |
| 217 | |
| 218 | // Especially includes server loop tasks to listen on sockets. Any error is considered fatal. |
| 219 | kj::TaskSet tasks; |
| 220 | |
| 221 | // Reports an exception thrown by a task in `tasks`. |
| 222 | void taskFailed(kj::Exception&& exception) override; |
| 223 | |
| 224 | // Tell all HttpServers to drain once the drainWhen promise resolves. |
| 225 | // This causes them to disconnect any connections that do not have a |
| 226 | // request in flight. |
| 227 | kj::Promise<void> handleDrain(kj::Promise<void> drainWhen); |
| 228 | |
| 229 | kj::Own<kj::TlsContext> makeTlsContext(config::TlsOptions::Reader conf); |
| 230 | kj::Promise<kj::Own<kj::NetworkAddress>> makeTlsNetworkAddress(config::TlsOptions::Reader conf, |
| 231 | kj::StringPtr addrStr, |
| 232 | kj::Maybe<kj::StringPtr> certificateHost, |
| 233 | uint defaultPort = 0); |
| 234 | |
| 235 | class HttpRewriter; |
| 236 | |
| 237 | kj::Own<Service> makeInvalidConfigService(); |
| 238 | kj::Own<Service> makeExternalService(kj::StringPtr name, |
| 239 | config::ExternalServer::Reader conf, |
| 240 | kj::HttpHeaderTable::Builder& headerTableBuilder); |
| 241 | kj::Own<Service> makeNetworkService(config::Network::Reader conf); |
| 242 | kj::Own<Service> makeDiskDirectoryService(kj::StringPtr name, |
| 243 | config::DiskDirectory::Reader conf, |
| 244 | kj::HttpHeaderTable::Builder& headerTableBuilder); |
| 245 | kj::Promise<kj::Own<Service>> makeWorker(kj::StringPtr name, |
| 246 | config::Worker::Reader conf, |
| 247 | capnp::List<config::Extension>::Reader extensions); |
| 248 | kj::Promise<kj::Own<Service>> makeService(config::Service::Reader conf, |
| 249 | kj::HttpHeaderTable::Builder& headerTableBuilder, |
| 250 | capnp::List<config::Extension>::Reader extensions); |
| 251 | |
| 252 | // Aborts all actors in this server except those in namespaces marked with `preventEviction`. |
| 253 | void abortAllActors(kj::Maybe<const kj::Exception&> reason); |
| 254 | |
| 255 | // Aborts all actors, cancels all alarms, and deletes all underlying storage for evictable |
| 256 | // namespaces. After this, DOs can be recreated with clean state. Useful for test isolation. |
| 257 | void deleteAllActors(kj::Maybe<const kj::Exception&> reason); |
| 258 | |
| 259 | // Can only be called in the link stage. |
| 260 | // |
| 261 | // May return a new object or may return a fake-own around a long-lived object. |
| 262 | kj::Own<Service> lookupService( |
| 263 | config::ServiceDesignator::Reader designator, kj::String errorContext); |
| 264 | |
| 265 | // Like lookupService() but looks up an actor class (especially for use as a facet class). |
| 266 | // Returns none on a config error. |
| 267 | kj::Own<ActorClass> lookupActorClass( |
| 268 | config::ServiceDesignator::Reader designator, kj::String errorContext); |
| 269 | |
| 270 | // Pretty similar to lookupService() and lookupActorClass(), but these callbacks are called by |
| 271 | // the `ChannelTokenHandler` when decoding tokens. |
| 272 | kj::Own<IoChannelFactory::SubrequestChannel> resolveEntrypoint( |
| 273 | kj::StringPtr serviceName, kj::Maybe<kj::StringPtr> entrypoint, Frankenvalue props) override; |
| 274 | kj::Own<IoChannelFactory::ActorClassChannel> resolveActorClass( |
| 275 | kj::StringPtr serviceName, kj::Maybe<kj::StringPtr> entrypoint, Frankenvalue props) override; |
| 276 | |
| 277 | kj::Array<byte> encodeChannelToken(IoChannelFactory::ChannelTokenUsage usage, |
| 278 | kj::StringPtr serviceName, |
| 279 | kj::Maybe<kj::StringPtr> entrypoint, |
| 280 | Frankenvalue& props); |
| 281 | |
| 282 | void decodeChannelToken(IoChannelFactory::ChannelTokenUsage usage, |
| 283 | kj::ArrayPtr<const byte> token, |
| 284 | kj::FunctionParam<void( |
| 285 | kj::StringPtr serviceName, kj::Maybe<kj::StringPtr> entrypoint, Frankenvalue props)> |
| 286 | callback); |
| 287 | |
| 288 | kj::Promise<void> listenHttp(kj::Own<kj::ConnectionReceiver> listener, |
| 289 | kj::Own<Service> service, |
| 290 | kj::StringPtr physicalProtocol, |
| 291 | kj::Own<HttpRewriter> rewriter); |
| 292 | |
| 293 | kj::Promise<void> listenTcp( |
| 294 | kj::Own<kj::ConnectionReceiver> listener, kj::Own<Service> service, kj::StringPtr addrStr); |
| 295 | |
| 296 | kj::Promise<void> listenDebugPort(kj::Own<kj::ConnectionReceiver> listener); |
| 297 | |
| 298 | class InvalidConfigService; |
| 299 | class InvalidConfigActorClass; |
| 300 | class ExternalHttpService; |
| 301 | class ExternalTcpService; |
| 302 | class NetworkService; |
| 303 | class DiskDirectoryService; |
| 304 | class WorkerService; |
| 305 | class WorkerEntrypointService; |
| 306 | class WorkerdBootstrapImpl; |
| 307 | class HttpListener; |
| 308 | class TcpListener; |
| 309 | class DebugPortListener; |
| 310 | |
| 311 | struct ErrorReporter; |
| 312 | struct ConfigErrorReporter; |
| 313 | struct DynamicErrorReporter; |
| 314 | struct WorkerDef; |
| 315 | kj::Promise<kj::Own<WorkerService>> makeWorkerImpl(kj::StringPtr name, |
| 316 | WorkerDef def, |
| 317 | capnp::List<config::Extension>::Reader extensions, |
| 318 | ErrorReporter& errorReporter); |
| 319 | |
| 320 | kj::Promise<void> startServices(jsg::V8System& v8System, |
| 321 | config::Config::Reader config, |
| 322 | kj::HttpHeaderTable::Builder& headerTableBuilder, |
| 323 | kj::ForkedPromise<void>& forkedDrainWhen); |
| 324 | |
| 325 | kj::Promise<void> listenOnSockets(config::Config::Reader config, |
| 326 | kj::HttpHeaderTable::Builder& headerTableBuilder, |
| 327 | kj::ForkedPromise<void>& forkedDrainWhen, |
| 328 | bool forTest = false); |
| 329 | |
| 330 | // Parsed socket protocol/TLS config. Extracted from the switch in listenOnSockets() to avoid |
| 331 | // goto-over-initialization inside a coroutine, which triggers a clang optimizer crash. |
| 332 | struct SocketTypeConfig { |
| 333 | uint defaultPort = 0; |
| 334 | config::HttpOptions::Reader httpOptions; |
| 335 | kj::Maybe<kj::Own<kj::TlsContext>> tls; |
| 336 | kj::StringPtr physicalProtocol; |
| 337 | }; |
| 338 | kj::Maybe<SocketTypeConfig> parseSocketType(config::Socket::Reader sock, kj::StringPtr name); |
| 339 | |
| 340 | void unlinkWorkerLoaders(); |
| 341 | |
| 342 | kj::Promise<void> preloadPython( |
| 343 | kj::StringPtr workerName, const WorkerDef& workerDef, ErrorReporter& errorReporter); |
| 344 | |
| 345 | friend struct FutureSubrequestChannel; |
| 346 | friend struct FutureActorClassChannel; |
| 347 | }; |
| 348 | |
| 349 | } // namespace workerd::server |