Skip to content
File

Blob: src/workerd/server/server.c++

250.9 KB
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#include "server.h"
6 
7#include "alarm-scheduler.h"
8#include "container-client.h"
9#include "pyodide.h"
10#include "workerd-api.h"
11 
12#include <workerd/api/actor-state.h>
13#include <workerd/api/analytics-engine.capnp.h>
14#include <workerd/api/pyodide/pyodide.h>
15#include <workerd/api/trace.h>
16#include <workerd/api/worker-rpc.h>
17#include <workerd/io/actor-cache.h>
18#include <workerd/io/actor-id.h>
19#include <workerd/io/actor-sqlite.h>
20#include <workerd/io/bundle-fs.h>
21#include <workerd/io/compatibility-date.h>
22#include <workerd/io/container.capnp.h>
23#include <workerd/io/hibernation-manager.h>
24#include <workerd/io/io-context.h>
25#include <workerd/io/limit-enforcer.h>
26#include <workerd/io/request-tracker.h>
27#include <workerd/io/trace-stream.h>
28#include <workerd/io/worker-entrypoint.h>
29#include <workerd/io/worker-fs.h>
30#include <workerd/io/worker-interface.h>
31#include <workerd/io/worker.h>
32#include <workerd/server/actor-id-impl.h>
33#include <workerd/server/facet-tree-index.h>
34#include <workerd/server/fallback-service.h>
35#include <workerd/util/exception.h>
36#include <workerd/util/http-util.h>
37#include <workerd/util/mimetype.h>
38#include <workerd/util/stream-utils.h>
39#include <workerd/util/use-perfetto-categories.h>
40#include <workerd/util/uuid.h>
41#include <workerd/util/websocket-error-handler.h>
42 
43#include <openssl/bio.h>
44#include <openssl/pem.h>
45 
46#include <capnp/compat/json.h>
47#include <capnp/message.h>
48#include <capnp/rpc-twoparty.h>
49#include <kj/compat/http.h>
50#include <kj/compat/tls.h>
51#include <kj/compat/url.h>
52#include <kj/debug.h>
53#include <kj/encoding.h>
54#include <kj/glob-filter.h>
55#include <kj/map.h>
56 
57#include <cstdlib>
58#include <ctime>
59 
60namespace workerd::server {
61 
62namespace {
63 
64struct PemData {
65 kj::String type;
66 kj::Array<byte> data;
67};
68 
69// Decode PEM format using OpenSSL helpers.
70static kj::Maybe<PemData> decodePem(kj::ArrayPtr<const char> text) {
71 // TODO(cleanup): Should this be part of the KJ TLS library? We don't technically use it for TLS.
72 // Maybe KJ should have a general crypto library that wraps OpenSSL?
73 
74 BIO* bio = BIO_new_mem_buf(const_cast<char*>(text.begin()), text.size());
75 KJ_DEFER(BIO_free(bio));
76 
77 class OpenSslDisposer: public kj::ArrayDisposer {
78 public:
79 void disposeImpl(void* firstElement,
80 size_t elementSize,
81 size_t elementCount,
82 size_t capacity,
83 void (*destroyElement)(void*)) const override {
84 OPENSSL_free(firstElement);
85 }
86 };
87 static constexpr OpenSslDisposer disposer;
88 
89 char* namePtr = nullptr;
90 char* headerPtr = nullptr;
91 byte* dataPtr = nullptr;
92 long dataLen = 0;
93 if (!PEM_read_bio(bio, &namePtr, &headerPtr, &dataPtr, &dataLen)) {
94 return kj::none;
95 }
96 kj::Array<char> nameArr(namePtr, strlen(namePtr) + 1, disposer);
97 KJ_DEFER(OPENSSL_free(headerPtr));
98 kj::Array<kj::byte> data(dataPtr, dataLen, disposer);
99 
100 return PemData{kj::String(kj::mv(nameArr)), kj::mv(data)};
101}
102 
103// Returns a time string in the format HTTP likes to use.
104static kj::String httpTime(kj::Date date) {
105 time_t time = (date - kj::UNIX_EPOCH) / kj::SECONDS;
106#if _WIN32
107 // `gmtime` is thread-safe on Windows: https://learn.microsoft.com/en-us/cpp/c-runtime-library/reference/gmtime-gmtime32-gmtime64?view=msvc-170#return-value
108 auto tm = *gmtime(&time);
109#else
110 struct tm tm;
111 KJ_ASSERT(gmtime_r(&time, &tm) == &tm);
112#endif
113 char buf[256]{};
114 size_t n = strftime(buf, sizeof(buf), "%a, %d %b %Y %H:%M:%S GMT", &tm);
115 KJ_ASSERT(n > 0);
116 return kj::heapString(buf, n);
117}
118 
119static kj::String escapeJsonString(kj::StringPtr text) {
120 static const char HEXDIGITS[] = "0123456789abcdef";
121 kj::Vector<char> escaped(text.size() + 1);
122 
123 for (char c: text) {
124 switch (c) {
125 case '"':
126 escaped.addAll("\\\""_kj);
127 break;
128 case '\\':
129 escaped.addAll("\\\\"_kj);
130 break;
131 case '\b':
132 escaped.addAll("\\b"_kj);
133 break;
134 case '\f':
135 escaped.addAll("\\f"_kj);
136 break;
137 case '\n':
138 escaped.addAll("\\n"_kj);
139 break;
140 case '\r':
141 escaped.addAll("\\r"_kj);
142 break;
143 case '\t':
144 escaped.addAll("\\t"_kj);
145 break;
146 default:
147 if (static_cast<uint8_t>(c) < 0x20) {
148 escaped.addAll("\\u00"_kj);
149 uint8_t c2 = c;
150 escaped.add(HEXDIGITS[c2 / 16]);
151 escaped.add(HEXDIGITS[c2 % 16]);
152 } else {
153 escaped.add(c);
154 }
155 break;
156 }
157 }
158 
159 return kj::str("\"", escaped.releaseAsArray(), "\"");
160}
161 
162template <typename T>
163static inline kj::Own<T> fakeOwn(T& ref) {
164 return kj::Own<T>(&ref, kj::NullDisposer::instance);
165}
166 
167void throwDynamicEntrypointTransferError() {
168 JSG_FAIL_REQUIRE(DOMDataCloneError,
169 "Entrypoints to dynamically-loaded workers cannot be transferred to other Workers, "
170 "because the system does not know how to reload this Worker from scratch. Instead, "
171 "have the parent Worker expose an entrypoint which constructs the dynamic worker "
172 "and forwards to it.");
173}
174 
175} // namespace
176 
177// =======================================================================================
178 
179Server::Server(kj::Filesystem& fs,
180 kj::Timer& timer,
181 const kj::MonotonicClock& monotonicClock,
182 kj::Network& network,
183 kj::EntropySource& entropySource,
184 Worker::LoggingOptions loggingOptions,
185 kj::Function<void(kj::String)> reportConfigError)
186 : fs(fs),
187 timer(timer),
188 monotonicClock(monotonicClock),
189 network(network),
190 entropySource(entropySource),
191 reportConfigError(kj::mv(reportConfigError)),
192 loggingOptions(loggingOptions),
193 memoryCacheProvider(kj::heap<api::MemoryCacheProvider>(timer)),
194 channelTokenHandler(*this),
195 tasks(*this) {}
196 
197struct Server::GlobalContext {
198 jsg::V8System& v8System;
199 capnp::ByteStreamFactory byteStreamFactory;
200 capnp::HttpOverCapnpFactory httpOverCapnpFactory;
201 ThreadContext threadContext;
202 kj::HttpHeaderTable& headerTable;
203 
204 GlobalContext(
205 Server& server, jsg::V8System& v8System, kj::HttpHeaderTable::Builder& headerTableBuilder)
206 : v8System(v8System),
207 httpOverCapnpFactory(
208 byteStreamFactory, headerTableBuilder, capnp::HttpOverCapnpFactory::LEVEL_2),
209 threadContext(server.timer,
210 server.entropySource,
211 headerTableBuilder,
212 httpOverCapnpFactory,
213 byteStreamFactory),
214 headerTable(headerTableBuilder.getFutureTable()) {}
215};
216 
217class Server::Service: public IoChannelFactory::SubrequestChannel {
218 public:
219 // Cross-links this service with other services. Must be called once before `startRequest()`.
220 virtual void link(Worker::ValidationErrorReporter& errorReporter) {}
221 
222 // Drops any cross-links created during link(). This called just before all the services are
223 // destroyed. An `Own<T>` cannot be destroyed unless the object it points to still exists, so
224 // we must clear all the `Own<Service>`s before we can actually destroy the `Service`s.
225 virtual void unlink() {}
226 
227 // Begin an incoming request. Returns a `WorkerInterface` object that will be used for one
228 // request then discarded.
229 virtual kj::Own<WorkerInterface> startRequest(
230 IoChannelFactory::SubrequestMetadata metadata) override = 0;
231 
232 // Returns true if the service exports the given handler, e.g. `fetch`, `scheduled`, etc.
233 virtual bool hasHandler(kj::StringPtr handlerName) = 0;
234 
235 // Return the service itself, or the underlying service if this instance wraps another service as
236 // with EntrypointService.
237 virtual Service* service() {
238 return this;
239 }
240 
241 // Implemented by EntrypointService for loopback ctx.exports entrypoints, to allow props to be
242 // specified.
243 virtual kj::Own<Service> forProps(Frankenvalue props) {
244 KJ_FAIL_REQUIRE("can't override props for this service");
245 }
246 
247 void requireAllowsTransfer() override {
248 // We consider all `Service` implementations to be safe to transfer, except for dynamic workers
249 // which we'll handle explicitly.
250 }
251};
252 
253class Server::ActorClass: public IoChannelFactory::ActorClassChannel {
254 public:
255 // The caller must call this before calling newActor(). If it returns a promise, then the
256 // caller must await the promise before calling other methods.
257 //
258 // In particular, this is needed with dynamically-loaded workers. The isolate may still be
259 // loading when the caller calls `getDurableObjectClass()` on it.
260 virtual kj::Maybe<kj::Promise<void>> whenReady() {
261 return kj::none;
262 }
263 
264 // Construct a new instance of the class. The parameters here are passed into `Worker::Actor`'s
265 // constructor.
266 virtual kj::Own<Worker::Actor> newActor(kj::Maybe<RequestTracker&> tracker,
267 Worker::Actor::Id actorId,
268 Worker::Actor::MakeActorCacheFunc makeActorCache,
269 Worker::Actor::MakeStorageFunc makeStorage,
270 kj::Own<Worker::Actor::Loopback> loopback,
271 kj::Maybe<kj::Own<Worker::Actor::HibernationManager>> manager,
272 kj::Maybe<rpc::Container::Client> container,
273 kj::Maybe<Worker::Actor::FacetManager&> facetManager) = 0;
274 
275 // Start a request on the actor. (The actor must have been created using newActor().)
276 virtual kj::Own<WorkerInterface> startRequest(
277 IoChannelFactory::SubrequestMetadata metadata, kj::Own<Worker::Actor> actor) = 0;
278 
279 virtual kj::Own<ActorClass> forProps(Frankenvalue props) {
280 KJ_FAIL_REQUIRE("can't override props for this actor class");
281 }
282};
283 
284Server::~Server() noexcept {
285 // This destructor is explicitly `noexcept` because if one of the `unlink()`s throws then we'd
286 // have a hard time avoiding a segfault later... and we're shutting down the server anyway so
287 // whatever, better to crash.
288 
289 // It's important to cancel all tasks before we start tearing down.
290 tasks.clear();
291 
292 // Unlink all the services, which should remove all refcount cycles.
293 unlinkWorkerLoaders();
294 for (auto& service: services) {
295 service.value->unlink();
296 }
297 
298 // Verify that unlinking actually eliminated cycles. Otherwise we have a memory leak -- and
299 // potentially use-after-free if we allow the `Server` to be destroyed while services still
300 // exist.
301 for (auto& service: services) {
302 KJ_ASSERT(
303 !service.value->isShared(), "service still has references after unlinking", service.key);
304 }
305}
306 
307// =======================================================================================
308 
309kj::Own<kj::TlsContext> Server::makeTlsContext(config::TlsOptions::Reader conf) {
310 kj::TlsContext::Options options;
311 
312 struct Attachments {
313 kj::Maybe<kj::TlsKeypair> keypair;
314 kj::Array<kj::TlsCertificate> trustedCerts;
315 };
316 auto attachments = kj::heap<Attachments>();
317 
318 if (conf.hasKeypair()) {
319 auto pairConf = conf.getKeypair();
320 options.defaultKeypair = attachments->keypair.emplace(
321 kj::TlsKeypair{.privateKey = kj::TlsPrivateKey(pairConf.getPrivateKey()),
322 .certificate = kj::TlsCertificate(pairConf.getCertificateChain())});
323 }
324 
325 options.verifyClients = conf.getRequireClientCerts();
326 options.useSystemTrustStore = conf.getTrustBrowserCas();
327 
328 auto trustList = conf.getTrustedCertificates();
329 if (trustList.size() > 0) {
330 attachments->trustedCerts = KJ_MAP(cert, trustList) { return kj::TlsCertificate(cert); };
331 options.trustedCertificates = attachments->trustedCerts;
332 }
333 
334 switch (conf.getMinVersion()) {
335 case config::TlsOptions::Version::GOOD_DEFAULT:
336 // Don't change.
337 goto validVersion;
338 case config::TlsOptions::Version::SSL3:
339 options.minVersion = kj::TlsVersion::SSL_3;
340 goto validVersion;
341 case config::TlsOptions::Version::TLS1_DOT0:
342 options.minVersion = kj::TlsVersion::TLS_1_0;
343 goto validVersion;
344 case config::TlsOptions::Version::TLS1_DOT1:
345 options.minVersion = kj::TlsVersion::TLS_1_1;
346 goto validVersion;
347 case config::TlsOptions::Version::TLS1_DOT2:
348 options.minVersion = kj::TlsVersion::TLS_1_2;
349 goto validVersion;
350 case config::TlsOptions::Version::TLS1_DOT3:
351 options.minVersion = kj::TlsVersion::TLS_1_3;
352 goto validVersion;
353 }
354 reportConfigError(kj::str("Encountered unknown TlsOptions::minVersion setting. Was the "
355 "config compiled with a newer version of the schema?"));
356 
357validVersion:
358 if (conf.hasCipherList()) {
359 options.cipherList = conf.getCipherList();
360 }
361 
362 return kj::heap<kj::TlsContext>(kj::mv(options));
363}
364 
365kj::Promise<kj::Own<kj::NetworkAddress>> Server::makeTlsNetworkAddress(
366 config::TlsOptions::Reader conf,
367 kj::StringPtr addrStr,
368 kj::Maybe<kj::StringPtr> certificateHost,
369 uint defaultPort) {
370 auto context = makeTlsContext(conf);
371 
372 KJ_IF_SOME(h, certificateHost) {
373 auto parsed = co_await network.parseAddress(addrStr, defaultPort);
374 co_return context->wrapAddress(kj::mv(parsed), h).attach(kj::mv(context));
375 }
376 
377 // Wrap the `Network` itself so we can use the TLS implementation's `parseAddress()` to extract
378 // the authority from the address.
379 auto tlsNetwork = context->wrapNetwork(network);
380 auto parsed = co_await network.parseAddress(addrStr, defaultPort);
381 co_return parsed.attach(kj::mv(context));
382}
383 
384// =======================================================================================
385 
386// Helper to apply config::HttpOptions.
387class Server::HttpRewriter {
388 // TODO(beta): Do we want to automatically add `Date`, `Server` (to outgoing responses),
389 // `User-Agent` (to outgoing requests), etc.?
390 
391 public:
392 HttpRewriter(
393 config::HttpOptions::Reader httpOptions, kj::HttpHeaderTable::Builder& headerTableBuilder)
394 : style(httpOptions.getStyle()),
395 requestInjector(httpOptions.getInjectRequestHeaders(), headerTableBuilder),
396 responseInjector(httpOptions.getInjectResponseHeaders(), headerTableBuilder) {
397 if (httpOptions.hasForwardedProtoHeader()) {
398 forwardedProtoHeader = headerTableBuilder.add(httpOptions.getForwardedProtoHeader());
399 }
400 if (httpOptions.hasCfBlobHeader()) {
401 cfBlobHeader = headerTableBuilder.add(httpOptions.getCfBlobHeader());
402 }
403 if (httpOptions.hasCapnpConnectHost()) {
404 capnpConnectHost = httpOptions.getCapnpConnectHost();
405 }
406 }
407 
408 bool hasCfBlobHeader() {
409 return cfBlobHeader != kj::none;
410 }
411 
412 bool needsRewriteRequest() {
413 return style == config::HttpOptions::Style::HOST || hasCfBlobHeader() ||
414 !requestInjector.empty();
415 }
416 
417 // Attach this to the promise returned by request().
418 struct Rewritten {
419 kj::Own<kj::HttpHeaders> headers;
420 kj::String ownUrl;
421 };
422 
423 Rewritten rewriteOutgoingRequest(
424 kj::StringPtr& url, const kj::HttpHeaders& headers, kj::Maybe<kj::StringPtr> cfBlobJson) {
425 Rewritten result{kj::heap(headers.cloneShallow()), nullptr};
426 
427 if (style == config::HttpOptions::Style::HOST) {
428 auto parsed = kj::Url::parse(url, kj::Url::HTTP_PROXY_REQUEST,
429 kj::Url::Options{.percentDecode = false, .allowEmpty = true});
430 result.headers->set(kj::HttpHeaderId::HOST, kj::mv(parsed.host));
431 KJ_IF_SOME(h, forwardedProtoHeader) {
432 result.headers->set(h, kj::mv(parsed.scheme));
433 }
434 url = result.ownUrl = parsed.toString(kj::Url::HTTP_REQUEST);
435 }
436 
437 KJ_IF_SOME(h, cfBlobHeader) {
438 KJ_IF_SOME(b, cfBlobJson) {
439 result.headers->setPtr(h, b);
440 } else {
441 result.headers->unset(h);
442 }
443 }
444 
445 requestInjector.apply(*result.headers);
446 
447 return result;
448 }
449 
450 kj::Maybe<Rewritten> rewriteIncomingRequest(kj::StringPtr& url,
451 kj::StringPtr physicalProtocol,
452 const kj::HttpHeaders& headers,
453 kj::Maybe<kj::String>& cfBlobJson) {
454 Rewritten result{kj::heap(headers.cloneShallow()), nullptr};
455 
456 if (style == config::HttpOptions::Style::HOST) {
457 auto parsed = kj::Url::parse(
458 url, kj::Url::HTTP_REQUEST, kj::Url::Options{.percentDecode = false, .allowEmpty = true});
459 parsed.host = kj::str(KJ_UNWRAP_OR_RETURN(headers.get(kj::HttpHeaderId::HOST), kj::none));
460 
461 KJ_IF_SOME(h, forwardedProtoHeader) {
462 KJ_IF_SOME(s, headers.get(h)) {
463 parsed.scheme = kj::str(s);
464 result.headers->unset(h);
465 }
466 }
467 
468 if (parsed.scheme == nullptr) parsed.scheme = kj::str(physicalProtocol);
469 
470 url = result.ownUrl = parsed.toString(kj::Url::HTTP_PROXY_REQUEST);
471 }
472 
473 KJ_IF_SOME(h, cfBlobHeader) {
474 KJ_IF_SOME(b, headers.get(h)) {
475 cfBlobJson = kj::str(b);
476 result.headers->unset(h);
477 }
478 }
479 
480 requestInjector.apply(*result.headers);
481 
482 return result;
483 }
484 
485 bool needsRewriteResponse() {
486 return !responseInjector.empty();
487 }
488 
489 void rewriteResponse(kj::HttpHeaders& headers) {
490 responseInjector.apply(headers);
491 }
492 
493 kj::Maybe<kj::StringPtr> getCapnpConnectHost() {
494 return capnpConnectHost;
495 }
496 
497 private:
498 config::HttpOptions::Style style;
499 kj::Maybe<kj::HttpHeaderId> forwardedProtoHeader;
500 kj::Maybe<kj::HttpHeaderId> cfBlobHeader;
501 kj::Maybe<kj::StringPtr> capnpConnectHost;
502 
503 class HeaderInjector {
504 public:
505 HeaderInjector(capnp::List<config::HttpOptions::Header>::Reader headers,
506 kj::HttpHeaderTable::Builder& headerTableBuilder)
507 : injectedHeaders(KJ_MAP(header, headers) {
508 InjectedHeader result;
509 result.id = headerTableBuilder.add(header.getName());
510 if (header.hasValue()) {
511 result.value = kj::str(header.getValue());
512 }
513 return result;
514 }) {}
515 
516 bool empty() {
517 return injectedHeaders.size() == 0;
518 }
519 
520 void apply(kj::HttpHeaders& headers) {
521 for (auto& header: injectedHeaders) {
522 KJ_IF_SOME(v, header.value) {
523 headers.setPtr(header.id, v);
524 } else {
525 headers.unset(header.id);
526 }
527 }
528 }
529 
530 private:
531 struct InjectedHeader {
532 kj::HttpHeaderId id;
533 kj::Maybe<kj::String> value;
534 };
535 kj::Array<InjectedHeader> injectedHeaders;
536 };
537 
538 HeaderInjector requestInjector;
539 HeaderInjector responseInjector;
540};
541 
542// =======================================================================================
543 
544// Service used when the service's config is invalid.
545class Server::InvalidConfigService final: public Service {
546 public:
547 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override {
548 JSG_FAIL_REQUIRE(Error, "Service cannot handle requests because its config is invalid.");
549 }
550 
551 bool hasHandler(kj::StringPtr handlerName) override {
552 return false;
553 }
554};
555 
556class Server::InvalidConfigActorClass final: public ActorClass {
557 public:
558 void requireAllowsTransfer() override {
559 // Can't get here because workerd would have failed to start.
560 KJ_UNREACHABLE;
561 }
562 
563 kj::Own<Worker::Actor> newActor(kj::Maybe<RequestTracker&> tracker,
564 Worker::Actor::Id actorId,
565 Worker::Actor::MakeActorCacheFunc makeActorCache,
566 Worker::Actor::MakeStorageFunc makeStorage,
567 kj::Own<Worker::Actor::Loopback> loopback,
568 kj::Maybe<kj::Own<Worker::Actor::HibernationManager>> manager,
569 kj::Maybe<rpc::Container::Client> container,
570 kj::Maybe<Worker::Actor::FacetManager&> facetManager) override {
571 JSG_FAIL_REQUIRE(
572 Error, "Cannot instantiate Durable Object class because its config is invalid.");
573 }
574 
575 kj::Own<WorkerInterface> startRequest(
576 IoChannelFactory::SubrequestMetadata metadata, kj::Own<Worker::Actor> actor) override {
577 // Can't get here because creating the actor would have required calling the other method.
578 KJ_UNREACHABLE;
579 }
580};
581 
582// Return a fake Own pointing to the singleton.
583kj::Own<Server::Service> Server::makeInvalidConfigService() {
584 return {invalidConfigServiceSingleton.get(), kj::NullDisposer::instance};
585}
586 
587// A NetworkAddress whose connect() method waits for a Promise<NetworkAddress> and then forwards
588// to it. Used by ExternalHttpService so that we don't have to wait for DNS lookup before the
589// server can start.
590class PromisedNetworkAddress final: public kj::NetworkAddress {
591 // TODO(cleanup): kj::Network should be extended with a new version of parseAddress() which does
592 // not do DNS lookup immediately, and therefore can return a NetworkAddress synchronously.
593 // In fact, this version should be designed to redo the DNS lookup periodically to see if it
594 // changed, which would be nice for workerd when the remote address may change over time.
595 public:
596 PromisedNetworkAddress(kj::Promise<kj::Own<kj::NetworkAddress>> promise)
597 : promise(promise.then([this](kj::Own<kj::NetworkAddress> result) { addr = kj::mv(result); })
598 .fork()) {}
599 
600 kj::Promise<kj::Own<kj::AsyncIoStream>> connect() override {
601 KJ_IF_SOME(a, addr) {
602 co_return co_await a.get()->connect();
603 } else {
604 co_await promise;
605 co_return co_await KJ_ASSERT_NONNULL(addr)->connect();
606 }
607 }
608 
609 kj::Promise<kj::AuthenticatedStream> connectAuthenticated() override {
610 KJ_IF_SOME(a, addr) {
611 co_return co_await a.get()->connectAuthenticated();
612 } else {
613 co_await promise;
614 co_return co_await KJ_ASSERT_NONNULL(addr)->connectAuthenticated();
615 }
616 }
617 
618 // We don't use any other methods, and they seem kinda annoying to implement.
619 kj::Own<kj::ConnectionReceiver> listen() override {
620 KJ_UNIMPLEMENTED("PromisedNetworkAddress::listen() not implemented");
621 }
622 kj::Own<kj::NetworkAddress> clone() override {
623 KJ_UNIMPLEMENTED("PromisedNetworkAddress::clone() not implemented");
624 }
625 kj::String toString() override {
626 KJ_UNIMPLEMENTED("PromisedNetworkAddress::toString() not implemented");
627 }
628 
629 private:
630 kj::ForkedPromise<void> promise;
631 kj::Maybe<kj::Own<kj::NetworkAddress>> addr;
632};
633 
634class Server::ExternalTcpService final: public Service, private WorkerInterface {
635 public:
636 ExternalTcpService(kj::Own<kj::NetworkAddress> addrParam): addr(kj::mv(addrParam)) {}
637 
638 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override {
639 return {this, kj::NullDisposer::instance};
640 }
641 
642 bool hasHandler(kj::StringPtr handlerName) override {
643 return handlerName == "fetch"_kj || handlerName == "connect"_kj;
644 }
645 
646 private:
647 kj::Own<kj::NetworkAddress> addr;
648 
649 kj::Promise<void> request(kj::HttpMethod method,
650 kj::StringPtr url,
651 const kj::HttpHeaders& headers,
652 kj::AsyncInputStream& requestBody,
653 kj::HttpService::Response& response) override {
654 throwUnsupported();
655 }
656 
657 kj::Promise<void> connect(kj::StringPtr host,
658 const kj::HttpHeaders& headers,
659 kj::AsyncIoStream& connection,
660 ConnectResponse& tunnel,
661 kj::HttpConnectSettings settings) override {
662 TRACE_EVENT("workerd", "ExternalTcpService::connect()", "host", host.cStr());
663 auto io_stream = co_await addr->connect();
664 
665 auto promises = kj::heapArrayBuilder<kj::Promise<void>>(2);
666 
667 promises.add(connection.pumpTo(*io_stream).then([&io_stream = *io_stream](uint64_t size) {
668 io_stream.shutdownWrite();
669 }));
670 
671 promises.add(io_stream->pumpTo(connection).then([&connection](uint64_t size) {
672 connection.shutdownWrite();
673 }));
674 
675 tunnel.accept(200, "OK", kj::HttpHeaders(kj::HttpHeaderTable{}));
676 
677 co_await kj::joinPromisesFailFast(promises.finish()).attach(kj::mv(io_stream));
678 }
679 
680 kj::Promise<void> prewarm(kj::StringPtr url) override {
681 return kj::READY_NOW;
682 }
683 kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override {
684 throwUnsupported();
685 }
686 kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override {
687 throwUnsupported();
688 }
689 kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override {
690 return event->notSupported();
691 }
692 
693 [[noreturn]] void throwUnsupported() {
694 JSG_FAIL_REQUIRE(Error, "External TCP servers don't support this event type.");
695 }
696};
697 
698// Service used when the service is configured as external HTTP service.
699class Server::ExternalHttpService final: public Service {
700 public:
701 ExternalHttpService(kj::Own<kj::NetworkAddress> addrParam,
702 kj::Own<HttpRewriter> rewriter,
703 kj::HttpHeaderTable& headerTable,
704 kj::Timer& timer,
705 kj::EntropySource& entropySource,
706 capnp::ByteStreamFactory& byteStreamFactory,
707 capnp::HttpOverCapnpFactory& httpOverCapnpFactory)
708 : addr(kj::mv(addrParam)),
709 webSocketErrorHandler(kj::heap<JsgifyWebSocketErrors>()),
710 inner(kj::newHttpClient(timer,
711 headerTable,
712 *addr,
713 {.entropySource = entropySource,
714 .webSocketCompressionMode = kj::HttpClientSettings::MANUAL_COMPRESSION,
715 .webSocketErrorHandler = *webSocketErrorHandler})),
716 serviceAdapter(kj::newHttpService(*inner)),
717 rewriter(kj::mv(rewriter)),
718 headerTable(headerTable),
719 byteStreamFactory(byteStreamFactory),
720 httpOverCapnpFactory(httpOverCapnpFactory) {}
721 
722 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override {
723 return kj::heap<WorkerInterfaceImpl>(*this, kj::mv(metadata));
724 }
725 
726 bool hasHandler(kj::StringPtr handlerName) override {
727 return handlerName == "fetch"_kj || handlerName == "connect"_kj;
728 }
729 
730 private:
731 kj::Own<kj::NetworkAddress> addr;
732 
733 kj::Own<JsgifyWebSocketErrors> webSocketErrorHandler;
734 kj::Own<kj::HttpClient> inner;
735 kj::Own<kj::HttpService> serviceAdapter;
736 
737 kj::Own<HttpRewriter> rewriter;
738 
739 kj::HttpHeaderTable& headerTable;
740 capnp::ByteStreamFactory& byteStreamFactory;
741 capnp::HttpOverCapnpFactory& httpOverCapnpFactory;
742 
743 struct CapnpClient {
744 kj::Own<kj::AsyncIoStream> connection;
745 capnp::TwoPartyClient rpcSystem;
746 
747 CapnpClient(kj::Own<kj::AsyncIoStream> connectionParam)
748 : connection(kj::mv(connectionParam)),
749 rpcSystem(*connection) {}
750 };
751 
752 // capnpClient is created on-demand when RPC is needed.
753 kj::Maybe<CapnpClient> capnpClient;
754 
755 // This task nulls out `capnpClient` when the connection is lost.
756 kj::Promise<void> clearCapnpClientTask = nullptr;
757 
758 // Get an WorkerdBootstrap representing the service on the other end of an HTTP connection. May
759 // reuse an existing connection, or form a new one over `client`.
760 rpc::WorkerdBootstrap::Client getOutgoingCapnp(kj::HttpClient& client) {
761 KJ_IF_SOME(c, capnpClient) {
762 return c.rpcSystem.bootstrap().castAs<rpc::WorkerdBootstrap>();
763 }
764 
765 // No existing client, need to create a new one.
766 kj::StringPtr host = KJ_UNWRAP_OR(rewriter->getCapnpConnectHost(),
767 { return JSG_KJ_EXCEPTION(FAILED, Error, "This ExternalServer not configured for RPC."); });
768 
769 auto req = client.connect(host, kj::HttpHeaders(headerTable), {});
770 auto& c = capnpClient.emplace(kj::mv(req.connection));
771 
772 // Arrange that when the connection is lost, we'll null out `capnpClient`. This ensures that
773 // on the next event, we'll attempt to reconnect.
774 //
775 // TODO(perf): Time out idle connections?
776 clearCapnpClientTask =
777 c.rpcSystem.onDisconnect().attach(kj::defer([this]() {
778 capnpClient = kj::none;
779 })).eagerlyEvaluate(nullptr);
780 
781 return c.rpcSystem.bootstrap().castAs<rpc::WorkerdBootstrap>();
782 }
783 
784 class WorkerInterfaceImpl final: public WorkerInterface, private kj::HttpService::Response {
785 public:
786 WorkerInterfaceImpl(ExternalHttpService& parent, IoChannelFactory::SubrequestMetadata metadata)
787 : parent(kj::addRef(parent)),
788 metadata(kj::mv(metadata)) {}
789 
790 kj::Promise<void> request(kj::HttpMethod method,
791 kj::StringPtr url,
792 const kj::HttpHeaders& headers,
793 kj::AsyncInputStream& requestBody,
794 kj::HttpService::Response& response) override {
795 TRACE_EVENT("workerd", "ExternalHttpServer::request()");
796 KJ_REQUIRE(wrappedResponse == kj::none, "object should only receive one request");
797 wrappedResponse = response;
798 if (parent->rewriter->needsRewriteRequest()) {
799 auto rewrite = parent->rewriter->rewriteOutgoingRequest(url, headers, metadata.cfBlobJson);
800 return parent->serviceAdapter->request(method, url, *rewrite.headers, requestBody, *this)
801 .attach(kj::mv(rewrite));
802 } else {
803 return parent->serviceAdapter->request(method, url, headers, requestBody, *this);
804 }
805 }
806 
807 kj::Promise<void> connect(kj::StringPtr host,
808 const kj::HttpHeaders& headers,
809 kj::AsyncIoStream& connection,
810 ConnectResponse& tunnel,
811 kj::HttpConnectSettings settings) override {
812 TRACE_EVENT("workerd", "ExternalHttpServer::connect()");
813 return parent->serviceAdapter->connect(host, headers, connection, tunnel, kj::mv(settings));
814 }
815 
816 kj::Promise<void> prewarm(kj::StringPtr url) override {
817 return kj::READY_NOW;
818 }
819 kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override {
820 throwUnsupported();
821 }
822 kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override {
823 throwUnsupported();
824 }
825 
826 kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override {
827 // We'll use capnp RPC for custom events.
828 auto bootstrap = parent->getOutgoingCapnp(*parent->inner);
829 auto dispatcher =
830 bootstrap.startEventRequest(capnp::MessageSize{4, 0}).send().getDispatcher();
831 return event
832 ->sendRpc(parent->httpOverCapnpFactory, parent->byteStreamFactory, kj::mv(dispatcher))
833 .attach(kj::mv(event));
834 }
835 
836 private:
837 kj::Own<ExternalHttpService> parent;
838 IoChannelFactory::SubrequestMetadata metadata;
839 kj::Maybe<kj::HttpService::Response&> wrappedResponse;
840 
841 [[noreturn]] void throwUnsupported() {
842 JSG_FAIL_REQUIRE(Error, "External HTTP servers don't support this event type.");
843 }
844 
845 kj::Own<kj::AsyncOutputStream> send(uint statusCode,
846 kj::StringPtr statusText,
847 const kj::HttpHeaders& headers,
848 kj::Maybe<uint64_t> expectedBodySize) override {
849 TRACE_EVENT("workerd", "ExternalHttpService::send()", "status", statusCode);
850 auto& response = KJ_ASSERT_NONNULL(wrappedResponse);
851 if (parent->rewriter->needsRewriteResponse()) {
852 auto rewrite = headers.cloneShallow();
853 parent->rewriter->rewriteResponse(rewrite);
854 return response.send(statusCode, statusText, rewrite, expectedBodySize);
855 } else {
856 return response.send(statusCode, statusText, headers, expectedBodySize);
857 }
858 }
859 
860 kj::Own<kj::WebSocket> acceptWebSocket(const kj::HttpHeaders& headers) override {
861 TRACE_EVENT("workerd", "ExternalHttpService::acceptWebSocket()");
862 auto& response = KJ_ASSERT_NONNULL(wrappedResponse);
863 if (parent->rewriter->needsRewriteResponse()) {
864 auto rewrite = headers.cloneShallow();
865 parent->rewriter->rewriteResponse(rewrite);
866 return response.acceptWebSocket(rewrite);
867 } else {
868 return response.acceptWebSocket(headers);
869 }
870 }
871 };
872};
873 
874kj::Own<Server::Service> Server::makeExternalService(kj::StringPtr name,
875 config::ExternalServer::Reader conf,
876 kj::HttpHeaderTable::Builder& headerTableBuilder) {
877 TRACE_EVENT("workerd", "Server::makeExternalService()", "name", name.cStr());
878 kj::StringPtr addrStr = nullptr;
879 kj::String ownAddrStr = nullptr;
880 
881 KJ_IF_SOME(override, externalOverrides.findEntry(name)) {
882 addrStr = ownAddrStr = kj::mv(override.value);
883 externalOverrides.erase(override);
884 } else if (conf.hasAddress()) {
885 addrStr = conf.getAddress();
886 } else {
887 reportConfigError(kj::str("External service \"", name,
888 "\" has no address in the config, so must be specified "
889 "on the command line with `--external-addr`."));
890 return makeInvalidConfigService();
891 }
892 
893 switch (conf.which()) {
894 case config::ExternalServer::HTTP: {
895 // We have to construct the rewriter upfront before waiting on any promises, since the
896 // HeaderTable::Builder is only available synchronously.
897 auto rewriter = kj::heap<HttpRewriter>(conf.getHttp(), headerTableBuilder);
898 auto addr = kj::heap<PromisedNetworkAddress>(network.parseAddress(addrStr, 80));
899 return kj::refcounted<ExternalHttpService>(kj::mv(addr), kj::mv(rewriter),
900 headerTableBuilder.getFutureTable(), timer, entropySource,
901 globalContext->byteStreamFactory, globalContext->httpOverCapnpFactory);
902 }
903 case config::ExternalServer::HTTPS: {
904 auto httpsConf = conf.getHttps();
905 kj::Maybe<kj::StringPtr> certificateHost;
906 if (httpsConf.hasCertificateHost()) {
907 certificateHost = httpsConf.getCertificateHost();
908 }
909 auto rewriter = kj::heap<HttpRewriter>(httpsConf.getOptions(), headerTableBuilder);
910 auto addr = kj::heap<PromisedNetworkAddress>(
911 makeTlsNetworkAddress(httpsConf.getTlsOptions(), addrStr, certificateHost, 443));
912 return kj::refcounted<ExternalHttpService>(kj::mv(addr), kj::mv(rewriter),
913 headerTableBuilder.getFutureTable(), timer, entropySource,
914 globalContext->byteStreamFactory, globalContext->httpOverCapnpFactory);
915 }
916 case config::ExternalServer::TCP: {
917 auto tcpConf = conf.getTcp();
918 auto addr = kj::heap<PromisedNetworkAddress>(network.parseAddress(addrStr, 80));
919 if (tcpConf.hasTlsOptions()) {
920 kj::Maybe<kj::StringPtr> certificateHost;
921 if (tcpConf.hasCertificateHost()) {
922 certificateHost = tcpConf.getCertificateHost();
923 }
924 addr = kj::heap<PromisedNetworkAddress>(
925 makeTlsNetworkAddress(tcpConf.getTlsOptions(), addrStr, certificateHost, 0));
926 }
927 return kj::refcounted<ExternalTcpService>(kj::mv(addr));
928 }
929 }
930 reportConfigError(kj::str("External service named \"", name,
931 "\" has unrecognized protocol. Was the config "
932 "compiled with a newer version of the schema?"));
933 return makeInvalidConfigService();
934}
935 
936// Service used when the service is configured as network service.
937class Server::NetworkService final: public Service, private WorkerInterface {
938 public:
939 NetworkService(kj::HttpHeaderTable& headerTable,
940 kj::Timer& timer,
941 kj::EntropySource& entropySource,
942 kj::Own<kj::Network> networkParam,
943 kj::Maybe<kj::Own<kj::Network>> tlsNetworkParam,
944 kj::Maybe<kj::SecureNetworkWrapper&> tlsContext)
945 : network(kj::mv(networkParam)),
946 tlsNetwork(kj::mv(tlsNetworkParam)),
947 webSocketErrorHandler(kj::heap<JsgifyWebSocketErrors>()),
948 inner(kj::newHttpClient(timer,
949 headerTable,
950 *network,
951 tlsNetwork,
952 {.entropySource = entropySource,
953 .webSocketCompressionMode = kj::HttpClientSettings::MANUAL_COMPRESSION,
954 .webSocketErrorHandler = *webSocketErrorHandler,
955 .tlsContext = tlsContext})),
956 serviceAdapter(kj::newHttpService(*inner)) {}
957 
958 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override {
959 return {this, kj::NullDisposer::instance};
960 }
961 
962 bool hasHandler(kj::StringPtr handlerName) override {
963 return handlerName == "fetch"_kj || handlerName == "connect"_kj;
964 }
965 
966 private:
967 kj::Own<kj::Network> network;
968 kj::Maybe<kj::Own<kj::Network>> tlsNetwork;
969 kj::Own<JsgifyWebSocketErrors> webSocketErrorHandler;
970 kj::Own<kj::HttpClient> inner;
971 kj::Own<kj::HttpService> serviceAdapter;
972 
973 kj::Promise<void> request(kj::HttpMethod method,
974 kj::StringPtr url,
975 const kj::HttpHeaders& headers,
976 kj::AsyncInputStream& requestBody,
977 kj::HttpService::Response& response) override {
978 TRACE_EVENT("workerd", "NetworkService::request()");
979 return serviceAdapter->request(method, url, headers, requestBody, response);
980 }
981 
982 kj::Promise<void> connect(kj::StringPtr host,
983 const kj::HttpHeaders& headers,
984 kj::AsyncIoStream& connection,
985 ConnectResponse& tunnel,
986 kj::HttpConnectSettings settings) override {
987 TRACE_EVENT("workerd", "NetworkService::connect()");
988 // This code is hit when the global `connect` function is called in a JS worker script.
989 // It represents a proxy-less TCP connection, which means we can simply defer the handling of
990 // the connection to the service adapter (likely NetworkHttpClient). Its behavior will be to
991 // connect directly to the host over TCP.
992 return serviceAdapter->connect(host, headers, connection, tunnel, kj::mv(settings));
993 }
994 
995 kj::Promise<void> prewarm(kj::StringPtr url) override {
996 return kj::READY_NOW;
997 }
998 kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override {
999 throwUnsupported();
1000 }
1001 kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override {
1002 throwUnsupported();
1003 }
1004 kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override {
1005 return event->notSupported();
1006 }
1007 
1008 [[noreturn]] void throwUnsupported() {
1009 JSG_FAIL_REQUIRE(Error, "External HTTP servers don't support this event type.");
1010 }
1011};
1012 
1013kj::Own<Server::Service> Server::makeNetworkService(config::Network::Reader conf) {
1014 TRACE_EVENT("workerd", "Server::makeNetworkService()");
1015 auto restrictedNetwork = network.restrictPeers( KJ_MAP(a, conf.getAllow()) -> kj::StringPtr {
1016 return a;
1017 }, KJ_MAP(a, conf.getDeny()) -> kj::StringPtr { return a; });
1018 
1019 kj::Maybe<kj::Own<kj::Network>> tlsNetwork;
1020 kj::Maybe<kj::SecureNetworkWrapper&> tlsContext;
1021 if (conf.hasTlsOptions()) {
1022 auto ownedTlsContext = makeTlsContext(conf.getTlsOptions());
1023 tlsContext = ownedTlsContext;
1024 tlsNetwork = ownedTlsContext->wrapNetwork(*restrictedNetwork).attach(kj::mv(ownedTlsContext));
1025 }
1026 
1027 return kj::refcounted<NetworkService>(globalContext->headerTable, timer, entropySource,
1028 kj::mv(restrictedNetwork), kj::mv(tlsNetwork), tlsContext);
1029}
1030 
1031// Service used when the service is configured as disk directory service.
1032class Server::DiskDirectoryService final: public Service, private WorkerInterface {
1033 public:
1034 DiskDirectoryService(config::DiskDirectory::Reader conf,
1035 kj::Own<const kj::Directory> dir,
1036 kj::HttpHeaderTable::Builder& headerTableBuilder)
1037 : writable(*dir),
1038 readable(kj::mv(dir)),
1039 headerTable(headerTableBuilder.getFutureTable()),
1040 hLastModified(headerTableBuilder.add("Last-Modified")),
1041 allowDotfiles(conf.getAllowDotfiles()) {}
1042 DiskDirectoryService(config::DiskDirectory::Reader conf,
1043 kj::Own<const kj::ReadableDirectory> dir,
1044 kj::HttpHeaderTable::Builder& headerTableBuilder)
1045 : readable(kj::mv(dir)),
1046 headerTable(headerTableBuilder.getFutureTable()),
1047 hLastModified(headerTableBuilder.add("Last-Modified")),
1048 allowDotfiles(conf.getAllowDotfiles()) {}
1049 
1050 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override {
1051 return {this, kj::NullDisposer::instance};
1052 }
1053 
1054 kj::Maybe<const kj::Directory&> getWritable() {
1055 return writable;
1056 }
1057 
1058 bool hasHandler(kj::StringPtr handlerName) override {
1059 return handlerName == "fetch"_kj;
1060 }
1061 
1062 private:
1063 kj::Maybe<const kj::Directory&> writable;
1064 kj::Own<const kj::ReadableDirectory> readable;
1065 kj::HttpHeaderTable& headerTable;
1066 kj::HttpHeaderId hLastModified;
1067 bool allowDotfiles;
1068 
1069 kj::Promise<void> request(kj::HttpMethod method,
1070 kj::StringPtr urlStr,
1071 const kj::HttpHeaders& requestHeaders,
1072 kj::AsyncInputStream& requestBody,
1073 kj::HttpService::Response& response) override {
1074 TRACE_EVENT("workerd", "DiskDirectoryService::request()", "url", urlStr.cStr());
1075 auto url = kj::Url::parse(urlStr);
1076 
1077 bool blockedPath = false;
1078 kj::Path path = nullptr;
1079 KJ_IF_SOME(exception,
1080 kj::runCatchingExceptions([&]() { path = kj::Path(url.path.releaseAsArray()); })) {
1081 (void)exception; // squash compiler warning about unused var
1082 // If the Path constructor throws, this path is not valid (e.g. it contains "..").
1083 blockedPath = true;
1084 }
1085 
1086 if (!blockedPath && !allowDotfiles) {
1087 for (auto& part: path) {
1088 if (part.startsWith(".")) {
1089 blockedPath = true;
1090 break;
1091 }
1092 }
1093 }
1094 
1095 if (method == kj::HttpMethod::GET || method == kj::HttpMethod::HEAD) {
1096 if (blockedPath) {
1097 co_return co_await response.sendError(404, "Not Found", headerTable);
1098 }
1099 
1100 auto file = KJ_UNWRAP_OR(readable->tryOpenFile(path),
1101 { co_return co_await response.sendError(404, "Not Found", headerTable); });
1102 
1103 auto meta = file->stat();
1104 
1105 switch (meta.type) {
1106 case kj::FsNode::Type::FILE: {
1107 // If this is a GET request with a Range header, return partial content if a single
1108 // satisfiable range is specified.
1109 // TODO(someday): consider supporting multiple ranges with multipart/byteranges
1110 kj::Maybe<kj::HttpByteRange> range;
1111 if (method == kj::HttpMethod::GET) {
1112 KJ_IF_SOME(header, requestHeaders.get(kj::HttpHeaderId::RANGE)) {
1113 KJ_SWITCH_ONEOF(kj::tryParseHttpRangeHeader(header.asArray(), meta.size)) {
1114 KJ_CASE_ONEOF(ranges, kj::Array<kj::HttpByteRange>) {
1115 KJ_ASSERT(ranges.size() > 0);
1116 if (ranges.size() == 1) range = ranges[0];
1117 }
1118 KJ_CASE_ONEOF(_, kj::HttpEverythingRange) {}
1119 KJ_CASE_ONEOF(_, kj::HttpUnsatisfiableRange) {
1120 kj::HttpHeaders headers(headerTable);
1121 headers.set(kj::HttpHeaderId::CONTENT_RANGE, kj::str("bytes */", meta.size));
1122 co_return co_await response.sendError(416, "Range Not Satisfiable", headers);
1123 }
1124 }
1125 }
1126 }
1127 
1128 kj::HttpHeaders headers(headerTable);
1129 headers.set(kj::HttpHeaderId::CONTENT_TYPE, MimeType::OCTET_STREAM.toString());
1130 headers.set(hLastModified, httpTime(meta.lastModified));
1131 
1132 // We explicitly set the Content-Length header because if we don't, and we were called
1133 // by a local Worker (without an actual HTTP connection in between), then the Worker
1134 // will not see a Content-Length header, but being able to query the content length
1135 // (especially with HEAD requests) is quite useful.
1136 // TODO(cleanup): Arguably the implementation of `fetch()` should be adjusted so that
1137 // if no `Content-Length` header is returned, but the body size is known via the KJ
1138 // HTTP API, then the header should be filled in automatically. Unclear if this is safe
1139 // to change without a compat flag.
1140 
1141 if (method == kj::HttpMethod::HEAD) {
1142 headers.set(kj::HttpHeaderId::CONTENT_LENGTH, kj::str(meta.size));
1143 response.send(200, "OK", headers, meta.size);
1144 co_return;
1145 } else KJ_IF_SOME(r, range) {
1146 KJ_ASSERT(r.start <= r.end);
1147 auto rangeSize = r.end - r.start + 1;
1148 headers.set(kj::HttpHeaderId::CONTENT_LENGTH, kj::str(rangeSize));
1149 headers.set(kj::HttpHeaderId::CONTENT_RANGE,
1150 kj::str("bytes ", r.start, "-", r.end, "/", meta.size));
1151 auto out = response.send(206, "Partial Content", headers, rangeSize);
1152 
1153 auto in = kj::heap<kj::FileInputStream>(*file, r.start);
1154 co_return co_await in->pumpTo(*out, rangeSize).ignoreResult();
1155 } else {
1156 headers.set(kj::HttpHeaderId::CONTENT_LENGTH, kj::str(meta.size));
1157 auto out = response.send(200, "OK", headers, meta.size);
1158 
1159 auto in = kj::heap<kj::FileInputStream>(*file);
1160 co_return co_await in->pumpTo(*out, meta.size).ignoreResult();
1161 }
1162 }
1163 case kj::FsNode::Type::DIRECTORY: {
1164 // Whoooops, we opened a directory. Back up and start over.
1165 
1166 auto dir = readable->openSubdir(path);
1167 
1168 kj::HttpHeaders headers(headerTable);
1169 headers.set(kj::HttpHeaderId::CONTENT_TYPE, MimeType::JSON.toString());
1170 headers.set(hLastModified, httpTime(meta.lastModified));
1171 
1172 // We intentionally don't provide the expected size here in order to reserve the right
1173 // to switch to streaming directory listing in the future.
1174 auto out = response.send(200, "OK", headers);
1175 
1176 if (method == kj::HttpMethod::HEAD) {
1177 co_return;
1178 } else {
1179 auto entries = dir->listEntries();
1180 kj::Vector<kj::String> jsonEntries(entries.size());
1181 for (auto& entry: entries) {
1182 if (!allowDotfiles && entry.name.startsWith(".")) {
1183 continue;
1184 }
1185 
1186 kj::StringPtr type = "other";
1187 switch (entry.type) {
1188 case kj::FsNode::Type::FILE:
1189 type = "file";
1190 break;
1191 case kj::FsNode::Type::DIRECTORY:
1192 type = "directory";
1193 break;
1194 case kj::FsNode::Type::SYMLINK:
1195 type = "symlink";
1196 break;
1197 case kj::FsNode::Type::BLOCK_DEVICE:
1198 type = "blockDevice";
1199 break;
1200 case kj::FsNode::Type::CHARACTER_DEVICE:
1201 type = "characterDevice";
1202 break;
1203 case kj::FsNode::Type::NAMED_PIPE:
1204 type = "namedPipe";
1205 break;
1206 case kj::FsNode::Type::SOCKET:
1207 type = "socket";
1208 break;
1209 case kj::FsNode::Type::OTHER:
1210 type = "other";
1211 break;
1212 }
1213 
1214 jsonEntries.add(
1215 kj::str("{\"name\":", escapeJsonString(entry.name), ",\"type\":\"", type, "\"}"));
1216 };
1217 
1218 auto content = kj::str('[', kj::strArray(jsonEntries, ","), ']');
1219 
1220 co_return co_await out->write(content.asBytes());
1221 }
1222 }
1223 default:
1224 co_return co_await response.sendError(406, "Not Acceptable", headerTable);
1225 }
1226 } else if (method == kj::HttpMethod::PUT) {
1227 auto& w = KJ_UNWRAP_OR(writable,
1228 { co_return co_await response.sendError(405, "Method Not Allowed", headerTable); });
1229 
1230 if (blockedPath || path.size() == 0) {
1231 co_return co_await response.sendError(403, "Unauthorized", headerTable);
1232 }
1233 
1234 auto replacer = w.replaceFile(
1235 path, kj::WriteMode::CREATE | kj::WriteMode::MODIFY | kj::WriteMode::CREATE_PARENT);
1236 auto stream = kj::heap<kj::FileOutputStream>(replacer->get());
1237 
1238 co_await requestBody.pumpTo(*stream);
1239 
1240 replacer->commit();
1241 kj::HttpHeaders headers(headerTable);
1242 response.send(204, "No Content", headers);
1243 co_return;
1244 } else if (method == kj::HttpMethod::DELETE) {
1245 auto& w = KJ_UNWRAP_OR(writable,
1246 { co_return co_await response.sendError(405, "Method Not Allowed", headerTable); });
1247 
1248 if (blockedPath || path.size() == 0) {
1249 co_return co_await response.sendError(403, "Unauthorized", headerTable);
1250 }
1251 
1252 auto found = w.tryRemove(path);
1253 
1254 kj::HttpHeaders headers(headerTable);
1255 if (found) {
1256 response.send(204, "No Content", headers);
1257 co_return;
1258 } else {
1259 co_return co_await response.sendError(404, "Not Found", headers);
1260 }
1261 } else {
1262 co_return co_await response.sendError(501, "Not Implemented", headerTable);
1263 }
1264 }
1265 
1266 kj::Promise<void> connect(kj::StringPtr host,
1267 const kj::HttpHeaders& headers,
1268 kj::AsyncIoStream& connection,
1269 kj::HttpService::ConnectResponse& response,
1270 kj::HttpConnectSettings settings) override {
1271 throwUnsupported();
1272 }
1273 kj::Promise<void> prewarm(kj::StringPtr url) override {
1274 return kj::READY_NOW;
1275 }
1276 kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override {
1277 throwUnsupported();
1278 }
1279 kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override {
1280 throwUnsupported();
1281 }
1282 kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override {
1283 return event->notSupported();
1284 }
1285 
1286 [[noreturn]] void throwUnsupported() {
1287 JSG_FAIL_REQUIRE(Error, "Disk directory services don't support this event type.");
1288 }
1289};
1290 
1291kj::Own<Server::Service> Server::makeDiskDirectoryService(kj::StringPtr name,
1292 config::DiskDirectory::Reader conf,
1293 kj::HttpHeaderTable::Builder& headerTableBuilder) {
1294 TRACE_EVENT("workerd", "Server::makeDiskDirectoryService()");
1295 kj::StringPtr pathStr = nullptr;
1296 kj::String ownPathStr;
1297 
1298 KJ_IF_SOME(override, directoryOverrides.findEntry(name)) {
1299 pathStr = ownPathStr = kj::mv(override.value);
1300 directoryOverrides.erase(override);
1301 } else if (conf.hasPath()) {
1302 pathStr = conf.getPath();
1303 } else {
1304 reportConfigError(kj::str("Directory \"", name,
1305 "\" has no path in the config, so must be specified on the "
1306 "command line with `--directory-path`."));
1307 return makeInvalidConfigService();
1308 }
1309 
1310 auto path = fs.getCurrentPath().evalNative(pathStr);
1311 
1312 if (conf.getWritable()) {
1313 auto openDir = KJ_UNWRAP_OR(fs.getRoot().tryOpenSubdir(kj::mv(path), kj::WriteMode::MODIFY), {
1314 reportConfigError(kj::str("Directory named \"", name, "\" not found: ", pathStr));
1315 return makeInvalidConfigService();
1316 });
1317 
1318 return kj::refcounted<DiskDirectoryService>(conf, kj::mv(openDir), headerTableBuilder);
1319 } else {
1320 auto openDir = KJ_UNWRAP_OR(fs.getRoot().tryOpenSubdir(kj::mv(path)), {
1321 reportConfigError(kj::str("Directory named \"", name, "\" not found: ", pathStr));
1322 return makeInvalidConfigService();
1323 });
1324 
1325 return kj::refcounted<DiskDirectoryService>(conf, kj::mv(openDir), headerTableBuilder);
1326 }
1327}
1328 
1329// =======================================================================================
1330 
1331// This class exists to update the InspectorService's table of isolates when a config
1332// has multiple services. The InspectorService exists on the stack of its own thread and
1333// initializes state that is bound to the thread, e.g. a http server and an event loop.
1334// This class provides a small thread-safe interface to the InspectorService so <name>:<isolate>
1335// mappings can be added after the InspectorService has started.
1336//
1337// The Cloudflare devtools only show the first service in workerd configuration. This service
1338// is always contains a users code. However, in packaging user code wrangler may add
1339// additional services that also have code. If using Chrome devtools to inspect a workerd,
1340// instance all services are visible and can be debugged.
1341class Server::InspectorServiceIsolateRegistrar final {
1342 public:
1343 InspectorServiceIsolateRegistrar() {}
1344 ~InspectorServiceIsolateRegistrar() noexcept(true);
1345 
1346 void registerIsolate(kj::StringPtr name, Worker::Isolate* isolate);
1347 
1348 KJ_DISALLOW_COPY_AND_MOVE(InspectorServiceIsolateRegistrar);
1349 
1350 private:
1351 void attach(const Server::InspectorService* anInspectorService) {
1352 *inspectorService.lockExclusive() = anInspectorService;
1353 }
1354 
1355 void detach() {
1356 *inspectorService.lockExclusive() = nullptr;
1357 }
1358 
1359 kj::MutexGuarded<const InspectorService*> inspectorService;
1360 friend class Server::InspectorService;
1361};
1362 
1363// Implements the interface for the devtools inspector protocol.
1364//
1365// The InspectorService is created when workerd serve is called using the -i option
1366// to define the inspector socket.
1367class Server::InspectorService final: public kj::HttpService, public kj::HttpServerErrorHandler {
1368 public:
1369 InspectorService(kj::Own<const kj::Executor> isolateThreadExecutor,
1370 kj::Timer& timer,
1371 kj::HttpHeaderTable::Builder& headerTableBuilder,
1372 InspectorServiceIsolateRegistrar& registrar)
1373 : isolateThreadExecutor(kj::mv(isolateThreadExecutor)),
1374 timer(timer),
1375 headerTable(headerTableBuilder.getFutureTable()),
1376 server(timer, headerTable, *this, kj::HttpServerSettings{.errorHandler = *this}),
1377 registrar(registrar) {
1378 registrar.attach(this);
1379 }
1380 
1381 ~InspectorService() {
1382 KJ_IF_SOME(r, registrar) {
1383 r.detach();
1384 }
1385 }
1386 
1387 void invalidateRegistrar() {
1388 registrar = kj::none;
1389 }
1390 
1391 kj::Promise<void> handleApplicationError(
1392 kj::Exception exception, kj::Maybe<kj::HttpService::Response&> response) override {
1393 if (exception.getType() == kj::Exception::Type::DISCONNECTED) {
1394 // Don't send a response, just close connection.
1395 co_return;
1396 }
1397 KJ_LOG(ERROR, kj::str("Uncaught exception: ", exception));
1398 KJ_IF_SOME(r, response) {
1399 co_return co_await r.sendError(500, "Internal Server Error", headerTable);
1400 }
1401 }
1402 
1403 kj::Promise<void> request(kj::HttpMethod method,
1404 kj::StringPtr url,
1405 const kj::HttpHeaders& headers,
1406 kj::AsyncInputStream& requestBody,
1407 kj::HttpService::Response& response) override {
1408 // The inspector protocol starts with the debug client sending ordinary HTTP GET requests
1409 // to /json/version and then to /json or /json/list. These must respond with valid JSON
1410 // documents that list the details of what isolates are available for inspection. Each
1411 // isolate must be listed separately. In the advertisement for each isolate is a URL
1412 // and a unique ID. The client will use the URL and ID to open a WebSocket request to
1413 // actually connect the debug session.
1414 kj::HttpHeaders responseHeaders(headerTable);
1415 if (headers.isWebSocket()) {
1416 KJ_IF_SOME(pos, url.findLast('/')) {
1417 auto id = url.slice(pos + 1);
1418 
1419 KJ_IF_SOME(isolate, isolates.find(id)) {
1420 // If getting the strong ref doesn't work it means that the Worker::Isolate
1421 // has already been cleaned up. We use a weak ref here in order to keep from
1422 // having the Worker::Isolate itself having to know anything at all about the
1423 // IsolateService and the registration process. So instead of having Isolate
1424 // explicitly clean up after itself we lazily evaluate the weak ref and clean
1425 // up when necessary.
1426 KJ_IF_SOME(ref, isolate->tryAddStrongRef()) {
1427 // When using --verbose, we'll output some logging to indicate when the
1428 // inspector client is attached/detached.
1429 KJ_LOG(INFO, kj::str("Inspector client attaching [", id, "]"));
1430 auto webSocket = response.acceptWebSocket(responseHeaders);
1431 kj::Duration timerOffset = 0 * kj::MILLISECONDS;
1432 try {
1433 co_return co_await ref->attachInspector(
1434 isolateThreadExecutor->addRef(), timer, timerOffset, *webSocket);
1435 } catch (...) {
1436 auto exception = kj::getCaughtExceptionAsKj();
1437 if (exception.getType() == kj::Exception::Type::DISCONNECTED) {
1438 // This likely just means that the inspector client was closed.
1439 // Nothing to do here but move along.
1440 KJ_LOG(INFO, "Inspector client detached"_kj);
1441 co_return;
1442 } else {
1443 // If it's any other kind of error, propagate it!
1444 kj::throwFatalException(kj::mv(exception));
1445 }
1446 }
1447 } else {
1448 // If we can't get a strong ref to the isolate here, it's been cleaned
1449 // up. The only thing we're going to do is clean up here and act like
1450 // nothing happened.
1451 isolates.erase(id);
1452 }
1453 }
1454 
1455 KJ_LOG(INFO, kj::str("Unknown worker session [", id, "]"));
1456 co_return co_await response.sendError(404, "Unknown worker session", responseHeaders);
1457 }
1458 
1459 // No / in url!? That's weird
1460 co_return co_await response.sendError(400, "Invalid request", responseHeaders);
1461 }
1462 
1463 // If the request is not a WebSocket request, it must be a GET to fetch details
1464 // about the implementation.
1465 if (method != kj::HttpMethod::GET) {
1466 co_return co_await response.sendError(501, "Unsupported Operation", responseHeaders);
1467 }
1468 
1469 if (url.endsWith("/json/version")) {
1470 responseHeaders.set(kj::HttpHeaderId::CONTENT_TYPE, MimeType::JSON.toString());
1471 auto content = kj::str("{\"Browser\": \"workerd\", \"Protocol-Version\": \"1.3\" }");
1472 auto out = response.send(200, "OK", responseHeaders, content.size());
1473 co_return co_await out->write(content.asBytes());
1474 } else if (url.endsWith("/json") || url.endsWith("/json/list") ||
1475 url.endsWith("/json/list?for_tab")) {
1476 responseHeaders.set(kj::HttpHeaderId::CONTENT_TYPE, MimeType::JSON.toString());
1477 
1478 auto baseWsUrl = KJ_UNWRAP_OR(headers.get(kj::HttpHeaderId::HOST),
1479 { co_return co_await response.sendError(400, "Bad Request", responseHeaders); });
1480 
1481 kj::Vector<kj::String> entries(isolates.size());
1482 kj::Vector<kj::String> toRemove;
1483 for (auto& entry: isolates) {
1484 // While we don't actually use the strong ref here we still attempt to acquire it
1485 // in order to determine if the isolate is actually still around. If the isolate
1486 // has been destroyed the weak ref will be cleared. We do it this way to keep from
1487 // having the Worker::Isolate know anything at all about the InspectorService.
1488 // We'll lazily clean up whenever we detect that the ref has been invalidated.
1489 //
1490 // TODO(cleanup): If we ever enable reloading of isolates for live services, we may
1491 // want to refactor this such that the WorkerService holds a handle to the registration
1492 // as opposed to using this lazy cleanup mechanism. For now, however, this is
1493 // sufficient.
1494 KJ_IF_SOME(ref, entry.value->tryAddStrongRef()) {
1495 (void)ref; // squash compiler warning about unused ref
1496 kj::Vector<kj::String> fields(9);
1497 fields.add(kj::str("\"id\":\"", entry.key, "\""));
1498 fields.add(kj::str("\"title\":\"workerd: worker ", entry.key, "\""));
1499 fields.add(kj::str("\"type\":\"node\""));
1500 fields.add(kj::str("\"description\":\"workerd worker\""));
1501 fields.add(kj::str("\"webSocketDebuggerUrl\":\"ws://", baseWsUrl, "/", entry.key, "\""));
1502 fields.add(kj::str(
1503 "\"devtoolsFrontendUrl\":\"devtools://devtools/bundled/js_app.html?experiments=true&v8only=true&ws=",
1504 baseWsUrl, "/\""));
1505 fields.add(kj::str(
1506 "\"devtoolsFrontendUrlCompat\":\"devtools://devtools/bundled/inspector.html?experiments=true&v8only=true&ws=",
1507 baseWsUrl, "/\""));
1508 fields.add(kj::str("\"faviconUrl\":\"https://workers.cloudflare.com/favicon.ico\""));
1509 fields.add(kj::str("\"url\":\"https://workers.dev\""));
1510 entries.add(kj::str('{', kj::strArray(fields, ","), '}'));
1511 } else {
1512 // If we're not able to get a reference to the isolate here, it's
1513 // been cleaned up and we should remove it from the list. We do this
1514 // after iterating to make sure we don't invalidate the iterator.
1515 toRemove.add(kj::str(entry.key));
1516 }
1517 }
1518 // Clean up if necessary
1519 for (auto& key: toRemove) {
1520 isolates.erase(key);
1521 }
1522 
1523 auto content = kj::str('[', kj::strArray(entries, ","), ']');
1524 
1525 auto out = response.send(200, "OK", responseHeaders, content.size());
1526 co_return co_await out->write(content.asBytes()).attach(kj::mv(content), kj::mv(out));
1527 }
1528 
1529 co_return co_await response.sendError(500, "Not yet implemented", responseHeaders);
1530 }
1531 
1532 kj::Promise<void> listen(kj::Own<kj::ConnectionReceiver> listener) {
1533 // Note that we intentionally do not make inspector connections be part of the usual drain()
1534 // procedure. Inspector connections are always long-lived WebSockets, and we do not want the
1535 // existence of such a connection to hold the server open. We do, however, want the connection
1536 // to stay open until all other requests are drained, for debugging purposes.
1537 //
1538 // Thus:
1539 // * We let connection loop tasks live on `HttpServer`'s own `TaskSet`, rather than our
1540 // server's main `TaskSet` which we wait to become empty on drain.
1541 // * We do not add this `HttpServer` to the server's `httpServers` list, so it will not receive
1542 // drain() requests. (However, our caller does cancel listening on the server port as soon
1543 // as we begin draining, since we may want new connections to go to a new instance of the
1544 // server.)
1545 co_return co_await server.listenHttp(*listener);
1546 }
1547 
1548 void registerIsolate(kj::StringPtr name, Worker::Isolate* isolate) {
1549 isolates.insert(kj::str(name), isolate->getWeakRef());
1550 }
1551 
1552 private:
1553 kj::Own<const kj::Executor> isolateThreadExecutor;
1554 kj::Timer& timer;
1555 kj::HttpHeaderTable& headerTable;
1556 kj::HashMap<kj::String, kj::Own<const Worker::Isolate::WeakIsolateRef>> isolates;
1557 kj::HttpServer server;
1558 kj::Maybe<InspectorServiceIsolateRegistrar&> registrar;
1559};
1560 
1561Server::InspectorServiceIsolateRegistrar::~InspectorServiceIsolateRegistrar() noexcept(true) {
1562 auto lockedInspectorService = this->inspectorService.lockExclusive();
1563 if (lockedInspectorService != nullptr) {
1564 auto is = const_cast<InspectorService*>(*lockedInspectorService);
1565 is->invalidateRegistrar();
1566 }
1567}
1568 
1569void Server::InspectorServiceIsolateRegistrar::registerIsolate(
1570 kj::StringPtr name, Worker::Isolate* isolate) {
1571 auto lockedInspectorService = this->inspectorService.lockExclusive();
1572 if (lockedInspectorService != nullptr) {
1573 auto is = const_cast<InspectorService*>(*lockedInspectorService);
1574 is->registerIsolate(name, isolate);
1575 }
1576}
1577 
1578// =======================================================================================
1579namespace {
1580class RequestObserverWithTracer final: public RequestObserver, public WorkerInterface {
1581 public:
1582 RequestObserverWithTracer(kj::Maybe<kj::Own<WorkerTracer>> tracer, kj::TaskSet& waitUntilTasks)
1583 : tracer(kj::mv(tracer)) {}
1584 
1585 ~RequestObserverWithTracer() noexcept(false) {
1586 KJ_IF_SOME(t, tracer) {
1587 // for a more precise end time, set the end timestamp now, if available
1588 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
1589 auto time = ioContext.now();
1590 t->recordTimestamp(time);
1591 }
1592 t->setOutcome(
1593 outcome, 0 * kj::MILLISECONDS /* cpu time */, 0 * kj::MILLISECONDS /* wall time */);
1594 }
1595 }
1596 
1597 WorkerInterface& wrapWorkerInterface(WorkerInterface& worker) override {
1598 if (tracer != kj::none) {
1599 inner = worker;
1600 return *this;
1601 }
1602 return worker;
1603 }
1604 
1605 void reportFailure(
1606 const kj::Exception& exception, FailureSource source = FailureSource::OTHER) override {
1607 outcome = RequestObserver::outcomeFromException(exception, source);
1608 }
1609 
1610 void setOutcome(EventOutcome newOutcome) override {
1611 outcome = newOutcome;
1612 }
1613 
1614 // WorkerInterface
1615 kj::Promise<void> request(kj::HttpMethod method,
1616 kj::StringPtr url,
1617 const kj::HttpHeaders& headers,
1618 kj::AsyncInputStream& requestBody,
1619 kj::HttpService::Response& response) override {
1620 try {
1621 SimpleResponseObserver responseWrapper(&fetchStatus, response);
1622 co_await KJ_ASSERT_NONNULL(inner).request(method, url, headers, requestBody, responseWrapper);
1623 } catch (...) {
1624 auto exception = kj::getCaughtExceptionAsKj();
1625 // Overloaded-type exceptions generally represent some resource exhaustion (i.e. not
1626 // necessarily an internal error) and correspond to HTTP error 503.
1627 if (exception.getType() == kj::Exception::Type::OVERLOADED) {
1628 fetchStatus = 503;
1629 } else {
1630 fetchStatus = 500;
1631 }
1632 reportFailure(exception);
1633 kj::throwFatalException(kj::mv(exception));
1634 }
1635 }
1636 
1637 kj::Promise<void> connect(kj::StringPtr host,
1638 const kj::HttpHeaders& headers,
1639 kj::AsyncIoStream& connection,
1640 ConnectResponse& response,
1641 kj::HttpConnectSettings settings) override {
1642 try {
1643 co_return co_await KJ_ASSERT_NONNULL(inner).connect(
1644 host, headers, connection, response, settings);
1645 } catch (...) {
1646 auto exception = kj::getCaughtExceptionAsKj();
1647 reportFailure(exception);
1648 kj::throwFatalException(kj::mv(exception));
1649 }
1650 }
1651 
1652 kj::Promise<void> prewarm(kj::StringPtr url) override {
1653 try {
1654 co_return co_await KJ_ASSERT_NONNULL(inner).prewarm(url);
1655 } catch (...) {
1656 auto exception = kj::getCaughtExceptionAsKj();
1657 reportFailure(exception);
1658 kj::throwFatalException(kj::mv(exception));
1659 }
1660 }
1661 
1662 kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override {
1663 try {
1664 co_return co_await KJ_ASSERT_NONNULL(inner).runScheduled(scheduledTime, cron);
1665 } catch (...) {
1666 auto exception = kj::getCaughtExceptionAsKj();
1667 reportFailure(exception);
1668 kj::throwFatalException(kj::mv(exception));
1669 }
1670 }
1671 
1672 kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override {
1673 try {
1674 co_return co_await KJ_ASSERT_NONNULL(inner).runAlarm(scheduledTime, retryCount);
1675 } catch (...) {
1676 auto exception = kj::getCaughtExceptionAsKj();
1677 reportFailure(exception);
1678 kj::throwFatalException(kj::mv(exception));
1679 }
1680 }
1681 
1682 kj::Promise<bool> test() override {
1683 try {
1684 co_return co_await KJ_ASSERT_NONNULL(inner).test();
1685 } catch (...) {
1686 auto exception = kj::getCaughtExceptionAsKj();
1687 reportFailure(exception);
1688 kj::throwFatalException(kj::mv(exception));
1689 }
1690 }
1691 
1692 kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override {
1693 try {
1694 co_return co_await KJ_ASSERT_NONNULL(inner).customEvent(kj::mv(event));
1695 } catch (...) {
1696 auto exception = kj::getCaughtExceptionAsKj();
1697 reportFailure(exception);
1698 kj::throwFatalException(kj::mv(exception));
1699 }
1700 }
1701 
1702 kj::Promise<kj::Maybe<kj::Date>> abandonAlarm(kj::Date scheduledTime) override {
1703 co_return co_await KJ_ASSERT_NONNULL(inner).abandonAlarm(scheduledTime);
1704 }
1705 
1706 private:
1707 kj::Maybe<kj::Own<WorkerTracer>> tracer;
1708 kj::Maybe<WorkerInterface&> inner;
1709 EventOutcome outcome = EventOutcome::OK;
1710 kj::uint fetchStatus = 0;
1711};
1712 
1713class SequentialSpanSubmitter final: public SpanSubmitter {
1714 public:
1715 SequentialSpanSubmitter(kj::Own<BaseTracer::WeakRef> weakTracer, kj::EntropySource& entropySource)
1716 : weakTracer(kj::mv(weakTracer)),
1717 entropySource(entropySource) {}
1718 void submitSpanClose(
1719 tracing::SpanId spanId, kj::Date startTime, kj::Date endTime, Span::TagMap&& tags) override {
1720 weakTracer->runIfAlive([&](BaseTracer& tracer) {
1721 tracing::SpanEndData spanEnd(spanId, endTime, kj::mv(tags));
1722 if (isPredictableModeForTest()) {
1723 startTime = spanEnd.endTime = kj::UNIX_EPOCH;
1724 }
1725 
1726 tracer.addSpanClose(kj::mv(spanEnd), startTime);
1727 });
1728 }
1729 
1730 bool submitSpanOpen(tracing::SpanId spanId,
1731 tracing::SpanId parentSpanId,
1732 kj::ConstString operationName,
1733 kj::Date startTime) override {
1734 bool submitted = false;
1735 weakTracer->runIfAlive([&](BaseTracer& tracer) {
1736 if (isPredictableModeForTest()) {
1737 startTime = kj::UNIX_EPOCH;
1738 }
1739 tracer.addSpanOpen(spanId, parentSpanId, kj::mv(operationName), startTime);
1740 submitted = true;
1741 });
1742 return submitted;
1743 }
1744 
1745 tracing::SpanId makeSpanId() override {
1746 if (isPredictableModeForTest()) {
1747 return tracing::SpanId(nextSpanId++);
1748 }
1749 return tracing::SpanId::fromEntropy(entropySource);
1750 }
1751 KJ_DISALLOW_COPY_AND_MOVE(SequentialSpanSubmitter);
1752 
1753 private:
1754 uint64_t nextSpanId = 1;
1755 kj::Own<BaseTracer::WeakRef> weakTracer;
1756 kj::EntropySource& entropySource;
1757};
1758 
1759// IsolateLimitEnforcer that enforces no limits.
1760class NullIsolateLimitEnforcer final: public IsolateLimitEnforcer {
1761 public:
1762 v8::Isolate::CreateParams getCreateParams() override {
1763 return {};
1764 }
1765 
1766 void customizeIsolate(v8::Isolate* isolate) override {}
1767 
1768 ActorCacheSharedLruOptions getActorCacheLruOptions() override {
1769 // TODO(someday): Make this configurable?
1770 return {.softLimit = 16 * (1ull << 20), // 16 MiB
1771 .hardLimit = 128 * (1ull << 20), // 128 MiB
1772 .staleTimeout = 30 * kj::SECONDS,
1773 .dirtyListByteLimit = 8 * (1ull << 20), // 8 MiB
1774 .maxKeysPerRpc = 128,
1775 
1776 // For now, we use `neverFlush` to implement in-memory-only actors.
1777 // See WorkerService::getActor().
1778 .neverFlush = true};
1779 }
1780 
1781 kj::Own<void> enterStartupJs(
1782 jsg::Lock& lock, kj::OneOf<kj::Exception, kj::Duration>&) const override {
1783 return {};
1784 }
1785 
1786 kj::Own<void> enterStartupPython(
1787 jsg::Lock& lock, kj::OneOf<kj::Exception, kj::Duration>&) const override {
1788 return {};
1789 }
1790 
1791 kj::Own<void> enterDynamicImportJs(
1792 jsg::Lock& lock, kj::OneOf<kj::Exception, kj::Duration>&) const override {
1793 return {};
1794 }
1795 
1796 kj::Own<void> enterLoggingJs(
1797 jsg::Lock& lock, kj::OneOf<kj::Exception, kj::Duration>&) const override {
1798 return {};
1799 }
1800 
1801 kj::Own<void> enterInspectorJs(
1802 jsg::Lock& loc, kj::OneOf<kj::Exception, kj::Duration>&) const override {
1803 return {};
1804 }
1805 
1806 void completedRequest(kj::StringPtr id) const override {}
1807 
1808 bool exitJs(jsg::Lock& lock) const override {
1809 return false;
1810 }
1811 
1812 void reportMetrics(IsolateObserver& isolateMetrics) const override {}
1813 
1814 kj::Maybe<size_t> checkPbkdfIterations(jsg::Lock& lock, size_t iterations) const override {
1815 // No limit on the number of iterations in workerd
1816 return kj::none;
1817 }
1818 
1819 bool hasExcessivelyExceededHeapLimit() const override {
1820 return false;
1821 }
1822 
1823 const TrackedWasmInstanceList& getTrackedWasmInstances() const override {
1824 return trackedWasmInstances;
1825 }
1826 
1827 private:
1828 TrackedWasmInstanceList trackedWasmInstances;
1829};
1830 
1831} // namespace
1832 
1833// Shared ErrorReporter base implemnetation. The logic to collect entrypoint information is the
1834// same regardless of where the code came from.
1835struct Server::ErrorReporter: public Worker::ValidationErrorReporter {
1836 // The `HashSet`s are the set of exported handlers, like `fetch`, `test`, etc.
1837 kj::HashMap<kj::String, kj::HashSet<kj::String>> namedEntrypoints;
1838 kj::Maybe<kj::HashSet<kj::String>> defaultEntrypoint;
1839 kj::HashSet<kj::String> actorClasses;
1840 kj::HashSet<kj::String> workflowClasses;
1841 
1842 void addEntrypoint(kj::Maybe<kj::StringPtr> exportName, kj::Array<kj::String> methods) override {
1843 kj::HashSet<kj::String> set;
1844 for (auto& method: methods) {
1845 set.insert(kj::mv(method));
1846 }
1847 KJ_IF_SOME(e, exportName) {
1848 namedEntrypoints.insert(kj::str(e), kj::mv(set));
1849 } else {
1850 defaultEntrypoint = kj::mv(set);
1851 }
1852 }
1853 
1854 void addActorClass(kj::StringPtr exportName) override {
1855 actorClasses.insert(kj::str(exportName));
1856 }
1857 
1858 void addWorkflowClass(kj::StringPtr exportName, kj::Array<kj::String> methods) override {
1859 // At runtime, we need to add it into the normal namedEntrypoints for Workflows to appear
1860 // in `WorkerService`. This is a different method compared to `addEntrypoint` because we need to
1861 // check for `WorkflowEntrypoint` inheritance at validation time.
1862 kj::HashSet<kj::String> set;
1863 for (auto& method: methods) {
1864 set.insert(kj::mv(method));
1865 }
1866 namedEntrypoints.insert(kj::str(exportName), kj::mv(set));
1867 workflowClasses.insert(kj::str(exportName));
1868 }
1869};
1870 
1871// Implementation of ErrorReporter specifically for reporting errors in the top-level workerd
1872// config.
1873struct Server::ConfigErrorReporter final: public ErrorReporter {
1874 ConfigErrorReporter(Server& server, kj::StringPtr name): server(server), name(name) {}
1875 
1876 Server& server;
1877 kj::StringPtr name;
1878 
1879 void addError(kj::String error) override {
1880 server.handleReportConfigError(kj::str("service ", name, ": ", error));
1881 }
1882};
1883 
1884// Implementation of ErrorReporter for dynamically-loaded Workers. We'll collect the errors and
1885// report them in an exception at the end.
1886struct Server::DynamicErrorReporter final: public ErrorReporter {
1887 kj::Vector<kj::String> errors;
1888 
1889 void addError(kj::String error) override {
1890 errors.add(kj::mv(error));
1891 }
1892 
1893 void throwIfErrors() {
1894 if (!errors.empty()) {
1895 JSG_FAIL_REQUIRE(Error, "Failed to start Worker:\n", kj::strArray(errors, "\n"));
1896 }
1897 }
1898};
1899 
1900class Server::WorkerService final: public Service,
1901 private kj::TaskSet::ErrorHandler,
1902 private IoChannelFactory,
1903 private TimerChannel,
1904 private LimitEnforcer {
1905 public:
1906 class ActorNamespace;
1907 
1908 // I/O channels, delivered when link() is called.
1909 struct LinkedIoChannels {
1910 kj::Array<kj::Own<IoChannelFactory::SubrequestChannel>> subrequest;
1911 kj::Array<kj::Maybe<ActorNamespace&>> actor; // null = configuration error
1912 kj::Array<kj::Own<ActorClass>> actorClass;
1913 kj::Maybe<kj::Own<IoChannelFactory::SubrequestChannel>> cache;
1914 kj::Maybe<const kj::Directory&> actorStorage;
1915 kj::Array<kj::Own<IoChannelFactory::SubrequestChannel>> tails;
1916 kj::Array<kj::Own<IoChannelFactory::SubrequestChannel>> streamingTails;
1917 kj::Array<kj::Rc<WorkerLoaderNamespace>> workerLoaders;
1918 kj::Maybe<kj::Network&> workerdDebugPortNetwork;
1919 };
1920 using LinkCallback =
1921 kj::Function<LinkedIoChannels(WorkerService&, Worker::ValidationErrorReporter&)>;
1922 using AbortActorsCallback = kj::Function<void(kj::Maybe<const kj::Exception&> reason)>;
1923 using DeleteActorsCallback = kj::Function<void(kj::Maybe<const kj::Exception&> reason)>;
1924 
1925 WorkerService(ChannelTokenHandler& channelTokenHandler,
1926 kj::Maybe<kj::StringPtr> serviceName,
1927 ThreadContext& threadContext,
1928 const kj::MonotonicClock& monotonicClock,
1929 kj::Own<const Worker> worker,
1930 kj::Maybe<kj::HashSet<kj::String>> defaultEntrypointHandlers,
1931 kj::HashMap<kj::String, kj::HashSet<kj::String>> namedEntrypoints,
1932 kj::HashSet<kj::String> actorClassEntrypointsParam,
1933 LinkCallback linkCallback,
1934 AbortActorsCallback abortActorsCallback,
1935 DeleteActorsCallback deleteActorsCallback,
1936 kj::Maybe<kj::String> dockerPathParam,
1937 kj::Maybe<kj::String> containerEgressInterceptorImageParam,
1938 bool isDynamic,
1939 kj::Maybe<kj::Function<void()>> abortIsolateCallback = kj::none)
1940 : channelTokenHandler(channelTokenHandler),
1941 serviceName(serviceName),
1942 threadContext(threadContext),
1943 monotonicClock(monotonicClock),
1944 ioChannels(kj::mv(linkCallback)),
1945 worker(kj::mv(worker)),
1946 defaultEntrypointHandlers(kj::mv(defaultEntrypointHandlers)),
1947 namedEntrypoints(kj::mv(namedEntrypoints)),
1948 actorClassEntrypoints(kj::mv(actorClassEntrypointsParam)),
1949 waitUntilTasks(*this),
1950 abortActorsCallback(kj::mv(abortActorsCallback)),
1951 deleteActorsCallback(kj::mv(deleteActorsCallback)),
1952 dockerPath(kj::mv(dockerPathParam)),
1953 containerEgressInterceptorImage(kj::mv(containerEgressInterceptorImageParam)),
1954 isDynamic(isDynamic),
1955 abortIsolateCallback(kj::mv(abortIsolateCallback)) {}
1956 
1957 // Call immediately after the constructor to set up `actorNamespaces`. This can't happen during
1958 // the constructor itself since it sets up cyclic references, which will throw an exception if
1959 // done during the constructor.
1960 void initActorNamespaces(
1961 const kj::HashMap<kj::String, ActorConfig>& actorClasses, kj::Network& network) {
1962 actorNamespaces.reserve(actorClasses.size());
1963 for (auto& entry: actorClasses) {
1964 if (!actorClassEntrypoints.contains(entry.key)) {
1965 KJ_LOG(WARNING,
1966 kj::str("A DurableObjectNamespace in the config referenced the class \"", entry.key,
1967 "\", but no such Durable Object class is exported from the worker. Please make "
1968 "sure the class name matches, it is exported, and the class extends "
1969 "'DurableObject'. Attempts to call to this Durable Object class will fail at "
1970 "runtime, but historically this was not a startup-time error. Future versions of "
1971 "workerd may make this a startup-time error."));
1972 }
1973 
1974 auto actorClass = kj::refcounted<ActorClassImpl>(*this, entry.key, Frankenvalue());
1975 auto ns = kj::heap<ActorNamespace>(kj::mv(actorClass), entry.value,
1976 kj::systemPreciseCalendarClock(), threadContext.getUnsafeTimer(),
1977 threadContext.getByteStreamFactory(), channelTokenHandler, network, dockerPath,
1978 containerEgressInterceptorImage, waitUntilTasks);
1979 actorNamespaces.insert(entry.key, kj::mv(ns));
1980 }
1981 }
1982 
1983 void requireAllowsTransfer() override {
1984 if (isDynamic) throwDynamicEntrypointTransferError();
1985 }
1986 
1987 kj::Maybe<kj::Own<Service>> getEntrypoint(kj::Maybe<kj::StringPtr> name, Frankenvalue props) {
1988 const kj::HashSet<kj::String>* handlers;
1989 KJ_IF_SOME(n, name) {
1990 KJ_IF_SOME(entry, namedEntrypoints.findEntry(n)) {
1991 name = entry.key; // replace with more-permanent string
1992 handlers = &entry.value;
1993 } else KJ_IF_SOME(className, actorClassEntrypoints.find(n)) {
1994 // TODO(soon): Restore this warning once miniflare no longer generates config that causes
1995 // it to log spuriously.
1996 //
1997 // KJ_LOG(WARNING,
1998 // kj::str("A ServiceDesignator in the config referenced the entrypoint \"", n,
1999 // "\", but this class does not extend 'WorkerEntrypoint'. Attempts to call this "
2000 // "entrypoint will fail at runtime, but historically this was not a startup-time "
2001 // "error. Future versions of workerd may make this a startup-time error."));
2002 
2003 static const kj::HashSet<kj::String> EMPTY_HANDLERS;
2004 name = className; // replace with more-permanent string
2005 handlers = &EMPTY_HANDLERS;
2006 } else {
2007 return kj::none;
2008 }
2009 } else {
2010 KJ_IF_SOME(d, defaultEntrypointHandlers) {
2011 handlers = &d;
2012 } else {
2013 // It would appear that there is no default export, therefore this refers to an entrypoint
2014 // that doesn't exist! However, this was historically allowed. For backwards-compatibility,
2015 // we preserve this behavior, by returning a reference to the WorkerService itself, whose
2016 // startRequest() will fail.
2017 //
2018 // What will happen if you invoke this entrypoint? Not what you think. Check out the
2019 // test case in server-test.c++ entitled "referencing non-extant default entrypoint is not
2020 // an error" for the sordid details.
2021 return kj::addRef(*this);
2022 }
2023 }
2024 return kj::refcounted<EntrypointService>(*this, name, kj::mv(props), *handlers);
2025 }
2026 
2027 // Like getEntrypoint() but used specifically to get the entrypoint for use in ctx.exports,
2028 // where it can be used raw (props are empty), or can be specialized with props.
2029 kj::Own<Service> getLoopbackEntrypoint(kj::Maybe<kj::StringPtr> name) {
2030 const kj::HashSet<kj::String>* handlers;
2031 KJ_IF_SOME(n, name) {
2032 KJ_IF_SOME(entry, namedEntrypoints.findEntry(n)) {
2033 name = entry.key; // replace with more-permanent string
2034 handlers = &entry.value;
2035 } else {
2036 KJ_FAIL_REQUIRE("getLoopbackEntrypoint() called for entrypoint that doesn't exist");
2037 }
2038 } else {
2039 KJ_IF_SOME(d, defaultEntrypointHandlers) {
2040 handlers = &d;
2041 } else {
2042 KJ_FAIL_REQUIRE("getLoopbackEntrypoint() called for entrypoint that doesn't exist");
2043 }
2044 }
2045 return kj::refcounted<EntrypointService>(*this, name, kj::none, *handlers);
2046 }
2047 
2048 kj::Maybe<kj::Own<ActorClass>> getActorClass(kj::Maybe<kj::StringPtr> name, Frankenvalue props) {
2049 KJ_IF_SOME(className, actorClassEntrypoints.find(KJ_UNWRAP_OR(name, return kj::none))) {
2050 return kj::refcounted<ActorClassImpl>(*this, className, kj::mv(props));
2051 } else {
2052 return kj::none;
2053 }
2054 }
2055 
2056 kj::Own<ActorClass> getLoopbackActorClass(kj::StringPtr name) {
2057 // Look up a more permanent class name string. (Also validates this is actually an export.)
2058 kj::StringPtr className = KJ_REQUIRE_NONNULL(actorClassEntrypoints.find(name),
2059 "getLoopbackActorClass() called for actor class that doesn't exist");
2060 
2061 return kj::refcounted<ActorClassImpl>(*this, className, kj::none);
2062 }
2063 
2064 bool hasDefaultEntrypoint() {
2065 return defaultEntrypointHandlers != kj::none;
2066 }
2067 
2068 kj::Array<kj::StringPtr> getEntrypointNames() {
2069 return KJ_MAP(e, namedEntrypoints) -> kj::StringPtr { return e.key; };
2070 }
2071 
2072 kj::Array<kj::StringPtr> getActorClassNames() {
2073 return KJ_MAP(name, actorClassEntrypoints) -> kj::StringPtr { return name; };
2074 }
2075 
2076 void link(Worker::ValidationErrorReporter& errorReporter) override {
2077 LinkCallback callback =
2078 kj::mv(KJ_REQUIRE_NONNULL(ioChannels.tryGet<LinkCallback>(), "already called link()"));
2079 auto linked = callback(*this, errorReporter);
2080 
2081 for (auto& ns: actorNamespaces) {
2082 ns.value->link(linked.actorStorage);
2083 }
2084 
2085 ioChannels = kj::mv(linked);
2086 }
2087 
2088 void unlink() override {
2089 // Need to remove all waited until tasks before destroying `ioChannels`
2090 waitUntilTasks.clear();
2091 
2092 // Need to tear down all actors before tearing down `ioChannels.actorStorage`.
2093 actorNamespaces.clear();
2094 
2095 // OK, now we can unlink.
2096 ioChannels = {};
2097 }
2098 
2099 kj::Maybe<ActorNamespace&> getActorNamespace(kj::StringPtr name) {
2100 KJ_IF_SOME(a, actorNamespaces.find(name)) {
2101 return *a;
2102 } else {
2103 return kj::none;
2104 }
2105 }
2106 
2107 kj::HashMap<kj::StringPtr, kj::Own<ActorNamespace>>& getActorNamespaces() {
2108 return actorNamespaces;
2109 }
2110 
2111 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override {
2112 return startRequest(kj::mv(metadata), kj::none, {});
2113 }
2114 
2115 bool hasHandler(kj::StringPtr handlerName) override {
2116 KJ_IF_SOME(h, defaultEntrypointHandlers) {
2117 return h.contains(handlerName);
2118 } else {
2119 return false;
2120 }
2121 }
2122 
2123 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata,
2124 kj::Maybe<kj::StringPtr> entrypointName,
2125 Frankenvalue props,
2126 kj::Maybe<kj::Own<Worker::Actor>> actor = kj::none,
2127 bool isTracer = false) {
2128 TRACE_EVENT("workerd", "Server::WorkerService::startRequest()");
2129 
2130 auto& channels = KJ_ASSERT_NONNULL(ioChannels.tryGet<LinkedIoChannels>());
2131 
2132 kj::Vector<kj::Own<WorkerInterface>> bufferedTailWorkers(channels.tails.size());
2133 kj::Vector<kj::Own<WorkerInterface>> streamingTailWorkers(channels.streamingTails.size());
2134 auto addWorkerIfNotRecursiveTracer = [this, isTracer](
2135 kj::Vector<kj::Own<WorkerInterface>>& workers,
2136 IoChannelFactory::SubrequestChannel& channel) {
2137 // Caution here... if the tail worker ends up having a circular dependency
2138 // on the worker we'll end up with an infinite loop trying to initialize.
2139 // We can test this directly but it's more difficult to test indirect
2140 // loops (dependency of dependency, etc). Here we're just going to keep
2141 // it simple and just check the direct dependency.
2142 // If service refers to an EntrypointService, we need to compare with the underlying
2143 // WorkerService to match this.
2144 auto& service = KJ_UNWRAP_OR(kj::dynamicDowncastIfAvailable<Service>(channel), {
2145 // Not a Service, probably not self-referential.
2146 workers.add(channel.startRequest({}));
2147 return;
2148 });
2149 
2150 if (service.service() == this) {
2151 if (!isTracer) {
2152 // This is a self-reference. Create a request with isTracer=true.
2153 KJ_IF_SOME(s, kj::dynamicDowncastIfAvailable<WorkerService>(service)) {
2154 workers.add(s.startRequest({}, kj::none, {}, kj::none, true));
2155 } else KJ_IF_SOME(s, kj::dynamicDowncastIfAvailable<EntrypointService>(service)) {
2156 workers.add(s.startRequest({}, true));
2157 } else {
2158 KJ_FAIL_ASSERT("Unexpected service type in recursive tail worker declaration");
2159 }
2160 } else {
2161 // Intentionally left empty to prevent infinite recursion with tail workers tailing
2162 // themselves
2163 }
2164 } else {
2165 workers.add(service.startRequest({}));
2166 }
2167 };
2168 
2169 // Do not add tracers for worker interfaces with the "test" entrypoint โ€“ we generally do not
2170 // need to trace the test event, although this is useful to test that span tracing works, so
2171 // we are not implementing a (more complex) mechanism to disable tracing for all test() events
2172 // here.
2173 if (entrypointName.orDefault("") != "test"_kj) {
2174 for (auto& service: channels.tails) {
2175 addWorkerIfNotRecursiveTracer(bufferedTailWorkers, *service);
2176 }
2177 for (auto& service: channels.streamingTails) {
2178 addWorkerIfNotRecursiveTracer(streamingTailWorkers, *service);
2179 }
2180 }
2181 
2182 kj::Maybe<kj::Own<WorkerTracer>> workerTracer = kj::none;
2183 
2184 if (!bufferedTailWorkers.empty() || !streamingTailWorkers.empty()) {
2185 // Setting up buffered tail workers support, but only if we actually have tail workers
2186 // configured.
2187 auto executionModel =
2188 actor == kj::none ? ExecutionModel::STATELESS : ExecutionModel::DURABLE_OBJECT;
2189 auto tailStreamWriter = tracing::initializeTailStreamWriter(
2190 streamingTailWorkers.releaseAsArray(), waitUntilTasks);
2191 auto trace = kj::refcounted<Trace>(kj::none /* stableId */, kj::none /* scriptName */,
2192 kj::none /* scriptVersion */, kj::none /* dispatchNamespace */, kj::none /* scriptId */,
2193 nullptr /* scriptTags */, mapCopyString(entrypointName), executionModel,
2194 kj::none /* durableObjectId */);
2195 kj::Own<WorkerTracer> tracer = kj::refcounted<WorkerTracer>(
2196 kj::none, kj::mv(trace), PipelineLogLevel::FULL, kj::none, kj::mv(tailStreamWriter));
2197 
2198 // When the tracer is complete, deliver traces to any buffered tail workers. We end up
2199 // creating two references to the WorkerTracer, one held by the observer and one that will be
2200 // passed to the IoContext. This ensures that the tracer lives long enough to receive all
2201 // events.
2202 if (!bufferedTailWorkers.empty()) {
2203 waitUntilTasks.add(tracer->onComplete().then(
2204 kj::coCapture([tailWorkers = bufferedTailWorkers.releaseAsArray()](
2205 kj::Own<Trace> trace) mutable -> kj::Promise<void> {
2206 for (auto& worker: tailWorkers) {
2207 auto event = kj::heap<workerd::api::TraceCustomEvent>(
2208 workerd::api::TraceCustomEvent::TYPE, kj::arr(kj::addRef(*trace)));
2209 co_await worker->customEvent(kj::mv(event)).ignoreResult();
2210 }
2211 co_return;
2212 })));
2213 }
2214 workerTracer = kj::mv(tracer);
2215 }
2216 
2217 KJ_IF_SOME(w, workerTracer) {
2218 w->setMakeUserRequestSpanFunc(
2219 [&w = *w, &entropySource = threadContext.getEntropySource()](
2220 tracing::TraceId traceId, kj::Maybe<tracing::TraceFlags> traceFlags) {
2221 return SpanParent(kj::refcounted<UserSpanObserver>(
2222 kj::refcounted<SequentialSpanSubmitter>(w.getWeakRef(), entropySource), kj::mv(traceId),
2223 traceFlags));
2224 });
2225 }
2226 kj::Own<RequestObserver> observer =
2227 kj::refcounted<RequestObserverWithTracer>(mapAddRef(workerTracer), waitUntilTasks);
2228 
2229 kj::Maybe<tracing::InvocationSpanContext> triggerContext;
2230 KJ_IF_SOME(ctx, metadata.userSpanParent.toSpanContext()) {
2231 KJ_IF_SOME(spanId, ctx.getSpanId()) {
2232 triggerContext = tracing::InvocationSpanContext(
2233 ctx.getTraceId(), tracing::TraceId::nullId, spanId, ctx.getTraceFlags());
2234 }
2235 }
2236 
2237 return newWorkerEntrypoint(threadContext, kj::atomicAddRef(*worker), entrypointName,
2238 kj::mv(props), kj::mv(actor), kj::Own<LimitEnforcer>(this, kj::NullDisposer::instance),
2239 {}, // ioContextDependency
2240 kj::Own<IoChannelFactory>(this, kj::NullDisposer::instance), kj::mv(observer),
2241 waitUntilTasks,
2242 true, // tunnelExceptions
2243 kj::mv(workerTracer), // workerTracer
2244 kj::mv(metadata.cfBlobJson),
2245 kj::none, // versionInfo
2246 kj::mv(triggerContext));
2247 }
2248 
2249 class ActorNamespace final {
2250 public:
2251 ActorNamespace(kj::Own<ActorClass> actorClass,
2252 const ActorConfig& config,
2253 const kj::Clock& clock,
2254 kj::Timer& timer,
2255 capnp::ByteStreamFactory& byteStreamFactory,
2256 ChannelTokenHandler& channelTokenHandler,
2257 kj::Network& dockerNetwork,
2258 kj::Maybe<kj::StringPtr> dockerPath,
2259 kj::Maybe<kj::StringPtr> containerEgressInterceptorImage,
2260 kj::TaskSet& waitUntilTasks)
2261 : actorClass(kj::mv(actorClass)),
2262 config(config),
2263 clock(clock),
2264 timer(timer),
2265 byteStreamFactory(byteStreamFactory),
2266 channelTokenHandler(channelTokenHandler),
2267 dockerNetwork(dockerNetwork),
2268 dockerPath(dockerPath),
2269 containerEgressInterceptorImage(containerEgressInterceptorImage),
2270 waitUntilTasks(waitUntilTasks) {}
2271 
2272 void link(kj::Maybe<const kj::Directory&> serviceActorStorage) {
2273 KJ_IF_SOME(dir, serviceActorStorage) {
2274 KJ_IF_SOME(d, config.tryGet<Durable>()) {
2275 this->actorStorage.emplace(dir.openSubdir(
2276 kj::Path({d.uniqueKey}), kj::WriteMode::CREATE | kj::WriteMode::MODIFY));
2277 }
2278 }
2279 
2280 KJ_IF_SOME(d, config.tryGet<Durable>()) {
2281 auto idFactory = kj::heap<ActorIdFactoryImpl>(d.uniqueKey);
2282 AlarmScheduler::GetActorFn getActor =
2283 [this, idFactory = kj::mv(idFactory)](
2284 kj::String idStr) mutable -> kj::Own<WorkerInterface> {
2285 Worker::Actor::Id id = idFactory->idFromString(kj::mv(idStr));
2286 auto actorContainer = this->getActorContainer(kj::mv(id));
2287 return newPromisedWorkerInterface(actorContainer->startRequest({}));
2288 };
2289 
2290 KJ_IF_SOME(as, this->actorStorage) {
2291 // Create per-namespace alarm scheduler backed by on-disk storage in the
2292 // namespace directory, alongside the per-actor .sqlite files.
2293 this->ownAlarmScheduler = kj::heap<AlarmScheduler>(
2294 clock, timer, as.vfs, kj::Path({"metadata.sqlite"}), kj::mv(getActor));
2295 } else {
2296 // No on-disk storage -- create an in-memory alarm scheduler.
2297 auto memDir = kj::newInMemoryDirectory(clock);
2298 auto vfs = kj::heap<SqliteDatabase::Vfs>(*memDir);
2299 this->ownAlarmScheduler = kj::heap<AlarmScheduler>(
2300 clock, timer, *vfs, kj::Path({"metadata.sqlite"}), kj::mv(getActor))
2301 .attach(kj::mv(vfs), kj::mv(memDir));
2302 }
2303 
2304 this->alarmScheduler = *KJ_ASSERT_NONNULL(ownAlarmScheduler);
2305 }
2306 }
2307 
2308 const ActorConfig& getConfig() {
2309 return config;
2310 }
2311 
2312 kj::Own<IoChannelFactory::ActorChannel> getActorChannel(Worker::Actor::Id id) {
2313 KJ_IF_SOME(doId, id.tryGet<kj::Own<ActorIdFactory::ActorId>>()) {
2314 KJ_IF_SOME(name, doId->getName()) {
2315 // To emulate production, we preserve the name on the id, but only if it's <= 1024 bytes.
2316 if (name.size() > 1024) {
2317 auto* idImpl = dynamic_cast<ActorIdFactoryImpl::ActorIdImpl*>(doId.get());
2318 KJ_ASSERT(idImpl != nullptr, "Unexpected ActorId type?");
2319 idImpl->clearName();
2320 }
2321 }
2322 }
2323 
2324 return kj::refcounted<ActorChannelImpl>(getActorContainer(kj::mv(id)));
2325 }
2326 
2327 class ActorContainer;
2328 using ActorMap = kj::HashMap<kj::StringPtr, kj::Own<ActorContainer>>;
2329 
2330 // ActorContainer mostly serves as a wrapper around Worker::Actor.
2331 // We use it to associate a HibernationManager with the Worker::Actor, since the
2332 // Worker::Actor can be destroyed during periods of prolonged inactivity.
2333 //
2334 // We use a RequestTracker to track strong references to this ActorContainer's Worker::Actor.
2335 // Once there are no Worker::Actor's left (excluding our own), `inactive()` is triggered and we
2336 // initiate the eviction of the Durable Object. If no requests arrive in the next 10 seconds,
2337 // the DO is evicted, otherwise we cancel the eviction task.
2338 class ActorContainer final: public RequestTracker::Hooks,
2339 public kj::Refcounted,
2340 public Worker::Actor::FacetManager {
2341 public:
2342 // Information which is needed before start() can be called, but may not be available yet
2343 // when the ActorContainer is constructed (especially in the case of facets).
2344 struct ClassAndId {
2345 kj::Own<ActorClass> actorClass;
2346 Worker::Actor::Id id;
2347 
2348 ClassAndId(kj::Own<ActorClass> actorClass, Worker::Actor::Id id)
2349 : actorClass(kj::mv(actorClass)),
2350 id(kj::mv(id)) {}
2351 };
2352 
2353 ActorContainer(kj::String key,
2354 ActorNamespace& ns,
2355 kj::Maybe<ActorContainer&> parent,
2356 kj::OneOf<ClassAndId, kj::Promise<ClassAndId>> classAndIdParam,
2357 kj::Timer& timer)
2358 : key(kj::mv(key)),
2359 tracker(kj::refcounted<RequestTracker>(*this)),
2360 ns(ns),
2361 root(parent.map([](ActorContainer& p) -> ActorContainer& { return p.root; })
2362 .orDefault(*this)),
2363 parent(parent),
2364 timer(timer),
2365 lastAccess(timer.now()) {
2366 KJ_SWITCH_ONEOF(classAndIdParam) {
2367 KJ_CASE_ONEOF(value, ClassAndId) {
2368 // `classAndId` is immediately available.
2369 classAndId = kj::mv(value);
2370 }
2371 KJ_CASE_ONEOF(promise, kj::Promise<ClassAndId>) {
2372 // We are receiving a promise for a `ClassAndId` to come later. Arrange to initialize
2373 // `classAndId` from the promise. Create a `ForkedPromise<void>` that resolves when
2374 // initialization is complete.
2375 classAndId = promise
2376 .then([this](ClassAndId value) {
2377 auto& forked = KJ_ASSERT_NONNULL(classAndId.tryGet<kj::ForkedPromise<void>>());
2378 if (!forked.hasBranches()) {
2379 // HACK: We're about to replace the ForkedPromise but it has no one waiting on it,
2380 // so we'd end up cancelling ourselves. Add a branch and detach it so this doesn't
2381 // happen.
2382 forked.addBranch().detach([](auto&&) {});
2383 }
2384 
2385 classAndId = kj::mv(value);
2386 }).fork();
2387 }
2388 }
2389 }
2390 
2391 ~ActorContainer() noexcept(false) {
2392 // Shutdown the tracker so we don't use active/inactive hooks anymore.
2393 tracker->shutdown();
2394 
2395 for (auto& facet: facets) {
2396 facet.value->abort(kj::none);
2397 }
2398 
2399 KJ_IF_SOME(a, actor) {
2400 // Unknown broken reason.
2401 auto reason = 0;
2402 a->shutdown(reason);
2403 }
2404 
2405 // Drop the container client reference
2406 // If setInactivityTimeout() was called, there's still a timer holding a reference
2407 // If not, this may be the last reference and the ContainerClient destructor will run
2408 containerClient = kj::none;
2409 }
2410 
2411 void active() override {
2412 // We're handling a new request, cancel the eviction promise.
2413 shutdownTask = kj::none;
2414 }
2415 
2416 void inactive() override {
2417 // Durable objects are evictable by default.
2418 bool isEvictable = true;
2419 KJ_SWITCH_ONEOF(ns.config) {
2420 KJ_CASE_ONEOF(c, Durable) {
2421 isEvictable = c.isEvictable;
2422 }
2423 KJ_CASE_ONEOF(c, Ephemeral) {
2424 isEvictable = c.isEvictable;
2425 }
2426 }
2427 if (isEvictable) {
2428 KJ_IF_SOME(a, actor) {
2429 KJ_IF_SOME(m, a->getHibernationManager()) {
2430 // The hibernation manager needs to survive actor eviction and be passed to the actor
2431 // constructor next time we create it.
2432 manager = m.addRef();
2433 }
2434 }
2435 shutdownTask =
2436 handleShutdown().eagerlyEvaluate([](kj::Exception&& e) { KJ_LOG(ERROR, e); });
2437 }
2438 }
2439 
2440 kj::StringPtr getKey() {
2441 return key;
2442 }
2443 RequestTracker& getTracker() {
2444 return *tracker;
2445 }
2446 kj::Maybe<kj::Own<Worker::Actor::HibernationManager>> tryGetManagerRef() {
2447 return manager.map(
2448 [&](kj::Own<Worker::Actor::HibernationManager>& m) { return kj::addRef(*m); });
2449 }
2450 void updateAccessTime() {
2451 lastAccess = timer.now();
2452 KJ_IF_SOME(p, parent) {
2453 p.updateAccessTime();
2454 }
2455 }
2456 kj::TimePoint getLastAccess() {
2457 return lastAccess;
2458 }
2459 
2460 bool hasClients() {
2461 // If anyone holds a reference to the container other than the actor map, then it must be
2462 // a client.
2463 if (isShared()) return true;
2464 for (auto& facet: facets) {
2465 if (facet.value->hasClients()) return true;
2466 }
2467 return false;
2468 }
2469 kj::Own<ActorContainer> addRef() {
2470 return kj::addRef(*this);
2471 }
2472 
2473 // Get the actor, starting it if it's not already running.
2474 kj::Promise<kj::Own<Worker::Actor>> getActor() {
2475 requireNotBroken();
2476 
2477 if (actor == kj::none) {
2478 KJ_IF_SOME(promise, classAndId.tryGet<kj::ForkedPromise<void>>()) {
2479 co_await promise;
2480 }
2481 
2482 auto& [actorClass, id] = KJ_ASSERT_NONNULL(classAndId.tryGet<ClassAndId>());
2483 
2484 KJ_IF_SOME(promise, actorClass->whenReady()) {
2485 co_await promise;
2486 }
2487 
2488 // A concurrent request could have started the actor, so check again.
2489 if (actor == kj::none) {
2490 start(actorClass, id);
2491 }
2492 }
2493 
2494 co_return KJ_ASSERT_NONNULL(actor)->addRef();
2495 }
2496 
2497 kj::Promise<kj::Own<WorkerInterface>> startRequest(
2498 IoChannelFactory::SubrequestMetadata metadata) {
2499 auto actor = co_await getActor();
2500 
2501 if (ns.cleanupTask == kj::none) {
2502 // Need to start the cleanup loop.
2503 ns.cleanupTask = ns.cleanupLoop();
2504 }
2505 
2506 // Since `getActor()` completed, `classAndId` must be resolved.
2507 auto& actorClass = KJ_ASSERT_NONNULL(classAndId.tryGet<ClassAndId>()).actorClass;
2508 
2509 co_return actorClass->startRequest(kj::mv(metadata), kj::mv(actor))
2510 .attach(kj::defer([self = kj::addRef(*this)]() mutable { self->updateAccessTime(); }));
2511 }
2512 
2513 // Abort this actor, shutting it down.
2514 //
2515 // It is the caller's responsibility to ensure that the aborted ActorContainer has been
2516 // removed from any maps that would cause it to receive further traffic, since any further
2517 // requests will be expected to fail. abort() does NOT attempt to remove the ActorContainer
2518 // from the parent facet map since at most call sites it makes more sense to handle this
2519 // directly.
2520 void abort(kj::Maybe<const kj::Exception&> reason) {
2521 if (brokenReason != kj::none) return;
2522 
2523 KJ_IF_SOME(a, actor) {
2524 // Unknown broken reason.
2525 a->shutdown(0, reason);
2526 }
2527 
2528 for (auto& facet: facets) {
2529 facet.value->abort(reason);
2530 }
2531 
2532 onBrokenTask = kj::none;
2533 shutdownTask = kj::none;
2534 manager = kj::none;
2535 tracker->shutdown();
2536 actor = kj::none;
2537 containerClient = kj::none;
2538 
2539 KJ_IF_SOME(r, reason) {
2540 brokenReason = r.clone();
2541 } else {
2542 brokenReason = JSG_KJ_EXCEPTION(FAILED, Error, "Actor aborted for uknown reason.");
2543 }
2544 }
2545 
2546 // Resets the actor's SQLite database while the connection is still open,
2547 // avoiding file-locking issues on Windows.
2548 void resetStorage() {
2549 KJ_IF_SOME(a, actor) {
2550 KJ_IF_SOME(cache, a->getPersistent()) {
2551 KJ_IF_SOME(db, cache.getSqliteDatabase()) {
2552 kj::runCatchingExceptions([&]() { db.reset(); });
2553 }
2554 }
2555 }
2556 }
2557 
2558 kj::Own<ActorContainer> getFacetContainer(
2559 kj::String childKey, kj::Function<kj::Promise<StartInfo>()> getStartInfo) {
2560 auto makeContainer = [&]() {
2561 auto promise = callFacetStartCallback(kj::mv(getStartInfo));
2562 return kj::refcounted<ActorContainer>(
2563 kj::mv(childKey), ns, *this, kj::mv(promise), timer);
2564 };
2565 
2566 bool isNew = false;
2567 
2568 auto& entry = facets.findOrCreateEntry(childKey, [&]() mutable {
2569 isNew = true;
2570 auto container = makeContainer();
2571 return ActorMap::Entry{container->getKey(), kj::mv(container)};
2572 });
2573 
2574 return entry.value->addRef();
2575 }
2576 
2577 uint getDepth() const override {
2578 KJ_IF_SOME(p, parent) {
2579 return 1 + p.getDepth();
2580 }
2581 return 0;
2582 }
2583 
2584 kj::Own<IoChannelFactory::ActorChannel> getFacet(
2585 kj::StringPtr name, kj::Function<kj::Promise<StartInfo>()> getStartInfo) override {
2586 auto facet = getFacetContainer(kj::str(name), kj::mv(getStartInfo));
2587 return kj::refcounted<ActorChannelImpl>(kj::mv(facet));
2588 }
2589 
2590 void abortFacet(kj::StringPtr name, kj::Exception reason) override {
2591 KJ_IF_SOME(entry, facets.findEntry(name)) {
2592 entry.value->abort(reason);
2593 facets.erase(entry);
2594 }
2595 }
2596 
2597 void deleteFacet(kj::StringPtr name) override {
2598 // First, abort any running facets.
2599 abortFacet(name, JSG_KJ_EXCEPTION(FAILED, Error, "Facet was deleted."));
2600 
2601 // Then delete the underlying storage.
2602 KJ_IF_SOME(as, ns.actorStorage) {
2603 // Note that if there's no facet index then there couldn't possibly be any child storage.
2604 KJ_IF_SOME(index, getFacetTreeIndexIfNotEmpty()) {
2605 uint childId = index.getId(getFacetId(), name);
2606 deleteDescendantStorage(*as.directory, childId);
2607 as.directory->remove(getSqlitePathForId(childId));
2608 }
2609 }
2610 }
2611 
2612 private:
2613 // The actor is constructed after the ActorContainer so it starts off empty.
2614 kj::Maybe<kj::Own<Worker::Actor>> actor;
2615 
2616 kj::String key;
2617 kj::Own<RequestTracker> tracker;
2618 ActorNamespace& ns;
2619 ActorContainer& root;
2620 kj::Maybe<ActorContainer&> parent;
2621 kj::Timer& timer;
2622 kj::TimePoint lastAccess;
2623 kj::Maybe<kj::Own<Worker::Actor::HibernationManager>> manager;
2624 kj::Maybe<kj::Promise<void>> shutdownTask;
2625 kj::Maybe<kj::Promise<void>> onBrokenTask;
2626 kj::Maybe<kj::Exception> brokenReason;
2627 
2628 // Reference to the ContainerClient (if container is enabled for this actor)
2629 kj::Maybe<kj::Own<ContainerClient>> containerClient;
2630 
2631 // If this is a `ForkedPromise<void>`, await the promise. When it has resolved, then
2632 // `classAndId` will have been replaced with the resolved `ClassAndId` value.
2633 kj::OneOf<ClassAndId, kj::ForkedPromise<void>> classAndId;
2634 
2635 // FacetTreeIndex for this actor. Only initialized on the root.
2636 kj::Maybe<kj::Own<FacetTreeIndex>> facetTreeIndex;
2637 
2638 // ID of this facet. Initialized when getFacetId() is first called.
2639 kj::Maybe<uint> facetId;
2640 
2641 ActorMap facets;
2642 
2643 // Get the facet ID for this facet. The root facet always has ID zero, but all other facets
2644 // need to be looked up in the index to make sure they are assigned consistent IDs.
2645 uint getFacetId() {
2646 KJ_IF_SOME(f, facetId) {
2647 return f;
2648 }
2649 
2650 ActorContainer& parent = KJ_UNWRAP_OR(this->parent, return 0);
2651 
2652 FacetTreeIndex& index = root.ensureFacetTreeIndex();
2653 return index.getId(parent.getFacetId(), key);
2654 }
2655 
2656 // Get the facet tree index, opening the file if it hasn't been opened yet, and creating it
2657 // if it hasn't been created yet.
2658 FacetTreeIndex& ensureFacetTreeIndex() {
2659 KJ_REQUIRE(parent == kj::none, "only 'root' may ensureFacetTreeIndex()");
2660 
2661 KJ_IF_SOME(i, facetTreeIndex) {
2662 return *i;
2663 } else {
2664 // Facet tree index hasn't been initialized yet. Do that now (opening the existing file,
2665 // or creating it if it doesn't exist).
2666 auto& as = KJ_REQUIRE_NONNULL(
2667 ns.actorStorage, "can't call getFacetId() when there's no backing storage");
2668 auto indexFile = as.directory->openFile(
2669 kj::Path({kj::str(key, ".facets")}), kj::WriteMode::CREATE | kj::WriteMode::MODIFY);
2670 return *facetTreeIndex.emplace(kj::heap<FacetTreeIndex>(kj::mv(indexFile)));
2671 }
2672 }
2673 
2674 // Like ensureFacetTreeIndex() but if the index doesn't exist on disk, return kj::none.
2675 kj::Maybe<FacetTreeIndex&> getFacetTreeIndexIfNotEmpty() {
2676 KJ_REQUIRE(parent == kj::none);
2677 
2678 KJ_IF_SOME(i, facetTreeIndex) {
2679 return *i;
2680 } else {
2681 // Facet tree index hasn't been initialized yet. If the file exists, open it. Otherwise,
2682 // assume empty and return none.
2683 auto& as = KJ_UNWRAP_OR(ns.actorStorage, return kj::none);
2684 auto indexFile = KJ_UNWRAP_OR(
2685 as.directory->tryOpenFile(kj::Path({kj::str(key, ".facets")}), kj::WriteMode::MODIFY),
2686 return kj::none);
2687 return *facetTreeIndex.emplace(kj::heap<FacetTreeIndex>(kj::mv(indexFile)));
2688 }
2689 }
2690 
2691 // Get the path to the facet's sqlite database, within the actor namespace directory.
2692 kj::Path getSqlitePathForId(uint id) {
2693 if (id == 0) {
2694 return kj::Path({kj::str(root.key, ".sqlite")});
2695 } else {
2696 return kj::Path({kj::str(root.key, '.', id, ".sqlite")});
2697 }
2698 }
2699 
2700 void deleteDescendantStorage(const kj::Directory& dir, uint parentId) {
2701 KJ_IF_SOME(index, getFacetTreeIndexIfNotEmpty()) {
2702 deleteDescendantStorage(dir, index, parentId);
2703 } else {
2704 // There's no index, so there must be no facets (other than the root).
2705 KJ_ASSERT(parentId == 0);
2706 }
2707 }
2708 
2709 void deleteDescendantStorage(const kj::Directory& dir, FacetTreeIndex& index, uint parentId) {
2710 index.forEachChild(parentId, [&](uint childId, kj::StringPtr childName) {
2711 deleteDescendantStorage(dir, index, childId);
2712 dir.remove(getSqlitePathForId(childId));
2713 });
2714 }
2715 
2716 void requireNotBroken() {
2717 KJ_IF_SOME(e, brokenReason) {
2718 kj::throwFatalException(e.clone());
2719 }
2720 }
2721 
2722 kj::Promise<void> monitorOnBroken(Worker::Actor& actor) {
2723 try {
2724 // It's possible for this to never resolve if the actor never breaks,
2725 // in which case the returned promise will just be canceled.
2726 co_await actor.onBroken();
2727 KJ_FAIL_ASSERT("actor.onBroken() resolved normally?");
2728 } catch (...) {
2729 brokenReason = kj::getCaughtExceptionAsKj();
2730 }
2731 
2732 for (auto& facet: facets) {
2733 facet.value->abort(brokenReason);
2734 }
2735 facets.clear();
2736 
2737 // HACK: Dropping the ActorContainer will delete onBrokenTask, cancelling ourselves. This
2738 // would crash. To avoid the problem, detach ourselves. This is safe because we know that
2739 // once we return there's nothing left for this promise to do anyway.
2740 KJ_ASSERT_NONNULL(onBrokenTask).detach([](kj::Exception&& e) {});
2741 
2742 // Hollow out the object, so that if it still has references, they won't keep these parts
2743 // alive. Since any further calls to `getActor()` will throw, we don't have to worry about
2744 // the actor being recreated.
2745 auto actorToDrop = kj::mv(this->actor);
2746 tracker->shutdown();
2747 auto managerToDrop = kj::mv(manager);
2748 
2749 // Note that we remove the entire ActorContainer from the map -- this drops the
2750 // HibernationManager so any connected hibernatable websockets will be disconnected.
2751 KJ_IF_SOME(p, parent) {
2752 p.facets.erase(key);
2753 } else {
2754 ns.actors.erase(key);
2755 }
2756 
2757 // WARNING: `this` MAY HAVE BEEN DELETED as a result of the above `erase()`. Do not access
2758 // it again here.
2759 }
2760 
2761 // Processes the eviction of the Durable Object and hibernates active websockets.
2762 kj::Promise<void> handleShutdown() {
2763 // After 10 seconds of inactivity, we destroy the Worker::Actor and hibernate any active
2764 // JS WebSockets.
2765 // TODO(someday): We could make this timeout configurable to make testing less burdensome.
2766 co_await timer.afterDelay(10 * kj::SECONDS);
2767 // Cancel the onBroken promise, since we're about to destroy the actor anyways and don't
2768 // want to trigger it.
2769 onBrokenTask = kj::none;
2770 KJ_IF_SOME(a, actor) {
2771 if (a->isShared()) {
2772 // Our ActiveRequest refcounting has broken somewhere. This is likely because we're
2773 // `addRef`-ing an actor that has had an ActiveRequest attached to its kj::Own (in other
2774 // words, the ActiveRequest count is less than it should be).
2775 //
2776 // Rather than dropping our actor and possibly ending up with split-brain,
2777 // we should opt out of the deferred proxy optimization and log the error to Sentry.
2778 KJ_LOG(ERROR,
2779 "Detected internal bug in hibernation: Durable Object has strong references "
2780 "when hibernation timeout expired.");
2781 
2782 co_return;
2783 }
2784 KJ_IF_SOME(m, manager) {
2785 auto& worker = a->getWorker();
2786 auto workerStrongRef = kj::atomicAddRef(worker);
2787 // Take an async lock, we can't use `takeAsyncLock(RequestObserver&)` since we don't
2788 // have an `IncomingRequest` at this point.
2789 //
2790 // Note that we do not have a race here because this is part of the `shutdownTask`
2791 // promise. If a new request comes in while we're waiting to get the lock then we will
2792 // cancel this promise.
2793 Worker::AsyncLock asyncLock = co_await worker.takeAsyncLockWithoutRequest(nullptr);
2794 workerStrongRef->runInLockScope(
2795 asyncLock, [&](Worker::Lock& lock) { m->hibernateWebSockets(lock); });
2796 }
2797 a->shutdown(
2798 0, KJ_EXCEPTION(DISCONNECTED, "broken.dropped; Actor freed due to inactivity"));
2799 }
2800 // Destroy the last strong Worker::Actor reference.
2801 actor = kj::none;
2802 
2803 // Drop our reference to the ContainerClient
2804 // If setInactivityTimeout() was called, the timer still holds a reference
2805 // so the container stays alive until the timeout expires
2806 containerClient = kj::none;
2807 }
2808 
2809 void start(kj::Own<ActorClass>& actorClass, Worker::Actor::Id& id) {
2810 KJ_REQUIRE(actor == nullptr);
2811 
2812 auto makeActorCache = [this](const ActorCache::SharedLru& sharedLru, OutputGate& outputGate,
2813 ActorCache::Hooks& hooks,
2814 SqliteObserver& sqliteObserver) mutable {
2815 return ns.config.tryGet<Durable>().map(
2816 [&](const Durable& d) -> kj::Own<ActorCacheInterface> {
2817 KJ_IF_SOME(as, ns.actorStorage) {
2818 kj::Own<ActorSqlite::Hooks> sqliteHooks;
2819 if (parent == kj::none) {
2820 KJ_IF_SOME(a, ns.alarmScheduler) {
2821 sqliteHooks = kj::heap<ActorSqliteHooks>(a, ActorKey{.actorId = key});
2822 } else {
2823 // No alarm scheduler available, use default hooks instance.
2824 sqliteHooks = fakeOwn(ActorSqlite::Hooks::getDefaultHooks());
2825 }
2826 } else {
2827 // TODO(someday): Support alarms in facets, somehow.
2828 sqliteHooks = fakeOwn(ActorSqlite::Hooks::getDefaultHooks());
2829 }
2830 
2831 uint selfId = getFacetId();
2832 auto path = getSqlitePathForId(selfId);
2833 auto db = kj::heap<SqliteDatabase>(
2834 as.vfs, kj::mv(path), kj::WriteMode::CREATE | kj::WriteMode::MODIFY);
2835 
2836 // Before we do anything, make sure the database is in WAL mode. We also need to
2837 // do this after reset() is used, so register a callback for that.
2838 db->run("PRAGMA journal_mode=WAL;");
2839 
2840 db->afterReset([this, &dir = *as.directory, selfId](SqliteDatabase& db) {
2841 db.run("PRAGMA journal_mode=WAL;");
2842 
2843 // reset() is used when the app called deleteAll(), in which case we also want to
2844 // delete all child facets.
2845 // TODO(someday): Arguably this should be transactional somehow so if we fail here
2846 // we don't leave the facets still there after the parent has already been reset.
2847 // But most filesystems do not support transactions, so we'd have to do something
2848 // like store a flag in the parent DB saying "reset pending" so that on a restart
2849 // we retry the deletions. Note that in production on SRS, this is actually
2850 // transactional -- there's only a problem when running locally with workerd.
2851 deleteDescendantStorage(dir, selfId);
2852 });
2853 
2854 return kj::heap<ActorSqlite>(kj::mv(db), outputGate,
2855 [](SpanParent) -> kj::Promise<void> { return kj::READY_NOW; }, *sqliteHooks)
2856 .attach(kj::mv(sqliteHooks));
2857 } else {
2858 // Create an ActorCache backed by a fake, empty storage. Elsewhere, we configure
2859 // ActorCache never to flush, so this effectively creates in-memory storage.
2860 return kj::heap<ActorCache>(
2861 newEmptyReadOnlyActorStorage(), sharedLru, outputGate, hooks);
2862 }
2863 });
2864 };
2865 
2866 bool enableSql = true;
2867 kj::Maybe<config::Worker::DurableObjectNamespace::ContainerOptions::Reader>
2868 containerOptions = kj::none;
2869 kj::Maybe<kj::StringPtr> uniqueKey;
2870 KJ_SWITCH_ONEOF(ns.config) {
2871 KJ_CASE_ONEOF(c, Durable) {
2872 enableSql = c.enableSql;
2873 containerOptions = c.containerOptions;
2874 uniqueKey = c.uniqueKey;
2875 }
2876 KJ_CASE_ONEOF(c, Ephemeral) {
2877 enableSql = c.enableSql;
2878 }
2879 }
2880 
2881 auto makeStorage =
2882 [enableSql = enableSql](jsg::Lock& js, const Worker::Api& api,
2883 ActorCacheInterface& actorCache) -> jsg::Ref<api::DurableObjectStorage> {
2884 return js.alloc<api::DurableObjectStorage>(
2885 js, IoContext::current().addObject(actorCache), enableSql);
2886 };
2887 
2888 auto loopback = kj::refcounted<Loopback>(*this);
2889 
2890 kj::Maybe<rpc::Container::Client> container = kj::none;
2891 KJ_IF_SOME(config, containerOptions) {
2892 KJ_ASSERT(config.hasImageName(), "Image name is required");
2893 auto imageName = config.getImageName();
2894 kj::String containerId;
2895 KJ_SWITCH_ONEOF(id) {
2896 KJ_CASE_ONEOF(globalId, kj::Own<ActorIdFactory::ActorId>) {
2897 containerId = globalId->toString();
2898 }
2899 KJ_CASE_ONEOF(existingId, kj::String) {
2900 containerId = kj::str(existingId);
2901 }
2902 }
2903 
2904 container = ns.getContainerClient(
2905 kj::str("workerd-", KJ_ASSERT_NONNULL(uniqueKey), "-", containerId), imageName);
2906 }
2907 
2908 auto actor = actorClass->newActor(getTracker(), Worker::Actor::cloneId(id),
2909 kj::mv(makeActorCache), kj::mv(makeStorage), kj::mv(loopback), tryGetManagerRef(),
2910 kj::mv(container), *this);
2911 onBrokenTask = monitorOnBroken(*actor);
2912 this->actor = kj::mv(actor);
2913 }
2914 
2915 // Helper coroutine to call `getStartInfo()`, the start callback for a facet, while making
2916 // sure the function stays alive until the returned promise resolves.
2917 static kj::Promise<ClassAndId> callFacetStartCallback(
2918 kj::Function<kj::Promise<StartInfo>()> getStartInfo) {
2919 auto info = co_await getStartInfo();
2920 co_return ClassAndId(info.actorClass.downcast<ActorClass>(), kj::mv(info.id));
2921 }
2922 };
2923 
2924 kj::Own<ActorContainer> getActorContainer(Worker::Actor::Id id) {
2925 kj::String key;
2926 
2927 KJ_SWITCH_ONEOF(id) {
2928 KJ_CASE_ONEOF(obj, kj::Own<ActorIdFactory::ActorId>) {
2929 KJ_REQUIRE(config.is<Durable>());
2930 key = obj->toString();
2931 }
2932 KJ_CASE_ONEOF(str, kj::String) {
2933 KJ_REQUIRE(config.is<Ephemeral>());
2934 key = kj::str(str);
2935 }
2936 }
2937 
2938 return actors
2939 .findOrCreate(key, [&]() mutable {
2940 auto container = kj::refcounted<ActorContainer>(kj::mv(key), *this, kj::none,
2941 ActorContainer::ClassAndId(kj::addRef(*actorClass), kj::mv(id)), timer);
2942 
2943 return kj::HashMap<kj::StringPtr, kj::Own<ActorContainer>>::Entry{
2944 container->getKey(), kj::mv(container)};
2945 })->addRef();
2946 }
2947 
2948 kj::Own<ContainerClient> getContainerClient(
2949 kj::StringPtr containerId, kj::StringPtr imageName) {
2950 KJ_IF_SOME(existingClient, containerClients.find(containerId)) {
2951 return existingClient->addRef();
2952 }
2953 
2954 // No existing container in the map, create a new one
2955 auto& dockerPathRef = KJ_ASSERT_NONNULL(
2956 dockerPath, "dockerPath must be defined to enable containers on this Durable Object.");
2957 
2958 // Grab a branch of any pending cleanup from a previous ContainerClient for this
2959 // container. If it exists, pass it to the container client so it knows that it has to sync.
2960 kj::Promise<void> previousCleanup = kj::READY_NOW;
2961 KJ_IF_SOME(state, containerCleanupState.find(containerId)) {
2962 previousCleanup = state.promise.addBranch();
2963 }
2964 
2965 // Upsert the cleanup state for this container ID. Replacing the
2966 // canceler auto-cancels any in-flight cleanup tasks from the previous
2967 // client's destructor. The generation counter is bumped on replacement
2968 // so the cleanup callback can detect stale ownership without relying
2969 // on raw pointer identity (which is vulnerable to address reuse).
2970 auto canceler = kj::heap<kj::Canceler>();
2971 uint64_t capturedGeneration = 0;
2972 containerCleanupState.upsert(kj::str(containerId),
2973 ContainerCleanupState{.canceler = kj::mv(canceler)},
2974 [&capturedGeneration](ContainerCleanupState& existing, ContainerCleanupState&& incoming) {
2975 existing.canceler = kj::mv(incoming.canceler);
2976 capturedGeneration = ++existing.generation;
2977 });
2978 
2979 // Cleanup callback: invoked from the ContainerClient destructor with the joined
2980 // with a cleanup promise
2981 kj::Function<void(kj::Promise<void>)> cleanupCallback =
2982 [this, containerId = kj::str(containerId), capturedGeneration](
2983 kj::Promise<void> cleanupPromise) mutable {
2984 KJ_IF_SOME(state, containerCleanupState.find(containerId)) {
2985 if (state.generation != capturedGeneration) {
2986 // A newer ContainerClient has replaced us already with another destructor.
2987 // drop the promise.
2988 return;
2989 }
2990 
2991 containerClients.erase(containerId);
2992 // Wrap with the canceler so a future client creation can cancel these
2993 // tasks
2994 auto cancellable =
2995 state.canceler->wrap(kj::mv(cleanupPromise)).catch_([](kj::Exception&&) {});
2996 
2997 auto forked = kj::mv(cancellable).fork();
2998 waitUntilTasks.add(forked.addBranch());
2999 state.promise = kj::mv(forked);
3000 }
3001 };
3002 
3003 auto client = kj::refcounted<ContainerClient>(byteStreamFactory, timer, dockerNetwork,
3004 kj::str(dockerPathRef), kj::str(containerId), kj::str(imageName),
3005 kj::str(KJ_ASSERT_NONNULL(containerEgressInterceptorImage,
3006 "containerEgressInterceptorImage must be configured for containers.")),
3007 waitUntilTasks, kj::mv(previousCleanup), kj::mv(cleanupCallback), channelTokenHandler);
3008 
3009 // Store raw pointer in map (does not own)
3010 containerClients.insert(kj::str(containerId), client.get());
3011 
3012 return kj::mv(client);
3013 }
3014 
3015 void abortAll(kj::Maybe<const kj::Exception&> reason) {
3016 for (auto& actor: actors) {
3017 actor.value->abort(reason);
3018 }
3019 actors.clear();
3020 }
3021 
3022 // Resets all actor databases, aborts all actors, and cancels all alarms so DOs
3023 // can be recreated with clean state.
3024 void deleteAll(kj::Maybe<const kj::Exception&> reason) {
3025 // Reset databases before aborting so connections are still open (avoids
3026 // Windows file-locking issues with deferred handle release).
3027 for (auto& actor: actors) {
3028 actor.value->resetStorage();
3029 }
3030 
3031 abortAll(reason);
3032 
3033 KJ_IF_SOME(scheduler, ownAlarmScheduler) {
3034 scheduler->deleteAll();
3035 }
3036 }
3037 
3038 private:
3039 kj::Own<ActorClass> actorClass;
3040 const ActorConfig& config;
3041 const kj::Clock& clock;
3042 
3043 struct ActorStorage {
3044 kj::Own<const kj::Directory> directory;
3045 SqliteDatabase::Vfs vfs;
3046 
3047 ActorStorage(kj::Own<const kj::Directory> directoryParam)
3048 : directory(kj::mv(directoryParam)),
3049 vfs(*directory) {}
3050 };
3051 
3052 // Note: The Vfs, actorStorage, and ownAlarmScheduler must not be torn down until all actors
3053 // have been torn down, so we declare them before `actors`.
3054 kj::Maybe<ActorStorage> actorStorage;
3055 kj::Maybe<kj::Own<AlarmScheduler>> ownAlarmScheduler;
3056 
3057 // Tracks the canceler and cleanup promise for a Docker container's lifecycle cleanup.
3058 // Useful to await on async calls of a ContainerClient destructor when the new
3059 // one appears before they've been resolved.
3060 struct ContainerCleanupState {
3061 // Canceler that wraps the promise fired in ~ContainerClient. Replacing
3062 // it cancels any pending cleanup, which resolves the promise immediately.
3063 kj::Own<kj::Canceler> canceler;
3064 
3065 // Forked cleanup promise. A branch is added to waitUntilTasks to keep the I/O alive,
3066 // and another branch is passed to the next ContainerClient so its status() can await.
3067 kj::ForkedPromise<void> promise = kj::Promise<void>(kj::READY_NOW).fork();
3068 
3069 // Monotonically increasing counter, bumped each time the canceler is replaced
3070 // via upsert. The cleanup callback captures the generation at creation time and
3071 // compares it to detect whether a newer ContainerClient has taken ownership,
3072 // avoiding a raw-pointer identity check that is vulnerable to address reuse.
3073 uint64_t generation = 0;
3074 };
3075 
3076 // Per-container cleanup state: canceler + forked cleanup promise.
3077 kj::HashMap<kj::String, ContainerCleanupState> containerCleanupState;
3078 
3079 // Map of container IDs to ContainerClients (for reconnection support with inactivity timeouts).
3080 // The map holds raw pointers (not ownership) - ContainerClients are owned by actors and timers.
3081 // When the last reference is dropped, the destructor removes the entry from this map.
3082 kj::HashMap<kj::String, ContainerClient*> containerClients;
3083 
3084 // If the actor is broken, we remove it from the map. However, if it's just evicted due to
3085 // inactivity, we keep the ActorContainer in the map but drop the Own<Worker::Actor>. When a new
3086 // request comes in, we recreate the Own<Worker::Actor>.
3087 ActorMap actors;
3088 
3089 kj::Maybe<kj::Promise<void>> cleanupTask;
3090 kj::Timer& timer;
3091 capnp::ByteStreamFactory& byteStreamFactory;
3092 ChannelTokenHandler& channelTokenHandler;
3093 kj::Network& dockerNetwork;
3094 kj::Maybe<kj::StringPtr> dockerPath;
3095 kj::Maybe<kj::StringPtr> containerEgressInterceptorImage;
3096 kj::TaskSet& waitUntilTasks;
3097 kj::Maybe<AlarmScheduler&> alarmScheduler;
3098 
3099 // Removes actors from `actors` after 70 seconds of last access.
3100 kj::Promise<void> cleanupLoop() {
3101 constexpr auto EXPIRATION = 70 * kj::SECONDS;
3102 
3103 // Don't bother running the loop if the config doesn't allow eviction.
3104 KJ_SWITCH_ONEOF(config) {
3105 KJ_CASE_ONEOF(c, Durable) {
3106 if (!c.isEvictable) co_return;
3107 }
3108 KJ_CASE_ONEOF(c, Ephemeral) {
3109 if (!c.isEvictable) co_return;
3110 }
3111 }
3112 
3113 while (true) {
3114 auto now = timer.now();
3115 actors.eraseAll([&](auto&, kj::Own<ActorContainer>& entry) {
3116 // Check getLastAccess() before hasClients() since it's faster.
3117 if ((now - entry->getLastAccess()) <= EXPIRATION) {
3118 // Used recently; don't evict.
3119 return false;
3120 }
3121 
3122 if (entry->hasClients()) {
3123 // There's still an active client; don't evict.
3124 return false;
3125 }
3126 
3127 // No clients and not used in a while, evict this actor.
3128 return true;
3129 });
3130 
3131 co_await timer.atTime(now + EXPIRATION);
3132 }
3133 }
3134 
3135 // Implements actor loopback, which is used by websocket hibernation to deliver events to the
3136 // actor from the websocket's read loop.
3137 class Loopback: public Worker::Actor::Loopback, public kj::Refcounted {
3138 public:
3139 Loopback(ActorContainer& actorContainer): actorContainer(actorContainer) {}
3140 
3141 kj::Own<WorkerInterface> getWorker(IoChannelFactory::SubrequestMetadata metadata) override {
3142 return newPromisedWorkerInterface(actorContainer.startRequest(kj::mv(metadata)));
3143 }
3144 
3145 kj::Own<Worker::Actor::Loopback> addRef() override {
3146 return kj::addRef(*this);
3147 }
3148 
3149 private:
3150 ActorContainer& actorContainer;
3151 };
3152 
3153 class ActorSqliteHooks final: public ActorSqlite::Hooks {
3154 public:
3155 ActorSqliteHooks(AlarmScheduler& alarmScheduler, ActorKey actor)
3156 : alarmScheduler(alarmScheduler),
3157 actor(actor) {}
3158 
3159 // We ignore the priorTask in workerd because everything should run synchronously.
3160 kj::Promise<void> scheduleRun(
3161 kj::Maybe<kj::Date> newAlarmTime, kj::Promise<void> priorTask) override {
3162 KJ_IF_SOME(scheduledTime, newAlarmTime) {
3163 alarmScheduler.setAlarm(actor, scheduledTime);
3164 } else {
3165 alarmScheduler.deleteAlarm(actor);
3166 }
3167 return kj::READY_NOW;
3168 }
3169 
3170 private:
3171 AlarmScheduler& alarmScheduler;
3172 ActorKey actor;
3173 };
3174 };
3175 
3176 private:
3177 class EntrypointService final: public Service {
3178 public:
3179 EntrypointService(WorkerService& worker,
3180 kj::Maybe<kj::StringPtr> entrypoint,
3181 kj::Maybe<Frankenvalue> props,
3182 const kj::HashSet<kj::String>& handlers)
3183 : worker(kj::addRef(worker)),
3184 entrypoint(entrypoint),
3185 handlers(handlers),
3186 props(kj::mv(props)) {}
3187 
3188 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override {
3189 return startRequest(kj::mv(metadata), false);
3190 }
3191 
3192 kj::Own<WorkerInterface> startRequest(
3193 IoChannelFactory::SubrequestMetadata metadata, bool isTracer) {
3194 Frankenvalue props;
3195 KJ_IF_SOME(p, this->props) {
3196 props = p.clone();
3197 } else {
3198 // Calling ctx.exports loopback without specifying props. Use empty props.
3199 }
3200 return worker->startRequest(kj::mv(metadata), entrypoint, kj::mv(props), kj::none, isTracer);
3201 }
3202 
3203 bool hasHandler(kj::StringPtr handlerName) override {
3204 return handlers.contains(handlerName);
3205 }
3206 
3207 // Return underlying WorkerService.
3208 virtual Service* service() override {
3209 return worker;
3210 }
3211 
3212 kj::Own<Service> forProps(Frankenvalue props) override {
3213 if (this->props != kj::none) {
3214 // This entrypoint is already specialized. Delegate to the default implementation (which
3215 // will throw an exception).
3216 return Service::forProps(kj::mv(props));
3217 }
3218 
3219 return kj::refcounted<EntrypointService>(*worker, entrypoint, kj::mv(props), handlers);
3220 }
3221 
3222 void requireAllowsTransfer() override {
3223 worker->requireAllowsTransfer();
3224 }
3225 
3226 kj::Array<byte> getToken(ChannelTokenUsage usage) override {
3227 worker->requireAllowsTransfer();
3228 
3229 // If requireAllowsTransfer() passed, then we are not dynamic so should have a service name.
3230 // Unspecialized loopback entrypoints are not serializable, so if we get here we must have
3231 // props.
3232 return worker->channelTokenHandler.encodeSubrequestChannelToken(
3233 usage, KJ_ASSERT_NONNULL(worker->serviceName), entrypoint, KJ_ASSERT_NONNULL(props));
3234 }
3235 
3236 private:
3237 kj::Own<WorkerService> worker;
3238 kj::Maybe<kj::StringPtr> entrypoint;
3239 const kj::HashSet<kj::String>& handlers;
3240 kj::Maybe<Frankenvalue> props;
3241 };
3242 
3243 class ActorClassImpl final: public ActorClass {
3244 public:
3245 ActorClassImpl(WorkerService& service, kj::StringPtr className, kj::Maybe<Frankenvalue> props)
3246 : service(kj::addRef(service)),
3247 className(className),
3248 props(kj::mv(props)) {}
3249 
3250 void requireAllowsTransfer() override {
3251 service->requireAllowsTransfer();
3252 }
3253 
3254 kj::Own<Worker::Actor> newActor(kj::Maybe<RequestTracker&> tracker,
3255 Worker::Actor::Id actorId,
3256 Worker::Actor::MakeActorCacheFunc makeActorCache,
3257 Worker::Actor::MakeStorageFunc makeStorage,
3258 kj::Own<Worker::Actor::Loopback> loopback,
3259 kj::Maybe<kj::Own<Worker::Actor::HibernationManager>> manager,
3260 kj::Maybe<rpc::Container::Client> container,
3261 kj::Maybe<Worker::Actor::FacetManager&> facetManager) override {
3262 TimerChannel& timerChannel = *service;
3263 
3264 // We define this event ID in the internal codebase, but to have WebSocket Hibernation
3265 // work for local development we need to pass an event type.
3266 static constexpr uint16_t hibernationEventTypeId = 8;
3267 
3268 Frankenvalue props;
3269 KJ_IF_SOME(p, this->props) {
3270 props = p.clone();
3271 } else {
3272 // Using ctx.exports class loopback without specifying props. Use empty props.
3273 }
3274 
3275 return kj::refcounted<Worker::Actor>(*service->worker, tracker, kj::mv(actorId), true,
3276 kj::mv(makeActorCache), className, kj::mv(props), kj::mv(makeStorage), kj::mv(loopback),
3277 timerChannel, kj::refcounted<ActorObserver>(), kj::mv(manager), hibernationEventTypeId,
3278 kj::mv(container), facetManager);
3279 }
3280 
3281 kj::Own<WorkerInterface> startRequest(
3282 IoChannelFactory::SubrequestMetadata metadata, kj::Own<Worker::Actor> actor) override {
3283 // The `props` parameter is empty here because props are not passed per-request, they are
3284 // passed at Actor construction time.
3285 return service->startRequest(kj::mv(metadata), className, {}, kj::mv(actor));
3286 }
3287 
3288 kj::Own<ActorClass> forProps(Frankenvalue props) override {
3289 if (this->props != kj::none) {
3290 // This entrypoint is already specialized. Delegate to the default implementation (which
3291 // will throw an exception).
3292 return ActorClass::forProps(kj::mv(props));
3293 }
3294 
3295 return kj::refcounted<ActorClassImpl>(*service, className, kj::mv(props));
3296 }
3297 
3298 kj::Array<byte> getToken(ChannelTokenUsage usage) override {
3299 service->requireAllowsTransfer();
3300 
3301 // If requireAllowsTransfer() passed, then we are not dynamic so should have a service name.
3302 // Unspecialized loopback entrypoints are not serializable, so if we get here we must have
3303 // props.
3304 return service->channelTokenHandler.encodeActorClassChannelToken(
3305 usage, KJ_ASSERT_NONNULL(service->serviceName), className, KJ_ASSERT_NONNULL(props));
3306 }
3307 
3308 private:
3309 kj::Own<WorkerService> service;
3310 kj::StringPtr className;
3311 kj::Maybe<Frankenvalue> props;
3312 };
3313 
3314 ChannelTokenHandler& channelTokenHandler;
3315 
3316 // This service's name as defined in the original config, or null if it's a dynamic isolate.
3317 // Used only for serialization.
3318 kj::Maybe<kj::StringPtr> serviceName;
3319 
3320 ThreadContext& threadContext;
3321 const kj::MonotonicClock& monotonicClock;
3322 
3323 // LinkedIoChannels owns the SqliteDatabase::Vfs, so make sure it is destroyed last.
3324 kj::OneOf<LinkCallback, LinkedIoChannels> ioChannels;
3325 
3326 kj::Own<const Worker> worker;
3327 kj::Maybe<kj::HashSet<kj::String>> defaultEntrypointHandlers;
3328 kj::HashMap<kj::String, kj::HashSet<kj::String>> namedEntrypoints;
3329 kj::HashSet<kj::String> actorClassEntrypoints;
3330 kj::HashMap<kj::StringPtr, kj::Own<ActorNamespace>> actorNamespaces;
3331 kj::TaskSet waitUntilTasks;
3332 AbortActorsCallback abortActorsCallback;
3333 DeleteActorsCallback deleteActorsCallback;
3334 kj::Maybe<kj::String> dockerPath;
3335 kj::Maybe<kj::String> containerEgressInterceptorImage;
3336 bool isDynamic;
3337 kj::Maybe<kj::Function<void()>> abortIsolateCallback;
3338 
3339 class ActorChannelImpl final: public IoChannelFactory::ActorChannel {
3340 public:
3341 ActorChannelImpl(kj::Own<ActorNamespace::ActorContainer> actorContainer)
3342 : actorContainer(kj::mv(actorContainer)) {}
3343 ~ActorChannelImpl() noexcept(false) {
3344 actorContainer->updateAccessTime();
3345 }
3346 
3347 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override {
3348 return newPromisedWorkerInterface(actorContainer->startRequest(kj::mv(metadata)));
3349 }
3350 
3351 private:
3352 kj::Own<ActorNamespace::ActorContainer> actorContainer;
3353 };
3354 
3355 // ---------------------------------------------------------------------------
3356 // implements kj::TaskSet::ErrorHandler
3357 
3358 void taskFailed(kj::Exception&& exception) override {
3359 KJ_LOG(ERROR, exception);
3360 }
3361 
3362 // ---------------------------------------------------------------------------
3363 // implements IoChannelFactory
3364 
3365 kj::Own<WorkerInterface> startSubrequest(uint channel, SubrequestMetadata metadata) override {
3366 auto& channels =
3367 KJ_REQUIRE_NONNULL(ioChannels.tryGet<LinkedIoChannels>(), "link() has not been called");
3368 
3369 KJ_REQUIRE(channel < channels.subrequest.size(), "invalid subrequest channel number");
3370 return channels.subrequest[channel]->startRequest(kj::mv(metadata));
3371 }
3372 
3373 capnp::Capability::Client getCapability(uint channel) override {
3374 KJ_FAIL_REQUIRE("no capability channels");
3375 }
3376 class CacheClientImpl final: public CacheClient {
3377 public:
3378 CacheClientImpl(
3379 IoChannelFactory::SubrequestChannel& cacheService, kj::HttpHeaderId cacheNamespaceHeader)
3380 : cacheService(kj::addRef(cacheService)),
3381 cacheNamespaceHeader(cacheNamespaceHeader) {}
3382 
3383 kj::Own<kj::HttpClient> getDefault(CacheClient::SubrequestMetadata metadata) override {
3384 return kj::heap<CacheHttpClientImpl>(*cacheService, cacheNamespaceHeader, kj::none,
3385 kj::mv(metadata.cfBlobJson), kj::mv(metadata.parentSpan));
3386 }
3387 
3388 kj::Own<kj::HttpClient> getNamespace(
3389 kj::StringPtr cacheName, CacheClient::SubrequestMetadata metadata) override {
3390 auto encodedName = kj::encodeUriComponent(cacheName);
3391 return kj::heap<CacheHttpClientImpl>(*cacheService, cacheNamespaceHeader, kj::mv(encodedName),
3392 kj::mv(metadata.cfBlobJson), kj::mv(metadata.parentSpan));
3393 }
3394 
3395 private:
3396 kj::Own<IoChannelFactory::SubrequestChannel> cacheService;
3397 kj::HttpHeaderId cacheNamespaceHeader;
3398 };
3399 
3400 class CacheHttpClientImpl final: public kj::HttpClient {
3401 public:
3402 CacheHttpClientImpl(IoChannelFactory::SubrequestChannel& parent,
3403 kj::HttpHeaderId cacheNamespaceHeader,
3404 kj::Maybe<kj::String> cacheName,
3405 kj::Maybe<kj::String> cfBlobJson,
3406 SpanParent parentSpan)
3407 : client(asHttpClient(parent.startRequest({kj::mv(cfBlobJson), kj::mv(parentSpan)}))),
3408 cacheName(kj::mv(cacheName)),
3409 cacheNamespaceHeader(cacheNamespaceHeader) {}
3410 
3411 Request request(kj::HttpMethod method,
3412 kj::StringPtr url,
3413 const kj::HttpHeaders& headers,
3414 kj::Maybe<uint64_t> expectedBodySize = kj::none) override {
3415 
3416 return client->request(method, url, addCacheNameHeader(headers, cacheName), expectedBodySize);
3417 }
3418 
3419 private:
3420 kj::Own<kj::HttpClient> client;
3421 kj::Maybe<kj::String> cacheName;
3422 kj::HttpHeaderId cacheNamespaceHeader;
3423 
3424 kj::HttpHeaders addCacheNameHeader(
3425 const kj::HttpHeaders& headers, kj::Maybe<kj::StringPtr> cacheName) {
3426 auto headersCopy = headers.cloneShallow();
3427 KJ_IF_SOME(name, cacheName) {
3428 headersCopy.setPtr(cacheNamespaceHeader, name);
3429 }
3430 
3431 return headersCopy;
3432 }
3433 };
3434 
3435 kj::Own<CacheClient> getCache() override {
3436 auto& channels =
3437 KJ_REQUIRE_NONNULL(ioChannels.tryGet<LinkedIoChannels>(), "link() has not been called");
3438 auto& cache = *JSG_REQUIRE_NONNULL(channels.cache, Error, "No Cache was configured");
3439 return kj::heap<CacheClientImpl>(cache, threadContext.getHeaderIds().cfCacheNamespace);
3440 }
3441 
3442 TimerChannel& getTimer() override {
3443 return *this;
3444 }
3445 
3446 kj::Promise<void> writeLogfwdr(
3447 uint channel, kj::FunctionParam<void(capnp::AnyPointer::Builder)> buildMessage) override {
3448 auto& context = IoContext::current();
3449 
3450 auto headers = kj::HttpHeaders(context.getHeaderTable());
3451 auto client = context.getHttpClient(channel, true, kj::none, "writeLogfwdr"_kjc);
3452 
3453 auto urlStr = kj::str("https://fake-host");
3454 
3455 capnp::MallocMessageBuilder requestMessage;
3456 auto requestBuilder = requestMessage.initRoot<capnp::AnyPointer>();
3457 
3458 buildMessage(requestBuilder);
3459 capnp::JsonCodec json;
3460 auto requestJson = json.encode(requestBuilder.getAs<api::AnalyticsEngineEvent>());
3461 
3462 co_await context.waitForOutputLocks();
3463 
3464 auto innerReq = client->request(kj::HttpMethod::POST, urlStr, headers, requestJson.size());
3465 auto request = attachToRequest(kj::mv(innerReq), kj::refcountedWrapper(kj::mv(client)));
3466 
3467 co_await request.body->write(requestJson.asBytes())
3468 .attach(kj::mv(requestJson), kj::mv(request.body));
3469 auto response = co_await request.response;
3470 
3471 KJ_REQUIRE(response.statusCode >= 200 && response.statusCode < 300,
3472 "writeLogfwdr request returned an error");
3473 co_await response.body->readAllBytes().attach(kj::mv(response.body)).ignoreResult();
3474 co_return;
3475 }
3476 
3477 kj::Own<SubrequestChannel> getSubrequestChannel(uint channel,
3478 kj::Maybe<Frankenvalue> props,
3479 kj::Maybe<VersionRequest> versionRequest) override {
3480 auto& channels =
3481 KJ_REQUIRE_NONNULL(ioChannels.tryGet<LinkedIoChannels>(), "link() has not been called");
3482 
3483 KJ_REQUIRE(channel < channels.subrequest.size(), "invalid subrequest channel number");
3484 
3485 SubrequestChannel& channelRef = *channels.subrequest[channel];
3486 
3487 KJ_IF_SOME(p, props) {
3488 // Requesting specialization of loopback (ctx.exports) entrypoint with props.
3489 auto& service = KJ_REQUIRE_NONNULL(kj::dynamicDowncastIfAvailable<Service>(channelRef),
3490 "referenced channel is not a loopback channel");
3491 return service.forProps(kj::mv(p));
3492 }
3493 
3494 return kj::addRef(channelRef);
3495 }
3496 
3497 kj::Own<ActorChannel> getGlobalActor(uint channel,
3498 const ActorIdFactory::ActorId& id,
3499 kj::Maybe<kj::String> locationHint,
3500 ActorGetMode mode,
3501 bool enableReplicaRouting,
3502 ActorRoutingMode routingMode,
3503 SpanParent parentSpan,
3504 kj::Maybe<ActorVersion> version) override {
3505 JSG_REQUIRE(mode == ActorGetMode::GET_OR_CREATE, Error,
3506 "workerd only supports GET_OR_CREATE mode for getting actor stubs");
3507 JSG_REQUIRE(!enableReplicaRouting, Error, "workerd does not support replica routing.");
3508 
3509 // Compile-time assert that we have considered every routing mode here.
3510 switch (routingMode) {
3511 case ActorRoutingMode::PRIMARY_ONLY:
3512 // Workerd-only configs only supports primaries anyway.
3513 break;
3514 case ActorRoutingMode::DEFAULT:
3515 // In workerd-only configs, DEFAULT means PRIMARY_ONLY.
3516 break;
3517 }
3518 
3519 auto& channels =
3520 KJ_REQUIRE_NONNULL(ioChannels.tryGet<LinkedIoChannels>(), "link() has not been called");
3521 
3522 KJ_REQUIRE(channel < channels.actor.size(), "invalid actor channel number");
3523 auto& ns = JSG_REQUIRE_NONNULL(
3524 channels.actor[channel], Error, "Actor namespace configuration was invalid.");
3525 KJ_REQUIRE(ns.getConfig().is<Durable>()); // should have been verified earlier
3526 return ns.getActorChannel(id.clone());
3527 }
3528 
3529 kj::Own<ActorChannel> getColoLocalActor(
3530 uint channel, kj::StringPtr id, SpanParent parentSpan) override {
3531 auto& channels =
3532 KJ_REQUIRE_NONNULL(ioChannels.tryGet<LinkedIoChannels>(), "link() has not been called");
3533 
3534 KJ_REQUIRE(channel < channels.actor.size(), "invalid actor channel number");
3535 auto& ns = JSG_REQUIRE_NONNULL(
3536 channels.actor[channel], Error, "Actor namespace configuration was invalid.");
3537 KJ_REQUIRE(ns.getConfig().is<Ephemeral>()); // should have been verified earlier
3538 return ns.getActorChannel(kj::str(id));
3539 }
3540 
3541 kj::Own<ActorClassChannel> getActorClass(uint channel, kj::Maybe<Frankenvalue> props) override {
3542 auto& channels =
3543 KJ_REQUIRE_NONNULL(ioChannels.tryGet<LinkedIoChannels>(), "link() has not been called");
3544 
3545 KJ_REQUIRE(channel < channels.actorClass.size(), "invalid actor class channel number");
3546 
3547 ActorClass& cls = *channels.actorClass[channel];
3548 
3549 KJ_IF_SOME(p, props) {
3550 return cls.forProps(kj::mv(p));
3551 }
3552 
3553 return kj::addRef(cls);
3554 }
3555 
3556 void abortAllActors(kj::Maybe<kj::Exception&> reason) override {
3557 abortActorsCallback(reason);
3558 }
3559 
3560 void deleteAllActors(kj::Maybe<kj::Exception&> reason) override {
3561 deleteActorsCallback(reason);
3562 }
3563 
3564 // For now, in workerd just abort the process for non-dynamic workers.
3565 void abortIsolate(kj::StringPtr reason) noexcept override {
3566 KJ_IF_SOME(cb, abortIsolateCallback) {
3567 // Removes the isolate from the isolates map.
3568 //
3569 // TODO: Should abort all outstanding calls to the isolate causing them to
3570 // throw the reason as the error.
3571 cb();
3572 } else {
3573 // Otherwise, abort the process. Throwing from a noexcept function will call
3574 // std::terminate, which produces a nicer error message than ::abort().
3575 if (reason == nullptr) {
3576 KJ_FAIL_REQUIRE("abortIsolate() called, terminating process");
3577 } else {
3578 KJ_FAIL_REQUIRE("abortIsolate() called, terminating process", reason);
3579 }
3580 }
3581 }
3582 
3583 kj::Own<WorkerStubChannel> loadIsolate(uint loaderChannel,
3584 kj::Maybe<kj::String> name,
3585 kj::Function<kj::Promise<DynamicWorkerSource>()> fetchSource) override;
3586 
3587 kj::Network& getWorkerdDebugPortNetwork() override {
3588 auto& channels =
3589 KJ_REQUIRE_NONNULL(ioChannels.tryGet<LinkedIoChannels>(), "link() has not been called");
3590 return KJ_REQUIRE_NONNULL(channels.workerdDebugPortNetwork,
3591 "workerdDebugPort binding is not enabled for this worker");
3592 }
3593 
3594 kj::Own<SubrequestChannel> subrequestChannelFromToken(
3595 ChannelTokenUsage usage, kj::ArrayPtr<const byte> token) override {
3596 return channelTokenHandler.decodeSubrequestChannelToken(usage, token);
3597 }
3598 
3599 kj::Own<ActorClassChannel> actorClassFromToken(
3600 ChannelTokenUsage usage, kj::ArrayPtr<const byte> token) override {
3601 return channelTokenHandler.decodeActorClassChannelToken(usage, token);
3602 }
3603 
3604 // ---------------------------------------------------------------------------
3605 // implements TimerChannel
3606 
3607 void syncTime() override {
3608 // Nothing to do
3609 }
3610 
3611 kj::Date now(kj::Maybe<kj::Date>) override {
3612 return kj::systemPreciseCalendarClock().now();
3613 }
3614 
3615 kj::Promise<void> atTime(kj::Date when) override {
3616 auto delay = when - now(kj::none);
3617 // We can't use `afterDelay(delay)` here because kj::Timer::afterDelay() is equivalent to
3618 // `atTime(timer.now() + delay)`, and kj::Timer::now() only advances when the event loop
3619 // polls for I/O. If JavaScript executed for a significant amount of time since the last
3620 // poll (e.g. compiling/running a script before the first setTimeout), timer.now() will be
3621 // stale and the delay will effectively be shortened by that staleness, causing the timer
3622 // to fire too early. Instead, we compute the target time using a fresh reading from the
3623 // monotonic clock so the delay is measured from the actual present.
3624 return threadContext.getUnsafeTimer().atTime(monotonicClock.now() + delay);
3625 }
3626 
3627 kj::Promise<void> afterLimitTimeout(kj::Duration t) override {
3628 return threadContext.getUnsafeTimer().afterDelay(t);
3629 }
3630 
3631 // ---------------------------------------------------------------------------
3632 // implements LimitEnforcer
3633 //
3634 // No limits are enforced.
3635 
3636 kj::Own<void> enterJs(jsg::Lock& lock, IoContext& context) override {
3637 return {};
3638 }
3639 void topUpActor() override {}
3640 void newSubrequest(bool isInHouse) override {}
3641 void newKvRequest(KvOpType op) override {}
3642 void newAnalyticsEngineRequest() override {}
3643 kj::Promise<void> limitDrain() override {
3644 return kj::NEVER_DONE;
3645 }
3646 kj::Promise<void> limitScheduled() override {
3647 return kj::NEVER_DONE;
3648 }
3649 kj::Duration getAlarmLimit() override {
3650 return 15 * kj::MINUTES;
3651 }
3652 size_t getBufferingLimit() override {
3653 return kj::maxValue;
3654 }
3655 kj::Maybe<EventOutcome> getLimitsExceeded() override {
3656 return kj::none;
3657 }
3658 kj::Promise<void> onLimitsExceeded() override {
3659 return kj::NEVER_DONE;
3660 }
3661 void setCpuLimitNearlyExceededCallback(kj::Function<void(void)> cb) override {}
3662 void requireLimitsNotExceeded() override {}
3663 void reportMetrics(RequestObserver& requestMetrics) override {}
3664 kj::Duration consumeTimeElapsedForPeriodicLogging() override {
3665 return 0 * kj::SECONDS;
3666 }
3667 size_t getSqliteMemoryUsage() const override {
3668 return 0;
3669 }
3670};
3671 
3672struct FutureSubrequestChannel {
3673 kj::OneOf<config::ServiceDesignator::Reader, kj::Own<IoChannelFactory::SubrequestChannel>>
3674 designator;
3675 kj::String errorContext;
3676 
3677 kj::Own<IoChannelFactory::SubrequestChannel> lookup(Server& server) && {
3678 KJ_SWITCH_ONEOF(designator) {
3679 KJ_CASE_ONEOF(conf, config::ServiceDesignator::Reader) {
3680 return server.lookupService(conf, kj::mv(errorContext));
3681 }
3682 KJ_CASE_ONEOF(channel, kj::Own<IoChannelFactory::SubrequestChannel>) {
3683 return kj::mv(channel);
3684 }
3685 }
3686 KJ_UNREACHABLE;
3687 }
3688};
3689 
3690struct FutureActorChannel {
3691 config::Worker::Binding::DurableObjectNamespaceDesignator::Reader designator;
3692 kj::String errorContext;
3693};
3694 
3695struct FutureActorClassChannel {
3696 kj::OneOf<config::ServiceDesignator::Reader, kj::Own<Server::ActorClass>> designator;
3697 kj::String errorContext;
3698 
3699 kj::Own<Server::ActorClass> lookup(Server& server) && {
3700 KJ_SWITCH_ONEOF(designator) {
3701 KJ_CASE_ONEOF(conf, config::ServiceDesignator::Reader) {
3702 return server.lookupActorClass(conf, kj::mv(errorContext));
3703 }
3704 KJ_CASE_ONEOF(channel, kj::Own<Server::ActorClass>) {
3705 return kj::mv(channel);
3706 }
3707 }
3708 KJ_UNREACHABLE;
3709 }
3710};
3711 
3712struct FutureWorkerLoaderChannel {
3713 kj::String name; // for error logging, not necessarily unique
3714 kj::Maybe<kj::String> id;
3715};
3716 
3717static kj::Maybe<WorkerdApi::Global> createBinding(kj::StringPtr workerName,
3718 config::Worker::Reader conf,
3719 config::Worker::Binding::Reader binding,
3720 Worker::ValidationErrorReporter& errorReporter,
3721 kj::Vector<FutureSubrequestChannel>& subrequestChannels,
3722 kj::Vector<FutureActorChannel>& actorChannels,
3723 kj::Vector<FutureActorClassChannel>& actorClassChannels,
3724 kj::Vector<FutureWorkerLoaderChannel>& workerLoaderChannels,
3725 bool& hasWorkerdDebugPortBinding,
3726 kj::HashMap<kj::String, kj::HashMap<kj::String, Server::ActorConfig>>& actorConfigs,
3727 bool experimental) {
3728 // creates binding object or returns null and reports an error
3729 using Global = WorkerdApi::Global;
3730 kj::StringPtr bindingName = binding.getName();
3731 TRACE_EVENT("workerd", "Server::WorkerService::createBinding()", "name", workerName.cStr(),
3732 "binding", bindingName.cStr());
3733 auto makeGlobal = [&](auto&& value) {
3734 return Global{.name = kj::str(bindingName), .value = kj::mv(value)};
3735 };
3736 
3737 auto errorContext = kj::str("Worker \"", workerName, "\"'s binding \"", bindingName, "\"");
3738 
3739 switch (binding.which()) {
3740 case config::Worker::Binding::UNSPECIFIED:
3741 errorReporter.addError(kj::str(errorContext, " does not specify any binding value."));
3742 return kj::none;
3743 
3744 case config::Worker::Binding::PARAMETER:
3745 KJ_UNIMPLEMENTED("TODO(beta): parameters");
3746 
3747 case config::Worker::Binding::TEXT:
3748 return makeGlobal(kj::str(binding.getText()));
3749 case config::Worker::Binding::DATA:
3750 return makeGlobal(kj::heapArray<byte>(binding.getData()));
3751 case config::Worker::Binding::JSON:
3752 return makeGlobal(Global::Json{kj::str(binding.getJson())});
3753 
3754 case config::Worker::Binding::WASM_MODULE:
3755 if (conf.isServiceWorkerScript()) {
3756 // Already handled earlier.
3757 } else {
3758 errorReporter.addError(kj::str(errorContext,
3759 " is a Wasm binding, but Wasm bindings are not allowed in "
3760 "modules-based scripts. Use Wasm modules instead."));
3761 }
3762 return kj::none;
3763 
3764 case config::Worker::Binding::CRYPTO_KEY: {
3765 auto keyConf = binding.getCryptoKey();
3766 Global::CryptoKey keyGlobal;
3767 
3768 switch (keyConf.which()) {
3769 case config::Worker::Binding::CryptoKey::RAW:
3770 keyGlobal.format = kj::str("raw");
3771 keyGlobal.keyData = kj::heapArray<kj::byte>(keyConf.getRaw());
3772 goto validFormat;
3773 case config::Worker::Binding::CryptoKey::HEX: {
3774 keyGlobal.format = kj::str("raw");
3775 auto decoded = kj::decodeHex(keyConf.getHex());
3776 if (decoded.hadErrors) {
3777 errorReporter.addError(
3778 kj::str("CryptoKey binding \"", binding.getName(), "\" contained invalid hex."));
3779 }
3780 keyGlobal.keyData = kj::Array<byte>(kj::mv(decoded));
3781 goto validFormat;
3782 }
3783 case config::Worker::Binding::CryptoKey::BASE64: {
3784 keyGlobal.format = kj::str("raw");
3785 auto decoded = kj::decodeBase64(keyConf.getBase64());
3786 if (decoded.hadErrors) {
3787 errorReporter.addError(
3788 kj::str("CryptoKey binding \"", binding.getName(), "\" contained invalid base64."));
3789 }
3790 keyGlobal.keyData = kj::Array<byte>(kj::mv(decoded));
3791 goto validFormat;
3792 }
3793 case config::Worker::Binding::CryptoKey::PKCS8: {
3794 keyGlobal.format = kj::str("pkcs8");
3795 auto pem = KJ_UNWRAP_OR(decodePem(keyConf.getPkcs8()), {
3796 errorReporter.addError(kj::str(
3797 "CryptoKey binding \"", binding.getName(), "\" contained invalid PEM format."));
3798 return kj::none;
3799 });
3800 if (pem.type != "PRIVATE KEY") {
3801 errorReporter.addError(kj::str("CryptoKey binding \"", binding.getName(),
3802 "\" contained wrong PEM type, "
3803 "expected \"PRIVATE KEY\" but got \"",
3804 pem.type, "\"."));
3805 return kj::none;
3806 }
3807 keyGlobal.keyData = kj::mv(pem.data);
3808 goto validFormat;
3809 }
3810 case config::Worker::Binding::CryptoKey::SPKI: {
3811 keyGlobal.format = kj::str("spki");
3812 auto pem = KJ_UNWRAP_OR(decodePem(keyConf.getSpki()), {
3813 errorReporter.addError(kj::str(
3814 "CryptoKey binding \"", binding.getName(), "\" contained invalid PEM format."));
3815 return kj::none;
3816 });
3817 if (pem.type != "PUBLIC KEY") {
3818 errorReporter.addError(kj::str("CryptoKey binding \"", binding.getName(),
3819 "\" contained wrong PEM type, "
3820 "expected \"PUBLIC KEY\" but got \"",
3821 pem.type, "\"."));
3822 return kj::none;
3823 }
3824 keyGlobal.keyData = kj::mv(pem.data);
3825 goto validFormat;
3826 }
3827 case config::Worker::Binding::CryptoKey::JWK:
3828 keyGlobal.format = kj::str("jwk");
3829 keyGlobal.keyData = Global::Json{kj::str(keyConf.getJwk())};
3830 goto validFormat;
3831 }
3832 errorReporter.addError(kj::str("Encountered unknown CryptoKey type for binding \"",
3833 binding.getName(), "\". Was the config compiled with a newer version of the schema?"));
3834 return kj::none;
3835 validFormat:
3836 
3837 auto algorithmConf = keyConf.getAlgorithm();
3838 switch (algorithmConf.which()) {
3839 case config::Worker::Binding::CryptoKey::Algorithm::NAME:
3840 keyGlobal.algorithm = Global::Json{escapeJsonString(algorithmConf.getName())};
3841 goto validAlgorithm;
3842 case config::Worker::Binding::CryptoKey::Algorithm::JSON:
3843 keyGlobal.algorithm = Global::Json{kj::str(algorithmConf.getJson())};
3844 goto validAlgorithm;
3845 }
3846 errorReporter.addError(kj::str("Encountered unknown CryptoKey algorithm type for binding \"",
3847 binding.getName(), "\". Was the config compiled with a newer version of the schema?"));
3848 return kj::none;
3849 validAlgorithm:
3850 
3851 keyGlobal.extractable = keyConf.getExtractable();
3852 keyGlobal.usages = KJ_MAP(usage, keyConf.getUsages()) { return kj::str(usage); };
3853 
3854 return makeGlobal(kj::mv(keyGlobal));
3855 return kj::none;
3856 }
3857 
3858 case config::Worker::Binding::SERVICE: {
3859 uint channel = static_cast<uint>(subrequestChannels.size()) +
3860 IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT;
3861 subrequestChannels.add(FutureSubrequestChannel{binding.getService(), kj::mv(errorContext)});
3862 return makeGlobal(
3863 Global::Fetcher{.channel = channel, .requiresHost = true, .isInHouse = false});
3864 }
3865 
3866 case config::Worker::Binding::DURABLE_OBJECT_NAMESPACE: {
3867 auto actorBinding = binding.getDurableObjectNamespace();
3868 const Server::ActorConfig* actorConfig;
3869 if (actorBinding.hasServiceName()) {
3870 auto& svcMap = KJ_UNWRAP_OR(actorConfigs.find(actorBinding.getServiceName()), {
3871 errorReporter.addError(kj::str(errorContext, " refers to a service \"",
3872 actorBinding.getServiceName(), "\", but no such service is defined."));
3873 return kj::none;
3874 });
3875 
3876 actorConfig = &KJ_UNWRAP_OR(svcMap.find(actorBinding.getClassName()), {
3877 errorReporter.addError(
3878 kj::str(errorContext, " refers to a Durable Object namespace named \"",
3879 actorBinding.getClassName(), "\" in service \"", actorBinding.getServiceName(),
3880 "\", but no such Durable Object namespace is defined by that service."));
3881 return kj::none;
3882 });
3883 } else {
3884 auto& localActorConfigs = KJ_ASSERT_NONNULL(actorConfigs.find(workerName));
3885 actorConfig = &KJ_UNWRAP_OR(localActorConfigs.find(actorBinding.getClassName()), {
3886 errorReporter.addError(kj::str(errorContext,
3887 " refers to a Durable Object namespace named \"", actorBinding.getClassName(),
3888 "\", but no such Durable Object namespace is defined "
3889 "by this Worker."));
3890 return kj::none;
3891 });
3892 }
3893 
3894 uint channel = static_cast<uint>(actorChannels.size());
3895 actorChannels.add(FutureActorChannel{actorBinding, kj::mv(errorContext)});
3896 
3897 KJ_SWITCH_ONEOF(*actorConfig) {
3898 KJ_CASE_ONEOF(durable, Server::Durable) {
3899 return makeGlobal(Global::DurableActorNamespace{
3900 .actorChannel = channel, .uniqueKey = durable.uniqueKey});
3901 }
3902 KJ_CASE_ONEOF(_, Server::Ephemeral) {
3903 return makeGlobal(Global::EphemeralActorNamespace{.actorChannel = channel});
3904 }
3905 }
3906 
3907 return kj::none;
3908 }
3909 
3910 case config::Worker::Binding::KV_NAMESPACE: {
3911 uint channel = static_cast<uint>(subrequestChannels.size()) +
3912 IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT;
3913 subrequestChannels.add(
3914 FutureSubrequestChannel{binding.getKvNamespace(), kj::mv(errorContext)});
3915 
3916 return makeGlobal(Global::KvNamespace{
3917 .subrequestChannel = channel, .bindingName = kj::str(binding.getName())});
3918 }
3919 
3920 case config::Worker::Binding::R2_BUCKET: {
3921 uint channel = static_cast<uint>(subrequestChannels.size()) +
3922 IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT;
3923 subrequestChannels.add(FutureSubrequestChannel{binding.getR2Bucket(), kj::mv(errorContext)});
3924 return makeGlobal(Global::R2Bucket{.subrequestChannel = channel,
3925 .bucket = kj::str(binding.getR2Bucket().getName()),
3926 .bindingName = kj::str(binding.getName())});
3927 }
3928 
3929 case config::Worker::Binding::OBSOLETE0:
3930 errorReporter.addError(kj::str(errorContext, " uses an obsolete binding type."));
3931 return kj::none;
3932 
3933 case config::Worker::Binding::QUEUE: {
3934 uint channel = static_cast<uint>(subrequestChannels.size()) +
3935 IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT;
3936 subrequestChannels.add(FutureSubrequestChannel{binding.getQueue(), kj::mv(errorContext)});
3937 
3938 return makeGlobal(Global::QueueBinding{.subrequestChannel = channel});
3939 }
3940 
3941 case config::Worker::Binding::WRAPPED: {
3942 auto wrapped = binding.getWrapped();
3943 kj::Vector<Global> innerGlobals;
3944 for (const auto& innerBinding: wrapped.getInnerBindings()) {
3945 KJ_IF_SOME(global,
3946 createBinding(workerName, conf, innerBinding, errorReporter, subrequestChannels,
3947 actorChannels, actorClassChannels, workerLoaderChannels, hasWorkerdDebugPortBinding,
3948 actorConfigs, experimental)) {
3949 innerGlobals.add(kj::mv(global));
3950 } else {
3951 // we've already communicated the error
3952 return kj::none;
3953 }
3954 }
3955 return makeGlobal(Global::Wrapped{
3956 .moduleName = kj::str(wrapped.getModuleName()),
3957 .entrypoint = kj::str(wrapped.getEntrypoint()),
3958 .innerBindings = innerGlobals.releaseAsArray(),
3959 });
3960 }
3961 
3962 case config::Worker::Binding::FROM_ENVIRONMENT: {
3963 const char* value = getenv(binding.getFromEnvironment().cStr());
3964 if (value == nullptr) {
3965 // TODO(cleanup): Maybe make a Global::Null? (Can't use nullptr_t in OneOf.) For now,
3966 // using JSON gets the job done hackily.
3967 return makeGlobal(Global::Json{kj::str("null")});
3968 } else {
3969 return makeGlobal(kj::str(value));
3970 }
3971 }
3972 
3973 case config::Worker::Binding::ANALYTICS_ENGINE: {
3974 if (!experimental) {
3975 errorReporter.addError(kj::str(
3976 "AnalyticsEngine bindings are an experimental feature which may change or go away in the future."
3977 "You must run workerd with `--experimental` to use this feature."));
3978 }
3979 
3980 uint channel = static_cast<uint>(subrequestChannels.size()) +
3981 IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT;
3982 subrequestChannels.add(
3983 FutureSubrequestChannel{binding.getAnalyticsEngine(), kj::mv(errorContext)});
3984 
3985 return makeGlobal(Global::AnalyticsEngine{
3986 .subrequestChannel = channel,
3987 .dataset = kj::str(binding.getAnalyticsEngine().getName()),
3988 .version = 0,
3989 });
3990 }
3991 case config::Worker::Binding::HYPERDRIVE: {
3992 uint channel = static_cast<uint>(subrequestChannels.size()) +
3993 IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT;
3994 subrequestChannels.add(
3995 FutureSubrequestChannel{binding.getHyperdrive().getDesignator(), kj::mv(errorContext)});
3996 return makeGlobal(Global::Hyperdrive{
3997 .subrequestChannel = channel,
3998 .database = kj::str(binding.getHyperdrive().getDatabase()),
3999 .user = kj::str(binding.getHyperdrive().getUser()),
4000 .password = kj::str(binding.getHyperdrive().getPassword()),
4001 .scheme = kj::str(binding.getHyperdrive().getScheme()),
4002 });
4003 }
4004 case config::Worker::Binding::UNSAFE_EVAL: {
4005 if (!experimental) {
4006 errorReporter.addError(kj::str("Unsafe eval is an experimental feature. ",
4007 "You must run workerd with `--experimental` to use this feature."));
4008 return kj::none;
4009 }
4010 return makeGlobal(Global::UnsafeEval{});
4011 }
4012 case config::Worker::Binding::MEMORY_CACHE: {
4013 if (!experimental) {
4014 errorReporter.addError(kj::str(
4015 "MemoryCache bindings are an experimental feature which may change or go away "
4016 "in the future. You must run workerd with `--experimental` to use this feature."));
4017 return kj::none;
4018 }
4019 auto cache = binding.getMemoryCache();
4020 // TODO(cleanup): Should we have some reasonable default for these so they can
4021 // be optional?
4022 if (!cache.hasLimits()) {
4023 errorReporter.addError(
4024 kj::str("MemoryCache bindings must specify limits. Please "
4025 "update the binding in the worker configuration and try again."));
4026 return kj::none;
4027 }
4028 Global::MemoryCache cacheCopy;
4029 // The id is optional. If provided, then multiple bindings with the same id will
4030 // share the same cache. Otherwise, a unique id is generated for the cache.
4031 if (cache.hasId()) {
4032 cacheCopy.cacheId = kj::str(cache.getId());
4033 }
4034 auto limits = cache.getLimits();
4035 cacheCopy.maxKeys = limits.getMaxKeys();
4036 cacheCopy.maxValueSize = limits.getMaxValueSize();
4037 cacheCopy.maxTotalValueSize = limits.getMaxTotalValueSize();
4038 return makeGlobal(kj::mv(cacheCopy));
4039 }
4040 
4041 case config::Worker::Binding::DURABLE_OBJECT_CLASS: {
4042 if (!experimental) {
4043 errorReporter.addError(kj::str(
4044 "Durable Object class bindings are an experimental feature which may change or go away "
4045 "in the future. You must run workerd with `--experimental` to use this feature."));
4046 return kj::none;
4047 }
4048 uint channel = actorClassChannels.size();
4049 actorClassChannels.add(
4050 FutureActorClassChannel{binding.getDurableObjectClass(), kj::mv(errorContext)});
4051 return makeGlobal(Global::ActorClass{.channel = channel});
4052 }
4053 
4054 case config::Worker::Binding::WORKER_LOADER: {
4055 if (!experimental) {
4056 errorReporter.addError(kj::str(
4057 "Worker loader bindings are an experimental feature which may change or go away "
4058 "in the future. You must run workerd with `--experimental` to use this feature."));
4059 return kj::none;
4060 }
4061 
4062 auto loaderConf = binding.getWorkerLoader();
4063 
4064 FutureWorkerLoaderChannel channel;
4065 if (loaderConf.hasId()) {
4066 channel.name = kj::str(loaderConf.getId());
4067 channel.id = kj::str(channel.name);
4068 } else {
4069 channel.name = kj::str(bindingName);
4070 }
4071 
4072 uint channelNumber = workerLoaderChannels.size();
4073 workerLoaderChannels.add(kj::mv(channel));
4074 return makeGlobal(Global::WorkerLoader{.channel = channelNumber});
4075 }
4076 
4077 case config::Worker::Binding::WORKERD_DEBUG_PORT: {
4078 if (!experimental) {
4079 errorReporter.addError(kj::str(
4080 "workerdDebugPort bindings are an experimental feature which may change or go away "
4081 "in the future. You must run workerd with `--experimental` to use this feature."));
4082 return kj::none;
4083 }
4084 
4085 hasWorkerdDebugPortBinding = true;
4086 return makeGlobal(Global::WorkerdDebugPort{});
4087 }
4088 }
4089 errorReporter.addError(kj::str(errorContext,
4090 "has unrecognized type. Was the config compiled with a newer version of "
4091 "the schema?"));
4092}
4093 
4094uint startInspector(
4095 kj::StringPtr inspectorAddress, Server::InspectorServiceIsolateRegistrar& registrar);
4096 
4097void Server::abortAllActors(kj::Maybe<const kj::Exception&> reason) {
4098 for (auto& service: services) {
4099 if (WorkerService* worker = dynamic_cast<WorkerService*>(&*service.value)) {
4100 for (auto& [className, ns]: worker->getActorNamespaces()) {
4101 bool isEvictable = true;
4102 KJ_SWITCH_ONEOF(ns->getConfig()) {
4103 KJ_CASE_ONEOF(c, Durable) {
4104 isEvictable = c.isEvictable;
4105 }
4106 KJ_CASE_ONEOF(c, Ephemeral) {
4107 isEvictable = c.isEvictable;
4108 }
4109 }
4110 if (isEvictable) ns->abortAll(reason);
4111 }
4112 }
4113 }
4114}
4115 
4116void Server::deleteAllActors(kj::Maybe<const kj::Exception&> reason) {
4117 for (auto& service: services) {
4118 if (WorkerService* worker = dynamic_cast<WorkerService*>(&*service.value)) {
4119 for (auto& [className, ns]: worker->getActorNamespaces()) {
4120 bool isEvictable = true;
4121 KJ_SWITCH_ONEOF(ns->getConfig()) {
4122 KJ_CASE_ONEOF(c, Durable) {
4123 isEvictable = c.isEvictable;
4124 }
4125 KJ_CASE_ONEOF(c, Ephemeral) {
4126 isEvictable = c.isEvictable;
4127 }
4128 }
4129 if (isEvictable) ns->deleteAll(reason);
4130 }
4131 }
4132 }
4133}
4134 
4135// WorkerDef is an intermediate representation of everything from `config::Worker::Reader` that
4136// `Server::makeWorkerImpl()` needs. Similar to `WorkerSource`, we factor out this intermediate
4137// representation so that we can potentially build it dynamically from input that isn't a
4138// workerd config file.
4139struct Server::WorkerDef {
4140 CompatibilityFlags::Reader featureFlags;
4141 WorkerSource source;
4142 kj::Maybe<kj::StringPtr> moduleFallback;
4143 const kj::HashMap<kj::String, ActorConfig>& localActorConfigs;
4144 bool isDynamic;
4145 
4146 FutureSubrequestChannel globalOutbound;
4147 kj::Maybe<FutureSubrequestChannel> cacheApiOutbound;
4148 kj::Vector<FutureSubrequestChannel> subrequestChannels;
4149 kj::Vector<FutureActorChannel> actorChannels;
4150 kj::Vector<FutureActorClassChannel> actorClassChannels;
4151 kj::Vector<FutureWorkerLoaderChannel> workerLoaderChannels;
4152 bool hasWorkerdDebugPortBinding = false;
4153 kj::Array<FutureSubrequestChannel> tails;
4154 kj::Array<FutureSubrequestChannel> streamingTails;
4155 
4156 // Dynamically-loaded isolates can't directly have storage, so for now I'm using a raw capnp
4157 // Reader here. A default-constructed Reader will have type `none` which is appropriate for
4158 // dynamically-loaded workers. Same story for ContainerEngine.
4159 config::Worker::DurableObjectStorage::Reader actorStorageConf;
4160 config::Worker::ContainerEngine::Reader containerEngineConf;
4161 
4162 // Similar to the `compileBindings` callback passed into `Worker`'s constructor, except that
4163 // `ctx.exports` is taken care of separately. This is provided as a callback since `env` is
4164 // constructed in a vastly different way for dynamically-loaded workers.
4165 kj::Function<void(jsg::Lock& lock, const Worker::Api& api, v8::Local<v8::Object> target)>
4166 compileBindings;
4167 
4168 // If the WorkerDef was created from a DymamicWorkerSource and that
4169 // source contains a clone of the source bundle, this will take ownership.
4170 kj::Maybe<kj::Own<void>> maybeOwnedSourceCode;
4171 
4172 // Callback invoked when abortIsolate() is called. Used by dynamic workers to remove
4173 // themselves from the loader's isolate map.
4174 kj::Maybe<kj::Function<void()>> abortIsolateCallback;
4175};
4176 
4177class Server::WorkerLoaderNamespace: public kj::Refcounted {
4178 public:
4179 WorkerLoaderNamespace(Server& server, kj::String namespaceName)
4180 : server(server),
4181 namespaceName(kj::mv(namespaceName)) {}
4182 
4183 void unlink() {
4184 for (auto& isolate: isolates) {
4185 isolate.value->unlink();
4186 }
4187 }
4188 
4189 kj::Own<WorkerStubChannel> loadIsolate(
4190 kj::Maybe<kj::String> name, kj::Function<kj::Promise<DynamicWorkerSource>()> fetchSource) {
4191 KJ_IF_SOME(n, name) {
4192 return isolates
4193 .findOrCreate(n,
4194 [&]() -> decltype(isolates)::Entry {
4195 // This name isn't actually used in any maps nor is it ever revealed back to the app, but it
4196 // may be used in error logs.
4197 auto isolateName = kj::str(namespaceName, ':', n);
4198 
4199 // On abort, remove the entry from this namespace's isolates map so
4200 // subsequent loadIsolate() calls with the same name will create a fresh
4201 // isolate.
4202 kj::Function<void()> onAborted = [this, mapKey = kj::str(n)]() { removeIsolate(mapKey); };
4203 
4204 return {.key = kj::mv(n),
4205 .value = kj::rc<WorkerStubImpl>(
4206 server, kj::mv(isolateName), kj::mv(onAborted), kj::mv(fetchSource))};
4207 })
4208 .addRef()
4209 .toOwn();
4210 } else {
4211 auto isolateName = kj::str(namespaceName, ":dynamic:", randomUUID(server.entropySource));
4212 return kj::rc<WorkerStubImpl>(server, kj::mv(isolateName), kj::none, kj::mv(fetchSource))
4213 .toOwn();
4214 }
4215 }
4216 
4217 void removeIsolate(kj::StringPtr name) {
4218 // This is called by abortIsolate()
4219 isolates.erase(name);
4220 }
4221 
4222 private:
4223 Server& server;
4224 kj::String namespaceName;
4225 
4226 class WorkerStubImpl;
4227 kj::HashMap<kj::String, kj::Rc<WorkerStubImpl>> isolates;
4228 
4229 class NullGlobalOutboundChannel: public IoChannelFactory::SubrequestChannel {
4230 public:
4231 kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override {
4232 JSG_FAIL_REQUIRE(Error,
4233 "This worker is not permitted to access the internet via global functions like fetch(). "
4234 "It must use capabilities (such as bindings in 'env') to talk to the outside world.");
4235 }
4236 
4237 void requireAllowsTransfer() override {
4238 // It's difficult to get here, because the null outbound is not normally something you can
4239 // reference. That said, it is possible to get a `Fetcher` representing the `next` outbound
4240 // by pulling it off an incoming `Request` object, and in practice that points to the same
4241 // thing as the null outbound. You could then try to transfer it.
4242 //
4243 // We disallow this for now because it's not clear why it would be needed. That said, if it
4244 // is needed for some reason, it wouldn't be hard to support. But we might want to change
4245 // the error message it throws from startRequest(), since the error would be somewhat
4246 // misleading after the channel has been transferred.
4247 JSG_FAIL_REQUIRE(DOMDataCloneError, "The null global outbound is not transferrable.");
4248 }
4249 };
4250 
4251 class WorkerStubImpl final: public WorkerStubChannel, public kj::Refcounted {
4252 public:
4253 WorkerStubImpl(Server& server,
4254 kj::String isolateName,
4255 kj::Maybe<kj::Function<void()>> onAborted,
4256 kj::Function<kj::Promise<DynamicWorkerSource>()> fetchSource)
4257 : onAborted(kj::mv(onAborted)),
4258 startupTask(start(server, kj::mv(isolateName), kj::mv(fetchSource)).fork()) {}
4259 
4260 ~WorkerStubImpl() {
4261 unlink();
4262 }
4263 
4264 void unlink() {
4265 KJ_IF_SOME(s, service) {
4266 s->unlink();
4267 }
4268 }
4269 
4270 kj::Own<IoChannelFactory::SubrequestChannel> getEntrypoint(
4271 kj::Maybe<kj::String> name, Frankenvalue props, kj::Maybe<ResourceLimits> limits) override {
4272 return kj::refcounted<SubrequestChannelImpl>(addRefToThis(), kj::mv(name), kj::mv(props));
4273 }
4274 
4275 kj::Own<IoChannelFactory::ActorClassChannel> getActorClass(
4276 kj::Maybe<kj::String> name, Frankenvalue props, kj::Maybe<ResourceLimits> limits) override {
4277 return kj::refcounted<ActorClassImpl>(addRefToThis(), kj::mv(name), kj::mv(props));
4278 }
4279 
4280 private:
4281 // Callback to remove the worker stub from the isolates map. None for
4282 // unnamed dynamic isolates.
4283 kj::Maybe<kj::Function<void()>> onAborted;
4284 
4285 kj::Maybe<kj::Own<WorkerService>> service; // null if still starting up
4286 kj::ForkedPromise<void> startupTask; // resolves when `service` is non-null
4287 
4288 void onAbortIsolate() {
4289 KJ_IF_SOME(cb, onAborted) {
4290 auto callback = kj::mv(cb);
4291 onAborted = kj::none;
4292 callback();
4293 }
4294 }
4295 
4296 kj::Promise<void> start(Server& server,
4297 kj::String isolateName,
4298 kj::Function<kj::Promise<DynamicWorkerSource>()> fetchSource) {
4299 auto source = co_await fetchSource();
4300 static const kj::HashMap<kj::String, ActorConfig> EMPTY_ACTOR_CONFIGS;
4301 
4302 // Rewrite the capabilities in `env` in order to build the I/O channel table.
4303 kj::Vector<FutureSubrequestChannel> subrequestChannels;
4304 kj::Vector<FutureActorClassChannel> actorClassChannels;
4305 source.env.rewriteCaps([&](kj::Own<Frankenvalue::CapTableEntry> entry) {
4306 if (auto channel = dynamic_cast<IoChannelFactory::SubrequestChannel*>(entry.get())) {
4307 uint channelNumber =
4308 subrequestChannels.size() + IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT;
4309 subrequestChannels.add(FutureSubrequestChannel{
4310 .designator = kj::addRef(*channel),
4311 .errorContext = kj::str("Worker's env"),
4312 });
4313 return kj::heap<IoChannelCapTableEntry>(
4314 IoChannelCapTableEntry::SUBREQUEST, channelNumber);
4315 } else if (auto channel = dynamic_cast<ActorClass*>(entry.get())) {
4316 uint channelNumber = actorClassChannels.size();
4317 actorClassChannels.add(FutureActorClassChannel{
4318 .designator = kj::addRef(*channel),
4319 .errorContext = kj::str("Worker's env"),
4320 });
4321 return kj::heap<IoChannelCapTableEntry>(
4322 IoChannelCapTableEntry::ACTOR_CLASS, channelNumber);
4323 } else {
4324 // Generally, it shouldn't be possible to get here, but just in case, let's at least
4325 // provide some sort of error, although it's a vague one.
4326 JSG_FAIL_REQUIRE(DOMDataCloneError,
4327 "Dynamic 'env' contains one or more objects that are not supported for use in "
4328 "'env', although they would be supported in 'props'.");
4329 }
4330 });
4331 
4332 WorkerDef def{
4333 .featureFlags = source.compatibilityFlags,
4334 .source = kj::mv(source.source),
4335 .moduleFallback = kj::none,
4336 .localActorConfigs = EMPTY_ACTOR_CONFIGS,
4337 .isDynamic = true,
4338 
4339 // clang-format off
4340 .globalOutbound{
4341 .designator = kj::mv(source.globalOutbound)
4342 .orDefault([]() { return kj::refcounted<NullGlobalOutboundChannel>(); }),
4343 .errorContext = kj::str("Worker's globalOutbound"),
4344 },
4345 
4346 .subrequestChannels = kj::mv(subrequestChannels),
4347 .actorClassChannels = kj::mv(actorClassChannels),
4348 
4349 .tails = KJ_MAP(tail, source.tails) -> FutureSubrequestChannel {
4350 return {
4351 .designator = kj::mv(tail),
4352 .errorContext = kj::str("Worker's tail"),
4353 };
4354 },
4355 .streamingTails = KJ_MAP(tail, source.streamingTails) -> FutureSubrequestChannel {
4356 return {
4357 .designator = kj::mv(tail),
4358 .errorContext = kj::str("Worker's streaming tail"),
4359 };
4360 },
4361 
4362 .compileBindings = [env = kj::mv(source.env)](
4363 jsg::Lock& js, const Worker::Api& api, v8::Local<v8::Object> target) mutable {
4364 env.populateJsObject(js, jsg::JsObject(target));
4365 },
4366 
4367 // Note here that we always keep the ownContent from the source, even if
4368 // ownContentIsRpcResponse is true. This is safe in workerd because we
4369 // are single-threaded here and we don't need to worry about the cross-thread
4370 // ownership issues. For the downstream use, however, we need to be careful
4371 // to not copy the ownContent if it is an RPC response.
4372 .maybeOwnedSourceCode = kj::mv(source.ownContent),
4373 // The callback is owned by the WorkerService, which is owned by `this`, so a raw
4374 // pointer is safe.
4375 .abortIsolateCallback = kj::Function<void()>([this]() { onAbortIsolate(); }),
4376 // clang-format on
4377 };
4378 
4379 DynamicErrorReporter errorReporter;
4380 
4381 auto service = co_await server.makeWorkerImpl(isolateName, kj::mv(def), {}, errorReporter);
4382 errorReporter.throwIfErrors();
4383 
4384 service->link(errorReporter);
4385 errorReporter.throwIfErrors();
4386 
4387 this->service = kj::mv(service);
4388 }
4389 
4390 class SubrequestChannelImpl final: public IoChannelFactory::SubrequestChannel {
4391 public:
4392 SubrequestChannelImpl(
4393 kj::Rc<WorkerStubImpl> isolate, kj::Maybe<kj::String> entrypointName, Frankenvalue props)
4394 : isolate(kj::mv(isolate)),
4395 entrypointName(kj::mv(entrypointName)),
4396 props(kj::mv(props)) {}
4397 
4398 kj::Own<WorkerInterface> startRequest(
4399 IoChannelFactory::SubrequestMetadata metadata) override {
4400 if (isolate->service == kj::none) {
4401 // Capture a refcounted reference rather than a raw `this` pointer so that the
4402 // SubrequestChannelImpl is kept alive until the startup task resolves, even if the
4403 // owning Fetcher is garbage-collected while the deferred promise is pending.
4404 return newPromisedWorkerInterface(isolate->startupTask.addBranch().then(
4405 [self = kj::addRef(*this), metadata = kj::mv(metadata)]() mutable {
4406 return self->startRequestImpl(kj::mv(metadata));
4407 }));
4408 } else {
4409 return startRequestImpl(kj::mv(metadata));
4410 }
4411 }
4412 
4413 void requireAllowsTransfer() override {
4414 throwDynamicEntrypointTransferError();
4415 }
4416 
4417 private:
4418 kj::Rc<WorkerStubImpl> isolate;
4419 kj::Maybe<kj::String> entrypointName;
4420 Frankenvalue props; // moved away when `entrypointService` is initialized
4421 
4422 kj::Maybe<kj::Own<Service>> entrypointService;
4423 
4424 kj::Own<WorkerInterface> startRequestImpl(IoChannelFactory::SubrequestMetadata metadata) {
4425 auto& service = KJ_ASSERT_NONNULL(isolate->service);
4426 if (entrypointService == kj::none) {
4427 entrypointService = service->getEntrypoint(entrypointName, kj::mv(props));
4428 }
4429 KJ_IF_SOME(ep, entrypointService) {
4430 // Attach a refcounted reference to `this` (SubrequestChannelImpl) to the returned
4431 // WorkerInterface. This keeps the SubrequestChannelImpl alive for the duration of
4432 // the request, which in turn keeps the WorkerStubImpl alive (via Rc), preventing
4433 // WorkerStubImpl::unlink() from destroying the WorkerService's I/O channels while
4434 // the request's IoContext still holds raw pointers to the WorkerService.
4435 //
4436 // Without this, if the JS Fetcher object is garbage-collected mid-request (e.g.
4437 // because it was a temporary expression), the SubrequestChannelImpl is destroyed,
4438 // the WorkerStubImpl refcount drops to zero, unlink() clears the WorkerService's
4439 // LinkedIoChannels, and the child worker's IoContext crashes accessing them.
4440 return ep->startRequest(kj::mv(metadata)).attach(kj::addRef(*this));
4441 } else {
4442 KJ_IF_SOME(en, entrypointName) {
4443 JSG_FAIL_REQUIRE(Error, "Worker has no such entrypoint: ", en);
4444 } else {
4445 JSG_FAIL_REQUIRE(Error, "Worker has no default entrypoint.");
4446 }
4447 }
4448 }
4449 };
4450 
4451 class ActorClassImpl final: public ActorClass {
4452 public:
4453 ActorClassImpl(
4454 kj::Rc<WorkerStubImpl> isolate, kj::Maybe<kj::String> entrypointName, Frankenvalue props)
4455 : isolate(kj::mv(isolate)),
4456 entrypointName(kj::mv(entrypointName)),
4457 props(kj::mv(props)) {}
4458 
4459 void requireAllowsTransfer() override {
4460 throwDynamicEntrypointTransferError();
4461 }
4462 
4463 kj::Maybe<kj::Promise<void>> whenReady() override {
4464 if (inner != kj::none) return kj::none;
4465 
4466 KJ_IF_SOME(service, isolate->service) {
4467 inner = service->getActorClass(entrypointName, kj::mv(props));
4468 return kj::none;
4469 }
4470 
4471 // Have to wait for the isolate to start up. Capture a refcounted reference rather than
4472 // a raw `this` pointer so that the ActorClassImpl stays alive until the startup task
4473 // resolves, even if the owning object is garbage-collected while waiting.
4474 return isolate->startupTask.addBranch().then([self = kj::addRef(*this)]() mutable {
4475 if (self->inner == kj::none) {
4476 self->inner = KJ_ASSERT_NONNULL(self->isolate->service)
4477 ->getActorClass(self->entrypointName, kj::mv(self->props));
4478 }
4479 });
4480 }
4481 
4482 kj::Own<Worker::Actor> newActor(kj::Maybe<RequestTracker&> tracker,
4483 Worker::Actor::Id actorId,
4484 Worker::Actor::MakeActorCacheFunc makeActorCache,
4485 Worker::Actor::MakeStorageFunc makeStorage,
4486 kj::Own<Worker::Actor::Loopback> loopback,
4487 kj::Maybe<kj::Own<Worker::Actor::HibernationManager>> manager,
4488 kj::Maybe<rpc::Container::Client> container,
4489 kj::Maybe<Worker::Actor::FacetManager&> facetManager) override {
4490 return getInner().newActor(tracker, kj::mv(actorId), kj::mv(makeActorCache),
4491 kj::mv(makeStorage), kj::mv(loopback), kj::mv(manager), kj::mv(container),
4492 facetManager);
4493 }
4494 
4495 kj::Own<WorkerInterface> startRequest(
4496 IoChannelFactory::SubrequestMetadata metadata, kj::Own<Worker::Actor> actor) override {
4497 return getInner().startRequest(kj::mv(metadata), kj::mv(actor));
4498 }
4499 
4500 private:
4501 kj::Rc<WorkerStubImpl> isolate;
4502 kj::Maybe<kj::String> entrypointName;
4503 Frankenvalue props; // moved away when `inner` is initialized
4504 
4505 kj::Maybe<kj::Own<ActorClass>> inner;
4506 
4507 ActorClass& getInner() {
4508 return *KJ_ASSERT_NONNULL(
4509 inner, "ActorClassChannel is not ready yet; should have awaited whenReady()");
4510 }
4511 };
4512 };
4513};
4514 
4515void Server::unlinkWorkerLoaders() {
4516 for (auto& loader: workerLoaderNamespaces) {
4517 loader.value->unlink();
4518 }
4519 for (auto& loader: anonymousWorkerLoaderNamespaces) {
4520 loader->unlink();
4521 }
4522}
4523 
4524kj::Own<WorkerStubChannel> Server::WorkerService::loadIsolate(uint loaderChannel,
4525 kj::Maybe<kj::String> name,
4526 kj::Function<kj::Promise<DynamicWorkerSource>()> fetchSource) {
4527 auto& channels =
4528 KJ_REQUIRE_NONNULL(ioChannels.tryGet<LinkedIoChannels>(), "link() has not been called");
4529 KJ_REQUIRE(loaderChannel < channels.workerLoaders.size(), "invalid worker loader channel number");
4530 
4531 return channels.workerLoaders[loaderChannel]->loadIsolate(kj::mv(name), kj::mv(fetchSource));
4532}
4533 
4534kj::Promise<kj::Own<Server::Service>> Server::makeWorker(kj::StringPtr name,
4535 config::Worker::Reader conf,
4536 capnp::List<config::Extension>::Reader extensions) {
4537 TRACE_EVENT("workerd", "Server::makeWorker()", "name", name.cStr());
4538 auto& localActorConfigs = KJ_ASSERT_NONNULL(actorConfigs.find(name));
4539 
4540 ConfigErrorReporter errorReporter(*this, name);
4541 
4542 capnp::MallocMessageBuilder arena;
4543 // TODO(beta): Factor out FeatureFlags from WorkerBundle.
4544 auto featureFlags = arena.initRoot<CompatibilityFlags>();
4545 
4546 KJ_IF_SOME(overrideDate, testCompatibilityDateOverride) {
4547 // When testCompatibilityDateOverride is set, the config must NOT specify compatibilityDate.
4548 if (conf.hasCompatibilityDate()) {
4549 errorReporter.addError(kj::str(
4550 "Worker specifies compatibilityDate but --compat-date was provided. "
4551 "When using --compat-date, workers must not specify compatibilityDate in the config. "
4552 "Use compatibilityFlags to enable/disable specific flags if needed."));
4553 }
4554 // Use FUTURE_FOR_TEST to allow any valid date (including far future like 2999-12-31)
4555 // without validation against CODE_VERSION or current date.
4556 compileCompatibilityFlags(overrideDate, conf.getCompatibilityFlags(), featureFlags,
4557 errorReporter, experimental, CompatibilityDateValidation::FUTURE_FOR_TEST);
4558 } else if (conf.hasCompatibilityDate()) {
4559 compileCompatibilityFlags(conf.getCompatibilityDate(), conf.getCompatibilityFlags(),
4560 featureFlags, errorReporter, experimental, CompatibilityDateValidation::CODE_VERSION);
4561 } else {
4562 errorReporter.addError(kj::str("Worker must specify compatibilityDate."));
4563 }
4564 
4565 kj::Vector<FutureSubrequestChannel> subrequestChannels;
4566 kj::Vector<FutureActorChannel> actorChannels;
4567 kj::Vector<FutureActorClassChannel> actorClassChannels;
4568 kj::Vector<FutureWorkerLoaderChannel> workerLoaderChannels;
4569 bool hasWorkerdDebugPortBinding = false;
4570 
4571 auto confBindings = conf.getBindings();
4572 kj::Vector<WorkerdApi::Global> globals(confBindings.size());
4573 for (auto binding: confBindings) {
4574 KJ_IF_SOME(global,
4575 createBinding(name, conf, binding, errorReporter, subrequestChannels, actorChannels,
4576 actorClassChannels, workerLoaderChannels, hasWorkerdDebugPortBinding, actorConfigs,
4577 experimental)) {
4578 globals.add(kj::mv(global));
4579 }
4580 }
4581 
4582 // Construct `WorkerDef` from `conf`.
4583 WorkerDef def{
4584 .featureFlags = featureFlags.asReader(),
4585 .source = WorkerdApi::extractSource(name, conf, featureFlags.asReader(), errorReporter),
4586 .moduleFallback = conf.hasModuleFallback() ? kj::some(conf.getModuleFallback()) : kj::none,
4587 .localActorConfigs = localActorConfigs,
4588 .isDynamic = false,
4589 
4590 .globalOutbound{
4591 .designator = conf.getGlobalOutbound(),
4592 .errorContext = kj::str("Worker \"", name, "\"'s globalOutbound"),
4593 },
4594 
4595 .cacheApiOutbound = conf.hasCacheApiOutbound()
4596 ? kj::some(FutureSubrequestChannel{
4597 .designator = conf.getCacheApiOutbound(),
4598 .errorContext = kj::str("Worker \"", name, "\"'s cacheApiOutbound"),
4599 })
4600 : kj::none,
4601 
4602 .subrequestChannels = kj::mv(subrequestChannels),
4603 .actorChannels = kj::mv(actorChannels),
4604 .actorClassChannels = kj::mv(actorClassChannels),
4605 .workerLoaderChannels = kj::mv(workerLoaderChannels),
4606 .hasWorkerdDebugPortBinding = hasWorkerdDebugPortBinding,
4607 
4608 // clang-format off
4609 .tails = KJ_MAP(tail, conf.getTails()) -> FutureSubrequestChannel {
4610 return {
4611 .designator = tail,
4612 .errorContext = kj::str("Worker \"", name, "\"'s tails"),
4613 };
4614 },
4615 
4616 .streamingTails = KJ_MAP(streamingTail, conf.getStreamingTails()) -> FutureSubrequestChannel {
4617 return {
4618 .designator = streamingTail,
4619 .errorContext = kj::str("Worker \"", name, "\"'s streaming tails"),
4620 };
4621 },
4622 
4623 .actorStorageConf = conf.getDurableObjectStorage(),
4624 .containerEngineConf = conf.getContainerEngine(),
4625 
4626 .compileBindings = [globals = kj::mv(globals)](
4627 jsg::Lock& lock, const Worker::Api& api, v8::Local<v8::Object> target) {
4628 return WorkerdApi::from(api).compileGlobals(lock, globals, target, 1);
4629 },
4630 // clang-format on
4631 };
4632 
4633 co_return co_await makeWorkerImpl(name, kj::mv(def), extensions, errorReporter);
4634}
4635 
4636kj::Promise<kj::Own<Server::WorkerService>> Server::makeWorkerImpl(kj::StringPtr name,
4637 WorkerDef def,
4638 capnp::List<config::Extension>::Reader extensions,
4639 ErrorReporter& errorReporter) {
4640 // Load Python artifacts if this is a Python worker
4641 co_await preloadPython(name, def, errorReporter);
4642 
4643 auto jsgobserver = kj::atomicRefcounted<JsgIsolateObserver>();
4644 auto observer = kj::atomicRefcounted<IsolateObserver>();
4645 auto limitEnforcer = kj::refcounted<NullIsolateLimitEnforcer>();
4646 
4647 // Create the FsMap that will be used to map known file system
4648 // roots to configurable locations.
4649 // TODO(node-fs): This is set up to allow users to configure the "mount"
4650 // points for known roots but we currently do not expose that in the
4651 // config. So for now this just uses the defaults.
4652 auto workerFs = newWorkerFileSystem(kj::heap<FsMap>(), getBundleDirectory(def.source));
4653 
4654 // TODO(soon): Either make python workers support the new module registry before
4655 // NMR is defaulted on, or disable NMR by default when python workers are enabled.
4656 // While NMR is experimental, we'll just throw an error if both are enabled.
4657 if (def.featureFlags.getPythonWorkers()) {
4658 KJ_REQUIRE(!def.featureFlags.getNewModuleRegistry(),
4659 "Python workers do not currently support the new ModuleRegistry implementation. "
4660 "Please disable the new ModuleRegistry feature flag to use Python workers.");
4661 }
4662 
4663 bool usingNewModuleRegistry = def.featureFlags.getNewModuleRegistry();
4664 kj::Maybe<kj::Arc<jsg::modules::ModuleRegistry>> newModuleRegistry;
4665 // TODO(soon): Python workers do not currently support the new module registry.
4666 if (usingNewModuleRegistry) {
4667 KJ_REQUIRE(experimental,
4668 "The new ModuleRegistry implementation is an experimental feature. "
4669 "You must run workerd with `--experimental` to use this feature.");
4670 
4671 // We use the same path for modules that the virtual file system uses.
4672 // For instance, if the user specifies a bundle path of "/foo/bar" and
4673 // there is a module in the bundle at "/foo/bar/baz.js", then the module's
4674 // import specifier url will be "file:///foo/bar/baz.js".
4675 const jsg::Url& bundleBase = workerFs->getBundleRoot();
4676 
4677 // In workerd the module registry is always associated with just a single
4678 // worker instance, so we initialize it here. In production, however, a
4679 // single instance may be shared across multiple replicas.
4680 kj::Maybe<kj::String> maybeFallbackService;
4681 KJ_IF_SOME(moduleFallback, def.moduleFallback) {
4682 maybeFallbackService = kj::str(moduleFallback);
4683 }
4684 
4685 using ArtifactBundler = workerd::api::pyodide::ArtifactBundler;
4686 auto isPythonWorker = def.featureFlags.getPythonWorkers();
4687 auto artifactBundler = isPythonWorker
4688 ? ArtifactBundler::makePackagesOnlyBundler(pythonConfig.pyodidePackageManager)
4689 : ArtifactBundler::makeDisabledBundler();
4690 
4691 newModuleRegistry = WorkerdApi::newWorkerdModuleRegistry(*jsgobserver,
4692 def.source.variant.tryGet<Worker::Script::ModulesSource>(), def.featureFlags, pythonConfig,
4693 bundleBase, extensions, kj::mv(maybeFallbackService), kj::mv(artifactBundler));
4694 }
4695 
4696 auto isolateGroup = v8::IsolateGroup::GetDefault();
4697 auto api = kj::heap<WorkerdApi>(globalContext->v8System, def.featureFlags, extensions,
4698 limitEnforcer->getCreateParams(), isolateGroup, kj::mv(jsgobserver), *memoryCacheProvider,
4699 pythonConfig);
4700 
4701 auto inspectorPolicy = Worker::Isolate::InspectorPolicy::DISALLOW;
4702 if (inspectorOverride != kj::none) {
4703 // For workerd, if the inspector is enabled, it is always fully trusted.
4704 inspectorPolicy = Worker::Isolate::InspectorPolicy::ALLOW_FULLY_TRUSTED;
4705 }
4706 Worker::LoggingOptions isolateLoggingOptions = loggingOptions;
4707 isolateLoggingOptions.consoleMode =
4708 def.source.variant.is<WorkerSource::ScriptSource>() && !usingNewModuleRegistry
4709 ? Worker::ConsoleMode::INSPECTOR_ONLY
4710 : loggingOptions.consoleMode;
4711 auto isolate = kj::atomicRefcounted<Worker::Isolate>(kj::mv(api), kj::mv(observer), name,
4712 kj::mv(limitEnforcer), inspectorPolicy, kj::mv(isolateLoggingOptions));
4713 
4714 // If we are using the inspector, we need to register the Worker::Isolate
4715 // with the inspector service.
4716 KJ_IF_SOME(isolateRegistrar, inspectorIsolateRegistrar) {
4717 isolateRegistrar->registerIsolate(name, isolate.get());
4718 }
4719 
4720 if (!usingNewModuleRegistry) {
4721 KJ_IF_SOME(moduleFallback, def.moduleFallback) {
4722 KJ_REQUIRE(experimental,
4723 "The module fallback service is an experimental feature. "
4724 "You must run workerd with `--experimental` to use the module fallback service.");
4725 // If the config has the moduleFallback option, then we are going to set up the ability
4726 // to load certain modules from a fallback service. This is generally intended for local
4727 // dev/testing purposes only.
4728 auto& apiIsolate = isolate->getApi();
4729 auto fallbackClient =
4730 kj::heap<workerd::fallback::FallbackServiceClient>(kj::str(moduleFallback));
4731 apiIsolate.setModuleFallbackCallback(
4732 [client = kj::mv(fallbackClient), featureFlags = apiIsolate.getFeatureFlags()](
4733 jsg::Lock& js, kj::StringPtr specifier, kj::Maybe<kj::String> referrer,
4734 jsg::CompilationObserver& observer, jsg::ModuleRegistry::ResolveMethod method,
4735 kj::Maybe<kj::StringPtr> rawSpecifier) mutable
4736 -> kj::Maybe<kj::OneOf<kj::String, jsg::ModuleRegistry::ModuleInfo>> {
4737 kj::HashMap<kj::StringPtr, kj::StringPtr> attributes;
4738 KJ_IF_SOME(moduleOrRedirect,
4739 client->tryResolve(workerd::fallback::Version::V1,
4740 method == jsg::ModuleRegistry::ResolveMethod::IMPORT
4741 ? workerd::fallback::ImportType::IMPORT
4742 : workerd::fallback::ImportType::REQUIRE,
4743 specifier, rawSpecifier.orDefault(nullptr), referrer.orDefault(kj::String()),
4744 attributes)) {
4745 KJ_SWITCH_ONEOF(moduleOrRedirect) {
4746 KJ_CASE_ONEOF(redirect, kj::String) {
4747 // If a string is returned, then the fallback service returned a 301 redirect.
4748 // The value is the specifier of the new target module.
4749 return kj::Maybe(kj::mv(redirect));
4750 }
4751 KJ_CASE_ONEOF(module, kj::Own<config::Worker::Module::Reader>) {
4752 KJ_IF_SOME(module,
4753 WorkerdApi::tryCompileModule(js, *module, observer, featureFlags)) {
4754 return kj::Maybe(kj::mv(module));
4755 }
4756 KJ_LOG(ERROR, "Fallback service does not support this module type", module->which());
4757 }
4758 }
4759 }
4760 
4761 return kj::none;
4762 });
4763 }
4764 }
4765 
4766 using ArtifactBundler = workerd::api::pyodide::ArtifactBundler;
4767 auto isPythonWorker = def.featureFlags.getPythonWorkers();
4768 auto artifactBundler = isPythonWorker
4769 ? ArtifactBundler::makePackagesOnlyBundler(pythonConfig.pyodidePackageManager)
4770 : ArtifactBundler::makeDisabledBundler();
4771 
4772 auto script = isolate->newScript(name, def.source, IsolateObserver::StartType::COLD,
4773 SpanParent(nullptr), workerFs.attach(kj::mv(def.maybeOwnedSourceCode)), false, errorReporter,
4774 kj::mv(artifactBundler), kj::mv(newModuleRegistry));
4775 
4776 using Global = WorkerdApi::Global;
4777 jsg::V8Ref<v8::Object> ctxExportsHandle = nullptr;
4778 auto compileBindings = [&](jsg::Lock& lock, const Worker::Api& api, v8::Local<v8::Object> target,
4779 v8::Local<v8::Object> ctxExports) {
4780 // We can't fill in ctx.exports yet because we need to run the validator first to discover
4781 // entrypoints, which we cannot do until after the Worker constructor completes. We are
4782 // permitted to hold a handle until then, though.
4783 ctxExportsHandle = lock.v8Ref(ctxExports);
4784 
4785 return def.compileBindings(lock, api, target);
4786 };
4787 auto worker = kj::atomicRefcounted<Worker>(kj::mv(script), kj::atomicRefcounted<WorkerObserver>(),
4788 kj::mv(compileBindings), IsolateObserver::StartType::COLD, SpanParent(nullptr),
4789 Worker::Lock::TakeSynchronously(kj::none), errorReporter);
4790 
4791 uint totalActorChannels = 0;
4792 
4793 worker->runInLockScope(Worker::Lock::TakeSynchronously(kj::none), [&](Worker::Lock& lock) {
4794 lock.validateHandlers(errorReporter);
4795 
4796 // Build `ctx.exports` based on the entrypoints reported by `validateHandlers()`.
4797 kj::Vector<Global> ctxExports(
4798 errorReporter.namedEntrypoints.size() + def.localActorConfigs.size());
4799 
4800 // Start numbering loopback channels for stateless entrypoints after the last subrequest
4801 // channel used by bindings.
4802 uint nextSubrequestChannel =
4803 def.subrequestChannels.size() + IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT;
4804 if (errorReporter.defaultEntrypoint != kj::none) {
4805 ctxExports.add(Global{.name = kj::str("default"),
4806 .value = Global::LoopbackServiceStub{.channel = nextSubrequestChannel++}});
4807 }
4808 for (auto& ep: errorReporter.namedEntrypoints) {
4809 // Workflow classes are treated as stateless entrypoints for runtime purposes, but should
4810 // NOT be reflected in ctx.exports.
4811 // TODO(someday): Currently Workflows must be given a name independent of their class name,
4812 // and the binding must reference that name. If the name were just the class name -- like
4813 // Durable Object namespaces -- then we could put a `Workflow` binding into `ctx.exports`.
4814 if (!errorReporter.workflowClasses.contains(ep.key)) {
4815 ctxExports.add(Global{.name = kj::str(ep.key),
4816 .value = Global::LoopbackServiceStub{.channel = nextSubrequestChannel++}});
4817 }
4818 }
4819 
4820 // Start numbering loopback channels for actor classes after the last actor channel and actor
4821 // class channel used by bindings. Note that every exported actor class will have a ctx.exports
4822 // entry, but only the ones that have storage configured will be namespace bindings; the others
4823 // will be simply actor class bindings, which can be used with facets. We will iterate over
4824 // the exported class names and cross-reference with the storage config. Note that if the
4825 // storage config contains a class name that isn't among the exports, we won't create a
4826 // ctx.exports entry for it (it wouldn't work anyway).
4827 uint nextActorChannel = def.actorChannels.size();
4828 uint nextActorClassChannel = def.actorClassChannels.size();
4829 for (auto& className: errorReporter.actorClasses) {
4830 uint actorClassChannel = nextActorClassChannel++;
4831 
4832 decltype(Global::value) value;
4833 KJ_IF_SOME(ns, def.localActorConfigs.find(className)) {
4834 // This class has storage attached. We'll create a loopback actor namespace binding.
4835 KJ_SWITCH_ONEOF(ns) {
4836 KJ_CASE_ONEOF(durable, Durable) {
4837 value = Global::LoopbackDurableActorNamespace{
4838 .actorChannel = nextActorChannel++,
4839 .uniqueKey = durable.uniqueKey,
4840 .classChannel = actorClassChannel,
4841 };
4842 }
4843 KJ_CASE_ONEOF(ephemeral, Ephemeral) {
4844 value = Global::LoopbackEphemeralActorNamespace{
4845 .actorChannel = nextActorChannel++,
4846 .classChannel = actorClassChannel,
4847 };
4848 }
4849 }
4850 } else {
4851 // No storage attached. We'll create an actual class binding (for use with facets).
4852 value = Global::LoopbackActorClass{.channel = actorClassChannel};
4853 }
4854 ctxExports.add(Global{.name = kj::str(className), .value = kj::mv(value)});
4855 }
4856 totalActorChannels = nextActorChannel;
4857 
4858 JSG_WITHIN_CONTEXT_SCOPE(lock, lock.getContext(), [&](jsg::Lock& js) {
4859 WorkerdApi::from(worker->getIsolate().getApi())
4860 .compileGlobals(lock, ctxExports, ctxExportsHandle.getHandle(js), 1);
4861 });
4862 
4863 // As an optimization, drop this now while we have the lock.
4864 { auto drop = kj::mv(ctxExportsHandle); }
4865 });
4866 
4867 // Extract abortIsolateCallback before moving def into linkCallback lambda
4868 auto abortIsolateCallback = kj::mv(def.abortIsolateCallback);
4869 
4870 auto linkCallback = [this, def = kj::mv(def), totalActorChannels](WorkerService& workerService,
4871 Worker::ValidationErrorReporter& errorReporter) mutable {
4872 WorkerService::LinkedIoChannels result;
4873 
4874 auto entrypointNames = workerService.getEntrypointNames();
4875 auto actorClassNames = workerService.getActorClassNames();
4876 
4877 auto services = kj::heapArrayBuilder<kj::Own<IoChannelFactory::SubrequestChannel>>(
4878 def.subrequestChannels.size() + IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT +
4879 entrypointNames.size() + workerService.hasDefaultEntrypoint());
4880 
4881 auto globalService = kj::mv(def.globalOutbound).lookup(*this);
4882 
4883 // Bind both "next" and "null" to the global outbound. (The difference between these is a
4884 // legacy artifact that no one should be depending on.)
4885 static_assert(IoContext::SPECIAL_SUBREQUEST_CHANNEL_COUNT == 2);
4886 services.add(kj::addRef(*globalService));
4887 services.add(kj::mv(globalService));
4888 
4889 for (auto& channel: def.subrequestChannels) {
4890 services.add(kj::mv(channel).lookup(*this));
4891 }
4892 
4893 // Link the ctx.exports self-referential channels. Note that it's important these are added
4894 // in exactyl the same order as the channels were allocated earlier when we compiled the
4895 // ctx.exports bindings.
4896 if (workerService.hasDefaultEntrypoint()) {
4897 services.add(workerService.getLoopbackEntrypoint(/*name=*/kj::none));
4898 }
4899 for (auto& ep: entrypointNames) {
4900 services.add(workerService.getLoopbackEntrypoint(ep));
4901 }
4902 
4903 result.subrequest = services.finish();
4904 
4905 // Set up actor class channels
4906 auto actorClasses = kj::heapArrayBuilder<kj::Own<ActorClass>>(
4907 def.actorClassChannels.size() + actorClassNames.size());
4908 
4909 for (auto& channel: def.actorClassChannels) {
4910 actorClasses.add(kj::mv(channel).lookup(*this));
4911 }
4912 
4913 auto linkedActorChannels =
4914 kj::heapArrayBuilder<kj::Maybe<WorkerService::ActorNamespace&>>(totalActorChannels);
4915 
4916 for (auto& channel: def.actorChannels) {
4917 WorkerService* targetService = &workerService;
4918 if (channel.designator.hasServiceName()) {
4919 auto& svc = KJ_UNWRAP_OR(this->services.find(channel.designator.getServiceName()), {
4920 // error was reported earlier
4921 linkedActorChannels.add(kj::none);
4922 continue;
4923 });
4924 targetService = dynamic_cast<WorkerService*>(svc.get());
4925 if (targetService == nullptr) {
4926 // error was reported earlier
4927 linkedActorChannels.add(kj::none);
4928 continue;
4929 }
4930 }
4931 
4932 // (If getActorNamespace() returns null, an error was reported earlier.)
4933 linkedActorChannels.add(targetService->getActorNamespace(channel.designator.getClassName()));
4934 };
4935 
4936 // Link the ctx.exports self-referential actor channels. Again, it's important that these
4937 // be added in the same order as before. kj::HashMap iteration order is deterministic, and
4938 // is exactly insertion order as long as no entries have been removed, so we can expect that
4939 // `workerService.getActorClassNames()` iterates in the same order as
4940 // `errorReporter.actorClasses` did earlier. As before, every exported class gets an actor
4941 // class channel, but only the ones with configured storage will also get namespace channels.
4942 auto& selfActorNamespaces = workerService.getActorNamespaces();
4943 for (auto& className: actorClassNames) {
4944 actorClasses.add(workerService.getLoopbackActorClass(className));
4945 KJ_IF_SOME(ns, selfActorNamespaces.find(className)) {
4946 linkedActorChannels.add(*ns);
4947 }
4948 }
4949 
4950 result.actor = linkedActorChannels.finish();
4951 result.actorClass = actorClasses.finish();
4952 
4953 KJ_IF_SOME(out, def.cacheApiOutbound) {
4954 result.cache = kj::mv(out).lookup(*this);
4955 }
4956 
4957 if (def.actorStorageConf.isLocalDisk()) {
4958 kj::StringPtr diskName = def.actorStorageConf.getLocalDisk();
4959 KJ_IF_SOME(svc, this->services.find(def.actorStorageConf.getLocalDisk())) {
4960 auto diskSvc = dynamic_cast<DiskDirectoryService*>(svc.get());
4961 if (diskSvc == nullptr) {
4962 errorReporter.addError(kj::str("durableObjectStorage config refers to the service \"",
4963 diskName, "\", but that service is not a local disk service."));
4964 } else KJ_IF_SOME(dir, diskSvc->getWritable()) {
4965 result.actorStorage = dir;
4966 } else {
4967 errorReporter.addError(
4968 kj::str("durableObjectStorage config refers to the disk service \"", diskName,
4969 "\", but that service is defined read-only."));
4970 }
4971 } else {
4972 errorReporter.addError(kj::str("durableObjectStorage config refers to a service \"",
4973 diskName, "\", but no such service is defined."));
4974 }
4975 }
4976 
4977 result.tails = KJ_MAP(tail, def.tails) { return kj::mv(tail).lookup(*this); };
4978 
4979 result.streamingTails = KJ_MAP(tail, def.streamingTails) { return kj::mv(tail).lookup(*this); };
4980 
4981 result.workerLoaders = KJ_MAP(il, def.workerLoaderChannels) {
4982 KJ_IF_SOME(id, il.id) {
4983 return workerLoaderNamespaces
4984 .findOrCreate(id, [&]() -> decltype(workerLoaderNamespaces)::Entry {
4985 return {
4986 .key = kj::mv(id),
4987 .value = kj::rc<WorkerLoaderNamespace>(*this, kj::mv(il.name)),
4988 };
4989 }).addRef();
4990 } else {
4991 return anonymousWorkerLoaderNamespaces
4992 .add(kj::rc<WorkerLoaderNamespace>(*this, kj::mv(il.name)))
4993 .addRef();
4994 }
4995 };
4996 
4997 if (def.hasWorkerdDebugPortBinding) {
4998 result.workerdDebugPortNetwork = network;
4999 }
5000 
5001 return result;
5002 };
5003 
5004 kj::Maybe<kj::String> dockerPath = kj::none;
5005 kj::Maybe<kj::String> containerEgressInterceptorImage = kj::none;
5006 switch (def.containerEngineConf.which()) {
5007 case config::Worker::ContainerEngine::NONE:
5008 // No container engine configured
5009 break;
5010 case config::Worker::ContainerEngine::LOCAL_DOCKER: {
5011 auto dockerConf = def.containerEngineConf.getLocalDocker();
5012 dockerPath = kj::str(dockerConf.getSocketPath());
5013 if (dockerConf.hasContainerEgressInterceptorImage()) {
5014 containerEgressInterceptorImage = kj::str(dockerConf.getContainerEgressInterceptorImage());
5015 }
5016 break;
5017 }
5018 }
5019 
5020 kj::Maybe<kj::StringPtr> serviceName;
5021 if (!def.isDynamic) serviceName = name;
5022 
5023 auto result =
5024 kj::refcounted<WorkerService>(channelTokenHandler, serviceName, globalContext->threadContext,
5025 monotonicClock, kj::mv(worker), kj::mv(errorReporter.defaultEntrypoint),
5026 kj::mv(errorReporter.namedEntrypoints), kj::mv(errorReporter.actorClasses),
5027 kj::mv(linkCallback), KJ_BIND_METHOD(*this, abortAllActors),
5028 KJ_BIND_METHOD(*this, deleteAllActors), kj::mv(dockerPath),
5029 kj::mv(containerEgressInterceptorImage), def.isDynamic, kj::mv(abortIsolateCallback));
5030 result->initActorNamespaces(def.localActorConfigs, network);
5031 co_return result;
5032}
5033 
5034// =======================================================================================
5035 
5036kj::Promise<kj::Own<Server::Service>> Server::makeService(config::Service::Reader conf,
5037 kj::HttpHeaderTable::Builder& headerTableBuilder,
5038 capnp::List<config::Extension>::Reader extensions) {
5039 kj::StringPtr name = conf.getName();
5040 
5041 switch (conf.which()) {
5042 case config::Service::UNSPECIFIED:
5043 reportConfigError(kj::str("Service named \"", name, "\" does not specify what to serve."));
5044 co_return makeInvalidConfigService();
5045 
5046 case config::Service::EXTERNAL:
5047 co_return makeExternalService(name, conf.getExternal(), headerTableBuilder);
5048 
5049 case config::Service::NETWORK:
5050 co_return makeNetworkService(conf.getNetwork());
5051 
5052 case config::Service::WORKER:
5053 co_return co_await makeWorker(name, conf.getWorker(), extensions);
5054 
5055 case config::Service::DISK:
5056 co_return makeDiskDirectoryService(name, conf.getDisk(), headerTableBuilder);
5057 }
5058 
5059 reportConfigError(kj::str("Service named \"", name,
5060 "\" has unrecognized type. Was the config compiled with a "
5061 "newer version of the schema?"));
5062 co_return makeInvalidConfigService();
5063}
5064 
5065void Server::taskFailed(kj::Exception&& exception) {
5066 fatalFulfiller->reject(kj::mv(exception));
5067}
5068 
5069kj::Own<Server::Service> Server::lookupService(
5070 config::ServiceDesignator::Reader designator, kj::String errorContext) {
5071 kj::StringPtr targetName = designator.getName();
5072 Service* service = KJ_UNWRAP_OR(services.find(targetName), {
5073 reportConfigError(kj::str(errorContext, " refers to a service \"", targetName,
5074 "\", but no such service is defined."));
5075 return kj::addRef(*invalidConfigServiceSingleton);
5076 });
5077 
5078 kj::Maybe<kj::StringPtr> entrypointName;
5079 if (designator.hasEntrypoint()) {
5080 entrypointName = designator.getEntrypoint();
5081 }
5082 
5083 auto props = [&]() -> Frankenvalue {
5084 auto props = designator.getProps();
5085 switch (props.which()) {
5086 case config::ServiceDesignator::Props::EMPTY:
5087 return {};
5088 case config::ServiceDesignator::Props::JSON:
5089 return Frankenvalue::fromJson(kj::str(props.getJson()));
5090 }
5091 reportConfigError(kj::str(errorContext,
5092 " has unrecognized props type. Was the config compiled with a "
5093 "newer version of the schema?"));
5094 return {};
5095 }();
5096 
5097 if (WorkerService* worker = dynamic_cast<WorkerService*>(service)) {
5098 KJ_IF_SOME(ep, worker->getEntrypoint(entrypointName, kj::mv(props))) {
5099 return kj::mv(ep);
5100 } else KJ_IF_SOME(ep, entrypointName) {
5101 reportConfigError(kj::str(errorContext, " refers to service \"", targetName,
5102 "\" with a named entrypoint \"", ep, "\", but \"", targetName,
5103 "\" has no such named entrypoint."));
5104 return kj::addRef(*invalidConfigServiceSingleton);
5105 } else {
5106 reportConfigError(kj::str(errorContext, " refers to service \"", targetName,
5107 "\", but does not specify an entrypoint, and the service does not have a "
5108 "default entrypoint."));
5109 return kj::addRef(*invalidConfigServiceSingleton);
5110 }
5111 } else {
5112 KJ_IF_SOME(ep, entrypointName) {
5113 reportConfigError(kj::str(errorContext, " refers to service \"", targetName,
5114 "\" with a named entrypoint \"", ep, "\", but \"", targetName,
5115 "\" is not a Worker, so does not have any named entrypoints."));
5116 } else if (!props.empty()) {
5117 reportConfigError(kj::str(errorContext, " refers to service \"", targetName,
5118 "\" and provides a `props` value, but \"", targetName,
5119 "\" is not a Worker, so cannot accept `props`"));
5120 }
5121 
5122 return kj::addRef(*service);
5123 }
5124}
5125 
5126kj::Own<Server::ActorClass> Server::lookupActorClass(
5127 config::ServiceDesignator::Reader designator, kj::String errorContext) {
5128 // TODO(cleanup): There's a lot of repeated code with lookupService(), should it be refactored?
5129 
5130 kj::StringPtr targetName = designator.getName();
5131 Service* service = KJ_UNWRAP_OR(services.find(targetName), {
5132 reportConfigError(kj::str(errorContext, " refers to a service \"", targetName,
5133 "\", but no such service is defined."));
5134 return kj::addRef(*invalidConfigActorClassSingleton);
5135 });
5136 
5137 kj::Maybe<kj::StringPtr> entrypointName;
5138 if (designator.hasEntrypoint()) {
5139 entrypointName = designator.getEntrypoint();
5140 }
5141 
5142 auto props = [&]() -> Frankenvalue {
5143 auto props = designator.getProps();
5144 switch (props.which()) {
5145 case config::ServiceDesignator::Props::EMPTY:
5146 return {};
5147 case config::ServiceDesignator::Props::JSON:
5148 return Frankenvalue::fromJson(kj::str(props.getJson()));
5149 }
5150 reportConfigError(kj::str(errorContext,
5151 " has unrecognized props type. Was the config compiled with a "
5152 "newer version of the schema?"));
5153 return {};
5154 }();
5155 
5156 if (WorkerService* worker = dynamic_cast<WorkerService*>(service)) {
5157 KJ_IF_SOME(ep, worker->getActorClass(entrypointName, kj::mv(props))) {
5158 return kj::mv(ep);
5159 } else KJ_IF_SOME(ep, entrypointName) {
5160 reportConfigError(kj::str(errorContext, " refers to service \"", targetName,
5161 "\" with a Durable Object entrypoint \"", ep, "\", but \"", targetName,
5162 "\" has no such exported entrypoint class."));
5163 return kj::addRef(*invalidConfigActorClassSingleton);
5164 } else {
5165 reportConfigError(kj::str(errorContext, " refers to service \"", targetName,
5166 "\", but does not specify an entrypoint, and the service does export a "
5167 "Durable Object class as its default entrypoint."));
5168 return kj::addRef(*invalidConfigActorClassSingleton);
5169 }
5170 } else {
5171 KJ_IF_SOME(ep, entrypointName) {
5172 reportConfigError(kj::str(errorContext, " refers to service \"", targetName,
5173 "\" with a named Durable Object entrypoint \"", ep, "\", but \"", targetName,
5174 "\" is not a Worker, so does not have any named entrypoints."));
5175 } else {
5176 reportConfigError(kj::str(errorContext, " refers to service \"", targetName,
5177 "\" as a Durable Object class, but \"", targetName,
5178 "\" is not a Worker, so cannot be used as a class."));
5179 }
5180 
5181 return kj::addRef(*invalidConfigActorClassSingleton);
5182 }
5183}
5184 
5185kj::Own<IoChannelFactory::SubrequestChannel> Server::resolveEntrypoint(
5186 kj::StringPtr serviceName, kj::Maybe<kj::StringPtr> entrypoint, Frankenvalue props) {
5187 auto& service = *JSG_REQUIRE_NONNULL(services.find(serviceName), Error,
5188 "Stub refers to a service that doesn't exist: ", serviceName);
5189 
5190 auto& worker = JSG_REQUIRE_NONNULL(kj::tryDowncast<WorkerService>(service), Error,
5191 "Stub refers to a service that is not a Worker: ", serviceName);
5192 
5193 return JSG_REQUIRE_NONNULL(worker.getEntrypoint(entrypoint, kj::mv(props)), Error,
5194 "Stub refers to a an entrypoint of the target service that doesn't exist: ",
5195 entrypoint.orDefault("default"));
5196}
5197 
5198kj::Own<IoChannelFactory::ActorClassChannel> Server::resolveActorClass(
5199 kj::StringPtr serviceName, kj::Maybe<kj::StringPtr> entrypoint, Frankenvalue props) {
5200 auto& service = *JSG_REQUIRE_NONNULL(services.find(serviceName), Error,
5201 "Stub refers to a service that doesn't exist: ", serviceName);
5202 
5203 auto& worker = JSG_REQUIRE_NONNULL(kj::tryDowncast<WorkerService>(service), Error,
5204 "Stub refers to a service that is not a Worker: ", serviceName);
5205 
5206 return JSG_REQUIRE_NONNULL(worker.getActorClass(entrypoint, kj::mv(props)), Error,
5207 "Stub refers to a an entrypoint of the target service that doesn't exist: ",
5208 entrypoint.orDefault("default"));
5209}
5210 
5211// =======================================================================================
5212 
5213class Server::WorkerdBootstrapImpl final: public rpc::WorkerdBootstrap::Server {
5214 public:
5215 WorkerdBootstrapImpl(kj::Own<IoChannelFactory::SubrequestChannel> service,
5216 capnp::HttpOverCapnpFactory& httpOverCapnpFactory)
5217 : service(kj::mv(service)),
5218 httpOverCapnpFactory(httpOverCapnpFactory) {}
5219 
5220 kj::Promise<void> startEvent(StartEventContext context) override {
5221 // Extract the optional cf blob from the RPC params and pass it along with the
5222 // service channel to EventDispatcherImpl. The cf blob will be included in
5223 // SubrequestMetadata when creating the WorkerInterface for HTTP events.
5224 kj::Maybe<kj::String> cfBlobJson;
5225 auto params = context.getParams();
5226 if (params.hasCfBlobJson()) {
5227 cfBlobJson = kj::str(params.getCfBlobJson());
5228 }
5229 context.initResults(capnp::MessageSize{4, 1})
5230 .setDispatcher(kj::heap<EventDispatcherImpl>(
5231 httpOverCapnpFactory, kj::addRef(*service), kj::mv(cfBlobJson)));
5232 return kj::READY_NOW;
5233 }
5234 
5235 private:
5236 kj::Own<IoChannelFactory::SubrequestChannel> service;
5237 capnp::HttpOverCapnpFactory& httpOverCapnpFactory;
5238 
5239 class EventDispatcherImpl final: public rpc::EventDispatcher::Server {
5240 public:
5241 EventDispatcherImpl(capnp::HttpOverCapnpFactory& httpOverCapnpFactory,
5242 kj::Own<IoChannelFactory::SubrequestChannel> service,
5243 kj::Maybe<kj::String> cfBlobJson)
5244 : httpOverCapnpFactory(httpOverCapnpFactory),
5245 service(kj::mv(service)),
5246 cfBlobJson(kj::mv(cfBlobJson)) {}
5247 
5248 kj::Promise<void> getHttpService(GetHttpServiceContext context) override {
5249 // Create WorkerInterface with cf blob metadata (if provided via startEvent).
5250 IoChannelFactory::SubrequestMetadata metadata;
5251 KJ_IF_SOME(cf, cfBlobJson) {
5252 metadata.cfBlobJson = kj::str(cf);
5253 }
5254 auto worker = getService()->startRequest(kj::mv(metadata));
5255 context.initResults(capnp::MessageSize{4, 1})
5256 .setHttp(httpOverCapnpFactory.kjToCapnp(kj::mv(worker)));
5257 return kj::READY_NOW;
5258 }
5259 
5260 kj::Promise<void> sendTraces(SendTracesContext context) override {
5261 auto traces =
5262 KJ_MAP(trace, context.getParams().getTraces()){ return kj::refcounted<Trace>(trace); };
5263 auto event = kj::heap<api::TraceCustomEvent>(api::TraceCustomEvent::TYPE, kj::mv(traces));
5264 auto worker = getWorker();
5265 auto result = co_await worker->customEvent(kj::mv(event));
5266 auto resp = context.getResults().getResult();
5267 resp.setOutcome(result.outcome);
5268 }
5269 
5270 kj::Promise<void> prewarm(PrewarmContext context) override {
5271 throwUnsupported();
5272 }
5273 
5274 kj::Promise<void> runScheduled(RunScheduledContext context) override {
5275 throwUnsupported();
5276 }
5277 
5278 kj::Promise<void> runAlarm(RunAlarmContext context) override {
5279 throwUnsupported();
5280 }
5281 
5282 kj::Promise<void> queue(QueueContext context) override {
5283 throwUnsupported();
5284 }
5285 
5286 kj::Promise<void> jsRpcSession(JsRpcSessionContext context) override {
5287 return api::JsRpcSessionCustomEvent::receiveRpc(context, getWorker());
5288 }
5289 
5290 kj::Promise<void> tailStreamSession(TailStreamSessionContext context) override {
5291 auto customEvent = kj::heap<tracing::TailStreamCustomEvent>();
5292 auto cap = customEvent->getCap();
5293 capnp::PipelineBuilder<TailStreamSessionResults> pipelineBuilder;
5294 pipelineBuilder.setTopLevel(cap);
5295 context.setPipeline(pipelineBuilder.build());
5296 context.getResults().setTopLevel(kj::mv(cap));
5297 
5298 auto worker = getWorker();
5299 auto result = co_await worker->customEvent(kj::mv(customEvent)).attach(kj::mv(worker));
5300 auto response = context.getResults();
5301 response.setResult(result.outcome);
5302 }
5303 
5304 private:
5305 capnp::HttpOverCapnpFactory& httpOverCapnpFactory;
5306 kj::Maybe<kj::Own<IoChannelFactory::SubrequestChannel>> service;
5307 kj::Maybe<kj::String> cfBlobJson;
5308 
5309 kj::Own<IoChannelFactory::SubrequestChannel> getService() {
5310 auto result =
5311 kj::mv(KJ_ASSERT_NONNULL(service, "EventDispatcher can only be used for one request"));
5312 service = kj::none;
5313 return result;
5314 }
5315 
5316 kj::Own<WorkerInterface> getWorker() {
5317 // For non-HTTP events (RPC, traces, etc.), create WorkerInterface with
5318 // empty metadata since there's no HTTP request to extract cf from.
5319 return getService()->startRequest({});
5320 }
5321 
5322 [[noreturn]] void throwUnsupported() {
5323 JSG_FAIL_REQUIRE(Error, "RPC connections don't yet support this event type.");
5324 }
5325 };
5326};
5327 
5328class Server::HttpListener final: public kj::Refcounted {
5329 public:
5330 HttpListener(Server& owner,
5331 kj::Own<kj::ConnectionReceiver> listener,
5332 kj::Own<Service> service,
5333 kj::StringPtr physicalProtocol,
5334 kj::Own<HttpRewriter> rewriter,
5335 kj::HttpHeaderTable& headerTable,
5336 kj::Timer& timer,
5337 capnp::HttpOverCapnpFactory& httpOverCapnpFactory)
5338 : owner(owner),
5339 listener(kj::mv(listener)),
5340 service(kj::mv(service)),
5341 headerTable(headerTable),
5342 timer(timer),
5343 httpOverCapnpFactory(httpOverCapnpFactory),
5344 physicalProtocol(physicalProtocol),
5345 rewriter(kj::mv(rewriter)) {}
5346 
5347 kj::Promise<void> run() {
5348 TRACE_EVENT("workerd", "HttpListener::run");
5349 for (;;) {
5350 kj::AuthenticatedStream stream = co_await listener->acceptAuthenticated();
5351 TRACE_EVENT("workerd", "HTTPListener handle connection");
5352 
5353 kj::Maybe<kj::String> cfBlobJson;
5354 if (!rewriter->hasCfBlobHeader()) {
5355 // Construct a cf blob describing the client identity.
5356 
5357 kj::PeerIdentity* peerId;
5358 
5359 KJ_IF_SOME(tlsId,
5360 kj::dynamicDowncastIfAvailable<kj::TlsPeerIdentity>(*stream.peerIdentity)) {
5361 peerId = &tlsId.getNetworkIdentity();
5362 
5363 // TODO(someday): Add client certificate info to the cf blob? At present, KJ only
5364 // supplies the common name, but that doesn't even seem to be one of the fields that
5365 // Cloudflare-hosted Workers receive. We should probably try to match those.
5366 } else {
5367 peerId = stream.peerIdentity;
5368 }
5369 
5370 KJ_IF_SOME(remote, kj::dynamicDowncastIfAvailable<kj::NetworkPeerIdentity>(*peerId)) {
5371 cfBlobJson = kj::str("{\"clientIp\": ", escapeJsonString(remote.toString()), "}");
5372 } else KJ_IF_SOME(local, kj::dynamicDowncastIfAvailable<kj::LocalPeerIdentity>(*peerId)) {
5373 auto creds = local.getCredentials();
5374 
5375 kj::Vector<kj::String> parts;
5376 KJ_IF_SOME(p, creds.pid) {
5377 parts.add(kj::str("\"clientPid\":", p));
5378 }
5379 KJ_IF_SOME(u, creds.uid) {
5380 parts.add(kj::str("\"clientUid\":", u));
5381 }
5382 
5383 cfBlobJson = kj::str("{", kj::strArray(parts, ","), "}");
5384 }
5385 }
5386 
5387 auto conn = kj::heap<Connection>(*this, kj::mv(cfBlobJson));
5388 
5389 static auto constexpr listen = [](kj::Own<HttpListener> self, kj::Own<Connection> conn,
5390 kj::Own<kj::AsyncIoStream> stream) -> kj::Promise<void> {
5391 try {
5392 co_await conn->listedHttp.httpServer.listenHttp(kj::mv(stream));
5393 } catch (...) {
5394 KJ_LOG(ERROR, kj::getCaughtExceptionAsKj());
5395 }
5396 };
5397 
5398 // Run the connection handler loop in the global task set, so that run() waits for open
5399 // connections to finish before returning, even if the listener loop is canceled. However,
5400 // do not consider exceptions from a specific connection to be fatal.
5401 owner.tasks.add(listen(kj::addRef(*this), kj::mv(conn), kj::mv(stream.stream)));
5402 }
5403 }
5404 
5405 private:
5406 Server& owner;
5407 kj::Own<kj::ConnectionReceiver> listener;
5408 kj::Own<Service> service;
5409 kj::HttpHeaderTable& headerTable;
5410 kj::Timer& timer;
5411 capnp::HttpOverCapnpFactory& httpOverCapnpFactory;
5412 kj::StringPtr physicalProtocol;
5413 kj::Own<HttpRewriter> rewriter;
5414 
5415 kj::Maybe<capnp::TwoPartyServer> capnpServer;
5416 
5417 kj::Promise<void> acceptCapnpConnection(kj::AsyncIoStream& conn) {
5418 KJ_IF_SOME(s, capnpServer) {
5419 return s.accept(conn);
5420 }
5421 
5422 // Capnp server not initialized. Create it now.
5423 auto& s = capnpServer.emplace(
5424 kj::heap<WorkerdBootstrapImpl>(kj::addRef(*service), httpOverCapnpFactory));
5425 return s.accept(conn);
5426 }
5427 
5428 struct Connection final: public kj::HttpService, public kj::HttpServerErrorHandler {
5429 Connection(HttpListener& parent, kj::Maybe<kj::String> cfBlobJson)
5430 : parent(parent),
5431 cfBlobJson(kj::mv(cfBlobJson)),
5432 webSocketErrorHandler(kj::heap<JsgifyWebSocketErrors>()),
5433 listedHttp(parent.owner,
5434 parent.timer,
5435 parent.headerTable,
5436 *this,
5437 kj::HttpServerSettings{.errorHandler = *this,
5438 .webSocketErrorHandler = *webSocketErrorHandler,
5439 .webSocketCompressionMode = kj::HttpServerSettings::MANUAL_COMPRESSION}) {}
5440 
5441 HttpListener& parent;
5442 kj::Maybe<kj::String> cfBlobJson;
5443 kj::Own<JsgifyWebSocketErrors> webSocketErrorHandler;
5444 ListedHttpServer listedHttp;
5445 
5446 class ResponseWrapper final: public kj::HttpService::Response {
5447 public:
5448 ResponseWrapper(kj::HttpService::Response& inner, HttpRewriter& rewriter)
5449 : inner(inner),
5450 rewriter(rewriter) {}
5451 
5452 kj::Own<kj::AsyncOutputStream> send(uint statusCode,
5453 kj::StringPtr statusText,
5454 const kj::HttpHeaders& headers,
5455 kj::Maybe<uint64_t> expectedBodySize = kj::none) override {
5456 TRACE_EVENT("workerd", "ResponseWrapper::send()");
5457 auto rewrite = headers.cloneShallow();
5458 rewriter.rewriteResponse(rewrite);
5459 return inner.send(statusCode, statusText, rewrite, expectedBodySize);
5460 }
5461 
5462 kj::Own<kj::WebSocket> acceptWebSocket(const kj::HttpHeaders& headers) override {
5463 TRACE_EVENT("workerd", "ResponseWrapper::acceptWebSocket()");
5464 auto rewrite = headers.cloneShallow();
5465 rewriter.rewriteResponse(rewrite);
5466 return inner.acceptWebSocket(rewrite);
5467 }
5468 
5469 private:
5470 kj::HttpService::Response& inner;
5471 HttpRewriter& rewriter;
5472 };
5473 
5474 // ---------------------------------------------------------------------------
5475 // implements kj::HttpService
5476 
5477 kj::Promise<void> request(kj::HttpMethod method,
5478 kj::StringPtr url,
5479 const kj::HttpHeaders& headers,
5480 kj::AsyncInputStream& requestBody,
5481 kj::HttpService::Response& response) override {
5482 TRACE_EVENT("workerd", "Connection:request()");
5483 IoChannelFactory::SubrequestMetadata metadata;
5484 metadata.cfBlobJson = mapCopyString(cfBlobJson);
5485 
5486 Response* wrappedResponse = &response;
5487 kj::Own<ResponseWrapper> ownResponse;
5488 if (parent.rewriter->needsRewriteResponse()) {
5489 wrappedResponse = ownResponse = kj::heap<ResponseWrapper>(response, *parent.rewriter);
5490 }
5491 
5492 if (parent.rewriter->needsRewriteRequest() || cfBlobJson != kj::none) {
5493 auto rewrite = KJ_UNWRAP_OR(parent.rewriter->rewriteIncomingRequest(
5494 url, parent.physicalProtocol, headers, metadata.cfBlobJson),
5495 { co_return co_await response.sendError(400, "Bad Request", parent.headerTable); });
5496 auto worker = parent.service->startRequest(kj::mv(metadata));
5497 co_return co_await worker->request(
5498 method, url, *rewrite.headers, requestBody, *wrappedResponse);
5499 } else {
5500 auto worker = parent.service->startRequest(kj::mv(metadata));
5501 co_return co_await worker->request(method, url, headers, requestBody, *wrappedResponse);
5502 }
5503 }
5504 
5505 kj::Promise<void> connect(kj::StringPtr host,
5506 const kj::HttpHeaders& headers,
5507 kj::AsyncIoStream& connection,
5508 ConnectResponse& response,
5509 kj::HttpConnectSettings settings) override {
5510 TRACE_EVENT("workerd", "Connection:connect()");
5511 KJ_IF_SOME(h, parent.rewriter->getCapnpConnectHost()) {
5512 if (h == host) {
5513 // Client is requesting to open a capnp session!
5514 response.accept(200, "OK", kj::HttpHeaders(parent.headerTable));
5515 co_return co_await parent.acceptCapnpConnection(connection);
5516 }
5517 }
5518 
5519 IoChannelFactory::SubrequestMetadata metadata;
5520 metadata.cfBlobJson = mapCopyString(cfBlobJson);
5521 
5522 auto worker = parent.service->startRequest(kj::mv(metadata));
5523 co_return co_await worker->connect(host, headers, connection, response, kj::mv(settings));
5524 }
5525 
5526 // ---------------------------------------------------------------------------
5527 // implements kj::HttpServerErrorHandler
5528 
5529 kj::Promise<void> handleApplicationError(
5530 kj::Exception exception, kj::Maybe<kj::HttpService::Response&> response) override {
5531 if (exception.getType() == kj::Exception::Type::DISCONNECTED) {
5532 // Don't send a response, just close connection.
5533 co_return;
5534 }
5535 KJ_LOG(ERROR, kj::str("Uncaught exception: ", exception));
5536 KJ_IF_SOME(r, response) {
5537 co_return co_await r.sendError(500, "Internal Server Error", parent.headerTable);
5538 }
5539 }
5540 };
5541};
5542 
5543class Server::TcpListener final: public kj::Refcounted {
5544 public:
5545 TcpListener(Server& owner,
5546 kj::Own<kj::ConnectionReceiver> listener,
5547 kj::Own<Service> service,
5548 kj::HttpHeaderTable& headerTable,
5549 kj::StringPtr addrStr)
5550 : owner(owner),
5551 listener(kj::mv(listener)),
5552 service(kj::mv(service)),
5553 headerTable(headerTable),
5554 addrStr(addrStr) {}
5555 
5556 kj::Promise<void> run() {
5557 TRACE_EVENT("workerd", "TcpListener::run");
5558 for (;;) {
5559 kj::AuthenticatedStream stream = co_await listener->acceptAuthenticated();
5560 TRACE_EVENT("workerd", "TcpListener handle connection");
5561 
5562 IoChannelFactory::SubrequestMetadata metadata;
5563 auto req = service->startRequest(kj::mv(metadata));
5564 auto response = kj::heap<ResponseWrapper>();
5565 kj::HttpHeaders headers(headerTable);
5566 owner.tasks.add(req->connect(addrStr, headers, *stream.stream, *response, {})
5567 .attach(kj::mv(stream.stream), kj::mv(response))
5568 .attach(kj::mv(req)));
5569 }
5570 }
5571 
5572 private:
5573 Server& owner;
5574 kj::Own<kj::ConnectionReceiver> listener;
5575 kj::Own<Service> service;
5576 kj::HttpHeaderTable& headerTable;
5577 kj::StringPtr addrStr;
5578 
5579 struct ResponseWrapper final: public kj::HttpService::ConnectResponse {
5580 void accept(
5581 uint statusCode, kj::StringPtr statusText, const kj::HttpHeaders& headers) override {
5582 // Ok.. we're accepting the connection... anything to do?
5583 }
5584 kj::Own<kj::AsyncOutputStream> reject(uint statusCode,
5585 kj::StringPtr statusText,
5586 const kj::HttpHeaders& headers,
5587 kj::Maybe<uint64_t> expectedBodySize = kj::none) override {
5588 // Doh... we're rejecting the connection... anything to do?
5589 return newNullOutputStream();
5590 }
5591 };
5592};
5593 
5594kj::Promise<void> Server::listenHttp(kj::Own<kj::ConnectionReceiver> listener,
5595 kj::Own<Service> service,
5596 kj::StringPtr physicalProtocol,
5597 kj::Own<HttpRewriter> rewriter) {
5598 auto obj =
5599 kj::refcounted<HttpListener>(*this, kj::mv(listener), kj::mv(service), physicalProtocol,
5600 kj::mv(rewriter), globalContext->headerTable, timer, globalContext->httpOverCapnpFactory);
5601 co_return co_await obj->run();
5602}
5603 
5604kj::Promise<void> Server::listenTcp(
5605 kj::Own<kj::ConnectionReceiver> listener, kj::Own<Service> service, kj::StringPtr addrStr) {
5606 auto obj = kj::refcounted<TcpListener>(
5607 *this, kj::mv(listener), kj::mv(service), globalContext->headerTable, addrStr);
5608 co_return co_await obj->run();
5609}
5610 
5611// =======================================================================================
5612// Debug port for exposing all services via RPC
5613 
5614class Server::DebugPortListener {
5615 public:
5616 DebugPortListener(Server& owner,
5617 kj::Own<kj::ConnectionReceiver> listener,
5618 capnp::HttpOverCapnpFactory& httpOverCapnpFactory)
5619 : owner(owner),
5620 listener(kj::mv(listener)),
5621 httpOverCapnpFactory(httpOverCapnpFactory) {}
5622 
5623 kj::Promise<void> run() {
5624 capnp::TwoPartyServer server(kj::heap<WorkerdDebugPortImpl>(&owner, httpOverCapnpFactory));
5625 co_return co_await server.listen(*listener);
5626 }
5627 
5628 private:
5629 Server& owner;
5630 kj::Own<kj::ConnectionReceiver> listener;
5631 capnp::HttpOverCapnpFactory& httpOverCapnpFactory;
5632 
5633 class WorkerdDebugPortImpl final: public rpc::WorkerdDebugPort::Server {
5634 public:
5635 WorkerdDebugPortImpl(
5636 workerd::server::Server* srvPtr, capnp::HttpOverCapnpFactory& httpOverCapnpFactory)
5637 : srv(*srvPtr),
5638 httpOverCapnpFactory(httpOverCapnpFactory) {}
5639 
5640 kj::Promise<void> getEntrypoint(GetEntrypointContext context) override {
5641 auto params = context.getParams();
5642 auto serviceName = params.getService();
5643 auto propsReader = params.getProps();
5644 
5645 // Look up the service.
5646 auto& serviceEntry = KJ_ASSERT_NONNULL(srv.services.find(serviceName),
5647 kj::str("jsg.Error: Worker \"", serviceName, "\" not found"));
5648 auto service = serviceEntry->service();
5649 
5650 // Convert props from Frankenvalue if provided
5651 Frankenvalue props;
5652 if (params.hasProps()) {
5653 props = Frankenvalue::fromCapnp(propsReader);
5654 }
5655 
5656 kj::Own<Service> targetService;
5657 
5658 // Try to cast to WorkerService to support entrypoints and props
5659 auto* workerService = dynamic_cast<WorkerService*>(service);
5660 if (workerService != nullptr) {
5661 // This is a WorkerService, use getEntrypoint which supports both entrypoints and props
5662 kj::Maybe<kj::StringPtr> maybeEntrypoint;
5663 if (params.hasEntrypoint()) {
5664 maybeEntrypoint = params.getEntrypoint();
5665 }
5666 
5667 targetService =
5668 KJ_ASSERT_NONNULL(workerService->getEntrypoint(maybeEntrypoint, kj::mv(props)),
5669 kj::str("jsg.Error: Worker does not export an entrypoint named \"",
5670 maybeEntrypoint.orDefault("(default)"), "\""));
5671 } else {
5672 // Not a WorkerService
5673 KJ_ASSERT(!params.hasEntrypoint(), "jsg.Error: Worker does not support named entrypoints");
5674 
5675 // Try to apply props if the service supports it
5676 if (params.hasProps()) {
5677 targetService = service->forProps(kj::mv(props));
5678 } else {
5679 // No props, just use the service as-is
5680 targetService = kj::addRef(*service);
5681 }
5682 }
5683 
5684 // Return a WorkerdBootstrap that wraps this service using the generic implementation.
5685 context.initResults(capnp::MessageSize{4, 1})
5686 .setEntrypoint(
5687 kj::heap<WorkerdBootstrapImpl>(kj::mv(targetService), httpOverCapnpFactory));
5688 return kj::READY_NOW;
5689 }
5690 
5691 kj::Promise<void> getActor(GetActorContext context) override {
5692 auto params = context.getParams();
5693 auto serviceName = params.getService();
5694 auto entrypointName = params.getEntrypoint();
5695 auto actorIdStr = params.getActorId();
5696 
5697 // Look up the service
5698 auto& serviceEntry = KJ_ASSERT_NONNULL(srv.services.find(serviceName),
5699 kj::str("jsg.Error: Worker \"", serviceName, "\" not found"));
5700 auto service = serviceEntry->service();
5701 
5702 // Try to cast to WorkerService
5703 auto* workerService = dynamic_cast<WorkerService*>(service);
5704 KJ_REQUIRE(workerService != nullptr, "jsg.Error: Worker does not support Durable Objects");
5705 
5706 // Look up the actor namespace
5707 auto& actorNamespace = KJ_ASSERT_NONNULL(workerService->getActorNamespace(entrypointName),
5708 kj::str("jsg.Error: Worker does not export a Durable Object class named \"",
5709 entrypointName, "\""));
5710 
5711 // Create an actor ID - use the namespace config to determine if it's durable or ephemeral
5712 Worker::Actor::Id actorId;
5713 KJ_SWITCH_ONEOF(actorNamespace.getConfig()) {
5714 KJ_CASE_ONEOF(c, Durable) {
5715 // Durable Object ID (hex-encoded SHA256 hash)
5716 auto decoded = kj::decodeHex(actorIdStr);
5717 KJ_REQUIRE(decoded.size() == SHA256_DIGEST_LENGTH,
5718 "Invalid Durable Object ID: expected 64 hex characters (32 bytes)", decoded.size());
5719 kj::Own<ActorIdFactory::ActorId> id =
5720 kj::heap<ActorIdFactoryImpl::ActorIdImpl>(decoded.begin(), kj::none);
5721 actorId = kj::mv(id);
5722 }
5723 KJ_CASE_ONEOF(c, Ephemeral) {
5724 // Ephemeral actor ID (plain string)
5725 actorId = kj::str(actorIdStr);
5726 }
5727 }
5728 
5729 // Wrap the actor channel using the generic WorkerdBootstrap implementation.
5730 context.initResults(capnp::MessageSize{4, 1})
5731 .setActor(kj::heap<WorkerdBootstrapImpl>(
5732 actorNamespace.getActorChannel(kj::mv(actorId)), httpOverCapnpFactory));
5733 return kj::READY_NOW;
5734 }
5735 
5736 private:
5737 workerd::server::Server& srv;
5738 capnp::HttpOverCapnpFactory& httpOverCapnpFactory;
5739 };
5740};
5741 
5742kj::Promise<void> Server::listenDebugPort(kj::Own<kj::ConnectionReceiver> listener) {
5743 DebugPortListener obj(*this, kj::mv(listener), globalContext->httpOverCapnpFactory);
5744 co_return co_await obj.run();
5745}
5746 
5747// =======================================================================================
5748// Server::run()
5749 
5750kj::Promise<void> Server::handleDrain(kj::Promise<void> drainWhen) {
5751 co_await drainWhen;
5752 TRACE_EVENT("workerd", "Server::handleDrain()");
5753 // Tell all HttpServers to drain. This causes them to disconnect any connections that don't
5754 // have a request in-flight.
5755 for (auto& httpServer: httpServers) {
5756 // The promise returned by `drain()` resolves when all connections have ended. But, we need
5757 // the promise returned by handleDrain() to resolve immediately when draining has started,
5758 // since that's what signals us to stop accepting incoming connections. So, we should not
5759 // co_await the promise returned by `drain()`. Technically, we don't actually have to wait
5760 // on it at all -- `drain()` returns the promise end of a promise-and-fulfiller, so simply
5761 // dropping it won't actually cancel anything. But since that's not documented in drain()'s
5762 // doc comment, we instead add the promise to `tasks` to be safe.
5763 tasks.add(httpServer.httpServer.drain());
5764 }
5765}
5766 
5767kj::Promise<void> Server::run(
5768 jsg::V8System& v8System, config::Config::Reader config, kj::Promise<void> drainWhen) {
5769 TRACE_EVENT("workerd", "Server.run");
5770 
5771 // Update logging settings from config (overridding structuredLogging when so)
5772 if (config.hasLogging()) {
5773 auto logging = config.getLogging();
5774 loggingOptions.structuredLogging = StructuredLogging(logging.getStructuredLogging());
5775 if (logging.hasStdoutPrefix()) {
5776 loggingOptions.stdoutPrefix = kj::ConstString(kj::str(logging.getStdoutPrefix()));
5777 }
5778 if (logging.hasStderrPrefix()) {
5779 loggingOptions.stderrPrefix = kj::ConstString(kj::str(logging.getStderrPrefix()));
5780 }
5781 } else {
5782 loggingOptions.structuredLogging = StructuredLogging(config.getStructuredLogging());
5783 }
5784 
5785 kj::HttpHeaderTable::Builder headerTableBuilder;
5786 globalContext = kj::heap<GlobalContext>(*this, v8System, headerTableBuilder);
5787 invalidConfigServiceSingleton = kj::refcounted<InvalidConfigService>();
5788 invalidConfigActorClassSingleton = kj::refcounted<InvalidConfigActorClass>();
5789 
5790 auto [fatalPromise, fatalFulfiller] = kj::newPromiseAndFulfiller<void>();
5791 this->fatalFulfiller = kj::mv(fatalFulfiller);
5792 
5793 auto forkedDrainWhen = handleDrain(kj::mv(drainWhen)).fork();
5794 
5795 co_await startServices(v8System, config, headerTableBuilder, forkedDrainWhen);
5796 
5797 auto listenPromise = listenOnSockets(config, headerTableBuilder, forkedDrainWhen);
5798 
5799 // We should have registered all headers synchronously. This is important because we want to
5800 // be able to start handling requests as soon as the services are available, even if some other
5801 // services take longer to get ready.
5802 auto ownHeaderTable = headerTableBuilder.build();
5803 
5804 co_return co_await listenPromise.exclusiveJoin(kj::mv(fatalPromise));
5805}
5806 
5807// Configure and start the inspector socket, returning the port the socket started on.
5808uint startInspector(
5809 kj::StringPtr inspectorAddress, Server::InspectorServiceIsolateRegistrar& registrar) {
5810 static constexpr uint UNASSIGNED_PORT = 0;
5811 static constexpr uint DEFAULT_PORT = 9229;
5812 kj::MutexGuarded<uint> inspectorPort(UNASSIGNED_PORT);
5813 
5814 // `startInspector()` is called on the Isolate thread. V8 requires CPU profiling to be started and
5815 // stopped on the same thread which executes JavaScript -- that is, the Isolate thread -- which
5816 // means we need to dispatch inspector messages on this thread. To help make that happen, we
5817 // capture this thread's kj::Executor here, and pass it into the InspectorService below. Later,
5818 // when the InspectorService receives a WebSocket connection, it calls
5819 // `Isolate::attachInspector()`, which uses the kj::Executor we create here to create a
5820 // XThreadNotifier and start a dispatch loop. The InspectorService reads subsequent WebSocket
5821 // inspector messages and feeds them to that dispatch loop via the XThreadNotifier.
5822 auto isolateThreadExecutor = kj::getCurrentThreadExecutor().addRef();
5823 
5824 // Start the InspectorService thread.
5825 kj::Thread thread([inspectorAddress, &inspectorPort, &registrar,
5826 isolateThreadExecutor = kj::mv(isolateThreadExecutor)]() mutable {
5827 kj::AsyncIoContext io = kj::setupAsyncIo();
5828 
5829 kj::HttpHeaderTable::Builder headerTableBuilder;
5830 
5831 // Create the special inspector service.
5832 auto inspectorService(kj::heap<Server::InspectorService>(
5833 kj::mv(isolateThreadExecutor), io.provider->getTimer(), headerTableBuilder, registrar));
5834 auto ownHeaderTable = headerTableBuilder.build();
5835 
5836 // Configure and start the inspector socket.
5837 
5838 auto& network = io.provider->getNetwork();
5839 
5840 // TODO(cleanup): There's an issue here that if listen fails, nothing notices. The
5841 // server will continue running but will no longer accept inspector connections.
5842 // This should be fixed by:
5843 // 1. Replacing the kj::NEVER_DONE with listen
5844 // 2. Making the thread's lambda `noexcept` so that if it throws the process crashes
5845 // 3. Probably also throw if listen completes without an exception (even if unlikely to
5846 // happen)
5847 auto listen = (kj::coCapture(
5848 [&network, &inspectorAddress, &inspectorPort, &inspectorService]() -> kj::Promise<void> {
5849 auto parsed = co_await network.parseAddress(inspectorAddress, DEFAULT_PORT);
5850 auto listener = parsed->listen();
5851 // EW-7716: Signal to thread that started the inspector service that the inspector is ready.
5852 *inspectorPort.lockExclusive() = listener->getPort();
5853 KJ_LOG(INFO, "Inspector is listening");
5854 co_await inspectorService->listen(kj::mv(listener));
5855 }))();
5856 
5857 kj::NEVER_DONE.wait(io.waitScope);
5858 });
5859 thread.detach();
5860 
5861 // EW-7716: Wait for the InspectorService instance to be initialized before proceeding.
5862 return inspectorPort.when([](const uint& port) { return port != UNASSIGNED_PORT; },
5863 [](const uint& port) { return port; });
5864}
5865 
5866kj::Promise<void> Server::preloadPython(
5867 kj::StringPtr workerName, const WorkerDef& workerDef, ErrorReporter& errorReporter) {
5868 if (workerDef.featureFlags.getPythonWorkers()) {
5869 auto pythonRelease = getPythonSnapshotRelease(workerDef.featureFlags);
5870 KJ_IF_SOME(release, pythonRelease) {
5871 auto version = getPythonBundleName(release);
5872 
5873 // Fetch the Pyodide bundle.
5874 co_await server::fetchPyodideBundle(pythonConfig, kj::mv(version), network, timer);
5875 
5876 // Preload Python packages.
5877 KJ_IF_SOME(modulesSource, workerDef.source.variant.tryGet<Worker::Script::ModulesSource>()) {
5878 if (modulesSource.isPython) {
5879 auto pythonRequirements = getPythonRequirements(modulesSource);
5880 
5881 // Store the packages in the package manager that is stored in the pythonConfig
5882 co_await server::fetchPyodidePackages(pythonConfig, pythonConfig.pyodidePackageManager,
5883 pythonRequirements, release, network, timer);
5884 }
5885 }
5886 }
5887 }
5888}
5889 
5890kj::Promise<void> Server::startServices(jsg::V8System& v8System,
5891 config::Config::Reader config,
5892 kj::HttpHeaderTable::Builder& headerTableBuilder,
5893 kj::ForkedPromise<void>& forkedDrainWhen) {
5894 // ---------------------------------------------------------------------------
5895 // Configure services
5896 TRACE_EVENT("workerd", "startServices");
5897 
5898 // First pass: Extract actor namespace configs.
5899 for (auto serviceConf: config.getServices()) {
5900 kj::StringPtr name = serviceConf.getName();
5901 kj::HashMap<kj::String, ActorConfig> serviceActorConfigs;
5902 
5903 if (serviceConf.isWorker()) {
5904 auto workerConf = serviceConf.getWorker();
5905 bool hadDurable = false;
5906 for (auto ns: workerConf.getDurableObjectNamespaces()) {
5907 switch (ns.which()) {
5908 case config::Worker::DurableObjectNamespace::UNIQUE_KEY:
5909 hadDurable = true;
5910 serviceActorConfigs.insert(kj::str(ns.getClassName()),
5911 Durable{.uniqueKey = kj::str(ns.getUniqueKey()),
5912 .isEvictable = !ns.getPreventEviction(),
5913 .enableSql = ns.getEnableSql(),
5914 .containerOptions = ns.hasContainer() ? kj::Maybe(ns.getContainer()) : kj::none});
5915 continue;
5916 case config::Worker::DurableObjectNamespace::EPHEMERAL_LOCAL:
5917 if (!experimental) {
5918 reportConfigError(kj::str(
5919 "Ephemeral objects (Durable Object namespaces with type 'ephemeralLocal') are an "
5920 "experimental feature which may change or go away in the future. You must run "
5921 "workerd with `--experimental` to use this feature."));
5922 }
5923 serviceActorConfigs.insert(kj::str(ns.getClassName()),
5924 Ephemeral{.isEvictable = !ns.getPreventEviction(), .enableSql = ns.getEnableSql()});
5925 continue;
5926 }
5927 reportConfigError(kj::str("Encountered unknown DurableObjectNamespace type in service \"",
5928 name, "\", class \"", ns.getClassName(),
5929 "\". Was the config compiled with a newer version "
5930 "of the schema?"));
5931 }
5932 
5933 switch (workerConf.getDurableObjectStorage().which()) {
5934 case config::Worker::DurableObjectStorage::NONE:
5935 if (hadDurable) {
5936 reportConfigError(kj::str("Worker service \"", name,
5937 "\" implements durable object classes but has "
5938 "`durableObjectStorage` set to `none`."));
5939 }
5940 goto validDurableObjectStorage;
5941 case config::Worker::DurableObjectStorage::IN_MEMORY:
5942 case config::Worker::DurableObjectStorage::LOCAL_DISK:
5943 goto validDurableObjectStorage;
5944 }
5945 reportConfigError(kj::str("Encountered unknown durableObjectStorage type in service \"", name,
5946 "\". Was the config compiled with a newer version of the schema?"));
5947 
5948 validDurableObjectStorage:
5949 if (workerConf.hasDurableObjectUniqueKeyModifier()) {
5950 // This should be implemented along with parameterized workers. It's not relevant
5951 // otherwise, but let's make sure no one sets it accidentally.
5952 KJ_UNIMPLEMENTED("durableObjectUniqueKeyModifier is not implemented yet");
5953 }
5954 }
5955 
5956 actorConfigs.upsert(kj::str(name), kj::mv(serviceActorConfigs), [&](auto&&...) {
5957 reportConfigError(kj::str("Config defines multiple services named \"", name, "\"."));
5958 });
5959 }
5960 
5961 // If we are using the inspector, we need to register the Worker::Isolate
5962 // with the inspector service.
5963 KJ_IF_SOME(inspectorAddress, inspectorOverride) {
5964 auto registrar = kj::heap<InspectorServiceIsolateRegistrar>();
5965 auto port = startInspector(inspectorAddress, *registrar);
5966 KJ_IF_SOME(stream, controlOverride) {
5967 auto message = kj::str("{\"event\":\"listen-inspector\",\"port\":", port, "}\n");
5968 try {
5969 stream->write(message.asBytes());
5970 } catch (kj::Exception& e) {
5971 KJ_LOG(ERROR, e);
5972 }
5973 }
5974 inspectorIsolateRegistrar = kj::mv(registrar);
5975 }
5976 
5977 // Second pass: Build services.
5978 for (auto serviceConf: config.getServices()) {
5979 kj::StringPtr name = serviceConf.getName();
5980 auto service = co_await makeService(serviceConf, headerTableBuilder, config.getExtensions());
5981 
5982 services.upsert(kj::str(name), kj::mv(service), [&](auto&&...) {
5983 reportConfigError(kj::str("Config defines multiple services named \"", name, "\"."));
5984 });
5985 }
5986 
5987 // Make the default "internet" service if it's not there already.
5988 services.findOrCreate("internet"_kj, [&]() {
5989 auto publicNetwork = network.restrictPeers({"public"_kj});
5990 
5991 kj::TlsContext::Options options;
5992 options.useSystemTrustStore = true;
5993 
5994 kj::Own<kj::TlsContext> tls = kj::heap<kj::TlsContext>(kj::mv(options));
5995 auto tlsNetwork = tls->wrapNetwork(*publicNetwork);
5996 
5997 // Attaching to refcounted NetworkService is safe since services map is long-lived
5998 auto service = kj::refcounted<NetworkService>(globalContext->headerTable, timer, entropySource,
5999 kj::mv(publicNetwork), kj::mv(tlsNetwork), *tls)
6000 .attachToThisReference(kj::mv(tls));
6001 
6002 return decltype(services)::Entry{kj::str("internet"_kj), kj::mv(service)};
6003 });
6004 
6005 // Third pass: Cross-link services.
6006 for (auto& service: services) {
6007 ConfigErrorReporter errorReporter(*this, service.key);
6008 service.value->link(errorReporter);
6009 }
6010}
6011 
6012kj::Maybe<Server::SocketTypeConfig> Server::parseSocketType(
6013 config::Socket::Reader sock, kj::StringPtr name) {
6014 switch (sock.which()) {
6015 case config::Socket::HTTP: {
6016 SocketTypeConfig result;
6017 result.defaultPort = 80;
6018 result.httpOptions = sock.getHttp();
6019 result.physicalProtocol = "http";
6020 return kj::mv(result);
6021 }
6022 case config::Socket::HTTPS: {
6023 auto https = sock.getHttps();
6024 SocketTypeConfig result;
6025 result.defaultPort = 443;
6026 result.httpOptions = https.getOptions();
6027 result.tls = makeTlsContext(https.getTlsOptions());
6028 result.physicalProtocol = "https";
6029 return kj::mv(result);
6030 }
6031 case config::Socket::TCP: {
6032 auto tcp = sock.getTcp();
6033 SocketTypeConfig result;
6034 if (tcp.hasTlsOptions()) {
6035 result.tls = makeTlsContext(tcp.getTlsOptions());
6036 }
6037 return kj::mv(result);
6038 }
6039 }
6040 reportConfigError(kj::str("Encountered unknown socket type in \"", name,
6041 "\". Was the config compiled with a newer version of the schema?"));
6042 return kj::none;
6043}
6044 
6045kj::Promise<void> Server::listenOnSockets(config::Config::Reader config,
6046 kj::HttpHeaderTable::Builder& headerTableBuilder,
6047 kj::ForkedPromise<void>& forkedDrainWhen,
6048 bool forTest) {
6049 // ---------------------------------------------------------------------------
6050 // Start sockets
6051 TRACE_EVENT("workerd", "listenOnSockets");
6052 for (auto sock: config.getSockets()) {
6053 kj::StringPtr name = sock.getName();
6054 kj::StringPtr addrStr = nullptr;
6055 kj::String ownAddrStr;
6056 kj::Maybe<kj::Own<kj::ConnectionReceiver>> listenerOverride;
6057 
6058 kj::Own<Service> service = lookupService(sock.getService(), kj::str("Socket \"", name, "\""));
6059 
6060 KJ_IF_SOME(override, socketOverrides.findEntry(name)) {
6061 KJ_SWITCH_ONEOF(override.value) {
6062 KJ_CASE_ONEOF(str, kj::String) {
6063 addrStr = ownAddrStr = kj::mv(str);
6064 break;
6065 }
6066 KJ_CASE_ONEOF(l, kj::Own<kj::ConnectionReceiver>) {
6067 listenerOverride = kj::mv(l);
6068 break;
6069 }
6070 }
6071 socketOverrides.erase(override);
6072 } else if (sock.hasAddress()) {
6073 addrStr = sock.getAddress();
6074 } else {
6075 reportConfigError(kj::str("Socket \"", name,
6076 "\" has no address in the config, so must be specified on the "
6077 "command line with `--socket-addr`."));
6078 continue;
6079 }
6080 
6081 auto maybeSocketConfig = parseSocketType(sock, name);
6082 if (maybeSocketConfig == kj::none) continue;
6083 auto& socketConfig = KJ_ASSERT_NONNULL(maybeSocketConfig);
6084 
6085 using PromisedReceived = kj::Promise<kj::Own<kj::ConnectionReceiver>>;
6086 PromisedReceived listener = nullptr;
6087 KJ_IF_SOME(l, listenerOverride) {
6088 listener = kj::mv(l);
6089 } else {
6090 listener = ([](kj::Promise<kj::Own<kj::NetworkAddress>> promise) -> PromisedReceived {
6091 auto parsed = co_await promise;
6092 co_return parsed->listen();
6093 })(network.parseAddress(addrStr, socketConfig.defaultPort));
6094 }
6095 
6096 KJ_IF_SOME(t, socketConfig.tls) {
6097 listener = ([](kj::Promise<kj::Own<kj::ConnectionReceiver>> promise,
6098 kj::Own<kj::TlsContext> tls) -> PromisedReceived {
6099 auto port = co_await promise;
6100 co_return tls->wrapPort(kj::mv(port)).attach(kj::mv(tls));
6101 })(kj::mv(listener), kj::mv(t));
6102 }
6103 
6104 // Need to create rewriter before waiting on anything since `headerTableBuilder` will no longer
6105 // be available later.
6106 auto rewriter = kj::heap<HttpRewriter>(socketConfig.httpOptions, headerTableBuilder);
6107 
6108 auto handle = kj::coCapture(
6109 [this, service = kj::mv(service), rewriter = kj::mv(rewriter),
6110 physicalProtocol = socketConfig.physicalProtocol, name,
6111 isHttp = sock.which() != config::Socket::TCP, addrStr](
6112 kj::Promise<kj::Own<kj::ConnectionReceiver>> promise) mutable -> kj::Promise<void> {
6113 if (isHttp) {
6114 TRACE_EVENT("workerd", "setup listenHttp");
6115 } else {
6116 TRACE_EVENT("workerd", "setup listenTcp");
6117 }
6118 
6119 auto listener = co_await promise;
6120 KJ_IF_SOME(stream, controlOverride) {
6121 auto message = kj::str("{\"event\":\"listen\",\"socket\":\"", name,
6122 "\",\"port\":", listener->getPort(), "}\n");
6123 try {
6124 stream->write(message.asBytes());
6125 } catch (kj::Exception& e) {
6126 KJ_LOG(ERROR, e);
6127 }
6128 }
6129 
6130 if (isHttp) {
6131 co_await listenHttp(kj::mv(listener), kj::mv(service), physicalProtocol, kj::mv(rewriter));
6132 } else {
6133 co_await listenTcp(kj::mv(listener), kj::mv(service), addrStr);
6134 }
6135 });
6136 tasks.add(handle(kj::mv(listener)).exclusiveJoin(forkedDrainWhen.addBranch()));
6137 }
6138 
6139 // Start debug port if configured
6140 KJ_IF_SOME(addr, debugPortOverride) {
6141 auto handle = kj::coCapture(
6142 [this, addr = kj::str(addr)](kj::ForkedPromise<void>& drain) mutable -> kj::Promise<void> {
6143 auto parsed = co_await network.parseAddress(addr, 0);
6144 auto listener = parsed->listen();
6145 
6146 KJ_IF_SOME(stream, controlOverride) {
6147 auto message = kj::str("{\"event\":\"listen\",\"socket\":\"debug-port"
6148 "\",\"port\":",
6149 listener->getPort(), "}\n");
6150 try {
6151 stream->write(message.asBytes());
6152 } catch (kj::Exception& e) {
6153 KJ_LOG(ERROR, e);
6154 }
6155 }
6156 
6157 co_await listenDebugPort(kj::mv(listener));
6158 });
6159 tasks.add(handle(forkedDrainWhen).exclusiveJoin(forkedDrainWhen.addBranch()));
6160 }
6161 
6162 for (auto& unmatched: socketOverrides) {
6163 reportConfigError(kj::str("Config did not define any socket named \"", unmatched.key,
6164 "\" to match the override "
6165 "provided on the command line."));
6166 }
6167 
6168 for (auto& unmatched: externalOverrides) {
6169 reportConfigError(kj::str("Config did not define any external service named \"", unmatched.key,
6170 "\" to match the "
6171 "override provided on the command line."));
6172 }
6173 
6174 for (auto& unmatched: directoryOverrides) {
6175 if (forTest && unmatched.key == "TEST_TMPDIR") {
6176 // Due to a historical bug, `workerd test` didn't check for the existence of unmatched
6177 // overrides, and our own tests became dependent on the ability to override TEST_TMPDIR
6178 // even if it was not used in the config. For now, we ignore this problem.
6179 //
6180 // TODO(cleanup): Figure out the right solution here.
6181 continue;
6182 }
6183 
6184 reportConfigError(kj::str("Config did not define any disk service named \"", unmatched.key,
6185 "\" to match the "
6186 "override provided on the command line."));
6187 }
6188 
6189 co_await tasks.onEmpty();
6190 
6191 // Give a chance for any errors to bubble up before we return success. In particular
6192 // Server::taskFailed() fulfills `fatalFulfiller`, which causes the server to exit with an error.
6193 // But the `TaskSet` may have become empty at the same time. We want the error to win the race
6194 // against the success.
6195 //
6196 // TODO(cleanup): A better solution would be for `TaskSet` to have a new variant of the
6197 // `onEmpty()` method like `onEmptyOrException()`, which propagates any exception thrown by
6198 // any task.
6199 co_await kj::yieldUntilQueueEmpty();
6200}
6201 
6202// =======================================================================================
6203// Server::test()
6204 
6205kj::Promise<bool> Server::test(jsg::V8System& v8System,
6206 config::Config::Reader config,
6207 kj::StringPtr servicePattern,
6208 kj::StringPtr entrypointPattern) {
6209 
6210 if (config.hasLogging()) {
6211 auto logging = config.getLogging();
6212 loggingOptions.structuredLogging = StructuredLogging(logging.getStructuredLogging());
6213 if (logging.hasStdoutPrefix()) {
6214 loggingOptions.stdoutPrefix = kj::ConstString(kj::str(logging.getStdoutPrefix()));
6215 }
6216 if (logging.hasStderrPrefix()) {
6217 loggingOptions.stderrPrefix = kj::ConstString(kj::str(logging.getStderrPrefix()));
6218 }
6219 } else {
6220 loggingOptions.structuredLogging = StructuredLogging(config.getStructuredLogging());
6221 }
6222 
6223 kj::HttpHeaderTable::Builder headerTableBuilder;
6224 globalContext = kj::heap<GlobalContext>(*this, v8System, headerTableBuilder);
6225 invalidConfigServiceSingleton = kj::refcounted<InvalidConfigService>();
6226 
6227 auto [fatalPromise, fatalFulfiller] = kj::newPromiseAndFulfiller<void>();
6228 this->fatalFulfiller = kj::mv(fatalFulfiller);
6229 
6230 auto forkedDrainWhen = kj::Promise<void>(kj::NEVER_DONE).fork();
6231 
6232 co_await startServices(v8System, config, headerTableBuilder, forkedDrainWhen);
6233 
6234 // Tests usually do not configure sockets, but they can, especially loopback sockets. Arrange
6235 // to wait on them. Crash if listening fails.
6236 auto listenPromise =
6237 listenOnSockets(config, headerTableBuilder, forkedDrainWhen,
6238 /* forTest = */ true)
6239 .eagerlyEvaluate([](kj::Exception&& e) noexcept { kj::throwFatalException(kj::mv(e)); });
6240 
6241 auto ownHeaderTable = headerTableBuilder.build();
6242 
6243 // TODO(someday): If the inspector is enabled, pause and wait for an inspector connection before
6244 // proceeding?
6245 
6246 kj::GlobFilter serviceGlob(servicePattern);
6247 kj::GlobFilter entrypointGlob(entrypointPattern);
6248 
6249 uint passCount = 0, failCount = 0;
6250 
6251 auto doTest = [&](Service& service, kj::StringPtr name) -> kj::Promise<void> {
6252 // TODO(soon): Better way of reporting test results, KJ_LOG is ugly. We should probably have
6253 // some sort of callback interface. It would be nice to report the exceptions thrown through
6254 // that interface too... can we? Use a tracer maybe?
6255 // HACK: We use DBG log level because INFO logging is optional, and warning/error would confuse
6256 // people. Note that server-test.c++ actually tests for this logging, so simply writing to
6257 // stderr wouldn't work.
6258 KJ_LOG(DBG, kj::str("[ TEST ] "_kj, name));
6259 auto req = service.startRequest({});
6260 auto start = monotonicClock.now();
6261 
6262 bool result = co_await req->test();
6263 if (result) {
6264 ++passCount;
6265 } else {
6266 ++failCount;
6267 }
6268 
6269 auto end = monotonicClock.now();
6270 auto duration = end - start;
6271 
6272 KJ_LOG(DBG, kj::str(result ? "[ PASS ] "_kj : "[ FAIL ] "_kj, name, " (", duration, ")"));
6273 };
6274 
6275 for (auto& service: services) {
6276 if (serviceGlob.matches(service.key)) {
6277 if (service.value->hasHandler("test"_kj) && entrypointGlob.matches("default"_kj)) {
6278 co_await doTest(*service.value, service.key);
6279 }
6280 
6281 if (WorkerService* worker = dynamic_cast<WorkerService*>(service.value.get())) {
6282 for (auto& name: worker->getEntrypointNames()) {
6283 if (entrypointGlob.matches(name)) {
6284 kj::Own<Service> ep = KJ_ASSERT_NONNULL(worker->getEntrypoint(name, /*props=*/{}));
6285 if (ep->hasHandler("test"_kj)) {
6286 co_await doTest(*ep, kj::str(service.key, ':', name));
6287 }
6288 }
6289 }
6290 }
6291 }
6292 }
6293 
6294 if (passCount + failCount == 0) {
6295 KJ_LOG(ERROR, "No tests found!");
6296 }
6297 
6298 co_return passCount > 0 && failCount == 0;
6299}
6300 
6301} // namespace workerd::server