Skip to content
File

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

62.1 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 <workerd/api/unsafe.h>
8#include <workerd/io/compatibility-date.capnp.h>
9#include <workerd/io/compatibility-date.h>
10#include <workerd/io/release-version.embed.h>
11#include <workerd/jsg/setup.h>
12#include <workerd/rust/cxx-integration/lib.rs.h>
13#include <workerd/server/cpp-capnp-schema.embed.h>
14#include <workerd/server/json-logger.h>
15#include <workerd/server/v8-platform-impl.h>
16#include <workerd/server/workerd-capnp-schema.embed.h>
17#include <workerd/server/workerd.capnp.h>
18#include <workerd/util/autogate.h>
19#include <workerd/util/entropy.h>
20 
21#include <errno.h>
22#include <fcntl.h>
23#ifdef __linux__
24#include <sys/mman.h>
25#include <sys/stat.h>
26#endif
27 
28#include <capnp/dynamic.h>
29#include <capnp/message.h>
30#include <capnp/schema-parser.h>
31#include <capnp/serialize.h>
32#include <kj/async-queue.h>
33#include <kj/encoding.h>
34#include <kj/filesystem.h>
35#include <kj/main.h>
36#include <kj/map.h>
37 
38#if _WIN32
39#include <windows.h>
40#include <winsock2.h>
41 
42#include <kj/async-win32.h>
43#include <kj/win32-api-version.h>
44#include <kj/windows-sanity.h>
45 
46#include <iostream>
47#else
48#include <sys/ioctl.h>
49#include <sys/socket.h>
50#include <sys/syscall.h>
51#include <unistd.h>
52 
53#include <kj/async-unix.h>
54#endif
55 
56#if __linux__
57#include <sys/inotify.h>
58#elif __APPLE__ || __FreeBSD__ || __OpenBSD__ || __NetBSD__ || __DragonFly__
59#define WORKERD_USE_KQUEUE_FOR_FILE_WATCHER 1
60#include <sys/event.h>
61#include <sys/time.h>
62#include <sys/types.h>
63#endif
64 
65#ifdef __GLIBC__
66#include <sys/auxv.h>
67#endif
68 
69#ifdef __APPLE__
70#include <crt_externs.h>
71#include <libproc.h>
72#define environ (*_NSGetEnviron())
73#endif
74 
75#include <workerd/util/use-perfetto-categories.h>
76 
77// since kj installs their global signal handlers
78// and exits with 1 Fuzzilli doesn't realize that an application crashed due to the signo.
79// Therefore, we install a handler before and just raise the signo
80#ifdef WORKERD_FUZZILLI
81 
82void signalHandler(int signo, siginfo_t* info, void* context) noexcept {
83 // inform reprl - remove debug output for clean testing
84 struct sigaction sa = {};
85 sa.sa_handler = SIG_DFL;
86 sigemptyset(&sa.sa_mask);
87 sa.sa_flags = 0;
88 sigaction(signo, &sa, nullptr);
89 raise(signo);
90}
91 
92void initSignalHandlers() {
93 struct sigaction action {};
94 action.sa_flags = SA_SIGINFO;
95 action.sa_sigaction = &signalHandler;
96 
97 for (auto signo: {SIGBUS, SIGFPE, SIGABRT, SIGILL, SIGTRAP, SIGSEGV}) {
98 KJ_SYSCALL(sigaction(signo, &action, nullptr));
99 }
100}
101#endif
102 
103namespace workerd::server {
104namespace {
105 
106static kj::StringPtr getVersionString() {
107 static const kj::String result = kj::str("workerd ", RELEASE_VERSION);
108 return result;
109}
110 
111// =======================================================================================
112 
113// For ASan's leak sanitizer, suppress warnings about leaks with stacks that include "unknown
114// modules". This suppression is adopted from the GN build and applies to addresses that LSan can't
115// symbolize or even map to a binary – perhaps JIT or snapshot-generated code in V8's case?
116// TODO(someday): Suppression is needed to get several python tests to pass under LSan. Investigate
117// if this is an actual leak (perhaps a bug in V8 itself since it is suppressed there?) at a later
118// time.
119#if __has_feature(address_sanitizer)
120extern "C" __attribute__((no_sanitize("address"))) __attribute__((visibility("default")))
121__attribute__((used)) const char*
122__lsan_default_suppressions() {
123 return "leak:<unknown module>\n";
124}
125#endif
126 
127// =======================================================================================
128 
129class EntropySourceImpl: public kj::EntropySource {
130 public:
131 void generate(kj::ArrayPtr<kj::byte> buffer) override {
132 getEntropy(buffer);
133 }
134};
135 
136// =======================================================================================
137// Some generic CLI helpers so that we can throw exceptions rather than return
138// kj::MainBuilder::Validity. Honestly I do not know how people put up with patterns like
139// Result<T, E>, it seems like such a slog.
140 
141class CliError {
142 public:
143 CliError(kj::String description): description(kj::mv(description)) {}
144 kj::String description;
145};
146 
147template <typename Func>
148auto cliMethod(Func&& func) {
149 return [func = kj::fwd<Func>(func)](auto&&... params) mutable -> kj::MainBuilder::Validity {
150 try {
151 func(kj::fwd<decltype(params)>(params)...);
152 return true;
153 } catch (CliError& e) {
154 return kj::mv(e.description);
155 }
156 };
157}
158 
159// Pass to MainBuilder when a function returning kj::MainBuilder::Validity is needed, implemented
160// by a method of this class.
161#define CLI_METHOD(name) cliMethod(KJ_BIND_METHOD(*this, name))
162 
163// Throws an exception that is caught and reported as a usage error.
164#define CLI_ERROR(...) throw CliError(kj::str(__VA_ARGS__))
165 
166constexpr capnp::ReaderOptions CONFIG_READER_OPTIONS = {
167 .traversalLimitInWords = kj::maxValue
168 // Configs can legitimately be very large and are not malicious, so use an effectively-infinite
169 // traversal limit.
170};
171 
172// =======================================================================================
173 
174#if __linux__
175 
176// Class which uses inotify to watch a set of files and alert when they change.
177class FileWatcher {
178 public:
179 FileWatcher(kj::UnixEventPort& port)
180 : inotifyFd(makeInotify()),
181 observer(port, inotifyFd, kj::UnixEventPort::FdObserver::OBSERVE_READ) {}
182 
183 bool isSupported() {
184 return true;
185 }
186 
187 void watch(kj::PathPtr path, kj::Maybe<const kj::ReadableFile&> file) {
188 // `file` is provided if available. The Linux implementation doesn't use it.
189 
190 auto pathStr = path.parent().toNativeString(true);
191 
192 int wd = watches.findOrCreate(pathStr, [&]() {
193 int wd;
194 uint32_t mask = IN_DELETE | IN_MODIFY | IN_MOVE | IN_CREATE;
195 KJ_SYSCALL(wd = inotify_add_watch(inotifyFd, pathStr.cStr(), mask));
196 return decltype(watches)::Entry{kj::mv(pathStr), wd};
197 });
198 
199 auto& files =
200 filesWatched.findOrCreate(wd, [&]() { return decltype(filesWatched)::Entry{wd, {}}; });
201 
202 files.upsert(kj::str(path.basename()[0]), [](auto&&...) {});
203 }
204 
205 kj::Promise<void> onChange() {
206 kj::byte buffer[4096]{};
207 
208 for (;;) {
209 ssize_t n;
210 KJ_NONBLOCKING_SYSCALL(n = read(inotifyFd, buffer, sizeof(buffer)));
211 
212 if (n < 0) {
213 // No more data to read.
214 co_await observer.whenBecomesReadable();
215 continue;
216 }
217 
218 kj::byte* ptr = buffer;
219 while (n > 0) {
220 KJ_ASSERT(n >= sizeof(struct inotify_event));
221 
222 auto& event = *reinterpret_cast<struct inotify_event*>(ptr);
223 size_t eventSize = sizeof(struct inotify_event) + event.len;
224 KJ_ASSERT(n >= eventSize);
225 KJ_ASSERT(eventSize % sizeof(void*) == 0);
226 ptr += eventSize;
227 n -= eventSize;
228 
229 if (event.len > 0 && event.name[0] != '\0') {
230 auto& watched = KJ_ASSERT_NONNULL(filesWatched.find(event.wd));
231 if (watched.find(kj::StringPtr(event.name)) != kj::none) {
232 // HIT! We saw a change.
233 co_return;
234 }
235 }
236 }
237 }
238 }
239 
240 private:
241 kj::OwnFd inotifyFd;
242 kj::UnixEventPort::FdObserver observer;
243 
244 kj::HashMap<kj::String, int> watches;
245 kj::HashMap<int, kj::HashSet<kj::String>> filesWatched;
246 
247 static kj::OwnFd makeInotify() {
248 return KJ_SYSCALL_FD(inotify_init1(IN_NONBLOCK | IN_CLOEXEC));
249 }
250};
251 
252#elif WORKERD_USE_KQUEUE_FOR_FILE_WATCHER
253 
254// Class which uses inotify to watch a set of files and alert when they change.
255//
256// This version uses kqueue to watch for changes in files. kqueue typically doesn't scale well
257// to watching whole directory trees, since it must keep a file descriptor open for each watched
258// file. However, for our use case, we don't really want to watch a directory tree anyway, we
259// want to watch the specific set of files which were opened while parsing the config. This is
260// not so bad, probably.
261//
262// Apple provides the FSEvents API as an alternative, but it seems way more complicated and I
263// can't tell if it would provide a real advantage. Plus, kqueue works on BSD systems.
264class FileWatcher {
265 public:
266 FileWatcher(kj::UnixEventPort& port)
267 : kqueueFd(makeKqueue()),
268 observer(port, kqueueFd, kj::UnixEventPort::FdObserver::OBSERVE_READ) {}
269 
270 bool isSupported() {
271 return true;
272 }
273 
274 void watch(kj::PathPtr path, kj::Maybe<const kj::ReadableFile&> file) {
275 KJ_IF_SOME(f, file) {
276 KJ_IF_SOME(fd, f.getFd()) {
277 // We need to duplicate the FD because the original will probably be closed later and
278 // closing the FD unregisters it from kqueue.
279 watchFd(KJ_SYSCALL_FD(dup(fd)));
280 return;
281 }
282 }
283 
284 // No existing file, open from disk.
285 watchFd(KJ_SYSCALL_FD(open(path.toNativeString(true).cStr(), O_RDONLY)));
286 }
287 
288 kj::Promise<void> onChange() {
289 for (;;) {
290 struct kevent event;
291 struct timespec timeout;
292 memset(&event, 0, sizeof(event));
293 memset(&timeout, 0, sizeof(timeout));
294 
295 int n;
296 KJ_SYSCALL(n = kevent(kqueueFd, nullptr, 0, &event, 1, &timeout));
297 
298 if (n == 0) {
299 // No events, wait for the kqueue to become readable indicating an event has been
300 // delivered.
301 co_await observer.whenBecomesReadable();
302 continue;
303 } else {
304 // We only pay attention to events that indicate changes in the first place, so there's
305 // no need to examine the event, it definitely means something changed.
306 co_return;
307 }
308 }
309 }
310 
311 private:
312 kj::OwnFd kqueueFd;
313 kj::UnixEventPort::FdObserver observer;
314 kj::Vector<kj::OwnFd> filesWatched;
315 
316 static kj::OwnFd makeKqueue() {
317 auto fd = KJ_SYSCALL_FD(kqueue());
318 KJ_SYSCALL(fcntl(fd, F_SETFD, FD_CLOEXEC));
319 return kj::mv(fd);
320 }
321 
322 void watchFd(kj::OwnFd fd) {
323 KJ_SYSCALL(fcntl(fd, F_SETFD, FD_CLOEXEC));
324 
325 struct kevent change;
326 memset(&change, 0, sizeof(change));
327 change.ident = fd.get();
328 change.filter = EVFILT_VNODE;
329 change.flags = EV_ADD | EV_CLEAR;
330 change.fflags = NOTE_WRITE | NOTE_EXTEND | NOTE_DELETE | NOTE_RENAME;
331 KJ_SYSCALL(kevent(kqueueFd, &change, 1, nullptr, 0, nullptr));
332 filesWatched.add(kj::mv(fd));
333 }
334};
335 
336#elif _WIN32
337 
338class FileWatcher {
339 public:
340 FileWatcher(kj::Win32EventPort& port) {}
341 
342 bool isSupported() {
343 return false;
344 }
345 
346 void watch(kj::PathPtr path, kj::Maybe<const kj::ReadableFile&> file) {}
347 
348 kj::Promise<void> onChange() {
349 return kj::NEVER_DONE;
350 }
351 
352 private:
353};
354 
355#else
356 
357// Dummy FileWatcher implementation for operating systems that aren't supported yet.
358class FileWatcher {
359 public:
360 FileWatcher(kj::UnixEventPort& port) {}
361 
362 bool isSupported() {
363 return false;
364 }
365 
366 void watch(kj::PathPtr path, kj::Maybe<const kj::ReadableFile&> file) {}
367 
368 kj::Promise<void> onChange() {
369 return kj::NEVER_DONE;
370 }
371 
372 private:
373};
374 
375#endif // #__linux__, #else
376 
377// =======================================================================================
378 
379kj::Maybe<kj::Own<capnp::SchemaFile>> tryImportBulitin(kj::StringPtr name);
380 
381// Callbacks for capnp::SchemaFileLoader. Implementing this interface lets us control import
382// resolution, which we want to do mainly so that we can set watches on all imported files.
383//
384// These callbacks also give us more control over error reporting, in particular the ability
385// to not throw an exception on the first error seen.
386class SchemaFileImpl final: public capnp::SchemaFile {
387 public:
388 class ErrorReporter {
389 public:
390 virtual void reportParsingError(
391 kj::StringPtr file, SourcePos start, SourcePos end, kj::StringPtr message) = 0;
392 };
393 
394 SchemaFileImpl(const kj::Directory& root,
395 kj::PathPtr current,
396 kj::Path fullPathParam,
397 kj::PathPtr basePath,
398 kj::ArrayPtr<const kj::Path> importPath,
399 kj::Own<const kj::ReadableFile> fileParam,
400 kj::Maybe<FileWatcher&> watcher,
401 ErrorReporter& errorReporter)
402 : root(root),
403 current(current),
404 fullPath(kj::mv(fullPathParam)),
405 basePath(basePath),
406 importPath(importPath),
407 file(kj::mv(fileParam)),
408 watcher(watcher),
409 errorReporter(errorReporter) {
410 if (fullPath.startsWith(current)) {
411 // Simplify display name by removing current directory prefix.
412 displayName = fullPath.slice(current.size(), fullPath.size()).toNativeString();
413 } else {
414 // Use full path.
415 displayName = fullPath.toNativeString(true);
416 }
417 
418 KJ_IF_SOME(w, watcher) {
419 w.watch(fullPath, *file);
420 }
421 }
422 
423 kj::StringPtr getDisplayName() const override {
424 return displayName;
425 }
426 
427 kj::Array<const char> readContent() const override {
428 uint64_t size = file->stat().size;
429 if (!size) {
430 return nullptr;
431 }
432 return file->mmap(0, file->stat().size).releaseAsChars();
433 }
434 
435 kj::Maybe<kj::Own<SchemaFile>> import(kj::StringPtr target) const override {
436 if (target.startsWith("/")) {
437 auto parsedPath = kj::Path::parse(target.slice(1));
438 for (auto& candidate: importPath) {
439 auto newFullPath = candidate.append(parsedPath);
440 
441 KJ_IF_SOME(newFile, root.tryOpenFile(newFullPath)) {
442 return kj::implicitCast<kj::Own<SchemaFile>>(kj::heap<SchemaFileImpl>(root, current,
443 kj::mv(newFullPath), candidate, importPath, kj::mv(newFile), watcher, errorReporter));
444 }
445 }
446 // No matching file found. Check if we have a builtin.
447 return tryImportBulitin(target);
448 } else {
449 auto relativeTo = fullPath.slice(basePath.size(), fullPath.size());
450 auto parsed = relativeTo.parent().eval(target);
451 auto newFullPath = basePath.append(parsed);
452 
453 KJ_IF_SOME(newFile, root.tryOpenFile(newFullPath)) {
454 return kj::implicitCast<kj::Own<SchemaFile>>(kj::heap<SchemaFileImpl>(root, current,
455 kj::mv(newFullPath), basePath, importPath, kj::mv(newFile), watcher, errorReporter));
456 } else {
457 return kj::none;
458 }
459 }
460 }
461 
462 bool operator==(const SchemaFile& other) const override {
463 if (auto downcasted = dynamic_cast<const SchemaFileImpl*>(&other)) {
464 return fullPath == downcasted->fullPath;
465 } else {
466 return false;
467 }
468 }
469 
470 size_t hashCode() const override {
471 return kj::hashCode(fullPath);
472 }
473 
474 void reportError(SourcePos start, SourcePos end, kj::StringPtr message) const override {
475 errorReporter.reportParsingError(displayName, start, end, message);
476 }
477 
478 private:
479 const kj::Directory& root;
480 kj::PathPtr current;
481 
482 // Full path from root of filesystem to the file.
483 kj::Path fullPath;
484 
485 // If this file was reached by scanning `importPath`, `basePath` is the particular import path
486 // directory that was used, otherwise it is empty. `basePath` is always a prefix of `fullPath`.
487 kj::PathPtr basePath;
488 
489 // Paths to search for absolute imports.
490 kj::ArrayPtr<const kj::Path> importPath;
491 
492 kj::Own<const kj::ReadableFile> file;
493 kj::String displayName;
494 
495 // Mutable because the SchemaParser interface forces us to make all our methods `const` so that
496 // parsing can happen on multiple threads, but we do not actually use multiple threads for
497 // parsing, so we're good.
498 mutable kj::Maybe<FileWatcher&> watcher;
499 
500 ErrorReporter& errorReporter;
501};
502 
503// A schema file whose text is embedded into the binary for convenience.
504//
505// TODO(someday): Could `capnp::SchemaParser` be updated such that it can use the compiled-in
506// schema nodes rather than re-parse the file from scratch? This is tricky as some information
507// is lost after compilation which is needed to compile dependents, e.g. aliases are erased.
508class BuiltinSchemaFileImpl final: public capnp::SchemaFile {
509 public:
510 BuiltinSchemaFileImpl(kj::StringPtr name, kj::StringPtr content): name(name), content(content) {}
511 
512 kj::StringPtr getDisplayName() const override {
513 return name;
514 }
515 
516 kj::Array<const char> readContent() const override {
517 return kj::Array<const char>(content.begin(), content.size(), kj::NullArrayDisposer::instance);
518 }
519 
520 kj::Maybe<kj::Own<SchemaFile>> import(kj::StringPtr target) const override {
521 return tryImportBulitin(target);
522 }
523 
524 bool operator==(const SchemaFile& other) const override {
525 if (auto downcasted = dynamic_cast<const BuiltinSchemaFileImpl*>(&other)) {
526 return downcasted->name == name;
527 } else {
528 return false;
529 }
530 }
531 
532 size_t hashCode() const override {
533 return kj::hashCode(name);
534 }
535 
536 void reportError(SourcePos start, SourcePos end, kj::StringPtr message) const override {
537 KJ_FAIL_ASSERT("parse error in built-in schema?", start.line, start.column, message);
538 }
539 
540 private:
541 kj::StringPtr name;
542 kj::StringPtr content;
543};
544 
545kj::Maybe<kj::Own<capnp::SchemaFile>> tryImportBulitin(kj::StringPtr name) {
546 if (name == "/capnp/c++.capnp") {
547 return kj::heap<BuiltinSchemaFileImpl>("/capnp/c++.capnp", CPP_CAPNP_SCHEMA);
548 } else if (name == "/workerd/workerd.capnp") {
549 return kj::heap<BuiltinSchemaFileImpl>("/workerd/workerd.capnp", WORKERD_CAPNP_SCHEMA);
550 } else {
551 return kj::none;
552 }
553}
554 
555// =======================================================================================
556 
557// A kj::Network implementation which wraps some other network and optionally (if enabled)
558// implements "loopback:" network addresses, which are expected to be serviced within the same
559// process. Loopback addresses are enabled only when running `workerd test`. The purpose is to
560// allow end-to-end testing of the network stack without creating a real external-facing socket.
561//
562// There is no use for loopback sockets in production since direct service bindings are more
563// efficient while solving the same problems.
564class NetworkWithLoopback final: public kj::Network {
565 public:
566 NetworkWithLoopback(kj::Network& inner, kj::AsyncIoProvider& ioProvider)
567 : inner(inner),
568 ioProvider(ioProvider),
569 loopbackEnabled(rootLoopbackEnabled) {}
570 
571 NetworkWithLoopback(
572 kj::Own<kj::Network> inner, kj::AsyncIoProvider& ioProvider, bool& loopbackEnabled)
573 : inner(*inner),
574 ownInner(kj::mv(inner)),
575 ioProvider(ioProvider),
576 loopbackEnabled(loopbackEnabled) {}
577 
578 // Call once to enable loopback addresses.
579 void enableLoopback() {
580 loopbackEnabled = true;
581 }
582 
583 kj::Promise<kj::Own<kj::NetworkAddress>> parseAddress(
584 kj::StringPtr addr, uint portHint = 0) override {
585 if (loopbackEnabled && addr.startsWith(PREFIX)) {
586 return kj::Own<kj::NetworkAddress>(kj::heap<LoopbackAddr>(*this, addr.slice(PREFIX.size())));
587 } else {
588 return inner.parseAddress(addr, portHint);
589 }
590 }
591 
592 kj::Own<kj::NetworkAddress> getSockaddr(const void* sockaddr, uint len) override {
593 return inner.getSockaddr(sockaddr, len);
594 }
595 
596 kj::Own<kj::Network> restrictPeers(kj::ArrayPtr<const kj::StringPtr> allow,
597 kj::ArrayPtr<const kj::StringPtr> deny = nullptr) override {
598 return kj::heap<NetworkWithLoopback>(
599 inner.restrictPeers(allow, deny), ioProvider, loopbackEnabled);
600 }
601 
602 private:
603 kj::Network& inner;
604 kj::Own<kj::Network> ownInner;
605 kj::AsyncIoProvider& ioProvider KJ_UNUSED;
606 bool rootLoopbackEnabled = false;
607 
608 // Reference to `rootLoopbackEnabled` of the root NetworkWithLoopback. All descendants
609 // (created using `restrictPeers()` will share the same flag value.
610 bool& loopbackEnabled;
611 
612 using ConnectionQueue = kj::ProducerConsumerQueue<kj::Own<kj::AsyncIoStream>>;
613 kj::HashMap<kj::String, kj::Own<ConnectionQueue>> loopbackQueues;
614 
615 ConnectionQueue& getLoopbackQueue(kj::StringPtr name) {
616 return *loopbackQueues.findOrCreate(name, [&]() {
617 return decltype(loopbackQueues)::Entry{
618 .key = kj::str(name),
619 .value = kj::heap<ConnectionQueue>(),
620 };
621 });
622 }
623 
624 static constexpr kj::StringPtr PREFIX = "loopback:"_kj;
625 
626 class LoopbackAddr final: public kj::NetworkAddress {
627 public:
628 LoopbackAddr(NetworkWithLoopback& parent, kj::StringPtr name)
629 : parent(parent),
630 name(kj::str(name)) {}
631 
632 kj::Promise<kj::Own<kj::AsyncIoStream>> connect() override {
633 // The purpose of loopback sockets is to actually test the network stack end-to-end. If
634 // people don't want to test the full stack, then they can create a direct service binding
635 // without going through a loopback socket.
636 //
637 // So, we create a real loopback socket here.
638 auto pipe = parent.ioProvider.newTwoWayPipe();
639 
640 parent.getLoopbackQueue(name).push(kj::mv(pipe.ends[0]));
641 return kj::mv(pipe.ends[1]);
642 }
643 
644 kj::Own<kj::ConnectionReceiver> listen() override {
645 return kj::heap<LoopbackReceiver>(parent.getLoopbackQueue(name));
646 }
647 
648 kj::Own<kj::NetworkAddress> clone() override {
649 return kj::heap<LoopbackAddr>(parent, name);
650 }
651 
652 kj::String toString() override {
653 return kj::str(PREFIX, name);
654 }
655 
656 private:
657 NetworkWithLoopback& parent;
658 kj::String name;
659 };
660 
661 class LoopbackReceiver final: public kj::ConnectionReceiver {
662 public:
663 LoopbackReceiver(ConnectionQueue& queue): queue(queue) {}
664 
665 kj::Promise<kj::Own<kj::AsyncIoStream>> accept() override {
666 return queue.pop();
667 }
668 
669 uint getPort() override {
670 return 0;
671 }
672 
673 private:
674 ConnectionQueue& queue;
675 };
676};
677 
678// =======================================================================================
679 
680class CliMain final: public SchemaFileImpl::ErrorReporter {
681 public:
682 CliMain(StructuredLoggingProcessContext& context, char** argv)
683 : context(context),
684 argv(argv),
685 server(kj::heap<Server>(*fs,
686 io.provider->getTimer(),
687 kj::systemPreciseMonotonicClock(),
688 network,
689 entropySource,
690 Worker::LoggingOptions(Worker::ConsoleMode::STDOUT),
691 [&](kj::String error) {
692 if (watcher == kj::none) {
693 // TODO(someday): Don't just fail on the first error, keep going in order to report
694 // additional errors. The tricky part is we don't currently have any signal of when
695 // the server has completely finished loading, and also we probably don't want to
696 // accept any connections on any of the sockets if the server is partially broken.
697 context.exitError(error);
698 } else {
699 // In --watch mode, we don't want to exit from errors, we want to wait until things
700 // change. It's OK if we try to serve requests despite brokenness since this is a
701 // development server.
702 hadErrors = true;
703 context.error(error);
704 }
705 })) {
706 KJ_IF_SOME(e, exeInfo) {
707 auto& exe = *e.file;
708 auto size = exe.stat().size;
709 KJ_ASSERT(size > sizeof(COMPILED_MAGIC_SUFFIX) + sizeof(uint64_t));
710 kj::byte magic[sizeof(COMPILED_MAGIC_SUFFIX)]{};
711 exe.read(size - sizeof(COMPILED_MAGIC_SUFFIX), magic);
712 if (kj::arrayPtr(magic) == kj::asBytes(COMPILED_MAGIC_SUFFIX)) {
713 // Oh! It appears we are running a compiled binary, it has a config appended to the end.
714 uint64_t configSize;
715 exe.read(size - sizeof(COMPILED_MAGIC_SUFFIX) - sizeof(uint64_t), kj::asBytes(configSize));
716 KJ_ASSERT(size - sizeof(COMPILED_MAGIC_SUFFIX) - sizeof(uint64_t) >
717 configSize * sizeof(capnp::word));
718 size_t offset = size - sizeof(COMPILED_MAGIC_SUFFIX) - sizeof(uint64_t) -
719 configSize * sizeof(capnp::word);
720 
721 auto mapping = exe.mmap(offset, configSize * sizeof(capnp::word));
722 KJ_ASSERT(reinterpret_cast<uintptr_t>(mapping.begin()) % sizeof(capnp::word) == 0,
723 "compiled-in config is not aligned correctly?");
724 
725 config = capnp::readMessageUnchecked<config::Config>(
726 reinterpret_cast<const capnp::word*>(mapping.begin()));
727 configOwner = kj::heap(kj::mv(mapping));
728 }
729 } else {
730 context.warning(
731 "Unable to find and open the program executable, so unable to determine if there is a "
732 "compiled-in config file. Proceeding on the assumption that there is not.");
733 }
734 
735 // We don't want to force people to specify top-level file IDs in `workerd` config files, as
736 // those IDs would be totally irrelevant.
737 schemaParser.setFileIdsRequired(false);
738 }
739 
740 kj::MainFunc getMain() {
741 if (config == kj::none) {
742 return kj::MainBuilder(
743 context, getVersionString(), "Runs the Workers JavaScript/Wasm runtime.")
744 .addSubCommand("serve", KJ_BIND_METHOD(*this, getServe), "run the server")
745 .addSubCommand(
746 "compile", KJ_BIND_METHOD(*this, getCompile), "create a self-contained binary")
747#ifdef WORKERD_FUZZILLI
748 .addSubCommand("fuzzilli", KJ_BIND_METHOD(*this, getFuzz), "run reprl for fuzzing")
749#endif
750 .addSubCommand("test", KJ_BIND_METHOD(*this, getTest), "run unit tests")
751 .addSubCommand("pyodide-lock", KJ_BIND_METHOD(*this, getPyodideLock),
752 "outputs the package lock file used by Pyodide")
753 .addSubCommand("make-pyodide-baseline-snapshot",
754 KJ_BIND_METHOD(*this, getMakePyodideBaselineSnapshot),
755 "Make a Pyodide baseline memory snapshot")
756 .build();
757 // TODO(someday):
758 // "validate": Loads the config and parses all the code to report errors, but then exits
759 // without serving anything.
760 // "explain": Produces human-friendly description of the config.
761 } else {
762 // We already have a config, meaning this must be a compiled binary.
763 auto builder = kj::MainBuilder(context, getVersionString(),
764 "Serve requests based on the compiled config.",
765 "This binary has an embedded configuration.");
766 return addServeOptions(builder);
767 }
768 }
769 
770 kj::MainBuilder& addConfigParsingOptionsNoConstName(kj::MainBuilder& builder) {
771 return builder
772 .addOptionWithArg({'I', "import-path"}, CLI_METHOD(addImportPath), "<dir>",
773 "Add <dir> to the list of directories searched for non-relative "
774 "imports in the config file (ones that start with a '/').")
775 .addOption({'b', "binary"},
776 [this]() {
777 binaryConfig = true;
778 return true;
779 },
780 "Specifies that the configuration file is an encoded binary Cap'n Proto "
781 "message, rather than the usual text format. This is particularly useful when "
782 "driving the server from higher-level tooling that automatically generates a "
783 "config.")
784 .expectArg("<config-file>", CLI_METHOD(parseConfigFile));
785 }
786 
787 kj::MainBuilder& addConfigParsingOptions(kj::MainBuilder& builder) {
788 return addConfigParsingOptionsNoConstName(builder).expectOptionalArg(
789 "<const-name>", CLI_METHOD(setConstName));
790 }
791 
792 kj::MainBuilder& addServeOrTestOptions(kj::MainBuilder& builder) {
793 return builder
794 .addOptionWithArg({'d', "directory-path"}, CLI_METHOD(overrideDirectory), "<name>=<path>",
795 "Override the directory named <name> to point to <path> instead of the "
796 "path specified in the config file.")
797 .addOptionWithArg({'e', "external-addr"}, CLI_METHOD(overrideExternal), "<name>=<addr>",
798 "Override the external service named <name> to connect to the address "
799 "<addr> instead of the address specified in the config file.")
800 .addOptionWithArg({'i', "inspector-addr"}, CLI_METHOD(enableInspector), "<addr>",
801 "Enable the inspector protocol to connect to the address <addr>.")
802#ifdef WORKERD_USE_PERFETTO
803 // TODO(later): In the future, we might want to enable providing a perfetto
804 // TraceConfig structure here rather than just the categories.
805 .addOptionWithArg({"p", "perfetto-trace"}, CLI_METHOD(enablePerfetto),
806 "<path>=<categories>", "Enable perfetto tracing output to the specified file.")
807#endif
808 .addOption({'w', "watch"}, CLI_METHOD(watch),
809 "Watch configuration files (and server binary) and reload if they change. "
810 "Useful for development, but not recommended in production.")
811 .addOption({"experimental"},
812 [this]() {
813 server->allowExperimental();
814 return true;
815 },
816 "Permit the use of experimental features which may break backwards "
817 "compatibility in a future release.")
818 .addOptionWithArg({"pyodide-package-disk-cache-dir"}, CLI_METHOD(setPackageDiskCacheDir),
819 "<path>",
820 "Use <path> as a disk cache to avoid repeatedly fetching packages from the internet. ")
821 .addOptionWithArg({"pyodide-bundle-disk-cache-dir"}, CLI_METHOD(setPyodideDiskCacheDir),
822 "<path>",
823 "Use <path> as a disk cache to avoid repeatedly fetching Pyodide bundles from the internet. ")
824 .addOption({"python-save-snapshot"},
825 [this]() {
826 server->setPythonCreateSnapshot();
827 return true;
828 }, "Save a dedicated snapshot to the disk cache")
829 .addOption({"python-save-baseline-snapshot"},
830 [this]() {
831 server->setPythonCreateBaselineSnapshot();
832 return true;
833 }, "Save a baseline snapshot to the disk cache")
834 .addOptionWithArg({"python-load-snapshot"}, CLI_METHOD(setPythonLoadSnapshot), "<path>",
835 "Load a snapshot from the python snapshot directory.")
836 .addOptionWithArg({"python-snapshot-dir"}, CLI_METHOD(setPythonSnapshotDirectory), "<path>",
837 "Set the snapshot snapshot directory.");
838 }
839 
840 kj::MainFunc addServeOptions(kj::MainBuilder& builder) {
841 return addServeOrTestOptions(builder)
842 .addOptionWithArg({'s', "socket-addr"}, CLI_METHOD(overrideSocketAddr), "<name>=<addr>",
843 "Override the socket named <name> to bind to the address <addr> instead "
844 "of the address specified in the config file.")
845 .addOptionWithArg({'S', "socket-fd"}, CLI_METHOD(overrideSocketFd), "<name>=<fd>",
846 "Override the socket named <name> to listen on the already-open socket "
847 "descriptor <fd> instead of the address specified in the config file.")
848 .addOptionWithArg({"control-fd"}, CLI_METHOD(enableControl), "<fd>",
849 "Enable sending of control messages on descriptor <fd>. Currently this "
850 "only reports the port each socket is listening on when ready.")
851 .addOptionWithArg({"debug-port"}, CLI_METHOD(enableDebugPort), "<addr>",
852 "Listen on the specified address for debug RPC connections. This exposes "
853 "a privileged interface that allows access to all services in the process. "
854 "For use by miniflare and local development only.")
855 .callAfterParsing(CLI_METHOD(serve))
856 .build();
857 }
858 
859 kj::MainFunc getServe() {
860 auto builder = kj::MainBuilder(context, getVersionString(), "Serve requests based on a config.",
861 "Serves requests based on the configuration specified in <config-file>.");
862 return addServeOptions(addConfigParsingOptions(builder));
863 }
864 
865 kj::MainFunc getPyodideLock() {
866 auto builder = kj::MainBuilder(
867 context, getVersionString(), "Outputs the package lock file used by Pyodide.");
868 return builder
869 .callAfterParsing([]() -> kj::MainBuilder::Validity {
870 static const PythonConfig config{
871 .packageDiskCacheRoot = kj::none,
872 .pyodideDiskCacheRoot = kj::none,
873 .createSnapshot = false,
874 .createBaselineSnapshot = false,
875 };
876 
877 capnp::MallocMessageBuilder message;
878 // TODO(EW-8977): Implement option to specify python worker flags.
879 auto features = message.getRoot<CompatibilityFlags>();
880 features.setPythonWorkers(true);
881 auto pythonRelease = KJ_ASSERT_NONNULL(getPythonSnapshotRelease(features));
882 
883 auto lock = KJ_ASSERT_NONNULL(api::pyodide::getPyodideLock(pythonRelease));
884 
885 printf("%s\n", lock.cStr());
886 fflush(stdout);
887 return true;
888 }).build();
889 }
890 
891 kj::MainFunc getTest() {
892 auto builder = kj::MainBuilder(context, getVersionString(), "Runs tests based on a config.",
893 "Runs tests for services defined in <config-file>. <filter>, if given, specifies "
894 "exactly which tests to run. It has one of the following formats:\n"
895 " <service-pattern>\n"
896 " <service-pattern>:<entrypoint-pattern>\n"
897 " <const-name>:<service-pattern>:<entrypoint-pattern>\n"
898 "<service-pattern> is a glob pattern matching names of services which should be tested. "
899 "If not specified, '*' is assumed (which matches all services). <entrypoint-pattern> "
900 "is a glob pattern matching entrypoints within each service which should be tested; "
901 "again, the default is '*'. <const-name> has the same meaning as for the `serve` "
902 "command (this is rarely used).\n"
903 "\n"
904 "Tests can be defined by exporting a function called `test` instead of (or in addition "
905 "to) `fetch`. Example:\n"
906 " export default {\n"
907 " async test(ctrl, env, ctx) {\n"
908 " if (1 + 1 != 2) {\n"
909 " throw new Error('math is broken!');\n"
910 " }\n"
911 " }\n"
912 " }\n"
913 "The test passes if the test function completes without throwing. Multiple tests can "
914 "be exported under different entrypoint names:\n"
915 " export let test1 = {\n"
916 " async test(ctrl, env, ctx) {\n"
917 " ...\n"
918 " }\n"
919 " }\n"
920 " export let test2 = {\n"
921 " async test(ctrl, env, ctx) {\n"
922 " ...\n"
923 " }\n"
924 " }\n");
925 return addServeOrTestOptions(addConfigParsingOptionsNoConstName(builder))
926 .addOption({"no-verbose"},
927 [this]() {
928 noVerbose = true;
929 return true;
930 },
931 "Disable INFO-level logging for this test. Otherwise, INFO logging is enabled by "
932 "default for tests in order to show uncaught exceptions, but it can be noisey.")
933 .addOption({"predictable"},
934 [this]() {
935 predictable = true;
936 return true;
937 },
938 "Enable predictable mode. This makes workerd behave more deterministically by using "
939 "pre-set values instead of random data or timestamps to facilitate testing.")
940 .addOption({"gc-stress"},
941 [this]() {
942 gcStress = true;
943 return true;
944 },
945 "Force a full V8 GC at each awaitIo continuation. "
946 "Detects KJ async objects on the JS heap without IoOwn wrapping. Very slow.")
947 .addOption({"all-autogates"},
948 [this]() {
949 allAutogates = true;
950 return true;
951 },
952 "Enable all autogates. This is useful for testing code paths that are guarded by "
953 "autogates.")
954 .addOptionWithArg({"compat-date"}, CLI_METHOD(setTestCompatDate), "<date>",
955 "Set the compatibility date for all workers. When specified, workers must NOT "
956 "specify compatibilityDate in the config. Use '0000-00-00' for oldest behavior "
957 "or '9999-12-31' for newest behavior.")
958 .expectOptionalArg("<filter>", CLI_METHOD(setTestFilter))
959 .callAfterParsing(CLI_METHOD(test))
960 .build();
961 }
962 
963 kj::MainFunc getFuzz() {
964 auto builder = kj::MainBuilder(context, getVersionString(),
965 "Creates a custom signal handler and depending on the config leverages Stdin.reprl() to communicate with fuzzilli.");
966 
967 return addServeOrTestOptions(addConfigParsingOptionsNoConstName(builder))
968 .callAfterParsing(CLI_METHOD(test))
969 .build();
970 }
971 
972 kj::MainFunc getCompile() {
973 auto builder = kj::MainBuilder(context, getVersionString(),
974 "Builds a self-contained binary from a config.",
975 "This parses a config file in the same manner as the \"serve\" command, but instead "
976 "of then running it, it outputs a new binary to stdout that embeds the config and all "
977 "associated Worker code and data as one self-contained unit. This binary may then "
978 "be executed on another system to run the config -- without any other files being "
979 "present on that system.");
980 return addConfigParsingOptions(builder)
981 .addOption({"config-only"},
982 [this]() {
983 configOnly = true;
984 return true;
985 },
986 "Only write the encoded binary config to stdout. Do not attach it to an executable. "
987 "The encoded config can be used as input to the \"serve\" command, without the need "
988 "for any other files to be present.")
989 .callAfterParsing(CLI_METHOD(compile))
990 .build();
991 }
992 
993 kj::MainFunc getMakePyodideBaselineSnapshot() {
994 server->allowExperimental();
995 server->setPythonCreateBaselineSnapshot();
996 auto builder =
997 kj::MainBuilder(context, getVersionString(), "Make a Pyodide baseline memory snapshot", "");
998 setPyodideDiskCacheDir(".");
999 return builder.expectArg("<python-version>", CLI_METHOD(parsePythonCompatFlag))
1000 .expectArg("<output-directory>", CLI_METHOD(setPackageDiskCacheDir))
1001 .callAfterParsing(CLI_METHOD(test))
1002 .build();
1003 }
1004 
1005 void addImportPath(kj::StringPtr pathStr) {
1006 auto path = fs->getCurrentPath().evalNative(pathStr);
1007 if (fs->getRoot().tryOpenSubdir(path) != kj::none) {
1008 importPath.add(kj::mv(path));
1009 } else {
1010 CLI_ERROR("No such directory.");
1011 }
1012 }
1013 
1014 struct Override {
1015 kj::String name;
1016 kj::StringPtr value;
1017 };
1018 Override parseOverride(kj::StringPtr str) {
1019 auto equalPos = KJ_UNWRAP_OR(str.findFirst('='), CLI_ERROR("Expected <name>=<value>"));
1020 return {kj::str(str.first(equalPos)), str.slice(equalPos + 1)};
1021 }
1022 
1023 void overrideSocketAddr(kj::StringPtr param) {
1024 auto [name, value] = parseOverride(param);
1025 server->overrideSocket(kj::mv(name), kj::str(value));
1026 }
1027 
1028#if _WIN32
1029 void validateSocketFd(uint fd, kj::StringPtr label) {
1030 int acceptcon = 0;
1031 int optlen = sizeof(acceptcon);
1032 int result = getsockopt(fd, SOL_SOCKET, SO_ACCEPTCONN, (char*)&acceptcon, &optlen);
1033 if (result == SOCKET_ERROR) {
1034 // https://learn.microsoft.com/en-us/windows/win32/api/winsock/nf-winsock-getsockopt#return-value
1035 switch (int error = WSAGetLastError()) {
1036 case WSAENOTSOCK:
1037 CLI_ERROR("File descriptor is not a socket.");
1038 case WSAENOPROTOOPT:
1039 // Some operating systems don't support SO_ACCEPTCONN; in that case just move on and
1040 // assume it is listening.
1041 break;
1042 default:
1043 KJ_FAIL_SYSCALL("getsockopt(fd, SOL_SOCKET, SO_ACCEPTCONN)", error);
1044 }
1045 } else if (!acceptcon) {
1046 CLI_ERROR("Socket for ", label, " is not listening.");
1047 }
1048 }
1049#else
1050 void validateSocketFd(uint fd, kj::StringPtr label) {
1051 int acceptcon = 0;
1052 socklen_t optlen = sizeof(acceptcon);
1053 KJ_SYSCALL_HANDLE_ERRORS(getsockopt(fd, SOL_SOCKET, SO_ACCEPTCONN, &acceptcon, &optlen)) {
1054 case EBADF:
1055 CLI_ERROR("File descriptor is not open.");
1056 case ENOTSOCK:
1057 CLI_ERROR("File descriptor is not a socket.");
1058 case ENOPROTOOPT:
1059 // Some operating systems don't support SO_ACCEPTCONN; in that case just move on and
1060 // assume it is listening.
1061 break;
1062 default:
1063 KJ_FAIL_SYSCALL("getsockopt(fd, SOL_SOCKET, SO_ACCEPTCONN)", error);
1064 }
1065 else {
1066 if (!acceptcon) {
1067 CLI_ERROR("Socket for ", label, " is not listening.");
1068 }
1069 }
1070 }
1071#endif
1072 
1073 void overrideSocketFd(kj::StringPtr param) {
1074 auto [name, value] = parseOverride(param);
1075 
1076 int fd = KJ_UNWRAP_OR(value.tryParseAs<uint>(),
1077 CLI_ERROR("Socket value must be a file descriptor (non-negative integer)."));
1078 
1079 validateSocketFd(fd, name);
1080 
1081 inheritedFds.add(fd);
1082 server->overrideSocket(kj::mv(name),
1083 io.lowLevelProvider->wrapListenSocketFd(fd, kj::LowLevelAsyncIoProvider::TAKE_OWNERSHIP));
1084 }
1085 
1086 void overrideDirectory(kj::StringPtr param) {
1087 auto [name, value] = parseOverride(param);
1088 server->overrideDirectory(kj::mv(name), kj::str(value));
1089 }
1090 
1091 void overrideExternal(kj::StringPtr param) {
1092 auto [name, value] = parseOverride(param);
1093 server->overrideExternal(kj::mv(name), kj::str(value));
1094 }
1095 
1096#ifdef WORKERD_USE_PERFETTO
1097 void enablePerfetto(kj::StringPtr param) {
1098 auto [name, value] = parseOverride(param);
1099 perfettoTraceDestination = kj::str(name);
1100 perfettoTraceCategories = kj::str(value);
1101 }
1102#endif
1103 
1104 void enableInspector(kj::StringPtr param) {
1105 server->enableInspector(kj::str(param));
1106 }
1107 
1108 void enableControl(kj::StringPtr param) {
1109 int fd = KJ_UNWRAP_OR(param.tryParseAs<uint>(),
1110 CLI_ERROR("Output value must be a file descriptor (non-negative integer)."));
1111 server->enableControl(fd);
1112 }
1113 
1114 void enableDebugPort(kj::StringPtr param) {
1115 server->enableDebugPort(kj::str(param));
1116 }
1117 
1118 void setPackageDiskCacheDir(kj::StringPtr pathStr) {
1119 kj::Path path = fs->getCurrentPath().eval(pathStr);
1120 kj::Maybe<kj::Own<const kj::Directory>> dir =
1121 fs->getRoot().tryOpenSubdir(path, kj::WriteMode::MODIFY);
1122 server->setPackageDiskCacheRoot(
1123 kj::mv(KJ_UNWRAP_OR(dir, CLI_ERROR("package disk cache dir must exist"))));
1124 }
1125 
1126 void setPyodideDiskCacheDir(kj::StringPtr pathStr) {
1127 kj::Path path = fs->getCurrentPath().eval(pathStr);
1128 kj::Maybe<kj::Own<const kj::Directory>> dir =
1129 fs->getRoot().tryOpenSubdir(path, kj::WriteMode::MODIFY);
1130 server->setPyodideDiskCacheRoot(kj::mv(dir));
1131 }
1132 
1133 void setPythonLoadSnapshot(kj::StringPtr pathStr) {
1134 server->setPythonLoadSnapshot(kj::str(pathStr));
1135 }
1136 void setPythonSnapshotDirectory(kj::StringPtr pathStr) {
1137 kj::Path path = fs->getCurrentPath().eval(pathStr);
1138 kj::Maybe<kj::Own<const kj::Directory>> dir =
1139 fs->getRoot().tryOpenSubdir(path, kj::WriteMode::MODIFY);
1140 server->setPythonSnapshotDirectory(kj::mv(dir));
1141 }
1142 
1143 void parsePythonCompatFlag(kj::StringPtr compatFlagStr) {
1144 auto builder = kj::heap<capnp::MallocMessageBuilder>();
1145 auto configBuilder = builder->initRoot<config::Config>();
1146 auto service = configBuilder.initServices(1)[0];
1147 service.setName("main");
1148 auto worker = service.initWorker();
1149 worker.setCompatibilityDate("2023-12-18");
1150 auto flags = worker.initCompatibilityFlags(2);
1151 flags.set(0, compatFlagStr);
1152 flags.set(1, "python_workers");
1153 auto mod = worker.initModules(1)[0];
1154 mod.setName("main.py");
1155 mod.setPythonModule("def test():\n pass");
1156 config = configBuilder.asReader();
1157 configOwner = kj::mv(builder);
1158 util::Autogate::initAutogate(getConfig().getAutogates());
1159 }
1160 
1161 void watch() {
1162#if _WIN32
1163 auto& w = watcher.emplace(io.win32EventPort);
1164#else
1165 auto& w = watcher.emplace(io.unixEventPort);
1166#endif
1167 if (!w.isSupported()) {
1168 CLI_ERROR("File watching is not yet implemented on your OS. Sorry! Pull requests welcome!");
1169 }
1170 
1171 KJ_IF_SOME(e, exeInfo) {
1172 w.watch(fs->getCurrentPath().eval(e.path), kj::none);
1173 } else {
1174 CLI_ERROR("Can't use --watch when we're unable to find our own executable.");
1175 }
1176 }
1177 
1178 void parseConfigFile(kj::StringPtr pathStr) {
1179 if (pathStr == "-") {
1180 // Read from stdin.
1181 
1182 if (!binaryConfig) {
1183 CLI_ERROR("Reading config from stdin is only allowed with --binary.");
1184 }
1185 
1186 // Can't use mmap() because it's probably not a file.
1187#if _WIN32
1188 auto handle = GetStdHandle(STD_INPUT_HANDLE);
1189 auto stream = kj::HandleInputStream(handle);
1190 auto reader = kj::heap<capnp::InputStreamMessageReader>(stream, CONFIG_READER_OPTIONS);
1191#else
1192 auto reader = kj::heap<capnp::StreamFdMessageReader>(STDIN_FILENO, CONFIG_READER_OPTIONS);
1193#endif
1194 config = reader->getRoot<config::Config>();
1195 configOwner = kj::mv(reader);
1196 } else {
1197 // Read file from disk.
1198 auto path = fs->getCurrentPath().evalNative(pathStr);
1199 auto file = KJ_UNWRAP_OR(fs->getRoot().tryOpenFile(path), CLI_ERROR("No such file."));
1200 
1201 // Use stat() to check that we have a file vs a directory which will fail to mmap
1202 auto metadata = file->stat();
1203 if (metadata.type != kj::FsNode::Type::FILE) {
1204 CLI_ERROR("Config path is not a file.");
1205 }
1206 
1207 if (binaryConfig) {
1208 // Interpret as binary config.
1209 auto mapping = file->mmap(0, file->stat().size);
1210 auto words = kj::arrayPtr(reinterpret_cast<const capnp::word*>(mapping.begin()),
1211 mapping.size() / sizeof(capnp::word));
1212 auto reader = kj::heap<capnp::FlatArrayMessageReader>(words, CONFIG_READER_OPTIONS)
1213 .attach(kj::mv(mapping));
1214 config = reader->getRoot<config::Config>();
1215 configOwner = kj::mv(reader);
1216 } else {
1217 // Interpret as schema file.
1218 schemaParser.loadCompiledTypeAndDependencies<config::Config>();
1219 
1220 parsedSchema = schemaParser.parseFile(kj::heap<SchemaFileImpl>(fs->getRoot(),
1221 fs->getCurrentPath(), kj::mv(path), nullptr, importPath, kj::mv(file), watcher, *this));
1222 
1223 // Construct a list of top-level constants of type `Config`. If there is exactly one,
1224 // we can use it by default.
1225 for (auto nested: parsedSchema.getAllNested()) {
1226 if (nested.getProto().isConst()) {
1227 auto constSchema = nested.asConst();
1228 auto type = constSchema.getType();
1229 if (type.isStruct() &&
1230 type.asStruct().getProto().getId() == capnp::typeId<config::Config>()) {
1231 topLevelConfigConstants.add(constSchema);
1232 }
1233 }
1234 }
1235 }
1236 }
1237 
1238 // We'll fail at getConfig() if there are multiple top level Config objects.
1239 // The error message says that you have to specify which config to use, but
1240 // it's not clear that there is any mechanism to do that.
1241 util::Autogate::initAutogate(getConfig().getAutogates());
1242 }
1243 
1244 void setConstName(kj::StringPtr name) {
1245 auto parent = parsedSchema;
1246 
1247 for (;;) {
1248 auto dotPos = KJ_UNWRAP_OR(name.findFirst('.'), break);
1249 auto parentName = name.first(dotPos);
1250 parent = KJ_UNWRAP_OR(parent.findNested(kj::str(parentName)),
1251 CLI_ERROR("No such constant is defined in the config file (the parent scope '",
1252 parentName, "' does not exist)."));
1253 name = name.slice(dotPos + 1);
1254 }
1255 
1256 auto node = KJ_UNWRAP_OR(parsedSchema.findNested(name),
1257 CLI_ERROR("No such constant is defined in the config file."));
1258 
1259 if (!node.getProto().isConst()) {
1260 CLI_ERROR("Symbol is not a constant.");
1261 }
1262 
1263 auto constSchema = node.asConst();
1264 auto type = constSchema.getType();
1265 if (!type.isStruct() || type.asStruct().getProto().getId() != capnp::typeId<config::Config>()) {
1266 CLI_ERROR("Constant is not of type 'Config'.");
1267 }
1268 
1269 config = constSchema.as<config::Config>();
1270 }
1271 
1272 void setTestFilter(kj::StringPtr filter) {
1273 kj::Vector<kj::String> parts;
1274 
1275 for (;;) {
1276 KJ_IF_SOME(pos, filter.findFirst(':')) {
1277 parts.add(kj::str(filter.first(pos)));
1278 filter = filter.slice(pos + 1);
1279 } else {
1280 parts.add(kj::str(filter));
1281 break;
1282 }
1283 }
1284 
1285 switch (parts.size()) {
1286 case 0:
1287 KJ_UNREACHABLE;
1288 case 1:
1289 testServicePattern = kj::mv(parts[0]);
1290 break;
1291 case 2:
1292 testServicePattern = kj::mv(parts[0]);
1293 testEntrypointPattern = kj::mv(parts[1]);
1294 break;
1295 case 3:
1296 setConstName(parts[0]);
1297 testServicePattern = kj::mv(parts[1]);
1298 testEntrypointPattern = kj::mv(parts[2]);
1299 break;
1300 default:
1301 CLI_ERROR("Too many colons.");
1302 }
1303 }
1304 
1305 void setTestCompatDate(kj::StringPtr date) {
1306 testCompatDate = kj::str(date);
1307 }
1308 
1309 void compile() {
1310 if (hadErrors) {
1311 // Errors were already reported with context.error(), so context.exit() will exit with a
1312 // non-zero code.
1313 context.exit();
1314 }
1315 
1316 config::Config::Reader config = getConfig();
1317 
1318#if _WIN32
1319 if (_isatty(_fileno(stdout))) {
1320#else
1321 if (isatty(STDOUT_FILENO)) {
1322#endif
1323 context.exitError(
1324 "Refusing to write binary to the terminal. Please use `>` to send the output to a file.");
1325 }
1326 
1327#if !_WIN32
1328 // Grab the inode info before we write anything.
1329 struct stat stats;
1330 KJ_SYSCALL(fstat(STDOUT_FILENO, &stats));
1331#endif
1332 
1333#if _WIN32
1334 kj::FdOutputStream out(_fileno(stdout));
1335#else
1336 kj::FdOutputStream out(STDOUT_FILENO);
1337#endif
1338 
1339 if (configOnly) {
1340 // Write just the config -- in normal message format -- to stdout.
1341 uint64_t size = config.totalSize().wordCount + 1;
1342 capnp::MallocMessageBuilder builder(size + 1);
1343 builder.setRoot(config);
1344 KJ_DASSERT(builder.getSegmentsForOutput().size() == 1);
1345 capnp::writeMessage(out, builder);
1346 } else {
1347 // Write an executable file to stdout by concatenating this executable, the config, and the
1348 // magic suffix. This takes advantage of the fact that you can append arbitrary stuff to an
1349 // ELF binary or Windows executable without affecting the ability to execute the program.
1350 
1351 // Copy the executable to the output.
1352 {
1353 auto& exe = KJ_UNWRAP_OR(exeInfo,
1354 CLI_ERROR(
1355 "Unable to find and open the program's own executable, so cannot produce a new "
1356 "binary with compiled-in config."));
1357 
1358 auto mapping = exe.file->mmap(0, exe.file->stat().size);
1359 out.write(mapping);
1360 
1361 // Pad to a word boundary if necessary.
1362 size_t n = mapping.size() % sizeof(capnp::word);
1363 if (n != 0) {
1364 kj::byte pad[sizeof(capnp::word)] = {0};
1365 out.write(kj::arrayPtr(pad).slice(n));
1366 }
1367 }
1368 
1369 // Now write the config, plus magic suffix. We're going to write the config as a
1370 // single-segment flat message, which makes it easier to consume.
1371 {
1372 uint64_t size = config.totalSize().wordCount + 1;
1373 static_assert(sizeof(uint64_t) + sizeof(COMPILED_MAGIC_SUFFIX) == sizeof(capnp::word) * 3);
1374 auto words = kj::heapArray<capnp::word>(size + 3);
1375 words.asBytes().fill(0);
1376 capnp::copyToUnchecked(config, words.first(size));
1377 
1378 memcpy(&words[words.size() - 3], &size, sizeof(size));
1379 memcpy(&words[words.size() - 2], COMPILED_MAGIC_SUFFIX, sizeof(COMPILED_MAGIC_SUFFIX));
1380 
1381 out.write(words.asBytes());
1382 }
1383 
1384#if !_WIN32
1385 // If we wrote a regular file, and it was empty before we started writing, then let's go ahead
1386 // and set the executable bit on the file.
1387 if (S_ISREG(stats.st_mode) && stats.st_size == 0) {
1388 // Add executable bit for all users who have read access.
1389 mode_t mode = stats.st_mode;
1390 if (mode & S_IRUSR) {
1391 mode |= S_IXUSR;
1392 }
1393 if (mode & S_IRGRP) {
1394 mode |= S_IXGRP;
1395 }
1396 if (mode & S_IROTH) {
1397 mode |= S_IXOTH;
1398 }
1399 KJ_SYSCALL(fchmod(STDOUT_FILENO, mode));
1400 }
1401#endif
1402 }
1403 }
1404 
1405 template <typename Func>
1406 void serveImpl(Func&& func) noexcept {
1407 if (hadErrors) {
1408 // Can't start, stuff is broken.
1409 KJ_IF_SOME(w, watcher) {
1410 // In --watch mode, it's annoying if the server exits and stops watching. Let's wait for
1411 // someone to fix the config.
1412 context.warning(
1413 "Can't start server due to config errors, waiting for config files to change...");
1414 waitForChanges(w).wait(io.waitScope);
1415 reloadFromConfigChange();
1416 } else {
1417 // Errors were reported earlier, so context.exit() will exit with a non-zero status.
1418 context.exit();
1419 }
1420 } else {
1421#ifdef WORKERD_USE_PERFETTO
1422 kj::Maybe<PerfettoSession> maybePerfettoSession;
1423 KJ_IF_SOME(dest, perfettoTraceDestination) {
1424 maybePerfettoSession =
1425 PerfettoSession(dest, kj::mv(perfettoTraceCategories).orDefault(kj::String()));
1426 }
1427#endif
1428 TRACE_EVENT("workerd", "serveImpl()");
1429 auto config = getConfig();
1430 
1431 // Configure structured logging in the process context
1432 if (config.hasLogging() ? config.getLogging().getStructuredLogging()
1433 : config.getStructuredLogging()) {
1434 context.enableStructuredLogging();
1435 }
1436 
1437 auto platform = jsg::defaultPlatform(0);
1438 WorkerdPlatform v8Platform(*platform);
1439 jsg::V8System v8System(v8Platform,
1440 KJ_MAP(flag, config.getV8Flags()) -> kj::StringPtr { return flag; }, platform.get());
1441 auto promise = func(v8System, config);
1442 KJ_IF_SOME(w, watcher) {
1443 promise = promise.exclusiveJoin(waitForChanges(w).then([this]() {
1444 // Watch succeeded.
1445 reloadFromConfigChange();
1446 }));
1447 }
1448 promise.wait(io.waitScope);
1449#ifdef WORKERD_USE_PERFETTO
1450 KJ_IF_SOME(perfettoSession, maybePerfettoSession) {
1451 auto dropMe = kj::mv(perfettoSession);
1452 maybePerfettoSession = kj::none;
1453 }
1454#endif
1455 
1456 if (getenv("KJ_CLEAN_SHUTDOWN") == nullptr) {
1457 context.exit();
1458 }
1459 
1460 // Server maintains a reference to the v8 platform. Clean up before destroying the platform.
1461 server = nullptr;
1462 }
1463 }
1464 
1465 void serve() noexcept {
1466 serveImpl([&](jsg::V8System& v8System, config::Config::Reader config) {
1467#if _WIN32
1468 return server->run(v8System, config);
1469#else
1470 return server->run(v8System, config,
1471 // Gracefully drain when SIGTERM is received.
1472 io.unixEventPort.onSignal(SIGTERM).ignoreResult());
1473#endif
1474 });
1475 }
1476 
1477 void test() {
1478 if (!noVerbose) {
1479 // Always turn on info logging when running tests so that uncaught exceptions are displayed.
1480 // TODO(beta): This can be removed once we improve our error logging story.
1481 kj::_::Debug::setLogLevel(kj::LogSeverity::INFO);
1482 }
1483 if (predictable) {
1484 setPredictableModeForTest();
1485 }
1486 if (gcStress) {
1487 setGcStressModeForTest();
1488 }
1489 if (allAutogates) {
1490 util::Autogate::initAllAutogates();
1491 }
1492 
1493 KJ_IF_SOME(compatDate, testCompatDate) {
1494 server->setTestCompatibilityDateOverride(kj::str(compatDate));
1495 }
1496 
1497 // Enable loopback sockets in tests only.
1498 network.enableLoopback();
1499 
1500 serveImpl([&](jsg::V8System& v8System, config::Config::Reader config) {
1501 return server
1502 ->test(v8System, config,
1503 testServicePattern.map([](auto& s) -> kj::StringPtr { return s; }).orDefault("*"_kj),
1504 testEntrypointPattern.map([](auto& s) -> kj::StringPtr {
1505 return s;
1506 }).orDefault("*"_kj))
1507 .then([this](bool result) -> kj::Promise<void> {
1508 if (!result) {
1509 context.error("Tests failed!");
1510 }
1511 
1512 if (watcher == kj::none) {
1513 return kj::READY_NOW;
1514 } else {
1515 // Pause forever waiting for watcher.
1516 return kj::NEVER_DONE;
1517 }
1518 });
1519 });
1520 }
1521 
1522#if _WIN32
1523 void reloadFromConfigChange() {
1524 KJ_UNREACHABLE("Watching is not yet implemented on Windows");
1525 }
1526#else
1527 [[noreturn]] void reloadFromConfigChange() {
1528 // Write extra spaces to fully overwrite the line that we wrote earlier with a CR but no LF:
1529 // "Noticed configuration change, reloading shortly...\r"
1530 context.warning("Reloading due to config change... ");
1531 for (auto fd: inheritedFds) {
1532 // Disable close-on-exec for inherited FDs so that the successor process can also inherit
1533 // them.
1534 KJ_SYSCALL(ioctl(fd, FIONCLEX));
1535 }
1536 bool missingBinary = false;
1537 for (;;) {
1538 KJ_SYSCALL_HANDLE_ERRORS(execve(KJ_ASSERT_NONNULL(exeInfo).path.cStr(), argv, environ)) {
1539 case ENOENT: {
1540 // Write a message
1541 // TODO(cleanup): Writing directly to stderr is super-hacky.
1542 if (!missingBinary) {
1543 context.warning("The server executable is missing! Waiting for it to reappear...\r");
1544 missingBinary = true;
1545 }
1546 sleep(1);
1547 break;
1548 }
1549 default:
1550 KJ_FAIL_SYSCALL("execve", error);
1551 }
1552 }
1553 }
1554#endif
1555 
1556 private:
1557 StructuredLoggingProcessContext& context;
1558 char** argv;
1559 
1560 bool binaryConfig = false;
1561 bool configOnly = false;
1562 bool noVerbose = false;
1563 bool predictable = false;
1564 bool gcStress = false;
1565 bool allAutogates = false;
1566 kj::Maybe<kj::String> testCompatDate;
1567 kj::Maybe<FileWatcher> watcher;
1568 
1569 kj::Own<kj::Filesystem> fs = kj::newDiskFilesystem();
1570 kj::AsyncIoContext io = kj::setupAsyncIo();
1571 NetworkWithLoopback network{io.provider->getNetwork(), *io.provider};
1572 EntropySourceImpl entropySource;
1573 
1574 kj::Vector<kj::Path> importPath;
1575 capnp::SchemaParser schemaParser;
1576 capnp::ParsedSchema parsedSchema;
1577 kj::Vector<capnp::ConstSchema> topLevelConfigConstants;
1578 
1579 kj::Own<void> configOwner; // backing object for `config`, if it's not `schemaParser`.
1580 kj::Maybe<config::Config::Reader> config;
1581 
1582 kj::Vector<int> inheritedFds;
1583 
1584 kj::Maybe<kj::String> testServicePattern;
1585 kj::Maybe<kj::String> testEntrypointPattern;
1586 
1587#ifdef WORKERD_USE_PERFETTO
1588 kj::Maybe<kj::String> perfettoTraceDestination;
1589 kj::Maybe<kj::String> perfettoTraceCategories;
1590#endif
1591 
1592 kj::Own<Server> server;
1593 
1594 // This is a randomly-generated 128-bit number that identifies when a binary has been compiled
1595 // with a specific config in order to run stand-alone.
1596 static constexpr uint64_t COMPILED_MAGIC_SUFFIX[2] = {// The layout of such a binary is:
1597 //
1598 // - Binary executable data (copy of the Workers Runtime binary).
1599 // - Padding to 8-byte boundary.
1600 // - Cap'n-Proto-encoded config.
1601 // - 8-byte size of config, counted in 8-byte words.
1602 // - 16-byte magic number COMPILED_MAGIC_SUFFIX.
1603 
1604 0xa69eda94d3cc02b5ull, 0xa3d977fdbf547d7full};
1605 
1606 struct ExeInfo {
1607 kj::String path;
1608 kj::Own<const kj::ReadableFile> file;
1609 };
1610 
1611#if _WIN32
1612 static kj::Maybe<ExeInfo> tryOpenExe(kj::Filesystem& fs, kj::StringPtr path) {
1613 // TODO(bug): Like with Unix below, we should probably use native CreateFile() here, but it has
1614 // sooooo many arguments, I don't want to deal with it.
1615 auto parsedPath = fs.getCurrentPath().evalNative(path);
1616 KJ_IF_SOME(file, fs.getRoot().tryOpenFile(parsedPath)) {
1617 return ExeInfo{kj::str(path), kj::mv(file)};
1618 }
1619 return kj::none;
1620 }
1621#else
1622 static kj::Maybe<ExeInfo> tryOpenExe(kj::Filesystem& fs, kj::StringPtr path) {
1623 // Use open() and not fs.getRoot().tryOpenFile() because we probably want to use true kernel
1624 // path resolution here, not KJ's logical path resolution.
1625 int fd = open(path.cStr(), O_RDONLY);
1626 if (fd < 0) {
1627 return kj::none;
1628 }
1629 return ExeInfo{kj::str(path), kj::newDiskFile(kj::OwnFd(fd))};
1630 }
1631#endif
1632 
1633 static kj::Maybe<ExeInfo> getExecFile(kj::ProcessContext& context, kj::Filesystem& fs) {
1634#ifdef __GLIBC__
1635 auto execfn = getauxval(AT_EXECFN);
1636 if (execfn != 0) {
1637 return tryOpenExe(fs, reinterpret_cast<const char*>(execfn));
1638 }
1639#endif
1640 
1641#if __linux__
1642 KJ_IF_SOME(link, fs.getRoot().tryReadlink(kj::Path({"proc", "self", "exe"}))) {
1643 return tryOpenExe(fs, link);
1644 }
1645#endif
1646 
1647#if __APPLE__
1648 // https://astojanov.github.io/blog/2011/09/26/pid-to-absolute-path.html
1649 pid_t pid = getpid();
1650 char pathbuf[PROC_PIDPATHINFO_MAXSIZE];
1651 if (proc_pidpath(pid, pathbuf, sizeof(pathbuf)) > 0) {
1652 return tryOpenExe(fs, pathbuf);
1653 }
1654#endif
1655 
1656#if _WIN32
1657 wchar_t pathbuf[MAX_PATH];
1658 int result = GetModuleFileNameW(NULL, pathbuf, MAX_PATH);
1659 if (result > 0) {
1660 auto decoded = kj::decodeWideString(kj::arrayPtr(pathbuf, result));
1661 KJ_ASSERT(!decoded.hadErrors);
1662 return tryOpenExe(fs, decoded);
1663 }
1664#endif
1665 
1666 // TODO(beta): Fall back to searching $PATH.
1667 return kj::none;
1668 }
1669 
1670 config::Config::Reader getConfig() {
1671 KJ_IF_SOME(c, config) {
1672 return c;
1673 } else {
1674 // The optional `<const-name>` parameter must not have been given -- otherwise we would have
1675 // a non-null `config` by this point. See if we can infer the correct constant...
1676 if (topLevelConfigConstants.empty()) {
1677 context.exitError(
1678 "The config file does not define any top-level constants of type 'Config'.");
1679 } else if (topLevelConfigConstants.size() == 1) {
1680 return config.emplace(topLevelConfigConstants[0].as<config::Config>());
1681 } else {
1682 auto names = KJ_MAP(cnst, topLevelConfigConstants) { return cnst.getShortDisplayName(); };
1683 // TODO: this error message says "you must specify which one to use".
1684 // This is not actually possible? Either fix the error message to say
1685 // **how** to specify which config object to use or tell user to define
1686 // exactly one top level Config constant.
1687 context.exitError(kj::str(
1688 "The config file defines multiple top-level constants of type 'Config', so you must "
1689 "specify which one to use. The options are: ",
1690 kj::strArray(names, ", ")));
1691 }
1692 }
1693 }
1694 
1695 kj::Maybe<ExeInfo> exeInfo = getExecFile(context, *fs);
1696 
1697 bool hadErrors = false;
1698 
1699 void reportParsingError(kj::StringPtr file,
1700 capnp::SchemaFile::SourcePos start,
1701 capnp::SchemaFile::SourcePos end,
1702 kj::StringPtr message) override {
1703 if (start.line == end.line && start.column < end.column) {
1704 context.error(kj::str(
1705 file, ":", start.line + 1, ":", start.column + 1, "-", end.column + 1, ": ", message));
1706 } else {
1707 context.error(kj::str(file, ":", start.line + 1, ":", start.column + 1, ": ", message));
1708 }
1709 
1710 hadErrors = true;
1711 }
1712 
1713#if _WIN32
1714 kj::Promise<void> waitForChanges(FileWatcher& watcher) {
1715 KJ_UNIMPLEMENTED("Watching is not yet implemented on Windows");
1716 }
1717#else
1718 // Wait for the FileWatcher to report a change, and then wait a moment for changes to settle
1719 // down, in case there's a bunch of changes all at once.
1720 kj::Promise<void> waitForChanges(FileWatcher& watcher) {
1721 co_await watcher.onChange();
1722 
1723 // Saw our first change!
1724 
1725 // Let the user know we saw the config change.
1726 // We don't include a newline but rather a carriage return so that when the next
1727 // line is written, this line disappears, to reduce noise.
1728 // TODO(cleanup): Writing directly to stderr is super-hacky.
1729 auto message = "Noticed configuration change, reloading shortly...\r"_kjb;
1730 kj::FdOutputStream(STDERR_FILENO).write(message);
1731 
1732 static auto const waitForResult = [](kj::Promise<void> promise,
1733 bool result = false) -> kj::Promise<bool> {
1734 co_await promise;
1735 co_return result;
1736 };
1737 
1738 for (;;) {
1739 auto nextChange = waitForResult(watcher.onChange());
1740 auto timeout =
1741 waitForResult(io.provider->getTimer().afterDelay(500 * kj::MILLISECONDS), true);
1742 bool sawTimeout = co_await nextChange.exclusiveJoin(kj::mv(timeout));
1743 
1744 // If we timed out, we end the loop. If we didn't time out, then we must have seen yet
1745 // another change, so we loop again with a new timeout.
1746 if (sawTimeout) break;
1747 }
1748 
1749 co_return;
1750 }
1751#endif
1752};
1753 
1754} // namespace
1755} // namespace workerd::server
1756 
1757int main(int argc, char* argv[]) {
1758 workerd::server::StructuredLoggingProcessContext context(argv[0]);
1759 
1760#if !_WIN32
1761 kj::UnixEventPort::captureSignal(SIGTERM);
1762#endif
1763 workerd::rust::cxx_integration::init();
1764 workerd::server::CliMain mainObject(context, argv);
1765 
1766#if defined(WORKERD_FUZZILLI) && defined(__linux__)
1767 initSignalHandlers();
1768#endif
1769 
1770 return ::kj::runMainAndExit(context, mainObject.getMain(), argc, argv);
1771}