Skip to content
File

Blob: src/workerd/server/server.h

cpp350 lines
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 
20namespace kj {
21class TlsContext;
22}
23 
24namespace workerd::jsg {
25class V8System;
26}
27 
28namespace workerd::server {
29 
30using 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.
36class 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