Skip to content
File

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

203.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/jsg/jsg-test.h>
8#include <workerd/jsg/setup.h>
9#include <workerd/util/autogate.h>
10#include <workerd/util/capnp-mock.h>
11 
12#include <capnp/compat/http-over-capnp.h>
13#include <capnp/rpc-twoparty.h>
14#include <kj/async-queue.h>
15#include <kj/encoding.h>
16#include <kj/test.h>
17 
18#include <cstdlib>
19#include <regex>
20 
21#if __linux__
22#include <unistd.h>
23#endif
24 
25namespace workerd::server {
26namespace {
27 
28#define KJ_FAIL_EXPECT_AT(location, ...) KJ_LOG_AT(ERROR, location, ##__VA_ARGS__);
29#define KJ_EXPECT_AT(cond, location, ...) \
30 if (auto _kjCondition = ::kj::_::MAGIC_ASSERT << cond) \
31 ; \
32 else \
33 KJ_FAIL_EXPECT_AT(location, "failed: expected " #cond, _kjCondition, ##__VA_ARGS__)
34 
35jsg::V8System v8System;
36// This can only be created once per process, so we have to put it at the top level.
37 
38const bool verboseLog = ([]() {
39 // TODO(beta): Improve uncaught exception reporting so that we don't have to do this.
40 kj::_::Debug::setLogLevel(kj::LogSeverity::INFO);
41 return true;
42})();
43 
44kj::Own<config::Config::Reader> parseConfig(kj::StringPtr text, kj::SourceLocation loc) {
45 capnp::MallocMessageBuilder builder;
46 auto root = builder.initRoot<config::Config>();
47 KJ_IF_SOME(exception, kj::runCatchingExceptions([&]() { TEXT_CODEC.decode(text, root); })) {
48 KJ_FAIL_REQUIRE_AT(loc, exception);
49 }
50 
51 util::Autogate::initAutogate(root.asReader().getAutogates());
52 
53 return capnp::clone(root.asReader());
54}
55 
56// Accept an indented block of text and remove the indentation. From each line of text, this will
57// remove a number of spaces up to the indentation of the first line.
58//
59// This is intended to allow multi-line raw text to be specified conveniently using C++11
60// `R"(blah)"` literal syntax, without the need to mess up indentation relative to the
61// surrounding code.
62kj::String operator""_blockquote(const char* str, size_t n) {
63 kj::StringPtr text(str, n);
64 
65 // Ignore a leading newline so that `R"(` can be placed on the line before the initial indent.
66 if (text.startsWith("\n")) {
67 text = text.slice(1);
68 }
69 
70 // Count indent size.
71 size_t indent = 0;
72 while (text.startsWith(" ")) {
73 text = text.slice(1);
74 ++indent;
75 }
76 
77 // Process lines.
78 kj::Vector<char> result;
79 while (text != nullptr) {
80 // Add data from this line.
81 auto nl = text.findFirst('\n').orDefault(text.size() - 1) + 1;
82 result.addAll(text.first(nl));
83 text = text.slice(nl);
84 
85 // Skip indent of next line, up to the expected indent size.
86 size_t seenIndent = 0;
87 while (seenIndent < indent && text.startsWith(" ")) {
88 text = text.slice(1);
89 ++seenIndent;
90 }
91 }
92 
93 result.add('\0');
94 return kj::String(result.releaseAsArray());
95}
96 
97class TestStream {
98 public:
99 TestStream(kj::WaitScope& ws, kj::Own<kj::AsyncIoStream> stream)
100 : ws(ws),
101 stream(kj::mv(stream)) {}
102 
103 void send(kj::StringPtr data, kj::SourceLocation loc = {}) {
104 stream->write(data.asBytes()).wait(ws);
105 }
106 void recv(kj::StringPtr expected, kj::SourceLocation loc = {}) {
107 auto actual = readAllAvailable();
108 if (actual == nullptr) {
109 KJ_FAIL_EXPECT_AT(loc, "message never received");
110 } else {
111 KJ_EXPECT_AT(actual == expected, loc);
112 }
113 }
114 void recvRegex(kj::StringPtr matcher, kj::SourceLocation loc = {}) {
115 auto actual = readAllAvailable();
116 if (actual == nullptr) {
117 KJ_FAIL_EXPECT_AT(loc, "message never received");
118 } else {
119 std::regex target(matcher.cStr());
120 KJ_EXPECT(std::regex_match(actual.cStr(), target), actual, matcher, loc);
121 }
122 }
123 
124 void recvWebSocket(kj::StringPtr expected, kj::SourceLocation loc = {}) {
125 auto actual = readWebSocketMessage();
126 KJ_EXPECT_AT(actual == expected, loc);
127 }
128 
129 void recvWebSocketRegex(kj::StringPtr matcher, kj::SourceLocation loc = {}) {
130 auto actual = readWebSocketMessage();
131 std::regex target(matcher.cStr());
132 KJ_EXPECT(std::regex_match(actual.cStr(), target), actual, matcher, loc);
133 }
134 
135 void recvWebSocketClose(int expectedCode) {
136 auto actual = readWebSocketMessage();
137 KJ_EXPECT(actual.size() >= 2);
138 int gotCode = (static_cast<uint8_t>(actual[0]) << 8) + static_cast<uint8_t>(actual[1]);
139 KJ_EXPECT(gotCode == expectedCode);
140 }
141 
142 void sendHttpGet(kj::StringPtr path, kj::SourceLocation loc = {}) {
143 send(kj::str("GET ", path,
144 " HTTP/1.1\n"
145 "Host: foo\n"
146 "\n"),
147 loc);
148 }
149 
150 void recvHttp200(kj::StringPtr expectedResponse, kj::SourceLocation loc = {}) {
151 recv(kj::str("HTTP/1.1 200 OK\n"
152 "Content-Length: ",
153 expectedResponse.size(),
154 "\n"
155 "Content-Type: text/plain;charset=UTF-8\n"
156 "\n",
157 expectedResponse),
158 loc);
159 }
160 
161 void httpGet200(kj::StringPtr path, kj::StringPtr expectedResponse, kj::SourceLocation loc = {}) {
162 sendHttpGet(path, loc);
163 recvHttp200(expectedResponse, loc);
164 }
165 
166 // Return true if the stream is at EOF.
167 bool isEof() {
168 if (premature != kj::none) {
169 // We still have unread data so we're definitely not at EOF.
170 return false;
171 }
172 
173 char c;
174 auto promise = stream->tryRead(&c, 1, 1);
175 if (!promise.poll(ws)) {
176 // Read didn't complete immediately. We have no data available, but we're not at EOF.
177 return false;
178 }
179 
180 size_t n = promise.wait(ws);
181 if (n == 0) {
182 return true;
183 } else {
184 // Oops, the stream had data available and we accidentally read a byte of it. Store that off
185 // to the side.
186 KJ_ASSERT(n == 1);
187 premature = c;
188 return false;
189 }
190 }
191 
192 void upgradeToWebSocket() {
193 send(R"(
194 GET / HTTP/1.1
195 Host: foo
196 Upgrade: websocket
197 Sec-WebSocket-Key: AAAAAAAAAAAAAAAAAAAAAA==
198 Sec-WebSocket-Version: 13
199 
200 )"_blockquote,
201 {});
202 
203 recv(R"(
204 HTTP/1.1 101 Switching Protocols
205 Connection: Upgrade
206 Upgrade: websocket
207 Sec-WebSocket-Accept: ICX+Yqv66kxgM0FcWaLWlFLwTAI=
208 
209 )"_blockquote,
210 {});
211 }
212 
213 kj::AsyncIoStream& getStream() {
214 return *stream;
215 }
216 
217 private:
218 kj::WaitScope& ws;
219 kj::Own<kj::AsyncIoStream> stream;
220 
221 // isEof() may prematurely read a character. Keep it off to the side for the next actual read.
222 kj::Maybe<char> premature;
223 
224 kj::String readAllAvailable() {
225 kj::Vector<char> buffer(256);
226 KJ_IF_SOME(p, premature) {
227 buffer.add(p);
228 }
229 
230 // Continuously try to read until there's nothing to read (or we've gone way past the size
231 // expected).
232 for (;;) {
233 size_t pos = buffer.size();
234 buffer.resize(kj::max(buffer.size() + 256, buffer.capacity()));
235 
236 auto promise = stream->tryRead(buffer.begin() + pos, 1, buffer.size() - pos);
237 if (!promise.poll(ws)) {
238 // A tryRead() of 1 byte didn't resolve, there must be no data to read.
239 buffer.resize(pos);
240 break;
241 }
242 size_t n = promise.wait(ws);
243 if (n == 0) {
244 buffer.resize(pos);
245 break;
246 }
247 
248 // Strip out `\r`s for convenience. We do this in-place...
249 for (size_t i: kj::range(pos, pos + n)) {
250 if (buffer[i] != '\r') {
251 buffer[pos++] = buffer[i];
252 }
253 }
254 buffer.resize(pos);
255 };
256 
257 buffer.add('\0');
258 return kj::String(buffer.releaseAsArray());
259 }
260 
261 kj::String readWebSocketMessage(size_t maxMessageSize = 1 << 24) {
262 // Reads a single, non-fragmented WebSocket message. Returns just the payload.
263 kj::Vector<uint8_t> header(256);
264 kj::Vector<uint8_t> mask(4);
265 
266 KJ_IF_SOME(p, premature) {
267 header.add(p);
268 premature = kj::Maybe<char>();
269 }
270 
271 tryRead(header, 2 - header.size(), "reading first two bytes of header");
272 bool masked = header[1] & 0x80;
273 size_t sevenBitPayloadLength = header[1] & 0x7f;
274 size_t realPayloadLength = sevenBitPayloadLength;
275 
276 if (sevenBitPayloadLength == 126) {
277 tryRead(header, 2, "reading 16-bit payload length");
278 realPayloadLength = (static_cast<size_t>(header[2]) << 8) + static_cast<size_t>(header[3]);
279 } else if (sevenBitPayloadLength == 127) {
280 tryRead(header, 8, "reading 64-bit payload length");
281 realPayloadLength = (static_cast<size_t>(header[2]) << 56) +
282 (static_cast<size_t>(header[3]) << 48) + (static_cast<size_t>(header[4]) << 40) +
283 (static_cast<size_t>(header[5]) << 32) + (static_cast<size_t>(header[6]) << 24) +
284 (static_cast<size_t>(header[7]) << 16) + (static_cast<size_t>(header[8]) << 8) +
285 (static_cast<size_t>(header[9]));
286 
287 KJ_REQUIRE(realPayloadLength <= maxMessageSize,
288 kj::str("Payload size too big (", realPayloadLength, " > ", maxMessageSize, ")"));
289 }
290 
291 if (masked) {
292 tryRead(mask, 4, "reading mask key");
293 // Currently we assume the mask is always 0, so its application is a no-op, hence we don't
294 // bother.
295 }
296 kj::Vector<char> payload(realPayloadLength + 1);
297 
298 tryRead(payload, realPayloadLength, "reading payload");
299 payload.add('\0');
300 return kj::String(payload.releaseAsArray());
301 }
302 
303 template <typename T>
304 void tryRead(kj::Vector<T>& buffer, size_t bytesToRead, kj::StringPtr what) {
305 static_assert(sizeof(T) == 1, "not byte-sized");
306 
307 size_t pos = buffer.size();
308 size_t bytesRead = 0;
309 buffer.resize(buffer.size() + bytesToRead);
310 while (bytesRead < bytesToRead) {
311 auto promise = stream->tryRead(buffer.begin() + pos, 1, buffer.size() - pos);
312 KJ_REQUIRE(promise.poll(ws), kj::str("No data available while ", what));
313 // A tryRead() of 1 byte didn't resolve, there must be no data to read.
314 
315 size_t n = promise.wait(ws);
316 KJ_REQUIRE(n > 0, kj::str("Not enough data while ", what));
317 bytesRead += n;
318 }
319 }
320};
321 
322class TestServer final: private kj::Filesystem, private kj::EntropySource, private kj::Clock {
323 public:
324 TestServer(kj::StringPtr configText,
325 Worker::ConsoleMode consoleMode = Worker::ConsoleMode::INSPECTOR_ONLY,
326 kj::SourceLocation loc = {})
327 : ws(loop),
328 config(parseConfig(configText, loc)),
329 root(kj::newInMemoryDirectory(*this)),
330 pwd(kj::Path({"current", "dir"})),
331 cwd(root->openSubdir(pwd, kj::WriteMode::CREATE | kj::WriteMode::CREATE_PARENT)),
332 timer(kj::origin<kj::TimePoint>()),
333 server(*this,
334 timer,
335 timer,
336 mockNetwork,
337 *this,
338 Worker::LoggingOptions(consoleMode),
339 [this](kj::String error) {
340 if (expectedErrors.startsWith(error) && expectedErrors[error.size()] == '\n') {
341 expectedErrors = expectedErrors.slice(error.size() + 1);
342 } else {
343 KJ_FAIL_EXPECT(error, expectedErrors);
344 }
345 }),
346 fakeDate(kj::UNIX_EPOCH),
347 mockNetwork(*this, {}, {}) {}
348 
349 ~TestServer() noexcept(false) {
350 for (auto& subq: subrequests) {
351 subq.value->rejectAll(KJ_EXCEPTION(FAILED, "test ended"));
352 }
353 
354 if (!unwindDetector.isUnwinding()) {
355 // Make sure any errors are reported.
356 KJ_IF_SOME(t, runTask) {
357 t.poll(ws);
358 }
359 }
360 }
361 
362 // Start the server. Call before connect().
363 void start(kj::Promise<void> drainWhen = kj::NEVER_DONE) {
364 KJ_REQUIRE(runTask == kj::none);
365 auto task =
366 server.run(v8System, *config, kj::mv(drainWhen)).eagerlyEvaluate([](kj::Exception&& e) {
367 KJ_FAIL_EXPECT(e);
368 });
369 KJ_EXPECT(!task.poll(ws));
370 runTask = kj::mv(task);
371 }
372 
373 // Call instead of `start()` when the config is expected to produce errors. The parameter is
374 // the expected list of errors messages, one per line.
375 void expectErrors(kj::StringPtr expected) {
376 expectedErrors = expected;
377 server.run(v8System, *config).poll(ws);
378 KJ_EXPECT(expectedErrors == nullptr, "some expected errors weren't seen");
379 }
380 
381 // Connect to the server on the given address. The string just has to match what is in the
382 // config; the actual connection is in-memory with no network involved.
383 TestStream connect(kj::StringPtr addr) {
384 return TestStream(ws, KJ_REQUIRE_NONNULL(sockets.find(addr), addr)->connect().wait(ws));
385 }
386 
387 // Try to connect to the address and return whether or not this connection attempt hangs,
388 // i.e. a listener exists but connections are not being accepted.
389 bool connectHangs(kj::StringPtr addr) {
390 return !KJ_REQUIRE_NONNULL(sockets.find(addr), addr)->connect().poll(ws);
391 }
392 
393 // Expect an incoming connection on the given address and from a network with the given
394 // allowed / denied peer list.
395 TestStream receiveSubrequest(kj::StringPtr addr,
396 kj::ArrayPtr<const kj::StringPtr> allowedPeers = nullptr,
397 kj::ArrayPtr<const kj::StringPtr> deniedPeers = nullptr,
398 kj::SourceLocation loc = {}) {
399 auto expectedFilter = peerFilterToString(allowedPeers, deniedPeers);
400 
401 auto promise = getSubrequestQueue(addr).pop();
402 KJ_ASSERT_AT(promise.poll(ws), loc, "never received expected subrequest", addr);
403 
404 auto info = promise.wait(ws);
405 auto actualFilter = info.peerFilter;
406 KJ_EXPECT_AT(actualFilter == expectedFilter, loc);
407 
408 auto pipe = kj::newTwoWayPipe();
409 info.fulfiller->fulfill(kj::mv(pipe.ends[0]));
410 return TestStream(ws, kj::mv(pipe.ends[1]));
411 }
412 
413 TestStream receiveInternetSubrequest(kj::StringPtr addr, kj::SourceLocation loc = {}) {
414 return receiveSubrequest(addr, {"public"_kj}, {}, loc);
415 }
416 
417 // Advance the timer through `seconds` seconds of virtual time.
418 void wait(size_t seconds) {
419 auto delayPromise = timer.afterDelay(seconds * kj::SECONDS).eagerlyEvaluate(nullptr);
420 while (!delayPromise.poll(ws)) {
421 // Since this test has no external I/O at all other than time, we know no events could
422 // possibly occur until the next timer event. So just advance directly to it and continue.
423 timer.advanceTo(KJ_ASSERT_NONNULL(timer.nextEvent()));
424 }
425 delayPromise.wait(ws);
426 }
427 
428 kj::WaitScope& getWaitScope() {
429 return ws;
430 }
431 
432 kj::EventLoop loop;
433 kj::WaitScope ws;
434 
435 kj::Own<config::Config::Reader> config;
436 kj::Own<const kj::Directory> root;
437 kj::Path pwd;
438 kj::Own<const kj::Directory> cwd;
439 kj::TimerImpl timer;
440 Server server;
441 
442 kj::Maybe<kj::Promise<void>> runTask;
443 kj::StringPtr expectedErrors;
444 
445 kj::Date fakeDate;
446 
447 private:
448 kj::UnwindDetector unwindDetector;
449 
450 // ---------------------------------------------------------------------------
451 // implements Filesystem
452 
453 const kj::Directory& getRoot() const override {
454 return *root;
455 }
456 const kj::Directory& getCurrent() const override {
457 return *cwd;
458 }
459 kj::PathPtr getCurrentPath() const override {
460 return pwd;
461 }
462 
463 // ---------------------------------------------------------------------------
464 // implements Network
465 
466 // Addresses that the server is listening on.
467 kj::HashMap<kj::String, kj::Own<kj::NetworkAddress>> sockets;
468 
469 class MockNetwork;
470 
471 struct SubrequestInfo {
472 kj::Own<kj::PromiseFulfiller<kj::Own<kj::AsyncIoStream>>> fulfiller;
473 kj::StringPtr peerFilter;
474 };
475 using SubrequestQueue = kj::ProducerConsumerQueue<SubrequestInfo>;
476 // Expected incoming connections and callbacks that should be used to handle them.
477 kj::HashMap<kj::String, kj::Own<SubrequestQueue>> subrequests;
478 
479 SubrequestQueue& getSubrequestQueue(kj::StringPtr addr) {
480 return *subrequests.findOrCreate(addr, [&]() -> decltype(subrequests)::Entry {
481 return {kj::str(addr), kj::heap<SubrequestQueue>()};
482 });
483 }
484 
485 static kj::String peerFilterToString(
486 kj::ArrayPtr<const kj::StringPtr> allow, kj::ArrayPtr<const kj::StringPtr> deny) {
487 if (allow == nullptr && deny == nullptr) {
488 return kj::str("(none)");
489 } else {
490 return kj::str("allow: [", kj::strArray(allow, ", "),
491 "], "
492 "deny: [",
493 kj::strArray(deny, ", "), "]");
494 }
495 }
496 
497 class MockAddress final: public kj::NetworkAddress {
498 public:
499 MockAddress(TestServer& test, kj::StringPtr peerFilter, kj::String address)
500 : test(test),
501 peerFilter(peerFilter),
502 address(kj::mv(address)) {}
503 
504 kj::Promise<kj::Own<kj::AsyncIoStream>> connect() override {
505 KJ_IF_SOME(addr, test.sockets.find(address)) {
506 // If someone is listening on this address, connect directly to them.
507 return addr->connect();
508 }
509 
510 auto [promise, fulfiller] = kj::newPromiseAndFulfiller<kj::Own<kj::AsyncIoStream>>();
511 
512 test.getSubrequestQueue(address).push({kj::mv(fulfiller), peerFilter});
513 
514 return kj::mv(promise);
515 }
516 kj::Own<kj::ConnectionReceiver> listen() override {
517 auto pipe = kj::newCapabilityPipe();
518 auto receiver = kj::heap<kj::CapabilityStreamConnectionReceiver>(*pipe.ends[0])
519 .attach(kj::mv(pipe.ends[0]));
520 auto sender = kj::heap<kj::CapabilityStreamNetworkAddress>(kj::none, *pipe.ends[1])
521 .attach(kj::mv(pipe.ends[1]));
522 test.sockets.insert(kj::str(address), kj::mv(sender));
523 return receiver;
524 }
525 kj::Own<kj::NetworkAddress> clone() override {
526 KJ_UNIMPLEMENTED("unused");
527 }
528 kj::String toString() override {
529 KJ_UNIMPLEMENTED("unused");
530 }
531 
532 private:
533 TestServer& test;
534 kj::StringPtr peerFilter;
535 kj::String address;
536 };
537 
538 class MockNetwork final: public kj::Network {
539 public:
540 MockNetwork(TestServer& test,
541 kj::ArrayPtr<const kj::StringPtr> allow,
542 kj::ArrayPtr<const kj::StringPtr> deny)
543 : test(test),
544 filter(peerFilterToString(allow, deny)) {}
545 
546 kj::Promise<kj::Own<kj::NetworkAddress>> parseAddress(
547 kj::StringPtr addr, uint portHint = 0) override {
548 return kj::Own<kj::NetworkAddress>(kj::heap<MockAddress>(test, filter, kj::str(addr)));
549 }
550 kj::Own<kj::NetworkAddress> getSockaddr(const void* sockaddr, uint len) override {
551 KJ_UNIMPLEMENTED("unused");
552 }
553 kj::Own<kj::Network> restrictPeers(
554 kj::ArrayPtr<const kj::StringPtr> allow, kj::ArrayPtr<const kj::StringPtr> deny) override {
555 KJ_ASSERT(filter == "(none)", "can't nest restrictPeers()");
556 return kj::heap<MockNetwork>(test, allow, deny);
557 }
558 
559 private:
560 TestServer& test;
561 kj::String filter;
562 };
563 
564 MockNetwork mockNetwork;
565 
566 // ---------------------------------------------------------------------------
567 // implements EntropySource
568 
569 void generate(kj::ArrayPtr<kj::byte> buffer) override {
570 kj::byte random = 4; // chosen by fair die roll by Randall Munroe in 2007.
571 // guaranteed to be random.
572 buffer.fill(random);
573 }
574 
575 // ---------------------------------------------------------------------------
576 // implements Clock
577 
578 kj::Date now() const override {
579 return fakeDate;
580 }
581};
582 
583// =======================================================================================
584// Test Workers
585 
586kj::String singleWorker(kj::StringPtr def) {
587 return kj::str(R"((
588 services = [
589 ( name = "hello",
590 worker = )"_kj,
591 def, R"(
592 )
593 ],
594 sockets = [
595 ( name = "main",
596 address = "test-addr",
597 service = "hello"
598 )
599 ]
600 ))"_kj);
601}
602 
603KJ_TEST("Server: serve basic Service Worker") {
604 TestServer test(singleWorker(R"((
605 compatibilityDate = "2022-08-17",
606 serviceWorkerScript =
607 `addEventListener("fetch", event => {
608 ` event.respondWith(new Response("Hello: " + event.request.url + "\n"));
609 `})
610 ))"_kj));
611 
612 test.start();
613 
614 auto conn = test.connect("test-addr");
615 
616 // Send a request, get a response.
617 conn.httpGet200("/", "Hello: http://foo/\n");
618 
619 // Send another request on the same connection, different path and host.
620 conn.send(R"(
621 GET /baz/qux?corge=grault HTTP/1.1
622 Host: bar
623 
624 )"_blockquote);
625 conn.recv(R"(
626 HTTP/1.1 200 OK
627 Content-Length: 39
628 Content-Type: text/plain;charset=UTF-8
629 
630 Hello: http://bar/baz/qux?corge=grault
631 )"_blockquote);
632 
633 // A request without `Host:` should 400.
634 conn.send(R"(
635 GET /baz/qux?corge=grault HTTP/1.1
636 
637 )"_blockquote);
638 conn.recv(R"(
639 HTTP/1.1 400 Bad Request
640 Content-Length: 11
641 
642 Bad Request)"_blockquote);
643}
644 
645KJ_TEST("Server: use service name as Service Worker origin") {
646 TestServer test(singleWorker(R"((
647 compatibilityDate = "2022-08-17",
648 serviceWorkerScript =
649 `addEventListener("fetch", event => {
650 ` event.respondWith(new Response(new Error("Doh!").stack));
651 `})
652 ))"_kj));
653 
654 test.start();
655 auto conn = test.connect("test-addr");
656 conn.httpGet200("/", R"(
657 Error: Doh!
658 at hello:2:34)"_blockquote);
659}
660 
661KJ_TEST("Server: serve basic modular Worker") {
662 TestServer test(singleWorker(R"((
663 compatibilityDate = "2022-08-17",
664 modules = [
665 ( name = "main.js",
666 esModule =
667 `export default {
668 ` async fetch(request) {
669 ` return new Response("Hello: " + request.url);
670 ` }
671 `}
672 )
673 ]
674 ))"_kj));
675 test.start();
676 auto conn = test.connect("test-addr");
677 conn.httpGet200("/", "Hello: http://foo/");
678}
679 
680KJ_TEST("Server: serve modular Worker with imports") {
681 TestServer test(singleWorker(R"((
682 compatibilityDate = "2022-08-17",
683 modules = [
684 ( name = "main.js",
685 esModule =
686 `import { MESSAGE as FOO } from "foo.js";
687 `import BAR from "bar.txt";
688 `import BAZ from "baz.bin";
689 `import QUX from "qux.json";
690 `import CORGE from "corge.js";
691 `import SQUARE_WASM from "square.wasm";
692 `const SQUARE = new WebAssembly.Instance(SQUARE_WASM, {});
693 `export default {
694 ` async fetch(request) {
695 ` return new Response([
696 ` FOO, BAR, new TextDecoder().decode(BAZ), QUX.message, CORGE.message,
697 ` "square.wasm says square(5) = " + SQUARE.exports.square(5)]
698 ` .join("\n"));
699 ` }
700 `}
701 ),
702 ( name = "foo.js",
703 esModule =
704 `export let MESSAGE = "Hello from foo.js"
705 ),
706 ( name = "bar.txt",
707 text = "Hello from bar.txt"
708 ),
709 ( name = "baz.bin",
710 data = "Hello from baz.bin"
711 ),
712 ( name = "qux.json",
713 json = `{"message": "Hello from qux.json"}
714 ),
715 ( name = "corge.js",
716 commonJsModule =
717 `module.exports.message = "Hello from corge.js";
718 ),
719 ( name = "square.wasm",
720 # Exports a function 'square(x)' that returns x^2.
721 wasm = 0x"00 61 73 6d 01 00 00 00 01 06 01 60 01 7f 01 7f
722 03 02 01 00 05 03 01 00 02 06 08 01 7f 01 41 80
723 88 04 0b 07 13 02 06 6d 65 6d 6f 72 79 02 00 06
724 73 71 75 61 72 65 00 00 0a 09 01 07 00 20 00 20
725 00 6c 0b"
726 )
727 ]
728 ))"_kj));
729 
730 test.start();
731 auto conn = test.connect("test-addr");
732 conn.httpGet200("/",
733 "Hello from foo.js\n"
734 "Hello from bar.txt\n"
735 "Hello from baz.bin\n"
736 "Hello from qux.json\n"
737 "Hello from corge.js\n"
738 "square.wasm says square(5) = 25");
739}
740 
741KJ_TEST("Server: compatibility dates") {
742 // The easiest flag to test is the presence of the global `navigator`.
743 auto selfNavigatorCheckerWorker = [](kj::StringPtr compatProperties) {
744 return singleWorker(kj::str(R"((
745 )",
746 compatProperties, R"(,
747 modules = [
748 ( name = "main.js",
749 esModule =
750 `export default {
751 ` async fetch(request) {
752 ` return new Response(!!self.navigator);
753 ` }
754 `}
755 )
756 ]
757 ))"_kj));
758 };
759 
760 {
761 TestServer test(selfNavigatorCheckerWorker("compatibilityDate = \"2022-08-17\""));
762 
763 test.start();
764 auto conn = test.connect("test-addr");
765 conn.httpGet200("/", "true");
766 }
767 
768 // In the past, the global wasn't there.
769 {
770 TestServer test(selfNavigatorCheckerWorker("compatibilityDate = \"2020-01-01\""));
771 
772 test.start();
773 auto conn = test.connect("test-addr");
774 conn.httpGet200("/", "false");
775 }
776 
777 // Disable using a flag instead of a date.
778 {
779 TestServer test(selfNavigatorCheckerWorker(
780 "compatibilityDate = \"2022-08-17\", compatibilityFlags = [\"no_global_navigator\"]"));
781 
782 test.start();
783 auto conn = test.connect("test-addr");
784 conn.httpGet200("/", "false");
785 }
786}
787 
788KJ_TEST("Server: compatibility dates are required") {
789 TestServer test(singleWorker(R"((
790 serviceWorkerScript =
791 `addEventListener("fetch", event => {
792 ` event.respondWith(new Response("Hello: " + event.request.url + "\n"));
793 `})
794 ))"_kj));
795 
796 test.expectErrors(R"(
797 service hello: Worker must specify compatibilityDate.
798 )"_blockquote);
799}
800 
801KJ_TEST("Server: value bindings") {
802#if _WIN32
803 _putenv("TEST_ENVIRONMENT_VAR=Hello from environment variable");
804#else
805 setenv("TEST_ENVIRONMENT_VAR", "Hello from environment variable", true);
806#endif
807 
808 TestServer test(singleWorker(R"((
809 compatibilityDate = "2022-08-17",
810 # (Must use Service Worker syntax to allow Wasm bindings.)
811 serviceWorkerScript =
812 `const SQUARE = new WebAssembly.Instance(BAZ, {});
813 `async function handle(request) {
814 ` let items = [];
815 ` items.push(FOO);
816 ` items.push(new TextDecoder().decode(BAR));
817 ` items.push("wasm says square(5) = " + SQUARE.exports.square(5));
818 ` items.push(QUX.message);
819 ` items.push(CORGE);
820 ` items.push("GRAULT is null? " + (GRAULT === null));
821 ` return new Response(items.join("\n"));
822 `}
823 `addEventListener("fetch", event => {
824 ` event.respondWith(handle(event.request));
825 `});
826 ,
827 bindings = [
828 ( name = "FOO", text = "Hello from text binding" ),
829 ( name = "BAR", data = "Hello from data binding" ),
830 ( name = "BAZ",
831 # Exports a function 'square(x)' that returns x^2.
832 wasmModule = 0x"00 61 73 6d 01 00 00 00 01 06 01 60 01 7f 01 7f
833 03 02 01 00 05 03 01 00 02 06 08 01 7f 01 41 80
834 88 04 0b 07 13 02 06 6d 65 6d 6f 72 79 02 00 06
835 73 71 75 61 72 65 00 00 0a 09 01 07 00 20 00 20
836 00 6c 0b"
837 ),
838 ( name = "QUX",
839 json = `{"message": "Hello from json binding"}
840 ),
841 ( name = "CORGE", fromEnvironment = "TEST_ENVIRONMENT_VAR" ),
842 ( name = "GRAULT", fromEnvironment = "TEST_NONEXISTENT_ENVIRONMENT_VAR" ),
843 ]
844 ))"_kj));
845 
846 test.start();
847 auto conn = test.connect("test-addr");
848 conn.httpGet200("/",
849 "Hello from text binding\n"
850 "Hello from data binding\n"
851 "wasm says square(5) = 25\n"
852 "Hello from json binding\n"
853 "Hello from environment variable\n"
854 "GRAULT is null? true");
855}
856 
857KJ_TEST("Server: WebCrypto bindings") {
858 TestServer test(singleWorker(R"((
859 compatibilityDate = "2022-08-17",
860 modules = [
861 ( name = "main.js",
862 esModule =
863 `function hex(buffer) {
864 ` return [...new Uint8Array(buffer)]
865 ` .map(x => x.toString(16).padStart(2, '0'))
866 ` .join('');
867 `}
868 `
869 `export default {
870 ` async fetch(request, env) {
871 ` let items = [];
872 `
873 ` let plaintext = new TextEncoder().encode("hello");
874 ` let sig = await crypto.subtle.sign({"name": "HMAC", "hash": "SHA-256"},
875 ` env.hmac, plaintext);
876 ` items.push("hmac signature is " + hex(sig));
877 ` let ver1 = await crypto.subtle.verify({"name": "HMAC", "hash": "SHA-256"},
878 ` env.hmac, sig, plaintext);
879 ` let ver2 = await crypto.subtle.verify({"name": "HMAC", "hash": "SHA-256"},
880 ` env.hmac, sig, new Uint8Array([12, 34]));
881 ` items.push("hmac verifications: " + ver1 + ", " + ver2);
882 ` items.push("hmac extractable? " + env.hmac.extractable);
883 `
884 ` let hexSig = await crypto.subtle.sign({"name": "HMAC", "hash": "SHA-256"},
885 ` env.hmacHex, plaintext);
886 ` let b64Sig = await crypto.subtle.sign({"name": "HMAC", "hash": "SHA-256"},
887 ` env.hmacBase64, plaintext);
888 ` let jwkSig = await crypto.subtle.sign({"name": "HMAC", "hash": "SHA-256"},
889 ` env.hmacJwk, plaintext);
890 ` items.push("hmac signature (hex key) is " + hex(hexSig));
891 ` items.push("hmac signature (base64 key) is " + hex(b64Sig));
892 ` items.push("hmac signature (jwk key) is " + hex(jwkSig));
893 `
894 ` try {
895 ` await crypto.subtle.verify({"name": "HMAC", "hash": "SHA-256"},
896 ` env.hmacHex, sig, plaintext);
897 ` items.push("verification with hmacHex was allowed");
898 ` } catch (err) {
899 ` items.push("verification with hmacHex was not allowed: " + err.message);
900 ` }
901 `
902 ` let ecsig = await crypto.subtle.sign(
903 ` {"name": "ECDSA", "namedCurve": "P-256", "hash": "SHA-256"},
904 ` env.ecPriv, plaintext);
905 ` let ecver = await crypto.subtle.verify(
906 ` {"name": "ECDSA", "namedCurve": "P-256", "hash": "SHA-256"},
907 ` env.ecPub, ecsig, plaintext);
908 ` items.push("ec verification: " + ecver);
909 ` items.push("ec extractable? " + env.ecPriv.extractable +
910 ` ", " + env.ecPub.extractable);
911 `
912 ` return new Response(items.join("\n"));
913 ` }
914 `}
915 )
916 ],
917 bindings = [
918 ( name = "hmac",
919 cryptoKey = (
920 raw = "testkey",
921 algorithm = (
922 json = `{"name": "HMAC", "hash": "SHA-256"}
923 ),
924 usages = [ sign, verify ]
925 )
926 ),
927 ( name = "hmacHex",
928 cryptoKey = (
929 hex = "746573746b6579",
930 algorithm = (
931 json = `{"name": "HMAC", "hash": "SHA-256"}
932 ),
933 usages = [ sign ]
934 )
935 ),
936 ( name = "hmacBase64",
937 cryptoKey = (
938 base64 = "dGVzdGtleQ==",
939 algorithm = (
940 json = `{"name": "HMAC", "hash": "SHA-256"}
941 ),
942 usages = [ sign ]
943 )
944 ),
945 ( name = "hmacJwk",
946 cryptoKey = (
947 jwk = `{"alg":"HS256","k":"dGVzdGtleQ","kty":"oct"}
948 ,
949 algorithm = (
950 json = `{"name": "HMAC", "hash": "SHA-256"}
951 ),
952 usages = [ sign ]
953 )
954 ),
955 
956 ( name = "ecPriv",
957 cryptoKey = (
958 pkcs8 =
959 `-----BEGIN PRIVATE KEY-----
960 `MIGHAgEAMBMGByqGSM49AgEGCCqGSM49AwEHBG0wawIBAQQgXB5SjGILYt4DxPho
961 `VUX/lMnLzpJD5R6Jl0bLCuRj8V2hRANCAAQ6pM4KrujAsw2xz0qA6l4DF/waMYVP
962 `QNOAakb+S9GwkOgrTbw6AYoawTaW68Vbwadfe2S02ya6yEKGyE3N56by
963 `-----END PRIVATE KEY-----
964 ,
965 algorithm = (
966 json = `{"name": "ECDSA", "namedCurve": "P-256"}
967 ),
968 usages = [ sign ]
969 )
970 ),
971 
972 ( name = "ecPub",
973 cryptoKey = (
974 spki =
975 `-----BEGIN PUBLIC KEY-----
976 `MFkwEwYHKoZIzj0CAQYIKoZIzj0DAQcDQgAEOqTOCq7owLMNsc9KgOpeAxf8GjGF
977 `T0DTgGpG/kvRsJDoK028OgGKGsE2luvFW8GnX3tktNsmushChshNzeem8g==
978 `-----END PUBLIC KEY-----
979 ,
980 algorithm = (
981 json = `{"name": "ECDSA", "namedCurve": "P-256"}
982 ),
983 usages = [ verify ],
984 extractable = true
985 )
986 )
987 ]
988 ))"_kj));
989 
990 test.start();
991 auto conn = test.connect("test-addr");
992 conn.httpGet200("/",
993 "hmac signature is 4a27693183b28d2616209d6ff5e77646af5fc06ea6affac37415995b07be2ddf\n"
994 "hmac verifications: true, false\n"
995 "hmac extractable? false\n"
996 "hmac signature (hex key) is "
997 "4a27693183b28d2616209d6ff5e77646af5fc06ea6affac37415995b07be2ddf\n"
998 "hmac signature (base64 key) is "
999 "4a27693183b28d2616209d6ff5e77646af5fc06ea6affac37415995b07be2ddf\n"
1000 "hmac signature (jwk key) is "
1001 "4a27693183b28d2616209d6ff5e77646af5fc06ea6affac37415995b07be2ddf\n"
1002 "verification with hmacHex was not allowed: "
1003 "Requested key usage \"verify\" does not match any usage listed in this CryptoKey.\n"
1004 "ec verification: true\n"
1005 "ec extractable? false, true");
1006}
1007 
1008KJ_TEST("Server: subrequest to default outbound") {
1009 TestServer test(singleWorker(R"((
1010 compatibilityDate = "2022-08-17",
1011 modules = [
1012 ( name = "main.js",
1013 esModule =
1014 `export default {
1015 ` async fetch(request, env) {
1016 ` let resp = await fetch("http://subhost/foo");
1017 ` let txt = await resp.text();
1018 ` return new Response(
1019 ` "sub X-Foo header: " + resp.headers.get("X-Foo") + "\n" +
1020 ` "sub body: " + txt);
1021 ` }
1022 `}
1023 )
1024 ]
1025 ))"_kj));
1026 
1027 test.start();
1028 auto conn = test.connect("test-addr");
1029 conn.sendHttpGet("/");
1030 
1031 auto subreq = test.receiveInternetSubrequest("subhost");
1032 subreq.recv(R"(
1033 GET /foo HTTP/1.1
1034 Host: subhost
1035 
1036 )"_blockquote);
1037 subreq.send(R"(
1038 HTTP/1.1 200 OK
1039 Content-Length: 6
1040 X-Foo: bar
1041 
1042 corge
1043 )"_blockquote);
1044 
1045 conn.recvHttp200(R"(
1046 sub X-Foo header: bar
1047 sub body: corge
1048 )"_blockquote);
1049}
1050 
1051KJ_TEST("Server: override 'internet' service") {
1052 TestServer test(R"((
1053 services = [
1054 ( name = "hello",
1055 worker = (
1056 compatibilityDate = "2022-08-17",
1057 modules = [
1058 ( name = "main.js",
1059 esModule =
1060 `export default {
1061 ` async fetch(request, env) {
1062 ` return fetch(request);
1063 ` }
1064 `}
1065 )
1066 ]
1067 )
1068 ),
1069 ( name = "internet",
1070 external = "proxy-host" )
1071 ],
1072 sockets = [
1073 ( name = "main",
1074 address = "test-addr",
1075 service = "hello"
1076 )
1077 ]
1078 ))"_kj);
1079 
1080 test.start();
1081 auto conn = test.connect("test-addr");
1082 conn.sendHttpGet("/");
1083 
1084 auto subreq = test.receiveSubrequest("proxy-host");
1085 subreq.recv(R"(
1086 GET / HTTP/1.1
1087 Host: foo
1088 
1089 )"_blockquote);
1090 subreq.send(R"(
1091 HTTP/1.1 200 OK
1092 Content-Length: 2
1093 Content-Type: text/plain;charset=UTF-8
1094 
1095 OK
1096 )"_blockquote);
1097 
1098 conn.recvHttp200("OK");
1099}
1100 
1101KJ_TEST("Server: override globalOutbound") {
1102 TestServer test(R"((
1103 services = [
1104 ( name = "hello",
1105 worker = (
1106 compatibilityDate = "2022-08-17",
1107 modules = [
1108 ( name = "main.js",
1109 esModule =
1110 `export default {
1111 ` async fetch(request, env) {
1112 ` return fetch(request);
1113 ` }
1114 `}
1115 )
1116 ],
1117 globalOutbound = "alternate-outbound"
1118 )
1119 ),
1120 ( name = "alternate-outbound",
1121 external = "proxy-host" )
1122 ],
1123 sockets = [
1124 ( name = "main",
1125 address = "test-addr",
1126 service = "hello"
1127 )
1128 ]
1129 ))"_kj);
1130 
1131 test.start();
1132 auto conn = test.connect("test-addr");
1133 conn.sendHttpGet("/");
1134 
1135 auto subreq = test.receiveSubrequest("proxy-host");
1136 subreq.recv(R"(
1137 GET / HTTP/1.1
1138 Host: foo
1139 
1140 )"_blockquote);
1141 subreq.send(R"(
1142 HTTP/1.1 200 OK
1143 Content-Length: 2
1144 Content-Type: text/plain;charset=UTF-8
1145 
1146 OK
1147 )"_blockquote);
1148 
1149 conn.recvHttp200("OK");
1150}
1151 
1152KJ_TEST("Server: connect() to default outbound") {
1153 TestServer test(singleWorker(R"((
1154 compatibilityDate = "2022-08-17",
1155 compatibilityFlags = ["nodejs_compat"],
1156 modules = [
1157 ( name = "main.js",
1158 esModule =
1159 `import { connect } from 'cloudflare:sockets';
1160 `import assert from 'node:assert';
1161 `
1162 `export default {
1163 ` async fetch(request, env) {
1164 ` let sock = connect("subhost:123");
1165 `
1166 ` let writer = sock.writable.getWriter();
1167 ` await writer.write(new TextEncoder().encode("hello"));
1168 ` await writer.close();
1169 `
1170 ` let reader = sock.readable.getReader();
1171 ` let chunk = await reader.read();
1172 ` assert.strictEqual(chunk.done, false);
1173 ` assert.strictEqual(new TextDecoder().decode(chunk.value), "goodbye");
1174 `
1175 ` await sock.close();
1176 ` return new Response("OK");
1177 ` }
1178 `}
1179 )
1180 ]
1181 ))"_kj));
1182 
1183 test.start();
1184 auto conn = test.connect("test-addr");
1185 conn.sendHttpGet("/");
1186 
1187 auto subreq = test.receiveInternetSubrequest("subhost:123");
1188 subreq.recv("hello");
1189 subreq.send("goodbye");
1190 
1191 conn.recvHttp200("OK");
1192}
1193 
1194KJ_TEST("Server: connect() with Worker as outbound, no connect_pass_though") {
1195 TestServer test(R"((
1196 services = [
1197 ( name = "hello",
1198 worker = (
1199 compatibilityDate = "2022-08-17",
1200 compatibilityFlags = ["nodejs_compat"],
1201 globalOutbound = "outbound-worker",
1202 modules = [
1203 ( name = "main.js",
1204 esModule =
1205 `import { connect } from 'cloudflare:sockets';
1206 `import assert from 'node:assert';
1207 `
1208 `export default {
1209 ` async fetch(request, env) {
1210 ` // TODO(bug): At present this throws synchronously, which seems like a bug in
1211 ` // the implementation of connect(): errors coming from the destination
1212 ` // service really ought to be async (in prod, they always will be), showing
1213 ` // up on the first read or write. At present, though, I'm not looking to
1214 ` // fix this bug.
1215 ` assert.throws(() => connect("subhost:123"), {
1216 ` name: "TypeError",
1217 ` message: "Incoming CONNECT on a worker not supported",
1218 ` });
1219 `
1220 ` return new Response("OK");
1221 ` }
1222 `}
1223 )
1224 ]
1225 )
1226 ),
1227 ( name = "outbound-worker",
1228 worker = (
1229 compatibilityDate = "2022-08-17",
1230 modules = [
1231 ( name = "main.js",
1232 esModule =
1233 `export default {
1234 ` async fetch(request, env) {
1235 ` throw new Error("HTTP not expected");
1236 ` }
1237 `}
1238 )
1239 ]
1240 )
1241 ),
1242 ],
1243 sockets = [
1244 ( name = "main",
1245 address = "test-addr",
1246 service = "hello"
1247 )
1248 ]
1249 ))"_kj);
1250 
1251 test.server.allowExperimental();
1252 test.start();
1253 auto conn = test.connect("test-addr");
1254 conn.sendHttpGet("/");
1255 
1256 conn.recvHttp200("OK");
1257}
1258 
1259KJ_TEST("Server: connect() with Worker as outbound, with connect_pass_though") {
1260 TestServer test(R"((
1261 services = [
1262 ( name = "hello",
1263 worker = (
1264 compatibilityDate = "2022-08-17",
1265 compatibilityFlags = ["nodejs_compat"],
1266 globalOutbound = "outbound-worker",
1267 modules = [
1268 ( name = "main.js",
1269 esModule =
1270 `import { connect } from 'cloudflare:sockets';
1271 `import assert from 'node:assert';
1272 `
1273 `export default {
1274 ` async fetch(request, env) {
1275 ` let sock = connect("subhost:123");
1276 `
1277 ` let writer = sock.writable.getWriter();
1278 ` await writer.write(new TextEncoder().encode("hello"));
1279 ` await writer.close();
1280 `
1281 ` let reader = sock.readable.getReader();
1282 ` let chunk = await reader.read();
1283 ` assert.strictEqual(chunk.done, false);
1284 ` assert.strictEqual(new TextDecoder().decode(chunk.value), "goodbye");
1285 `
1286 ` await sock.close();
1287 ` return new Response("OK");
1288 ` }
1289 `}
1290 )
1291 ]
1292 )
1293 ),
1294 ( name = "outbound-worker",
1295 worker = (
1296 compatibilityDate = "2022-08-17",
1297 compatibilityFlags = ["connect_pass_through"],
1298 modules = [
1299 ( name = "main.js",
1300 esModule =
1301 `export default {
1302 ` async fetch(request, env) {
1303 ` throw new Error("HTTP not expected");
1304 ` }
1305 `}
1306 )
1307 ]
1308 )
1309 ),
1310 ],
1311 sockets = [
1312 ( name = "main",
1313 address = "test-addr",
1314 service = "hello"
1315 )
1316 ]
1317 ))"_kj);
1318 
1319 test.server.allowExperimental();
1320 test.start();
1321 auto conn = test.connect("test-addr");
1322 conn.sendHttpGet("/");
1323 
1324 auto subreq = test.receiveInternetSubrequest("subhost:123");
1325 subreq.recv("hello");
1326 subreq.send("goodbye");
1327 
1328 conn.recvHttp200("OK");
1329}
1330 
1331KJ_TEST("Server: capability bindings") {
1332 TestServer test(R"((
1333 services = [
1334 ( name = "hello",
1335 worker = (
1336 compatibilityDate = "2022-08-17",
1337 modules = [
1338 ( name = "main.js",
1339 esModule =
1340 `export default {
1341 ` async fetch(request, env) {
1342 ` let items = [];
1343 ` items.push(await (await env.fetcher.fetch("http://foo")).text());
1344 ` items.push(await env.kv.get("bar"));
1345 ` items.push(await (await env.r2.get("baz")).text());
1346 ` await env.queue.send("hello");
1347 ` items.push("Hello from Queue\n");
1348 ` const connection = await env.hyperdrive.connect();
1349 ` const encoded = new TextEncoder().encode("hyperdrive-test");
1350 ` await connection.writable.getWriter().write(new Uint8Array(encoded));
1351 ` items.push(`Hello from Hyperdrive(${env.hyperdrive.user})\n`);
1352 ` return new Response(items.join(""));
1353 ` }
1354 `}
1355 )
1356 ],
1357 bindings = [
1358 ( name = "fetcher",
1359 service = "service-outbound"
1360 ),
1361 ( name = "kv",
1362 kvNamespace = "kv-outbound"
1363 ),
1364 ( name = "r2",
1365 r2Bucket = "r2-outbound"
1366 ),
1367 ( name = "queue",
1368 queue = "queue-outbound"
1369 ),
1370 ( name = "hyperdrive",
1371 hyperdrive = (
1372 designator = "hyperdrive-outbound",
1373 database = "test-db",
1374 user = "test-user",
1375 password = "test-password",
1376 scheme = "postgresql"
1377 )
1378 )
1379 ]
1380 )
1381 ),
1382 ( name = "service-outbound", external = "service-host" ),
1383 ( name = "kv-outbound", external = "kv-host" ),
1384 ( name = "r2-outbound", external = "r2-host" ),
1385 ( name = "queue-outbound", external = "queue-host" ),
1386 ( name = "hyperdrive-outbound", external = (
1387 address = "hyperdrive-host",
1388 tcp = ()
1389 ))
1390 ],
1391 sockets = [
1392 ( name = "main",
1393 address = "test-addr",
1394 service = "hello"
1395 )
1396 ]
1397 ))"_kj);
1398 
1399 test.start();
1400 auto conn = test.connect("test-addr");
1401 conn.sendHttpGet("/");
1402 
1403 {
1404 auto subreq = test.receiveSubrequest("service-host");
1405 subreq.recv(R"(
1406 GET / HTTP/1.1
1407 Host: foo
1408 
1409 )"_blockquote);
1410 subreq.send(R"(
1411 HTTP/1.1 200 OK
1412 Content-Length: 16
1413 Content-Type: text/plain;charset=UTF-8
1414 
1415 Hello from HTTP
1416 )"_blockquote);
1417 }
1418 
1419 {
1420 auto subreq = test.receiveSubrequest("kv-host");
1421 subreq.recv(R"(
1422 GET /bar?urlencoded=true HTTP/1.1
1423 Host: fake-host
1424 CF-KV-FLPROD-405: https://fake-host/bar?urlencoded=true
1425 
1426 )"_blockquote);
1427 subreq.send(R"(
1428 HTTP/1.1 200 OK
1429 Content-Length: 14
1430 
1431 Hello from KV
1432 )"_blockquote);
1433 }
1434 
1435 {
1436 auto subreq = test.receiveSubrequest("r2-host");
1437 subreq.recv(R"(
1438 GET / HTTP/1.1
1439 Host: fake-host
1440 CF-R2-Request: {"version":1,"method":"get","object":"baz"}
1441 
1442 )"_blockquote);
1443 subreq.send(R"(
1444 HTTP/1.1 200 OK
1445 Content-Length: 16
1446 CF-R2-Metadata-Size: 2
1447 
1448 {}Hello from R2
1449 )"_blockquote);
1450 }
1451 
1452 {
1453 auto subreq = test.receiveSubrequest("queue-host");
1454 // We use a regex match to avoid dealing with the non-text characters in the POST body (which
1455 // may change as v8 serialization versions change over time).
1456 subreq.recvRegex(R"(
1457 POST /message HTTP/1.1
1458 Content-Length: 9
1459 Host: fake-host
1460 Content-Type: application/octet-stream
1461 
1462 .+hello)"_blockquote);
1463 subreq.send(R"(
1464 HTTP/1.1 200 OK
1465 Content-Type: application/json
1466 Content-Length: 27
1467 
1468 {"metadata":{"metrics":{}}}
1469 )"_blockquote);
1470 }
1471 
1472 {
1473 auto subreq = test.receiveSubrequest("hyperdrive-host");
1474 subreq.recv("hyperdrive-test");
1475 }
1476 conn.recvHttp200(R"(
1477 Hello from HTTP
1478 Hello from KV
1479 Hello from R2
1480 Hello from Queue
1481 Hello from Hyperdrive(test-user)
1482 )"_blockquote);
1483}
1484 
1485KJ_TEST("Server: cyclic bindings") {
1486 TestServer test(R"((
1487 services = [
1488 ( name = "service1",
1489 worker = (
1490 compatibilityDate = "2022-08-17",
1491 modules = [
1492 ( name = "main.js",
1493 esModule =
1494 `export default {
1495 ` async fetch(request, env) {
1496 ` if (request.url.endsWith("/done")) {
1497 ` return new Response("!");
1498 ` } else {
1499 ` let resp2 = await env.service2.fetch(request);
1500 ` let text = await resp2.text();
1501 ` return new Response("Hello " + text);
1502 ` }
1503 ` }
1504 `}
1505 )
1506 ],
1507 bindings = [(name = "service2", service = "service2")]
1508 )
1509 ),
1510 ( name = "service2",
1511 worker = (
1512 compatibilityDate = "2022-08-17",
1513 modules = [
1514 ( name = "main.js",
1515 esModule =
1516 `export default {
1517 ` async fetch(request, env) {
1518 ` let resp2 = await env.service1.fetch("http://foo/done");
1519 ` let text = await resp2.text();
1520 ` return new Response("World" + text);
1521 ` }
1522 `}
1523 )
1524 ],
1525 bindings = [(name = "service1", service = "service1")]
1526 )
1527 ),
1528 ],
1529 sockets = [
1530 ( name = "main",
1531 address = "test-addr",
1532 service = "service1"
1533 )
1534 ]
1535 ))"_kj);
1536 
1537 test.start();
1538 auto conn = test.connect("test-addr");
1539 conn.httpGet200("/", "Hello World!");
1540}
1541 
1542KJ_TEST("Server: named entrypoints") {
1543 TestServer test(R"((
1544 services = [
1545 ( name = "hello",
1546 worker = (
1547 compatibilityDate = "2022-08-17",
1548 modules = [
1549 ( name = "main.js",
1550 esModule =
1551 `export default {
1552 ` async fetch(request, env) {
1553 ` return new Response("hello from default entrypoint");
1554 ` }
1555 `}
1556 `export let foo = {
1557 ` async fetch(request, env) {
1558 ` return new Response("hello from foo entrypoint");
1559 ` }
1560 `}
1561 `export let bar = {
1562 ` async fetch(request, env) {
1563 ` return new Response("hello from bar entrypoint");
1564 ` }
1565 `}
1566 `
1567 `// Also export some symbols that aren't valid entrypoints, but we should still
1568 `// be allowed to point sockets at them. (Sending any actual requests to them
1569 `// will still fail.)
1570 `export let invalidObj = {}; // no handlers
1571 `export let invalidArray = [1, 2];
1572 `export let invalidMap = new Map();
1573 )
1574 ]
1575 )
1576 ),
1577 ],
1578 sockets = [
1579 ( name = "main", address = "test-addr", service = "hello" ),
1580 ( name = "alt1", address = "foo-addr", service = (name = "hello", entrypoint = "foo")),
1581 ( name = "alt2", address = "bar-addr", service = (name = "hello", entrypoint = "bar")),
1582 
1583 ( name = "invalid1", address = "invalid1-addr",
1584 service = (name = "hello", entrypoint = "invalidObj")),
1585 ( name = "invalid2", address = "invalid2-addr",
1586 service = (name = "hello", entrypoint = "invalidArray")),
1587 ( name = "invalid3", address = "invalid3-addr",
1588 service = (name = "hello", entrypoint = "invalidMap")),
1589 ]
1590 ))"_kj);
1591 
1592 test.start();
1593 
1594 {
1595 auto conn = test.connect("test-addr");
1596 conn.httpGet200("/", "hello from default entrypoint");
1597 }
1598 
1599 {
1600 auto conn = test.connect("foo-addr");
1601 conn.httpGet200("/", "hello from foo entrypoint");
1602 }
1603 
1604 {
1605 auto conn = test.connect("bar-addr");
1606 conn.httpGet200("/", "hello from bar entrypoint");
1607 }
1608}
1609 
1610KJ_TEST("Server: invalid entrypoint") {
1611 TestServer test(R"((
1612 services = [
1613 ( name = "hello",
1614 worker = (
1615 compatibilityDate = "2022-08-17",
1616 modules = [
1617 ( name = "main.js",
1618 esModule =
1619 `export default {
1620 ` async fetch(request, env) {
1621 ` return env.svc.fetch(request);
1622 ` }
1623 `}
1624 )
1625 ],
1626 bindings = [(name = "svc", service = (name = "hello", entrypoint = "bar"))],
1627 )
1628 ),
1629 ],
1630 sockets = [
1631 ( name = "main", address = "test-addr", service = "hello" ),
1632 ( name = "alt1", address = "foo-addr", service = (name = "hello", entrypoint = "foo")),
1633 ]
1634 ))"_kj);
1635 
1636 test.expectErrors(
1637 "Worker \"hello\"'s binding \"svc\" refers to service \"hello\" with a named entrypoint "
1638 "\"bar\", but \"hello\" has no such named entrypoint.\n"
1639 "Socket \"alt1\" refers to service \"hello\" with a named entrypoint \"foo\", but \"hello\" "
1640 "has no such named entrypoint.\n");
1641}
1642 
1643KJ_TEST("Server: referencing non-extant default entrypoint is not an error") {
1644 // For historical reasons, it's not a config error to refer to to the default entrypoint of
1645 // a service that has no default export.
1646 TestServer test(R"((
1647 services = [
1648 ( name = "hello",
1649 worker = (
1650 compatibilityDate = "2022-08-17",
1651 modules = [
1652 ( name = "main.js",
1653 esModule =
1654 `export let alt = {
1655 ` async fetch(request, env) {
1656 ` return new Response("OK");
1657 ` }
1658 `}
1659 )
1660 ],
1661 )
1662 ),
1663 ],
1664 sockets = [
1665 ( name = "main", address = "test-addr", service = "hello" ),
1666 ]
1667 ))"_kj);
1668 test.start();
1669 
1670 // A request will still fail at runtime, but we shouldn't have seen startup/config errors.
1671 auto conn = test.connect("test-addr");
1672 conn.sendHttpGet("/");
1673 
1674 // Due to the Deep Magic (bugs) going back to the dawn of Module Workers, if an HTTP request is
1675 // delivered to the default entrypoint of a module worker that has no default export, then the
1676 // system will fall back to calling event handlers registered with addEventListener("fetch").
1677 //
1678 // There is a magic deeper still in which, due to mistakes introduced in the stillness and the
1679 // darkness before Module Workers dawned, if none of those event listeners call
1680 // `event.respondWith()` (perhaps because *there are no event listeners*), then the request falls
1681 // back to default handling, in which it simply passes through to fetch() and makes a subrequest.
1682 //
1683 // So... we expect... a subrequest...
1684 {
1685 auto subreq = test.receiveSubrequest("foo", {"public"});
1686 subreq.recv(R"(
1687 GET / HTTP/1.1
1688 Host: foo
1689 
1690 )"_blockquote);
1691 subreq.send(R"(
1692 HTTP/1.1 200 OK
1693 Content-Length: 3
1694 
1695 wat)"_blockquote);
1696 }
1697 
1698 conn.recv(R"(
1699 HTTP/1.1 200 OK
1700 Content-Length: 3
1701 
1702 wat)"_blockquote);
1703}
1704 
1705KJ_TEST("Server: referencing DO class as entrypoint is not an error") {
1706 // For historical reasons, it's not a config error to refer to an actor class as a stateless
1707 // entrypoint.
1708 TestServer test(R"((
1709 services = [
1710 ( name = "hello",
1711 worker = (
1712 compatibilityDate = "2022-08-17",
1713 modules = [
1714 ( name = "main.js",
1715 esModule =
1716 `import { DurableObject } from "cloudflare:workers"
1717 `
1718 `export class SomeActor extends DurableObject {}
1719 `
1720 `export default {
1721 ` async fetch(request, env) {
1722 ` return new Response("OK");
1723 ` }
1724 `}
1725 )
1726 ],
1727 )
1728 ),
1729 ],
1730 sockets = [
1731 ( name = "main",
1732 address = "test-addr",
1733 service = (name = "hello", entrypoint = "SomeActor")
1734 ),
1735 ]
1736 ))"_kj);
1737 
1738 // We see a log warning at config time, but config otherwise completes successfully.
1739 {
1740 // TODO(soon): Restore this warning once miniflare no longer generates config that causes
1741 // it to log spuriously.
1742 //
1743 // KJ_EXPECT_LOG(WARNING,
1744 // "A ServiceDesignator in the config referenced the entrypoint \"SomeActor\", but this "
1745 // "class does not extend 'WorkerEntrypoint'. Attempts to call this entrypoint will "
1746 // "fail at runtime, but historically this was not a startup-time error. Future "
1747 // "versions of workerd may make this a startup-time error.");
1748 test.start();
1749 }
1750 
1751 // However, a request will still fail at runtime.
1752 KJ_EXPECT_LOG(ERROR, "worker is not an actor but class name was requested");
1753 KJ_EXPECT_LOG(INFO, "Unable to get exported handler");
1754 KJ_EXPECT_LOG(ERROR, "Unable to get exported handler");
1755 
1756 auto conn = test.connect("test-addr");
1757 conn.sendHttpGet("/");
1758 conn.recv(R"(
1759 HTTP/1.1 500 Internal Server Error
1760 Connection: close
1761 Content-Length: 21
1762 
1763 Internal Server Error)"_blockquote);
1764}
1765 
1766KJ_TEST("Server: exporting a DO class as the default export is not an error") {
1767 // For historical reasons, it's not a config error to export a DO class as the default
1768 // entrypoint. It doesn't work at runtime, but it's not a config error.
1769 TestServer test(R"((
1770 services = [
1771 ( name = "hello",
1772 worker = (
1773 compatibilityDate = "2022-08-17",
1774 modules = [
1775 ( name = "main.js",
1776 esModule =
1777 `import { DurableObject } from "cloudflare:workers"
1778 `
1779 `export default class extends DurableObject {
1780 ` async fetch(request) {
1781 ` return new Response("this should not be called");
1782 ` }
1783 `}
1784 )
1785 ],
1786 )
1787 ),
1788 ],
1789 sockets = [
1790 ( name = "main",
1791 address = "test-addr",
1792 service = "hello"
1793 ),
1794 ]
1795 ))"_kj);
1796 
1797 // We see a log error at config time, but config otherwise completes successfully.
1798 {
1799 KJ_EXPECT_LOG(ERROR,
1800 "Exported actor class as default entrypoint. This doesn't work, but historically "
1801 "did not produce a startup-time error.");
1802 test.start();
1803 }
1804 
1805 // Note that there is no way to actually configure the default export as a DO class since
1806 // `className` is non-optional in both `DurableObjectNamespace` and
1807 // `DurableObjectNamespaceDesignator`.
1808 //
1809 // We can, however, try to send a stateless request to the default entrypoint and see what
1810 // happens!
1811 //
1812 // Since the runtime does not believe there is any (stateless) entrypoint exported as the
1813 // default entrypoint, if you try to send a request to it, it behaves the same as if there were
1814 // no `export default` at all.
1815 //
1816 // The behavior of this is quite strange. See the comment in the earlier test:
1817 //
1818 // KJ_TEST("Server: referencing non-extant default entrypoint is not an error")
1819 auto conn = test.connect("test-addr");
1820 conn.sendHttpGet("/");
1821 
1822 {
1823 auto subreq = test.receiveSubrequest("foo", {"public"});
1824 subreq.recv(R"(
1825 GET / HTTP/1.1
1826 Host: foo
1827 
1828 )"_blockquote);
1829 subreq.send(R"(
1830 HTTP/1.1 200 OK
1831 Content-Length: 3
1832 
1833 wat)"_blockquote);
1834 }
1835 
1836 conn.recv(R"(
1837 HTTP/1.1 200 OK
1838 Content-Length: 3
1839 
1840 wat)"_blockquote);
1841}
1842 
1843KJ_TEST("Server: configuring a DO namespace with no class export is not an error") {
1844 // For historical reasons, it's not a config error to configure a DO namespace when there is
1845 // no corresponding class export.
1846 TestServer test(R"((
1847 services = [
1848 ( name = "hello",
1849 worker = (
1850 compatibilityDate = "2022-08-17",
1851 modules = [
1852 ( name = "main.js",
1853 esModule =
1854 `export default {
1855 ` async fetch(request, env) {
1856 ` return env.ns.get(env.ns.newUniqueId()).fetch(request);
1857 ` //return new Response("OK");
1858 ` }
1859 `}
1860 )
1861 ],
1862 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
1863 durableObjectNamespaces = [
1864 ( className = "MyActorClass",
1865 uniqueKey = "mykey",
1866 )
1867 ],
1868 durableObjectStorage = (inMemory = void)
1869 )
1870 ),
1871 ],
1872 sockets = [
1873 ( name = "main",
1874 address = "test-addr",
1875 service = "hello"
1876 ),
1877 ]
1878 ))"_kj);
1879 
1880 // We see a log warning at config time, but config otherwise completes successfully.
1881 {
1882 KJ_EXPECT_LOG(WARNING,
1883 "A DurableObjectNamespace in the config referenced the class \"MyActorClass\", but "
1884 "no such Durable Object class is exported from the worker. Please make sure the "
1885 "class name matches, it is exported, and the class extends 'DurableObject'. "
1886 "Attempts to call to this Durable Object class will fail at runtime, but historically "
1887 "this was not a startup-time error. Future versions of workerd may make this a "
1888 "startup-time error.");
1889 test.start();
1890 }
1891 
1892 // However, a request will still fail at runtime.
1893 KJ_EXPECT_LOG(ERROR, "no such actor class");
1894 KJ_EXPECT_LOG(INFO, "internal error");
1895 KJ_EXPECT_LOG(INFO, "internal error");
1896 KJ_EXPECT_LOG(ERROR, "internal error");
1897 
1898 auto conn = test.connect("test-addr");
1899 conn.sendHttpGet("/");
1900 conn.recv(R"(
1901 HTTP/1.1 500 Internal Server Error
1902 Connection: close
1903 Content-Length: 21
1904 
1905 Internal Server Error)"_blockquote);
1906}
1907 
1908KJ_TEST("Server: call queue handler on service binding") {
1909 TestServer test(R"((
1910 services = [
1911 ( name = "service1",
1912 worker = (
1913 compatibilityDate = "2022-08-17",
1914 compatibilityFlags = ["service_binding_extra_handlers"],
1915 modules = [
1916 ( name = "main.js",
1917 esModule =
1918 `export default {
1919 ` async fetch(request, env) {
1920 ` let result = await env.service2.queue("queueName1", [
1921 ` {id: "1", timestamp: 12345, body: "my message", attempts: 1},
1922 ` {id: "msg2", timestamp: 23456, body: 22, attempts: 2},
1923 ` ]);
1924 ` return new Response(`queue outcome: ${result.outcome}, ackAll: ${result.ackAll}`);
1925 ` }
1926 `}
1927 )
1928 ],
1929 bindings = [(name = "service2", service = "service2")]
1930 )
1931 ),
1932 ( name = "service2",
1933 worker = (
1934 compatibilityDate = "2022-08-17",
1935 modules = [
1936 ( name = "main.js",
1937 esModule =
1938 `export default {
1939 ` async fetch(request, env) {
1940 ` throw new Error("unimplemented");
1941 ` },
1942 ` async queue(event) {
1943 ` if (event.queue == "queueName1" &&
1944 ` event.messages.length == 2 &&
1945 ` event.messages[0].id == "1" &&
1946 ` event.messages[0].timestamp.getTime() == 12345 &&
1947 ` event.messages[0].body == "my message" &&
1948 ` event.messages[0].attempts == 1 &&
1949 ` event.messages[1].id == "msg2" &&
1950 ` event.messages[1].timestamp.getTime() == 23456 &&
1951 ` event.messages[1].body == 22 &&
1952 ` event.messages[1].attempts == 2) {
1953 ` event.ackAll();
1954 ` return;
1955 ` }
1956 ` throw new Error("messages didn't match expectations: " + JSON.stringify(event.messages));
1957 ` }
1958 `}
1959 )
1960 ]
1961 )
1962 ),
1963 ],
1964 sockets = [
1965 ( name = "main",
1966 address = "test-addr",
1967 service = "service1"
1968 )
1969 ]
1970 ))"_kj);
1971 
1972 test.server.allowExperimental();
1973 test.start();
1974 auto conn = test.connect("test-addr");
1975 conn.httpGet200("/", "queue outcome: ok, ackAll: true");
1976}
1977 
1978KJ_TEST("Server: Durable Objects (in memory)") {
1979 TestServer test(R"((
1980 services = [
1981 ( name = "hello",
1982 worker = (
1983 compatibilityDate = "2022-08-17",
1984 modules = [
1985 ( name = "main.js",
1986 esModule =
1987 `export default {
1988 ` async fetch(request, env) {
1989 ` let id = env.ns.idFromName(request.url)
1990 ` let actor = env.ns.get(id)
1991 ` return await actor.fetch(request)
1992 ` }
1993 `}
1994 `export class MyActorClass {
1995 ` constructor(state, env) {
1996 ` this.storage = state.storage;
1997 ` this.id = state.id;
1998 ` if (this.id.constructor.name != "DurableObjectId") {
1999 ` throw new Error("durable ID should be type DurableObjectId, " +
2000 ` `got: ${this.id.constructor.name}`);
2001 ` }
2002 ` if (typeof this.id.name !== "string" || this.id.name.length === 0) {
2003 ` throw new Error("ctx.id.name should be a non-empty string for " +
2004 ` `named DOs, got: ${JSON.stringify(this.id.name)}`);
2005 ` }
2006 ` }
2007 ` async fetch(request) {
2008 ` let count = (await this.storage.get("foo")) || 0;
2009 ` this.storage.put("foo", count + 1);
2010 ` return new Response(this.id + ": " + request.url + " " + count);
2011 ` }
2012 `}
2013 )
2014 ],
2015 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2016 durableObjectNamespaces = [
2017 ( className = "MyActorClass",
2018 uniqueKey = "mykey",
2019 )
2020 ],
2021 durableObjectStorage = (inMemory = void)
2022 )
2023 ),
2024 ],
2025 sockets = [
2026 ( name = "main",
2027 address = "test-addr",
2028 service = "hello"
2029 )
2030 ]
2031 ))"_kj);
2032 
2033 test.start();
2034 auto conn = test.connect("test-addr");
2035 conn.httpGet200(
2036 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 0");
2037 conn.httpGet200(
2038 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 1");
2039 conn.httpGet200(
2040 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 2");
2041 conn.httpGet200(
2042 "/bar", "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79: http://foo/bar 0");
2043 conn.httpGet200(
2044 "/bar", "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79: http://foo/bar 1");
2045 conn.httpGet200(
2046 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 3");
2047 conn.httpGet200(
2048 "/bar", "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79: http://foo/bar 2");
2049}
2050 
2051KJ_TEST("Server: Durable Objects keep ctx.id.name undefined for unique IDs") {
2052 TestServer test(R"((
2053 services = [
2054 ( name = "hello",
2055 worker = (
2056 compatibilityDate = "2022-08-17",
2057 modules = [
2058 ( name = "main.js",
2059 esModule =
2060 `export default {
2061 ` async fetch(request, env) {
2062 ` let actor = env.ns.get(env.ns.newUniqueId())
2063 ` return await actor.fetch(request)
2064 ` }
2065 `}
2066 `export class MyActorClass {
2067 ` constructor(state, env) {
2068 ` this.name = state.id.name;
2069 ` }
2070 ` async fetch(request) {
2071 ` return new Response(String(this.name));
2072 ` }
2073 `}
2074 )
2075 ],
2076 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2077 durableObjectNamespaces = [
2078 ( className = "MyActorClass",
2079 uniqueKey = "mykey",
2080 )
2081 ],
2082 durableObjectStorage = (inMemory = void)
2083 )
2084 ),
2085 ],
2086 sockets = [
2087 ( name = "main",
2088 address = "test-addr",
2089 service = "hello"
2090 )
2091 ]
2092 ))"_kj);
2093 
2094 test.start();
2095 auto conn = test.connect("test-addr");
2096 conn.httpGet200("/", "undefined");
2097}
2098 
2099KJ_TEST("Server: Durable Objects retain ctx.id.name for short names") {
2100 TestServer test(R"((
2101 services = [
2102 ( name = "hello",
2103 worker = (
2104 compatibilityDate = "2022-08-17",
2105 modules = [
2106 ( name = "main.js",
2107 esModule =
2108 `const name = "retained-name-123"
2109 `export default {
2110 ` async fetch(request, env) {
2111 ` let actor = env.ns.get(env.ns.idFromName(name))
2112 ` return await actor.fetch(request)
2113 ` }
2114 `}
2115 `export class MyActorClass {
2116 ` constructor(state, env) {
2117 ` this.name = state.id.name;
2118 ` }
2119 ` async fetch(request) {
2120 ` return new Response(String(this.name));
2121 ` }
2122 `}
2123 )
2124 ],
2125 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2126 durableObjectNamespaces = [
2127 ( className = "MyActorClass",
2128 uniqueKey = "mykey",
2129 )
2130 ],
2131 durableObjectStorage = (inMemory = void)
2132 )
2133 ),
2134 ],
2135 sockets = [
2136 ( name = "main",
2137 address = "test-addr",
2138 service = "hello"
2139 )
2140 ]
2141 ))"_kj);
2142 
2143 test.start();
2144 auto conn = test.connect("test-addr");
2145 conn.httpGet200("/", "retained-name-123");
2146}
2147 
2148KJ_TEST("Server: Durable Objects drop ctx.id.name for long names") {
2149 TestServer test(R"((
2150 services = [
2151 ( name = "hello",
2152 worker = (
2153 compatibilityDate = "2022-08-17",
2154 modules = [
2155 ( name = "main.js",
2156 esModule =
2157 `export default {
2158 ` async fetch(request, env) {
2159 ` let actor = env.ns.get(env.ns.idFromName("a".repeat(1025)))
2160 ` return await actor.fetch(request)
2161 ` }
2162 `}
2163 `export class MyActorClass {
2164 ` constructor(state, env) {
2165 ` this.name = state.id.name;
2166 ` }
2167 ` async fetch(request) {
2168 ` return new Response(String(this.name));
2169 ` }
2170 `}
2171 )
2172 ],
2173 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2174 durableObjectNamespaces = [
2175 ( className = "MyActorClass",
2176 uniqueKey = "mykey",
2177 )
2178 ],
2179 durableObjectStorage = (inMemory = void)
2180 )
2181 ),
2182 ],
2183 sockets = [
2184 ( name = "main",
2185 address = "test-addr",
2186 service = "hello"
2187 )
2188 ]
2189 ))"_kj);
2190 
2191 test.start();
2192 auto conn = test.connect("test-addr");
2193 conn.httpGet200("/", "undefined");
2194}
2195 
2196KJ_TEST("Server: Simultaneous requests to a DO that hasn't started don't cause split brain") {
2197 TestServer test(R"((
2198 services = [
2199 ( name = "hello",
2200 worker = (
2201 compatibilityDate = "2025-04-01",
2202 modules = [
2203 ( name = "main.js",
2204 esModule =
2205 `import {DurableObject} from "cloudflare:workers"
2206 `export default {
2207 ` async fetch(request, env) {
2208 ` let id = env.ns.idFromName(request.url)
2209 ` let actor = env.ns.get(id)
2210 ` let promise1 = actor.increment()
2211 ` let promise2 = actor.increment()
2212 ` let promise3 = actor.increment()
2213 ` return new Response(`${await promise1} ${await promise2} ${await promise3}`)
2214 ` }
2215 `}
2216 `export class Counter extends DurableObject {
2217 ` counter = 0;
2218 ` async increment() {
2219 ` return this.counter++;
2220 ` }
2221 `}
2222 )
2223 ],
2224 bindings = [(name = "ns", durableObjectNamespace = "Counter")],
2225 durableObjectNamespaces = [
2226 ( className = "Counter",
2227 uniqueKey = "mykey",
2228 )
2229 ],
2230 durableObjectStorage = (inMemory = void)
2231 )
2232 ),
2233 ],
2234 sockets = [
2235 ( name = "main",
2236 address = "test-addr",
2237 service = "hello"
2238 )
2239 ]
2240 ))"_kj);
2241 
2242 test.start();
2243 auto conn = test.connect("test-addr");
2244 conn.httpGet200("/", "0 1 2");
2245}
2246 
2247KJ_TEST("Server: Broken DO stays broken until stub replaced") {
2248 TestServer test(R"((
2249 services = [
2250 ( name = "hello",
2251 worker = (
2252 compatibilityDate = "2025-04-01",
2253 modules = [
2254 ( name = "main.js",
2255 esModule =
2256 `import {DurableObject} from "cloudflare:workers"
2257 `export default {
2258 ` async fetch(request, env) {
2259 ` let id = env.ns.idFromName(request.url)
2260 ` let actor = env.ns.get(id)
2261 ` let i1 = await actor.increment()
2262 ` try { await actor.abort() } catch {}
2263 ` try {
2264 ` let i2 = await actor.increment();
2265 ` throw new Error(`expected error from broken stub, got ${i2}`);
2266 ` } catch (err) {
2267 ` if (!err.message.includes("test abort reason")) {
2268 ` throw err
2269 ` }
2270 ` }
2271 ` actor = env.ns.get(id)
2272 ` let i3 = await actor.increment()
2273 ` return new Response(`${i1} ${i3}`)
2274 ` }
2275 `}
2276 `export class Counter extends DurableObject {
2277 ` counter = 0;
2278 ` async increment() {
2279 ` return this.counter++;
2280 ` }
2281 ` async abort() {
2282 ` this.ctx.abort(new Error("test abort reason"));
2283 ` }
2284 `}
2285 )
2286 ],
2287 bindings = [(name = "ns", durableObjectNamespace = "Counter")],
2288 durableObjectNamespaces = [
2289 ( className = "Counter",
2290 uniqueKey = "mykey",
2291 )
2292 ],
2293 durableObjectStorage = (inMemory = void)
2294 )
2295 ),
2296 ],
2297 sockets = [
2298 ( name = "main",
2299 address = "test-addr",
2300 service = "hello"
2301 )
2302 ]
2303 ))"_kj);
2304 
2305 test.start();
2306 
2307 auto conn = test.connect("test-addr");
2308 conn.httpGet200("/", "0 0");
2309}
2310 
2311KJ_TEST("Server: Durable Objects (on disk)") {
2312 kj::StringPtr config = R"((
2313 services = [
2314 ( name = "hello",
2315 worker = (
2316 compatibilityDate = "2022-08-17",
2317 modules = [
2318 ( name = "main.js",
2319 esModule =
2320 `export default {
2321 ` async fetch(request, env) {
2322 ` let id = env.ns.idFromName(request.url)
2323 ` let actor = env.ns.get(id)
2324 ` return await actor.fetch(request)
2325 ` }
2326 `}
2327 `export class MyActorClass {
2328 ` constructor(state, env) {
2329 ` this.storage = state.storage;
2330 ` this.id = state.id;
2331 ` if (this.id.constructor.name != "DurableObjectId") {
2332 ` throw new Error("durable ID should be type DurableObjectId, " +
2333 ` `got: ${this.id.constructor.name}`);
2334 ` }
2335 ` }
2336 ` async fetch(request) {
2337 ` let count = (await this.storage.get("foo")) || 0;
2338 ` this.storage.put("foo", count + 1);
2339 ` return new Response(this.id + ": " + request.url + " " + count);
2340 ` }
2341 `}
2342 )
2343 ],
2344 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2345 durableObjectNamespaces = [
2346 ( className = "MyActorClass",
2347 uniqueKey = "mykey",
2348 )
2349 ],
2350 durableObjectStorage = (localDisk = "my-disk")
2351 )
2352 ),
2353 ( name = "my-disk",
2354 disk = (
2355 path = "../../var/do-storage",
2356 writable = true,
2357 )
2358 ),
2359 ],
2360 sockets = [
2361 ( name = "main",
2362 address = "test-addr",
2363 service = "hello"
2364 )
2365 ]
2366 ))"_kj;
2367 
2368 // Create a directory outside of the test scope which we can use across multiple TestServers.
2369 auto dir = kj::newInMemoryDirectory(kj::nullClock());
2370 
2371 {
2372 TestServer test(config);
2373 
2374 // Link our directory into the test filesystem.
2375 test.root->transfer(kj::Path({"var"_kj, "do-storage"_kj}),
2376 kj::WriteMode::CREATE | kj::WriteMode::CREATE_PARENT, *dir, nullptr,
2377 kj::TransferMode::LINK);
2378 
2379 test.start();
2380 auto conn = test.connect("test-addr");
2381 conn.httpGet200(
2382 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 0");
2383 conn.httpGet200(
2384 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 1");
2385 conn.httpGet200(
2386 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 2");
2387 conn.httpGet200("/bar",
2388 "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79: http://foo/bar 0");
2389 conn.httpGet200("/bar",
2390 "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79: http://foo/bar 1");
2391 conn.httpGet200(
2392 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 3");
2393 conn.httpGet200("/bar",
2394 "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79: http://foo/bar 2");
2395 
2396 // The storage directory contains .sqlite and .sqlite-wal files for both objects, plus the
2397 // per-namespace metadata.sqlite (alarm scheduler) and its WAL file. Note that the `-shm`
2398 // files are missing because SQLite doesn't actually tell the VFS to create these as separate
2399 // files, it leaves it up to the VFS to decide how shared memory works, and our KJ-wrapping
2400 // VFS currently doesn't put this in SHM files. If we were using a real disk directory,
2401 // though, they would be there.
2402 KJ_EXPECT(dir->openSubdir(kj::Path({"mykey"}))->listNames().size() == 6);
2403 KJ_EXPECT(dir->exists(kj::Path(
2404 {"mykey", "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79.sqlite"})));
2405 KJ_EXPECT(dir->exists(kj::Path(
2406 {"mykey", "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79.sqlite-wal"})));
2407 KJ_EXPECT(dir->exists(kj::Path(
2408 {"mykey", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234.sqlite"})));
2409 KJ_EXPECT(dir->exists(kj::Path(
2410 {"mykey", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234.sqlite-wal"})));
2411 KJ_EXPECT(dir->exists(kj::Path({"mykey", "metadata.sqlite"})));
2412 KJ_EXPECT(dir->exists(kj::Path({"mykey", "metadata.sqlite-wal"})));
2413 }
2414 
2415 // Having torn everything down, the WAL files should be gone.
2416 KJ_EXPECT(dir->openSubdir(kj::Path({"mykey"}))->listNames().size() == 3);
2417 KJ_EXPECT(dir->exists(kj::Path(
2418 {"mykey", "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79.sqlite"})));
2419 KJ_EXPECT(dir->exists(kj::Path(
2420 {"mykey", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234.sqlite"})));
2421 
2422 // Let's start a new server and verify it can load the files from disk.
2423 {
2424 TestServer test(config);
2425 
2426 // Link our directory into the test filesystem.
2427 test.root->transfer(kj::Path({"var"_kj, "do-storage"_kj}),
2428 kj::WriteMode::CREATE | kj::WriteMode::CREATE_PARENT, *dir, nullptr,
2429 kj::TransferMode::LINK);
2430 
2431 test.start();
2432 auto conn = test.connect("test-addr");
2433 conn.httpGet200(
2434 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 4");
2435 conn.httpGet200(
2436 "/", "59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234: http://foo/ 5");
2437 conn.httpGet200("/bar",
2438 "02b496f65dd35cbac90e3e72dc5a398ee93926ea4a3821e26677082d2e6f9b79: http://foo/bar 3");
2439 }
2440}
2441 
2442KJ_TEST("Server: Durable Object alarm persistence (on disk)") {
2443 kj::StringPtr config = R"((
2444 services = [
2445 ( name = "hello",
2446 worker = (
2447 compatibilityDate = "2024-01-01",
2448 modules = [
2449 ( name = "main.js",
2450 esModule =
2451 `export default {
2452 ` async fetch(request, env) {
2453 ` let id = env.ns.idFromName("alarm-actor")
2454 ` let actor = env.ns.get(id)
2455 ` return await actor.fetch(request)
2456 ` }
2457 `}
2458 `export class MyActorClass {
2459 ` constructor(state, env) {
2460 ` this.storage = state.storage;
2461 ` }
2462 ` async fetch(request) {
2463 ` let url = new URL(request.url);
2464 ` if (url.pathname === "/set") {
2465 ` let time = parseInt(url.searchParams.get("t"));
2466 ` await this.storage.setAlarm(time);
2467 ` return new Response("alarm set to " + time);
2468 ` } else if (url.pathname === "/get") {
2469 ` let alarm = await this.storage.getAlarm();
2470 ` return new Response("alarm=" + alarm);
2471 ` } else {
2472 ` return new Response("unknown path", {status: 404});
2473 ` }
2474 ` }
2475 ` async alarm() {}
2476 `}
2477 )
2478 ],
2479 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2480 durableObjectNamespaces = [
2481 ( className = "MyActorClass",
2482 uniqueKey = "alarmkey",
2483 )
2484 ],
2485 durableObjectStorage = (localDisk = "my-disk")
2486 )
2487 ),
2488 ( name = "my-disk",
2489 disk = (
2490 path = "../../var/do-storage",
2491 writable = true,
2492 )
2493 ),
2494 ],
2495 sockets = [
2496 ( name = "main",
2497 address = "test-addr",
2498 service = "hello"
2499 )
2500 ]
2501 ))"_kj;
2502 
2503 auto dir = kj::newInMemoryDirectory(kj::nullClock());
2504 
2505 // A far-future alarm time (won't fire during the test).
2506 kj::StringPtr alarmTime = "4102444800000";
2507 
2508 {
2509 TestServer test(config);
2510 test.root->transfer(kj::Path({"var"_kj, "do-storage"_kj}),
2511 kj::WriteMode::CREATE | kj::WriteMode::CREATE_PARENT, *dir, nullptr,
2512 kj::TransferMode::LINK);
2513 
2514 test.start();
2515 auto conn = test.connect("test-addr");
2516 
2517 conn.httpGet200(kj::str("/set?t=", alarmTime), kj::str("alarm set to ", alarmTime));
2518 conn.httpGet200("/get", kj::str("alarm=", alarmTime));
2519 }
2520 
2521 // Verify metadata.sqlite exists on disk in the namespace directory.
2522 KJ_EXPECT(dir->exists(kj::Path({"alarmkey", "metadata.sqlite"})));
2523 
2524 // Start a new server and verify the alarm is still there.
2525 {
2526 TestServer test(config);
2527 test.root->transfer(kj::Path({"var"_kj, "do-storage"_kj}),
2528 kj::WriteMode::CREATE | kj::WriteMode::CREATE_PARENT, *dir, nullptr,
2529 kj::TransferMode::LINK);
2530 
2531 test.start();
2532 auto conn = test.connect("test-addr");
2533 
2534 conn.httpGet200("/get", kj::str("alarm=", alarmTime));
2535 }
2536}
2537 
2538KJ_TEST("Server: Ephemeral Objects") {
2539 TestServer test(R"((
2540 services = [
2541 ( name = "hello",
2542 worker = (
2543 compatibilityDate = "2022-08-17",
2544 modules = [
2545 ( name = "main.js",
2546 esModule =
2547 `export default {
2548 ` async fetch(request, env) {
2549 ` let actor = env.ns.get(request.url)
2550 ` return await actor.fetch(request)
2551 ` }
2552 `}
2553 `export class MyActorClass {
2554 ` constructor(state, env) {
2555 ` if (state.storage) throw new Error("storage shouldn't be present");
2556 ` this.id = state.id;
2557 ` if (typeof this.id != "string") {
2558 ` throw new Error("ephemeral ID should be type string, " +
2559 ` `got: ${this.id.constructor.name}`);
2560 ` }
2561 ` this.count = 0;
2562 ` }
2563 ` async fetch(request) {
2564 ` return new Response(this.id + ": " + request.url + " " + this.count++);
2565 ` }
2566 `}
2567 )
2568 ],
2569 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2570 durableObjectNamespaces = [
2571 ( className = "MyActorClass",
2572 ephemeralLocal = void,
2573 )
2574 ],
2575 durableObjectStorage = (inMemory = void)
2576 )
2577 ),
2578 ],
2579 sockets = [
2580 ( name = "main",
2581 address = "test-addr",
2582 service = "hello"
2583 )
2584 ]
2585 ))"_kj);
2586 
2587 test.server.allowExperimental();
2588 test.start();
2589 auto conn = test.connect("test-addr");
2590 conn.httpGet200("/", "http://foo/: http://foo/ 0");
2591 conn.httpGet200("/", "http://foo/: http://foo/ 1");
2592 conn.httpGet200("/", "http://foo/: http://foo/ 2");
2593 conn.httpGet200("/bar", "http://foo/bar: http://foo/bar 0");
2594 conn.httpGet200("/bar", "http://foo/bar: http://foo/bar 1");
2595 conn.httpGet200("/", "http://foo/: http://foo/ 3");
2596 conn.httpGet200("/bar", "http://foo/bar: http://foo/bar 2");
2597}
2598 
2599KJ_TEST("Server: Durable Objects (ephemeral) eviction") {
2600 TestServer test(R"((
2601 services = [
2602 ( name = "hello",
2603 worker = (
2604 compatibilityDate = "2023-08-17",
2605 modules = [
2606 ( name = "main.js",
2607 esModule =
2608 `export default {
2609 ` async fetch(request, env) {
2610 ` let id = env.ns.idFromName("59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234");
2611 ` let obj = env.ns.get(id)
2612 ` if (request.url.endsWith("/setup")) {
2613 ` return await obj.fetch("http://example.com/setup");
2614 ` } else if (request.url.endsWith("/check")) {
2615 ` try {
2616 ` return await obj.fetch("http://example.com/check");
2617 ` } catch(e) {
2618 ` throw e;
2619 ` }
2620 ` } else if (request.url.endsWith("/checkEvicted")) {
2621 ` return await obj.fetch("http://example.com/checkEvicted");
2622 ` }
2623 ` return new Response("Invalid Route!")
2624 ` }
2625 `}
2626 `export class MyActorClass {
2627 ` constructor(state, env) {
2628 ` this.defaultMessage = false; // Set to true on first "setup" request
2629 ` }
2630 ` async fetch(request) {
2631 ` if (request.url.endsWith("/setup")) {
2632 ` // Request 1, set defaultMessage, will remain true as long as actor is live.
2633 ` this.defaultMessage = true;
2634 ` return new Response("OK");
2635 ` } else if (request.url.endsWith("/check")) {
2636 ` // Request 2, assert that actor is still in alive (defaultMessage is still true).
2637 ` if (this.defaultMessage) {
2638 ` // Actor is still alive and we did not re-run the constructor
2639 ` return new Response("OK");
2640 ` }
2641 ` throw new Error("Error: Actor was evicted!");
2642 ` } else if (request.url.endsWith("/checkEvicted")) {
2643 ` // Final request (3), check if the defaultMessage has been set to false,
2644 ` // indicating the actor was evicted
2645 ` if (!this.defaultMessage) {
2646 ` // Actor was evicted and we re-ran the constructor!
2647 ` return new Response("OK");
2648 ` }
2649 ` throw new Error("Error: Actor was not evicted! We were still alive.");
2650 ` }
2651 ` }
2652 `}
2653 )
2654 ],
2655 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2656 durableObjectNamespaces = [
2657 ( className = "MyActorClass",
2658 uniqueKey = "mykey",
2659 )
2660 ],
2661 durableObjectStorage = (inMemory = void)
2662 )
2663 ),
2664 ],
2665 sockets = [
2666 ( name = "main",
2667 address = "test-addr",
2668 service = "hello"
2669 )
2670 ]
2671 ))"_kj);
2672 
2673 test.start();
2674 auto conn = test.connect("test-addr");
2675 conn.httpGet200("/setup", "OK");
2676 conn.httpGet200("/check", "OK");
2677 
2678 // Force hibernation by waiting 10 seconds.
2679 test.wait(10);
2680 // Need a second connection because of 5 second HTTP timeout.
2681 auto connTwo = test.connect("test-addr");
2682 connTwo.httpGet200("/checkEvicted", "OK");
2683}
2684 
2685KJ_TEST("Server: Durable Objects (ephemeral) prevent eviction") {
2686 TestServer test(R"((
2687 services = [
2688 ( name = "hello",
2689 worker = (
2690 compatibilityDate = "2023-08-17",
2691 modules = [
2692 ( name = "main.js",
2693 esModule =
2694 `export default {
2695 ` async fetch(request, env) {
2696 ` let id = env.ns.idFromName("59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234");
2697 ` let obj = env.ns.get(id);
2698 ` if (request.url.endsWith("/setup")) {
2699 ` return await obj.fetch("http://example.com/setup");
2700 ` } else if (request.url.endsWith("/assertNotEvicted")) {
2701 ` try {
2702 ` return await obj.fetch("http://example.com/assertNotEvicted");
2703 ` } catch(e) {
2704 ` throw e;
2705 ` }
2706 ` }
2707 ` return new Response("Invalid Route!")
2708 ` }
2709 `}
2710 `export class MyActorClass {
2711 ` constructor(state, env) {
2712 ` this.defaultMessage = false; // Set to true on first "setup" request
2713 ` }
2714 ` async fetch(request) {
2715 ` if (request.url.endsWith("/setup")) {
2716 ` // Request 1, set defaultMessage, will remain true as long as actor is live.
2717 ` this.defaultMessage = true;
2718 ` return new Response("OK");
2719 ` } else if (request.url.endsWith("/assertNotEvicted")) {
2720 ` // Request 2, assert that actor is still in alive (defaultMessage is still true).
2721 ` if (this.defaultMessage) {
2722 ` // Actor is still alive and we did not re-run the constructor
2723 ` return new Response("OK");
2724 ` }
2725 ` throw new Error("Error: Actor was evicted!");
2726 ` }
2727 ` }
2728 `}
2729 )
2730 ],
2731 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2732 durableObjectNamespaces = [
2733 ( className = "MyActorClass",
2734 uniqueKey = "mykey",
2735 preventEviction = true,
2736 )
2737 ],
2738 durableObjectStorage = (inMemory = void)
2739 )
2740 ),
2741 ],
2742 sockets = [
2743 ( name = "main",
2744 address = "test-addr",
2745 service = "hello"
2746 )
2747 ]
2748 ))"_kj);
2749 
2750 test.start();
2751 auto conn = test.connect("test-addr");
2752 conn.httpGet200("/setup", "OK");
2753 conn.httpGet200("/assertNotEvicted", "OK");
2754 
2755 // Attempt to force hibernation by waiting 10 seconds.
2756 test.wait(10);
2757 // Need a second connection because of 5 second HTTP timeout.
2758 auto connTwo = test.connect("test-addr");
2759 connTwo.httpGet200("/assertNotEvicted", "OK");
2760}
2761 
2762KJ_TEST("Server: Durable Object evictions when callback scheduled") {
2763 kj::StringPtr config = R"((
2764 services = [
2765 ( name = "hello",
2766 worker = (
2767 compatibilityDate = "2023-08-17",
2768 modules = [
2769 ( name = "main.js",
2770 esModule =
2771 `export default {
2772 ` async fetch(request, env) {
2773 ` let id = env.ns.idFromName("59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234");
2774 ` let obj = env.ns.get(id)
2775 ` return await obj.fetch(request.url);
2776 ` }
2777 `}
2778 `export class MyActorClass {
2779 ` constructor(state, env) {
2780 ` this.defaultMessage = false; // Set to true on first "setup" request
2781 ` this.storage = state.storage;
2782 ` this.count = 0;
2783 ` }
2784 ` async fetch(request) {
2785 ` if (request.url.endsWith("/15Seconds")) {
2786 ` // Schedule a callback to run in 15 seconds.
2787 ` // The DO should NOT be evicted by the inactivity timeout before this runs.
2788 ` this.defaultMessage = true;
2789 ` let id = setInterval(() => { clearInterval(id); }, 15000);
2790 ` return new Response("OK");
2791 ` } else if (request.url.endsWith("/20Seconds")) {
2792 ` // Schedule a callback to run every 20 seconds.
2793 ` // The DO should expire after 70 seconds.
2794 ` this.defaultMessage = true;
2795 ` this.count = 0;
2796 ` await this.storage.put("count", this.count);
2797 ` let id = setInterval(() => {
2798 ` // Increment number of times we ran this.
2799 ` this.count += 1;
2800 ` this.storage.put("count", this.count);
2801 ` }, 20000);
2802 ` return new Response("OK");
2803 ` } else if (request.url.endsWith("/assertActive")) {
2804 ` // Assert that actor is still in alive (defaultMessage is still true).
2805 ` if (this.defaultMessage) {
2806 ` // Actor is still alive and we did not re-run the constructor
2807 ` return new Response("OK");
2808 ` }
2809 ` throw new Error("Error: Actor was evicted!");
2810 ` } else if (request.url.endsWith("/assertEvicted")) {
2811 ` // Check if the defaultMessage has been set to false,
2812 ` // indicating the actor was evicted
2813 ` if (!this.defaultMessage) {
2814 ` // Actor was evicted and we re-ran the constructor!
2815 ` return new Response("OK");
2816 ` }
2817 ` throw new Error("Error: Actor was not evicted! We were still alive.");
2818 ` } else if (request.url.endsWith("/assertEvictedAndCount")) {
2819 ` // Check if the defaultMessage has been set to false,
2820 ` // indicating the actor was evicted
2821 ` if (!this.defaultMessage) {
2822 ` var count = await this.storage.get("count");
2823 ` if (!(4 < count && count < 8)) {
2824 ` // Something must have gone wrong. We have a 70 sec expiration,
2825 ` // and worst case is it takes ~140 seconds to evict. The callback runs
2826 ` // every 20 seconds, so it has to be evicted before the 8th callback.
2827 ` throw new Error(`Callback ran ${count} times, expected between 4 to 8!`);
2828 ` }
2829 ` // Actor was evicted and we had the right count!
2830 ` return new Response("OK");
2831 ` }
2832 ` throw new Error("Error: Actor was not evicted! We were still alive.");
2833 ` }
2834 ` }
2835 `}
2836 )
2837 ],
2838 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2839 durableObjectNamespaces = [
2840 ( className = "MyActorClass",
2841 uniqueKey = "mykey",
2842 )
2843 ],
2844 durableObjectStorage = (localDisk = "my-disk")
2845 )
2846 ),
2847 ( name = "my-disk",
2848 disk = (
2849 path = "../../var/do-storage",
2850 writable = true,
2851 )
2852 ),
2853 ],
2854 sockets = [
2855 ( name = "main",
2856 address = "test-addr",
2857 service = "hello"
2858 )
2859 ]
2860 ))"_kj;
2861 
2862 // Create a directory outside of the test scope which we can use across multiple TestServers.
2863 auto dir = kj::newInMemoryDirectory(kj::nullClock());
2864 {
2865 TestServer test(config);
2866 // Link our directory into the test filesystem.
2867 test.root->transfer(kj::Path({"var"_kj, "do-storage"_kj}),
2868 kj::WriteMode::CREATE | kj::WriteMode::CREATE_PARENT, *dir, nullptr,
2869 kj::TransferMode::LINK);
2870 
2871 test.start();
2872 auto conn = test.connect("test-addr");
2873 // Setup a callback that will run in 15 seconds.
2874 // This callback should prevent the DO from being evicted.
2875 conn.httpGet200("/15Seconds", "OK");
2876 
2877 // If we weren't waiting on anything, the DO would be evicted after 10 seconds,
2878 // however, it will actually be evicted in 25 seconds (15 seconds until setInterval is cleared +
2879 // 10 seconds for inactivity timer).
2880 
2881 test.wait(15);
2882 // The `setInterval()` will be cleared around now. Let's verify that we didn't get evicted.
2883 
2884 // Need a new connection because of 5 second HTTP timeout.
2885 auto connTwo = test.connect("test-addr");
2886 connTwo.httpGet200("/assertActive", "OK");
2887 
2888 // Force hibernation by waiting at least 10 seconds since we haven't scheduled any new work.
2889 test.wait(10);
2890 
2891 // Need a new connection because of 5 second HTTP timeout.
2892 auto connThree = test.connect("test-addr");
2893 connThree.httpGet200("/assertEvicted", "OK");
2894 
2895 // Now we know we aren't evicting DOs early if they have future work scheduled. Next, let's
2896 // ensure we ARE evicting DOs if there are no connected clients for 70 seconds.
2897 // Note that the `/20seconds` path calls setInterval to run every 20 seconds, and never clears.
2898 auto connFour = test.connect("test-addr");
2899 connFour.httpGet200("/20Seconds", "OK");
2900 // It's unlikely, but the worst case is the cleanupLoop checks just before the 70 sec expiration,
2901 // and has to wait another 70 seconds before trying to remove again. We'll wait for 142 seconds
2902 // to account for this.
2903 test.wait(142);
2904 
2905 auto connFive = test.connect("test-addr");
2906 connFive.httpGet200("/assertEvictedAndCount", "OK");
2907 }
2908}
2909 
2910KJ_TEST("Server: Durable Objects websocket") {
2911 TestServer test(R"((
2912 services = [
2913 ( name = "hello",
2914 worker = (
2915 compatibilityDate = "2023-08-17",
2916 modules = [
2917 ( name = "main.js",
2918 esModule =
2919 `export default {
2920 ` async fetch(request, env) {
2921 ` let id = env.ns.idFromName("59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234");
2922 ` let obj = env.ns.get(id)
2923 ` return await obj.fetch(request);
2924 ` }
2925 `}
2926 `
2927 `export class MyActorClass {
2928 ` constructor(state) {}
2929 `
2930 ` async fetch(request) {
2931 ` let pair = new WebSocketPair();
2932 ` let ws = pair[1]
2933 ` ws.accept();
2934 `
2935 ` ws.addEventListener("message", (m) => {
2936 ` ws.send(m.data);
2937 ` });
2938 ` ws.addEventListener("close", (c) => {
2939 ` ws.close(c.code, c.reason);
2940 ` });
2941 `
2942 ` return new Response(null, {status: 101, statusText: "Switching Protocols", webSocket: pair[0]});
2943 ` }
2944 `}
2945 )
2946 ],
2947 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
2948 durableObjectNamespaces = [
2949 ( className = "MyActorClass",
2950 uniqueKey = "mykey",
2951 )
2952 ],
2953 durableObjectStorage = (inMemory = void)
2954 )
2955 ),
2956 ],
2957 sockets = [
2958 ( name = "main",
2959 address = "test-addr",
2960 service = "hello"
2961 )
2962 ]
2963 ))"_kj);
2964 
2965 test.start();
2966 auto wsConn = test.connect("test-addr");
2967 wsConn.upgradeToWebSocket();
2968 constexpr kj::StringPtr expectedOne = "Hello"_kj;
2969 constexpr kj::StringPtr expectedTwo = "There"_kj;
2970 // \x81\x05 are part of the websocket frame.
2971 // \x81 is 10000001 -- leftmost bit implies this is the final frame, rightmost implies text data.
2972 // \x05 says the payload length is 5.
2973 wsConn.send(kj::str("\x81\x05", expectedOne));
2974 wsConn.send(kj::str("\x81\x05", expectedTwo));
2975 wsConn.recvWebSocket(expectedOne);
2976 wsConn.recvWebSocket(expectedTwo);
2977 
2978 // Force hibernation by waiting 10 seconds.
2979 test.wait(10);
2980 wsConn.send(kj::str("\x81\x05", expectedOne));
2981 wsConn.send(kj::str("\x81\x05", expectedTwo));
2982 wsConn.recvWebSocket(expectedOne);
2983 wsConn.recvWebSocket(expectedTwo);
2984}
2985 
2986KJ_TEST("Server: Durable Objects websocket hibernation") {
2987 TestServer test(R"((
2988 services = [
2989 ( name = "hello",
2990 worker = (
2991 compatibilityDate = "2023-08-17",
2992 modules = [
2993 ( name = "main.js",
2994 esModule =
2995 `export default {
2996 ` async fetch(request, env) {
2997 ` let id = env.ns.idFromName("59002eb8cf872e541722977a258a12d6a93bbe8192b502e1c0cb250aa91af234");
2998 ` let obj = env.ns.get(id)
2999 `
3000 ` // 1. Create a websocket (request 1)
3001 ` // 2. Use websocket once
3002 ` // 3. Let actor hibernate
3003 ` // 4. Wake actor by sending new request (request 2)
3004 ` // - This confirms we get back hibernation manager.
3005 ` // 5. Use websocket once
3006 ` // 6. Let actor hibernate
3007 ` // 7. Wake actor by using websocket
3008 ` // - This confirms we get back hibernation manager.
3009 ` // 8. Use websocket once
3010 ` try {
3011 ` return await obj.fetch(request);
3012 ` } catch (err) {
3013 ` if (request.url.endsWith("/abort")) {
3014 ` // expected
3015 ` return new Response("OK");
3016 ` } else {
3017 ` throw err;
3018 ` }
3019 ` }
3020 ` }
3021 `}
3022 `
3023 `export class MyActorClass {
3024 ` constructor(state) {
3025 ` this.state = state;
3026 ` // If reqCount is 0, then the actor's constructor has run.
3027 ` // This implies we're starting up, so either this is the first request or we were evicted.
3028 ` this.reqCount = 0;
3029 ` }
3030 `
3031 ` async fetch(request) {
3032 ` if (request.url.endsWith("/")) {
3033 ` // Request 1, accept a websocket.
3034 ` let pair = new WebSocketPair(true);
3035 ` let ws = pair[1];
3036 ` this.state.acceptWebSocket(ws);
3037 `
3038 ` this.reqCount += 1;
3039 ` if (this.reqCount != 1) {
3040 ` throw new Error(`Expected request count of 1 but got ${this.reqCount}`);
3041 ` }
3042 ` return new Response(null, {status: 101, statusText: "Switching Protocols", webSocket: pair[0]});
3043 ` } else if (request.url.endsWith("/wakeUpAndCheckWS")) {
3044 ` // Request 2, wake actor and check if WS available.
3045 ` let allWebsockets = this.state.getWebSockets();
3046 ` for (const ws of allWebsockets) {
3047 ` ws.send("Hello! Just woke up from a nap.");
3048 ` }
3049 `
3050 ` this.reqCount += 1;
3051 ` if (this.reqCount != 1) {
3052 ` throw new Error(`Expected request count of 1 but got ${this.reqCount}`);
3053 ` }
3054 `
3055 ` return new Response("OK");
3056 ` } else if (request.url.endsWith("/abort")) {
3057 ` this.state.abort("test abort message");
3058 ` }
3059 ` return new Error("Unknown path!");
3060 ` }
3061 `
3062 ` async webSocketMessage(ws, msg) {
3063 ` if (msg == "Regular message.") {
3064 ` ws.send("Regular response.");
3065 ` } else if (msg == "Confirm actor was evicted.") {
3066 ` // Called when waking from hibernation due to inbound websocket message.
3067 ` if (this.reqCount == 0) {
3068 ` ws.send("OK")
3069 ` } else {
3070 ` ws.send(`[ FAILURE ] - reqCount was ${this.reqCount} so actor wasn't evicted`);
3071 ` }
3072 ` }
3073 ` }
3074 `
3075 ` async webSocketClose(ws, code, reason, wasClean) {
3076 ` if (code == 1006) {
3077 ` if (reason != "WebSocket disconnected without sending Close frame.") {
3078 ` throw new Error(`Got abnormal closure with unexpected reason: ${reason}`);
3079 ` }
3080 ` if (wasClean) {
3081 ` throw new Error("Got abnormal closure but wasClean was true!");
3082 ` }
3083 ` } else if (code != 1234) {
3084 ` throw new Error(`Expected close code 1234, got ${code}`);
3085 ` } else if (reason != "OK") {
3086 ` throw new Error(`Expected close reason "OK", got ${reason}`);
3087 ` } else {
3088 ` ws.close(4321, "KO");
3089 ` }
3090 ` }
3091 `
3092 ` async webSocketError(ws, error) {
3093 ` console.log(`Encountered error: ${error}`);
3094 ` throw new Error(error);
3095 ` }
3096 `}
3097 
3098 )
3099 ],
3100 bindings = [(name = "ns", durableObjectNamespace = "MyActorClass")],
3101 durableObjectNamespaces = [
3102 ( className = "MyActorClass",
3103 uniqueKey = "mykey",
3104 )
3105 ],
3106 durableObjectStorage = (inMemory = void)
3107 )
3108 ),
3109 ],
3110 sockets = [
3111 ( name = "main",
3112 address = "test-addr",
3113 service = "hello"
3114 )
3115 ]
3116 ))"_kj);
3117 
3118 test.start();
3119 auto wsConn = test.connect("test-addr");
3120 wsConn.upgradeToWebSocket();
3121 // 1. Make hibernatable ws and use it.
3122 constexpr kj::StringPtr message = "Regular message."_kj;
3123 constexpr kj::StringPtr response = "Regular response."_kj;
3124 wsConn.send(kj::str("\x81\x10", message));
3125 wsConn.recvWebSocket(response);
3126 
3127 // 2. Hibernate
3128 test.wait(10);
3129 // 3. Use normal connection and read from ws.
3130 {
3131 auto conn = test.connect("test-addr");
3132 conn.httpGet200("/wakeUpAndCheckWS", "OK"_kj);
3133 }
3134 constexpr kj::StringPtr unpromptedResponse = "Hello! Just woke up from a nap."_kj;
3135 wsConn.recvWebSocket(unpromptedResponse);
3136 
3137 // 4. Hibernate again
3138 test.wait(10);
3139 
3140 // 5. Wake up by sending a message.
3141 constexpr kj::StringPtr confirmEviction = "Confirm actor was evicted."_kj;
3142 constexpr kj::StringPtr evicted = "OK"_kj;
3143 wsConn.send(kj::str("\x81\x1a", confirmEviction));
3144 wsConn.recvWebSocket(evicted);
3145 
3146 // 6. Hibernate again
3147 test.wait(10);
3148 
3149 // 7. Wake up the actor and have it abort itself. This should disconnect the WebSocket, even
3150 // though the WebSocket itself is still hibernated.
3151 KJ_EXPECT_LOG(INFO, "Error: test abort message");
3152 KJ_EXPECT_LOG(INFO, "other end of WebSocketPipe was destroyed");
3153 {
3154 auto conn = test.connect("test-addr");
3155 conn.httpGet200("/abort", "OK"_kj);
3156 }
3157 
3158 KJ_EXPECT(wsConn.isEof());
3159}
3160 
3161KJ_TEST("Server: tail workers") {
3162 TestServer test(R"((
3163 services = [
3164 ( name = "hello",
3165 worker = (
3166 compatibilityDate = "2024-11-01",
3167 modules = [
3168 ( name = "main.js",
3169 esModule =
3170 `export default {
3171 ` async fetch(req, env, ctx) {
3172 ` console.log("foo", "bar");
3173 ` console.log("baz");
3174 ` return new Response("OK");
3175 ` }
3176 `}
3177 )
3178 ],
3179 tails = ["tail", "tail2"],
3180 )
3181 ),
3182 ( name = "tail",
3183 worker = (
3184 compatibilityDate = "2024-11-01",
3185 modules = [
3186 ( name = "main.js",
3187 esModule =
3188 `export default {
3189 ` async tail(req, env, ctx) {
3190 ` await fetch("http://tail", {
3191 ` method: "POST",
3192 ` body: JSON.stringify(req[0].logs.map(log => log.message))
3193 ` });
3194 ` }
3195 `}
3196 )
3197 ],
3198 )
3199 ),
3200 ( name = "tail2",
3201 worker = (
3202 compatibilityDate = "2024-11-01",
3203 modules = [
3204 ( name = "main.js",
3205 esModule =
3206 `export default {
3207 ` async tail(req, env, ctx) {
3208 ` await fetch("http://tail2/" + req[0].logs.length);
3209 ` }
3210 `}
3211 )
3212 ],
3213 )
3214 ),
3215 ],
3216 sockets = [
3217 ( name = "main",
3218 address = "test-addr",
3219 service = "hello"
3220 )
3221 ]
3222 ))"_kj);
3223 
3224 test.start();
3225 auto conn = test.connect("test-addr");
3226 conn.sendHttpGet("/");
3227 conn.recvHttp200("OK");
3228 
3229 auto subreq = test.receiveInternetSubrequest("tail");
3230 subreq.recv(R"(
3231 POST / HTTP/1.1
3232 Content-Length: 23
3233 Host: tail
3234 Content-Type: text/plain;charset=UTF-8
3235 
3236 [["foo","bar"],["baz"]])"_blockquote);
3237 
3238 auto subreq2 = test.receiveInternetSubrequest("tail2");
3239 subreq2.recv(R"(
3240 GET /2 HTTP/1.1
3241 Host: tail2
3242 
3243 )"_blockquote);
3244 
3245 subreq.send(R"(
3246 HTTP/1.1 200 OK
3247 Content-Length: 0
3248 
3249 )"_blockquote);
3250 
3251 subreq2.send(R"(
3252 HTTP/1.1 200 OK
3253 Content-Length: 0
3254 
3255 )"_blockquote);
3256}
3257 
3258// =======================================================================================
3259// Test HttpOptions on receive
3260 
3261KJ_TEST("Server: serve proxy requests") {
3262 TestServer test(R"((
3263 services = [
3264 ( name = "hello",
3265 worker = (
3266 compatibilityDate = "2022-08-17",
3267 serviceWorkerScript =
3268 `addEventListener("fetch", event => {
3269 ` event.respondWith(new Response("Hello: " + event.request.url + "\n"));
3270 `})
3271 )
3272 )
3273 ],
3274 sockets = [
3275 ( name = "main",
3276 address = "test-addr",
3277 service = "hello",
3278 http = (style = proxy)
3279 )
3280 ]
3281 ))"_kj);
3282 
3283 test.start();
3284 
3285 auto conn = test.connect("test-addr");
3286 
3287 // Send a proxy-style request. No `Host:` header!
3288 conn.send(R"(
3289 GET http://foo/bar HTTP/1.1
3290 
3291 )"_blockquote);
3292 conn.recv(R"(
3293 HTTP/1.1 200 OK
3294 Content-Length: 22
3295 Content-Type: text/plain;charset=UTF-8
3296 
3297 Hello: http://foo/bar
3298 )"_blockquote);
3299}
3300 
3301KJ_TEST("Server: forwardedProtoHeader") {
3302 TestServer test(R"((
3303 services = [
3304 ( name = "hello",
3305 worker = (
3306 compatibilityDate = "2022-08-17",
3307 serviceWorkerScript =
3308 `addEventListener("fetch", event => {
3309 ` event.respondWith(new Response("Hello: " + event.request.url + "\n"));
3310 `})
3311 )
3312 )
3313 ],
3314 sockets = [
3315 ( name = "main",
3316 address = "test-addr",
3317 service = "hello",
3318 http = (forwardedProtoHeader = "Test-Proto")
3319 )
3320 ]
3321 ))"_kj);
3322 
3323 test.start();
3324 
3325 auto conn = test.connect("test-addr");
3326 
3327 // Send a request with a forwarded proto header.
3328 conn.send(R"(
3329 GET /bar HTTP/1.1
3330 Host: foo
3331 tEsT-pRoTo: baz
3332 
3333 )"_blockquote);
3334 conn.recv(R"(
3335 HTTP/1.1 200 OK
3336 Content-Length: 21
3337 Content-Type: text/plain;charset=UTF-8
3338 
3339 Hello: baz://foo/bar
3340 )"_blockquote);
3341 
3342 // Send a request without one.
3343 conn.send(R"(
3344 GET /bar HTTP/1.1
3345 Host: foo
3346 
3347 )"_blockquote);
3348 conn.recv(R"(
3349 HTTP/1.1 200 OK
3350 Content-Length: 22
3351 Content-Type: text/plain;charset=UTF-8
3352 
3353 Hello: http://foo/bar
3354 )"_blockquote);
3355}
3356 
3357KJ_TEST("Server: cfBlobHeader") {
3358 TestServer test(R"((
3359 services = [
3360 ( name = "hello",
3361 worker = (
3362 compatibilityDate = "2022-08-17",
3363 serviceWorkerScript =
3364 `addEventListener("fetch", event => {
3365 ` if (event.request.cf) {
3366 ` event.respondWith(new Response("cf.foo = " + event.request.cf.foo + "\n"));
3367 ` } else {
3368 ` event.respondWith(new Response("cf is null\n"));
3369 ` }
3370 `})
3371 )
3372 )
3373 ],
3374 sockets = [
3375 ( name = "main",
3376 address = "test-addr",
3377 service = "hello",
3378 http = (cfBlobHeader = "CF-Blob")
3379 )
3380 ]
3381 ))"_kj);
3382 
3383 test.start();
3384 
3385 auto conn = test.connect("test-addr");
3386 
3387 // Send a request with a CF blob.
3388 conn.send(R"(
3389 GET / HTTP/1.1
3390 Host: bar
3391 cF-bLoB: {"foo": "hello"}
3392 
3393 )"_blockquote);
3394 conn.recv(R"(
3395 HTTP/1.1 200 OK
3396 Content-Length: 15
3397 Content-Type: text/plain;charset=UTF-8
3398 
3399 cf.foo = hello
3400 )"_blockquote);
3401 
3402 // Send a request without one
3403 conn.send(R"(
3404 GET / HTTP/1.1
3405 Host: bar
3406 
3407 )"_blockquote);
3408 conn.recv(R"(
3409 HTTP/1.1 200 OK
3410 Content-Length: 11
3411 Content-Type: text/plain;charset=UTF-8
3412 
3413 cf is null
3414 )"_blockquote);
3415}
3416 
3417KJ_TEST("Server: inject headers on incoming request/response") {
3418 TestServer test(R"((
3419 services = [
3420 ( name = "hello",
3421 worker = (
3422 compatibilityDate = "2022-08-17",
3423 serviceWorkerScript =
3424 `addEventListener("fetch", event => {
3425 ` let text = [...event.request.headers]
3426 ` .map(([k,v]) => { return `${k}: ${v}\n` }).join("");
3427 ` event.respondWith(new Response(text));
3428 `})
3429 )
3430 )
3431 ],
3432 sockets = [
3433 ( name = "main",
3434 address = "test-addr",
3435 service = "hello",
3436 http = (
3437 injectRequestHeaders = [
3438 (name = "Foo", value = "oof"),
3439 (name = "Bar", value = "rab"),
3440 ],
3441 injectResponseHeaders = [
3442 (name = "Baz", value = "zab"),
3443 (name = "Qux", value = "xuq"),
3444 ]
3445 )
3446 )
3447 ]
3448 ))"_kj);
3449 
3450 test.start();
3451 
3452 auto conn = test.connect("test-addr");
3453 
3454 // Send a request, check headers.
3455 conn.send(R"(
3456 GET / HTTP/1.1
3457 Host: example.com
3458 
3459 )"_blockquote);
3460 conn.recv(R"(
3461 HTTP/1.1 200 OK
3462 Content-Length: 36
3463 Content-Type: text/plain;charset=UTF-8
3464 Baz: zab
3465 Qux: xuq
3466 
3467 bar: rab
3468 foo: oof
3469 host: example.com
3470 )"_blockquote);
3471}
3472 
3473KJ_TEST("Server: drain incoming HTTP connections") {
3474 TestServer test(singleWorker(R"((
3475 compatibilityDate = "2022-08-17",
3476 serviceWorkerScript =
3477 `addEventListener("fetch", event => {
3478 ` event.respondWith(new Response("hello"));
3479 `})
3480 ))"_kj));
3481 
3482 auto paf = kj::newPromiseAndFulfiller<void>();
3483 
3484 test.start(kj::mv(paf.promise));
3485 
3486 auto conn = test.connect("test-addr");
3487 auto conn2 = test.connect("test-addr");
3488 
3489 // Send a request on each connection, get a response.
3490 conn.httpGet200("/", "hello");
3491 conn2.httpGet200("/", "hello");
3492 
3493 // Send a partial request on conn2.
3494 conn2.send("GET");
3495 
3496 // No EOF yet.
3497 KJ_EXPECT(!conn.isEof());
3498 KJ_EXPECT(!conn2.isEof());
3499 
3500 // Drain the server.
3501 paf.fulfiller->fulfill();
3502 
3503 // Now we get EOF on conn.
3504 KJ_EXPECT(conn.isEof());
3505 
3506 // But conn2 is still open.
3507 KJ_EXPECT(!conn2.isEof());
3508 
3509 // New connections shouldn't be accepted at this point.
3510 KJ_EXPECT(test.connectHangs("test-addr"));
3511 
3512 // Finish the request on conn2.
3513 conn2.send(" / HTTP/1.1\nHost: foo\n\n");
3514 
3515 // We receive a response with Connection: close
3516 conn2.recv(R"(
3517 HTTP/1.1 200 OK
3518 Connection: close
3519 Content-Length: 5
3520 Content-Type: text/plain;charset=UTF-8
3521 
3522 hello)"_blockquote);
3523 
3524 // And then the connection is, in fact, closed.
3525 KJ_EXPECT(conn2.isEof());
3526}
3527 
3528// =======================================================================================
3529// Test alternate service types
3530//
3531// We're going to stop using JavaScript here because it's not really helping. We can directly
3532// connect a socket to a non-Worker service.
3533 
3534KJ_TEST("Server: network outbound with allow/deny") {
3535 TestServer test(R"((
3536 services = [
3537 (name = "hello", network = (allow = ["foo", "bar"], deny = ["baz", "qux"]))
3538 ],
3539 sockets = [
3540 (name = "main", address = "test-addr", service = "hello")
3541 ]
3542 ))"_kj);
3543 
3544 test.start();
3545 
3546 auto conn = test.connect("test-addr");
3547 
3548 conn.sendHttpGet("/path");
3549 
3550 {
3551 auto subreq = test.receiveSubrequest("foo", {"foo", "bar"}, {"baz", "qux"});
3552 subreq.recv(R"(
3553 GET /path HTTP/1.1
3554 Host: foo
3555 
3556 )"_blockquote);
3557 subreq.send(R"(
3558 HTTP/1.1 200 OK
3559 Content-Length: 2
3560 Content-Type: text/plain;charset=UTF-8
3561 
3562 OK)"_blockquote);
3563 }
3564 
3565 conn.recvHttp200("OK");
3566}
3567 
3568KJ_TEST("Server: external server") {
3569 TestServer test(R"((
3570 services = [
3571 (name = "hello", external = "ext-addr")
3572 ],
3573 sockets = [
3574 (name = "main", address = "test-addr", service = "hello")
3575 ]
3576 ))"_kj);
3577 
3578 test.start();
3579 
3580 auto conn = test.connect("test-addr");
3581 
3582 conn.sendHttpGet("/path");
3583 
3584 {
3585 auto subreq = test.receiveSubrequest("ext-addr");
3586 subreq.recv(R"(
3587 GET /path HTTP/1.1
3588 Host: foo
3589 
3590 )"_blockquote);
3591 subreq.send(R"(
3592 HTTP/1.1 200 OK
3593 Content-Length: 2
3594 Content-Type: text/plain;charset=UTF-8
3595 
3596 OK)"_blockquote);
3597 }
3598 
3599 conn.recvHttp200("OK");
3600}
3601 
3602KJ_TEST("Server: external server proxy style") {
3603 TestServer test(R"((
3604 services = [
3605 (name = "hello", external = (address = "ext-addr", http = (style = proxy)))
3606 ],
3607 sockets = [
3608 (name = "main", address = "test-addr", service = "hello")
3609 ]
3610 ))"_kj);
3611 
3612 test.start();
3613 
3614 auto conn = test.connect("test-addr");
3615 
3616 conn.sendHttpGet("/path");
3617 
3618 {
3619 auto subreq = test.receiveSubrequest("ext-addr");
3620 subreq.recv(R"(
3621 GET http://foo/path HTTP/1.1
3622 Host: foo
3623 
3624 )"_blockquote);
3625 subreq.send(R"(
3626 HTTP/1.1 200 OK
3627 Content-Length: 2
3628 Content-Type: text/plain;charset=UTF-8
3629 
3630 OK)"_blockquote);
3631 }
3632 
3633 conn.recvHttp200("OK");
3634}
3635 
3636KJ_TEST("Server: external server forwarded-proto") {
3637 TestServer test(R"((
3638 services = [
3639 (name = "hello", external = (address = "ext-addr", http = (forwardedProtoHeader = "X-Proto")))
3640 ],
3641 sockets = [
3642 (name = "main", address = "test-addr", service = "hello", http = (style = proxy))
3643 ]
3644 ))"_kj);
3645 
3646 test.start();
3647 
3648 auto conn = test.connect("test-addr");
3649 
3650 conn.send(R"(
3651 GET https://foo/path HTTP/1.1
3652 
3653 )"_blockquote);
3654 
3655 {
3656 auto subreq = test.receiveSubrequest("ext-addr");
3657 subreq.recv(R"(
3658 GET /path HTTP/1.1
3659 Host: foo
3660 X-Proto: https
3661 
3662 )"_blockquote);
3663 subreq.send(R"(
3664 HTTP/1.1 200 OK
3665 Content-Length: 2
3666 Content-Type: text/plain;charset=UTF-8
3667 
3668 OK)"_blockquote);
3669 }
3670 
3671 conn.recvHttp200("OK");
3672}
3673 
3674KJ_TEST("Server: external server inject headers") {
3675 TestServer test(R"((
3676 services = [
3677 ( name = "hello",
3678 external = (
3679 address = "ext-addr",
3680 http = (
3681 injectRequestHeaders = [
3682 (name = "Foo", value = "oof"),
3683 (name = "Bar", value = "rab"),
3684 ],
3685 injectResponseHeaders = [
3686 (name = "Baz", value = "zab"),
3687 (name = "Qux", value = "xuq"),
3688 ]
3689 )
3690 )
3691 )
3692 ],
3693 sockets = [
3694 (name = "main", address = "test-addr", service = "hello")
3695 ]
3696 ))"_kj);
3697 
3698 test.start();
3699 
3700 auto conn = test.connect("test-addr");
3701 
3702 conn.sendHttpGet("/path");
3703 
3704 {
3705 auto subreq = test.receiveSubrequest("ext-addr");
3706 subreq.recv(R"(
3707 GET /path HTTP/1.1
3708 Host: foo
3709 Foo: oof
3710 Bar: rab
3711 
3712 )"_blockquote);
3713 subreq.send(R"(
3714 HTTP/1.1 200 OK
3715 Content-Length: 2
3716 Content-Type: text/plain;charset=UTF-8
3717 
3718 OK)"_blockquote);
3719 }
3720 
3721 conn.recv(R"(
3722 HTTP/1.1 200 OK
3723 Content-Length: 2
3724 Content-Type: text/plain;charset=UTF-8
3725 Baz: zab
3726 Qux: xuq
3727 
3728 OK)"_blockquote);
3729}
3730 
3731KJ_TEST("Server: external server cf blob header") {
3732 TestServer test(R"((
3733 services = [
3734 ( name = "hello",
3735 worker = (
3736 compatibilityDate = "2022-08-17",
3737 modules = [
3738 ( name = "main.js",
3739 esModule =
3740 `export default {
3741 ` async fetch(request, env) {
3742 ` return env.ext.fetch("http://ext/path2", {cf: {hello: "world"}});
3743 ` }
3744 `}
3745 )
3746 ],
3747 bindings = [(name = "ext", service = "ext")]
3748 )
3749 ),
3750 (name = "ext", external = (address = "ext-addr", http = (cfBlobHeader = "CF-Blob")))
3751 ],
3752 sockets = [
3753 (name = "main", address = "test-addr", service = "hello")
3754 ]
3755 ))"_kj);
3756 
3757 test.start();
3758 
3759 auto conn = test.connect("test-addr");
3760 
3761 conn.sendHttpGet("/path");
3762 
3763 {
3764 auto subreq = test.receiveSubrequest("ext-addr");
3765 subreq.recv(R"(
3766 GET /path2 HTTP/1.1
3767 Host: ext
3768 CF-Blob: {"hello":"world"}
3769 
3770 )"_blockquote);
3771 subreq.send(R"(
3772 HTTP/1.1 200 OK
3773 Content-Length: 2
3774 Content-Type: text/plain;charset=UTF-8
3775 
3776 OK)"_blockquote);
3777 }
3778 
3779 conn.recv(R"(
3780 HTTP/1.1 200 OK
3781 Content-Length: 2
3782 Content-Type: text/plain;charset=UTF-8
3783 
3784 OK)"_blockquote);
3785}
3786 
3787KJ_TEST("Server: disk service") {
3788 TestServer test(R"((
3789 services = [
3790 (name = "hello", disk = "../../frob/blah")
3791 ],
3792 sockets = [
3793 (name = "main", address = "test-addr", service = "hello")
3794 ]
3795 ))"_kj);
3796 
3797 auto mode = kj::WriteMode::CREATE | kj::WriteMode::CREATE_PARENT;
3798 auto dir = test.root->openSubdir(kj::Path({"frob"_kj, "blah"_kj}), mode);
3799 test.fakeDate =
3800 kj::UNIX_EPOCH + 2 * kj::DAYS + 5 * kj::HOURS + 18 * kj::MINUTES + 23 * kj::SECONDS;
3801 dir->openFile(kj::Path({"foo.txt"}), mode)->writeAll("hello from foo.txt\n");
3802 dir->openFile(kj::Path({"numbers.txt"}), mode)->writeAll("0123456789\n");
3803 test.fakeDate = kj::UNIX_EPOCH + 400 * kj::DAYS + 2 * kj::HOURS + 52 * kj::MINUTES +
3804 9 * kj::SECONDS + 163 * kj::MILLISECONDS;
3805 dir->openFile(kj::Path({"bar.txt"}), mode)->writeAll("hello from bar.txt\n");
3806 test.fakeDate = kj::UNIX_EPOCH;
3807 dir->openFile(kj::Path({"baz", "qux.txt"}), mode)->writeAll("hello from qux.txt\n");
3808 dir->openFile(kj::Path({".dot"}), mode)->writeAll("this is a dotfile\n");
3809 dir->openFile(kj::Path({".dotdir", "foo"}), mode)->writeAll("this is a dotfile\n");
3810 
3811 test.start();
3812 
3813 auto conn = test.connect("test-addr");
3814 
3815 conn.sendHttpGet("/foo.txt");
3816 conn.recv(R"(
3817 HTTP/1.1 200 OK
3818 Content-Length: 19
3819 Content-Type: application/octet-stream
3820 Last-Modified: Sat, 03 Jan 1970 05:18:23 GMT
3821 
3822 hello from foo.txt
3823 )"_blockquote);
3824 
3825 conn.sendHttpGet("/bar.txt");
3826 conn.recv(R"(
3827 HTTP/1.1 200 OK
3828 Content-Length: 19
3829 Content-Type: application/octet-stream
3830 Last-Modified: Fri, 05 Feb 1971 02:52:09 GMT
3831 
3832 hello from bar.txt
3833 )"_blockquote);
3834 
3835 conn.sendHttpGet("/baz/qux.txt");
3836 conn.recv(R"(
3837 HTTP/1.1 200 OK
3838 Content-Length: 19
3839 Content-Type: application/octet-stream
3840 Last-Modified: Thu, 01 Jan 1970 00:00:00 GMT
3841 
3842 hello from qux.txt
3843 )"_blockquote);
3844 
3845 // TODO(beta): Test listing a directory. Unfortunately it doesn't work against the in-memory
3846 // filesystem right now.
3847 //
3848 // conn.sendHttpGet("/");
3849 
3850 // HEAD returns no content.
3851 conn.send(R"(
3852 HEAD /numbers.txt HTTP/1.1
3853 Host: foo
3854 
3855 )"_blockquote);
3856 conn.recv(R"(
3857 HTTP/1.1 200 OK
3858 Content-Length: 11
3859 Content-Type: application/octet-stream
3860 Last-Modified: Sat, 03 Jan 1970 05:18:23 GMT
3861 
3862 )"_blockquote);
3863 
3864 // GET with single range returns partial content.
3865 conn.send(R"(
3866 GET /numbers.txt HTTP/1.1
3867 Host: foo
3868 Range: bytes=3-5
3869 
3870 )"_blockquote);
3871 conn.recv(R"(
3872 HTTP/1.1 206 Partial Content
3873 Content-Length: 3
3874 Content-Type: application/octet-stream
3875 Content-Range: bytes 3-5/11
3876 Last-Modified: Sat, 03 Jan 1970 05:18:23 GMT
3877 
3878 345)"_blockquote);
3879 
3880 // GET with single covering range returns full content.
3881 conn.send(R"(
3882 GET /numbers.txt HTTP/1.1
3883 Host: foo
3884 Range: bytes=-50
3885 
3886 )"_blockquote);
3887 conn.recv(R"(
3888 HTTP/1.1 200 OK
3889 Content-Length: 11
3890 Content-Type: application/octet-stream
3891 Last-Modified: Sat, 03 Jan 1970 05:18:23 GMT
3892 
3893 0123456789
3894 )"_blockquote);
3895 
3896 // GET with many ranges returns full content.
3897 conn.send(R"(
3898 GET /numbers.txt HTTP/1.1
3899 Host: foo
3900 Range: bytes=1-3, 6-8
3901 
3902 )"_blockquote);
3903 conn.recv(R"(
3904 HTTP/1.1 200 OK
3905 Content-Length: 11
3906 Content-Type: application/octet-stream
3907 Last-Modified: Sat, 03 Jan 1970 05:18:23 GMT
3908 
3909 0123456789
3910 )"_blockquote);
3911 
3912 // GET with unsatisfiable range.
3913 conn.send(R"(
3914 GET /numbers.txt HTTP/1.1
3915 Host: foo
3916 Range: bytes=20-30
3917 
3918 )"_blockquote);
3919 conn.recv(R"(
3920 HTTP/1.1 416 Range Not Satisfiable
3921 Content-Length: 21
3922 Content-Range: bytes */11
3923 
3924 Range Not Satisfiable)"_blockquote);
3925 
3926 // File not found...
3927 conn.sendHttpGet("/no-such-file.txt");
3928 conn.recv(R"(
3929 HTTP/1.1 404 Not Found
3930 Content-Length: 9
3931 
3932 Not Found)"_blockquote);
3933 
3934 // Directory not found...
3935 conn.sendHttpGet("/no-such-dir/file.txt");
3936 conn.recv(R"(
3937 HTTP/1.1 404 Not Found
3938 Content-Length: 9
3939 
3940 Not Found)"_blockquote);
3941 
3942 // PUT is denied because not writable.
3943 conn.send(R"(
3944 PUT /corge.txt HTTP/1.1
3945 Host: foo
3946 Content-Length: 6
3947 
3948 corge
3949 )"_blockquote);
3950 conn.recv(R"(
3951 HTTP/1.1 405 Method Not Allowed
3952 Content-Length: 18
3953 
3954 Method Not Allowed)"_blockquote);
3955 
3956 // DELETE is denied because not writable.
3957 conn.send(R"(
3958 DELETE /corge.txt HTTP/1.1
3959 Host: foo
3960 
3961 )"_blockquote);
3962 conn.recv(R"(
3963 HTTP/1.1 405 Method Not Allowed
3964 Content-Length: 18
3965 
3966 Method Not Allowed)"_blockquote);
3967 
3968 // POST is denied because invalid method.
3969 conn.send(R"(
3970 POST /corge.txt HTTP/1.1
3971 Host: foo
3972 Content-Length: 6
3973 
3974 corge
3975 )"_blockquote);
3976 conn.recv(R"(
3977 HTTP/1.1 501 Not Implemented
3978 Content-Length: 15
3979 
3980 Not Implemented)"_blockquote);
3981 
3982 // Dotfile access is denied.
3983 conn.sendHttpGet("/.dot");
3984 conn.recv(R"(
3985 HTTP/1.1 404 Not Found
3986 Content-Length: 9
3987 
3988 Not Found)"_blockquote);
3989 
3990 // Dotfile directory access is denied.
3991 conn.sendHttpGet("/.dotdir/foo");
3992 conn.recv(R"(
3993 HTTP/1.1 404 Not Found
3994 Content-Length: 9
3995 
3996 Not Found)"_blockquote);
3997}
3998 
3999KJ_TEST("Server: disk service writable") {
4000 TestServer test(R"((
4001 services = [
4002 (name = "hello", disk = (path = "../../frob/blah", writable = true))
4003 ],
4004 sockets = [
4005 (name = "main", address = "test-addr", service = "hello")
4006 ]
4007 ))"_kj);
4008 
4009 auto mode = kj::WriteMode::CREATE | kj::WriteMode::CREATE_PARENT;
4010 auto dir = test.root->openSubdir(kj::Path({"frob"_kj, "blah"_kj}), mode);
4011 dir->openFile(kj::Path({"existing.txt"}), mode)->writeAll("replace me!");
4012 
4013 test.start();
4014 
4015 auto conn = test.connect("test-addr");
4016 
4017 // Write a file.
4018 conn.send(R"(
4019 PUT /newfile.txt HTTP/1.1
4020 Host: foo
4021 Content-Length: 6
4022 
4023 corge
4024 )"_blockquote);
4025 conn.recv(R"(
4026 HTTP/1.1 204 No Content
4027 
4028 )"_blockquote);
4029 
4030 // Read it back.
4031 KJ_EXPECT(dir->openFile(kj::Path({"newfile.txt"}))->readAllText() == "corge\n");
4032 
4033 // Delete it.
4034 conn.send(R"(
4035 DELETE /newfile.txt HTTP/1.1
4036 Host: foo
4037 
4038 )"_blockquote);
4039 conn.recv(R"(
4040 HTTP/1.1 204 No Content
4041 
4042 )"_blockquote);
4043 KJ_EXPECT(!dir->exists(kj::Path({"newfile.txt"})));
4044 
4045 // Delete a non-existent file.
4046 conn.send(R"(
4047 DELETE /notfound.txt HTTP/1.1
4048 Host: foo
4049 
4050 )"_blockquote);
4051 conn.recv(R"(
4052 HTTP/1.1 404 Not Found
4053 Content-Length: 9
4054 
4055 Not Found)"_blockquote);
4056 
4057 // Replace a file.
4058 conn.send(R"(
4059 PUT /existing.txt HTTP/1.1
4060 Host: foo
4061 Content-Length: 7
4062 
4063 grault
4064 )"_blockquote);
4065 conn.recv(R"(
4066 HTTP/1.1 204 No Content
4067 
4068 )"_blockquote);
4069 
4070 // Read it back.
4071 KJ_EXPECT(dir->openFile(kj::Path({"existing.txt"}))->readAllText() == "grault\n");
4072 
4073 // Write a file to a new directory.
4074 conn.send(R"(
4075 PUT /newdir/newfile.txt HTTP/1.1
4076 Host: foo
4077 Content-Length: 7
4078 
4079 garply
4080 )"_blockquote);
4081 conn.recv(R"(
4082 HTTP/1.1 204 No Content
4083 
4084 )"_blockquote);
4085 
4086 // Read it back.
4087 KJ_EXPECT(dir->openFile(kj::Path({"newdir", "newfile.txt"}))->readAllText() == "garply\n");
4088 
4089 // Delete the new directory.
4090 conn.send(R"(
4091 DELETE /newdir/ HTTP/1.1
4092 Host: foo
4093 
4094 )"_blockquote);
4095 conn.recv(R"(
4096 HTTP/1.1 204 No Content
4097 
4098 )"_blockquote);
4099 KJ_EXPECT(!dir->exists(kj::Path({"newdir"})));
4100 
4101 // POST is denied because invalid method.
4102 conn.send(R"(
4103 POST /corge.txt HTTP/1.1
4104 Host: foo
4105 Content-Length: 6
4106 
4107 waldo
4108 )"_blockquote);
4109 conn.recv(R"(
4110 HTTP/1.1 501 Not Implemented
4111 Content-Length: 15
4112 
4113 Not Implemented)"_blockquote);
4114 
4115 // Dotfile write access is denied.
4116 conn.send(R"(
4117 PUT /.dot HTTP/1.1
4118 Host: foo
4119 Content-Length: 6
4120 
4121 waldo
4122 )"_blockquote);
4123 conn.recv(R"(
4124 HTTP/1.1 403 Unauthorized
4125 Content-Length: 12
4126 
4127 Unauthorized)"_blockquote);
4128 
4129 // Dotfile directory write access is denied.
4130 conn.send(R"(
4131 PUT /.dotdir/foo HTTP/1.1
4132 Host: foo
4133 Content-Length: 6
4134 
4135 waldo
4136 )"_blockquote);
4137 conn.recv(R"(
4138 HTTP/1.1 403 Unauthorized
4139 Content-Length: 12
4140 
4141 Unauthorized)"_blockquote);
4142 
4143 // Dotfile delete access is denied.
4144 conn.send(R"(
4145 DELETE /.dot HTTP/1.1
4146 Host: foo
4147 
4148 )"_blockquote);
4149 conn.recv(R"(
4150 HTTP/1.1 403 Unauthorized
4151 Content-Length: 12
4152 
4153 Unauthorized)"_blockquote);
4154 
4155 // Root write is denied.
4156 conn.send(R"(
4157 PUT / HTTP/1.1
4158 Host: foo
4159 Content-Length: 6
4160 
4161 corge
4162 )"_blockquote);
4163 conn.recv(R"(
4164 HTTP/1.1 403 Unauthorized
4165 Content-Length: 12
4166 
4167 Unauthorized)"_blockquote);
4168 
4169 // Root delete is denied.
4170 conn.send(R"(
4171 DELETE / HTTP/1.1
4172 Host: foo
4173 
4174 )"_blockquote);
4175 conn.recv(R"(
4176 HTTP/1.1 403 Unauthorized
4177 Content-Length: 12
4178 
4179 Unauthorized)"_blockquote);
4180}
4181 
4182KJ_TEST("Server: disk service allow dotfiles") {
4183 TestServer test(R"((
4184 services = [
4185 (name = "hello", disk = (path = "../../frob", writable = true, allowDotfiles = true))
4186 ],
4187 sockets = [
4188 (name = "main", address = "test-addr", service = "hello")
4189 ]
4190 ))"_kj);
4191 
4192 auto mode = kj::WriteMode::CREATE | kj::WriteMode::CREATE_PARENT;
4193 auto dir = test.root->openSubdir(kj::Path({"frob"_kj}), mode);
4194 
4195 // Put a file at root that shouldn't be accessible.
4196 test.root->openFile(kj::Path({"secret"}), mode)->writeAll("this is super-secret");
4197 
4198 test.start();
4199 
4200 auto conn = test.connect("test-addr");
4201 
4202 conn.send(R"(
4203 PUT /.dot HTTP/1.1
4204 Host: foo
4205 Content-Length: 6
4206 
4207 waldo
4208 )"_blockquote);
4209 conn.recv(R"(
4210 HTTP/1.1 204 No Content
4211 
4212 )"_blockquote);
4213 
4214 KJ_EXPECT(dir->openFile(kj::Path({".dot"}))->readAllText() == "waldo\n");
4215 
4216 conn.sendHttpGet("/.dot");
4217 conn.recv(R"(
4218 HTTP/1.1 200 OK
4219 Content-Length: 6
4220 Content-Type: application/octet-stream
4221 Last-Modified: Thu, 01 Jan 1970 00:00:00 GMT
4222 
4223 waldo
4224 )"_blockquote);
4225 
4226 conn.sendHttpGet("/../secret");
4227 conn.recv(R"(
4228 HTTP/1.1 404 Not Found
4229 Content-Length: 9
4230 
4231 Not Found)"_blockquote);
4232 conn.sendHttpGet("/%2e%2e/secret");
4233 conn.recv(R"(
4234 HTTP/1.1 404 Not Found
4235 Content-Length: 9
4236 
4237 Not Found)"_blockquote);
4238 
4239 conn.send(R"(
4240 PUT /../secret HTTP/1.1
4241 Host: foo
4242 Content-Length: 5
4243 
4244 evil
4245 )"_blockquote);
4246 conn.recv(R"(
4247 HTTP/1.1 204 No Content
4248 
4249 )"_blockquote);
4250 // This actually wrote to /secret, because URL parsing simply ignores leading "../".
4251 KJ_EXPECT(dir->openFile(kj::Path({"secret"}))->readAllText() == "evil\n");
4252 KJ_EXPECT(test.root->openFile(kj::Path({"secret"}))->readAllText() == "this is super-secret");
4253 
4254 conn.send(R"(
4255 PUT /%2e%2e/secret HTTP/1.1
4256 Host: foo
4257 Content-Length: 5
4258 
4259 evil
4260 )"_blockquote);
4261 conn.recv(R"(
4262 HTTP/1.1 403 Unauthorized
4263 Content-Length: 12
4264 
4265 Unauthorized)"_blockquote);
4266 // This didn't work.
4267 KJ_EXPECT(test.root->openFile(kj::Path({"secret"}))->readAllText() == "this is super-secret");
4268}
4269 
4270// =======================================================================================
4271// Test Cache API
4272 
4273KJ_TEST("Server: If no cache service is defined, access to the cache API should error") {
4274 TestServer test(singleWorker(R"((
4275 compatibilityDate = "2022-08-17",
4276 modules = [
4277 ( name = "test.js",
4278 esModule =
4279 `export default {
4280 ` async fetch(request) {
4281 ` try {
4282 ` return new Response(await caches.default.match(request))
4283 ` } catch (e) {return new Response(e.message)}
4284 `
4285 ` }
4286 `}
4287 )
4288 ]
4289 ))"_kj));
4290 
4291 test.start();
4292 auto conn = test.connect("test-addr");
4293 conn.httpGet200("/", "No Cache was configured");
4294}
4295 
4296KJ_TEST("Server: cached response") {
4297 TestServer test(R"((
4298 services = [
4299 ( name = "hello",
4300 worker = (
4301 cacheApiOutbound = "cache-outbound",
4302 compatibilityDate = "2022-08-17",
4303 modules = [
4304 ( name = "main.js",
4305 esModule =
4306 `export default {
4307 ` async fetch(request, env, ctx) {
4308 ` const cache = caches.default;
4309 ` let response = await cache.match(request);
4310 ` return response ?? new Response('not cached');
4311 ` }
4312 `}
4313 )
4314 ]
4315 )
4316 ),
4317 ( name = "cache-outbound", external = "cache-host" ),
4318 ],
4319 sockets = [
4320 ( name = "main",
4321 address = "test-addr",
4322 service = "hello"
4323 )
4324 ]
4325 ))"_kj);
4326 
4327 test.start();
4328 auto conn = test.connect("test-addr");
4329 conn.sendHttpGet("/");
4330 
4331 {
4332 auto subreq = test.receiveSubrequest("cache-host");
4333 subreq.recv(R"(
4334 GET / HTTP/1.1
4335 Host: foo
4336 Cache-Control: only-if-cached
4337 
4338 )"_blockquote);
4339 subreq.send(R"(
4340 HTTP/1.1 200 OK
4341 CF-Cache-Status: HIT
4342 Content-Length: 6
4343 
4344 cached)"_blockquote);
4345 }
4346 
4347 conn.recv(R"(
4348 HTTP/1.1 200 OK
4349 Content-Length: 6
4350 CF-Cache-Status: HIT
4351 
4352 cached)"_blockquote);
4353}
4354 
4355KJ_TEST("Server: cache name is passed through to service") {
4356 TestServer test(R"((
4357 services = [
4358 ( name = "hello",
4359 worker = (
4360 cacheApiOutbound = "cache-outbound",
4361 compatibilityDate = "2022-08-17",
4362 modules = [
4363 ( name = "main.js",
4364 esModule =
4365 `export default {
4366 ` async fetch(request, env, ctx) {
4367 ` const cache = await caches.open('test-cache');
4368 ` let response = await cache.match(request);
4369 ` return response ?? new Response('not cached');
4370 ` }
4371 `}
4372 )
4373 ]
4374 )
4375 ),
4376 ( name = "cache-outbound", external = "cache-host" ),
4377 ],
4378 sockets = [
4379 ( name = "main",
4380 address = "test-addr",
4381 service = "hello"
4382 )
4383 ]
4384 ))"_kj);
4385 
4386 test.start();
4387 auto conn = test.connect("test-addr");
4388 conn.sendHttpGet("/");
4389 
4390 {
4391 auto subreq = test.receiveSubrequest("cache-host");
4392 subreq.recv(R"(
4393 GET / HTTP/1.1
4394 Host: foo
4395 Cache-Control: only-if-cached
4396 CF-Cache-Namespace: test-cache
4397 
4398 )"_blockquote);
4399 subreq.send(R"(
4400 HTTP/1.1 200 OK
4401 CF-Cache-Status: HIT
4402 Content-Length: 6
4403 
4404 cached)"_blockquote);
4405 }
4406 
4407 conn.recv(R"(
4408 HTTP/1.1 200 OK
4409 Content-Length: 6
4410 CF-Cache-Status: HIT
4411 
4412 cached)"_blockquote);
4413}
4414 
4415// =======================================================================================
4416// Test the test command
4417 
4418KJ_TEST("Server: cache name is passed through to service") {
4419 kj::StringPtr config = R"((
4420 services = [
4421 ( name = "hello",
4422 worker = (
4423 compatibilityDate = "2022-08-17",
4424 modules = [
4425 ( name = "main.js",
4426 esModule =
4427 `export default {
4428 ` async test(controller, env, ctx) {}
4429 `}
4430 `export let fail = {
4431 ` async test(controller, env, ctx) {
4432 ` throw new Error("ded");
4433 ` }
4434 `}
4435 `export let nonTest = {
4436 ` async fetch(req, env, ctx) {
4437 ` return new Response("ok");
4438 ` }
4439 `}
4440 )
4441 ]
4442 )
4443 ),
4444 ( name = "another",
4445 worker = (
4446 compatibilityDate = "2022-08-17",
4447 modules = [
4448 ( name = "main.js",
4449 esModule =
4450 `export default {
4451 ` async test(controller, env, ctx) {
4452 ` console.log(env.MESSAGE);
4453 ` }
4454 `}
4455 )
4456 ],
4457 bindings = [
4458 ( name = "MESSAGE", text = "other test" ),
4459 ]
4460 )
4461 ),
4462 ],
4463 sockets = [
4464 ( name = "main",
4465 address = "test-addr",
4466 service = "hello"
4467 )
4468 ]
4469 ))"_kj;
4470 
4471 {
4472 TestServer test(config);
4473 KJ_EXPECT_LOG(DBG, "[ TEST ] hello");
4474 KJ_EXPECT_LOG(DBG, "[ PASS ] hello");
4475 KJ_EXPECT(test.server.test(v8System, *test.config, "hello", "default").wait(test.ws));
4476 }
4477 
4478 {
4479 TestServer test(config);
4480 KJ_EXPECT_LOG(DBG, "[ TEST ] hello:fail");
4481 KJ_EXPECT_LOG(INFO, "Error: ded");
4482 KJ_EXPECT_LOG(DBG, "[ FAIL ] hello:fail");
4483 KJ_EXPECT(!test.server.test(v8System, *test.config, "hello", "fail").wait(test.ws));
4484 }
4485 
4486 {
4487 TestServer test(config);
4488 KJ_EXPECT_LOG(DBG, "[ TEST ] hello");
4489 KJ_EXPECT_LOG(DBG, "[ PASS ] hello");
4490 KJ_EXPECT_LOG(DBG, "[ TEST ] hello:fail");
4491 KJ_EXPECT_LOG(INFO, "Error: ded");
4492 KJ_EXPECT_LOG(DBG, "[ FAIL ] hello:fail");
4493 KJ_EXPECT(!test.server.test(v8System, *test.config, "hello", "*").wait(test.ws));
4494 }
4495 
4496 {
4497 TestServer test(config);
4498 KJ_EXPECT_LOG(DBG, "[ TEST ] hello");
4499 KJ_EXPECT_LOG(DBG, "[ PASS ] hello");
4500 KJ_EXPECT_LOG(DBG, "[ TEST ] another");
4501 KJ_EXPECT_LOG(INFO, "other test");
4502 KJ_EXPECT_LOG(DBG, "[ PASS ] another");
4503 KJ_EXPECT(test.server.test(v8System, *test.config, "*", "default").wait(test.ws));
4504 }
4505 
4506 {
4507 TestServer test(config);
4508 KJ_EXPECT_LOG(DBG, "[ TEST ] hello");
4509 KJ_EXPECT_LOG(DBG, "[ PASS ] hello");
4510 KJ_EXPECT_LOG(DBG, "[ TEST ] hello:fail");
4511 KJ_EXPECT_LOG(INFO, "Error: ded");
4512 KJ_EXPECT_LOG(DBG, "[ FAIL ] hello:fail");
4513 KJ_EXPECT_LOG(DBG, "[ TEST ] another");
4514 KJ_EXPECT_LOG(INFO, "other test");
4515 KJ_EXPECT_LOG(DBG, "[ PASS ] another");
4516 KJ_EXPECT(!test.server.test(v8System, *test.config, "*", "*").wait(test.ws));
4517 }
4518}
4519 
4520// =======================================================================================
4521 
4522KJ_TEST("Server: JS RPC over HTTP connections") {
4523 // Test that we can send RPC over an ExternalServer pointing back to our own loopback socket,
4524 // as long as both are configured with a `capnpConnectHost`.
4525 
4526 TestServer test(R"((
4527 services = [
4528 ( name = "hello",
4529 worker = (
4530 compatibilityDate = "2024-02-23",
4531 compatibilityFlags = ["experimental"],
4532 modules = [
4533 ( name = "main.js",
4534 esModule =
4535 `import {WorkerEntrypoint} from "cloudflare:workers";
4536 `export default {
4537 ` async fetch(request, env) {
4538 ` return new Response("got: " + await env.OUT.frob(3, 11));
4539 ` }
4540 `}
4541 `export class MyRpc extends WorkerEntrypoint {
4542 ` async frob(a, b) { return a * b + 2; }
4543 `}
4544 )
4545 ],
4546 bindings = [( name = "OUT", service = "outbound")]
4547 )
4548 ),
4549 (name = "outbound", external = (address = "loopback", http = (capnpConnectHost = "cappy")))
4550 ],
4551 sockets = [
4552 ( name = "main", address = "test-addr", service = "hello" ),
4553 ( name = "alt1", address = "loopback",
4554 service = (name = "hello", entrypoint = "MyRpc"),
4555 http = (capnpConnectHost = "cappy")),
4556 ]
4557 ))"_kj);
4558 
4559 test.server.allowExperimental();
4560 test.start();
4561 
4562 auto conn = test.connect("test-addr");
4563 conn.httpGet200("/", "got: 35");
4564}
4565 
4566KJ_TEST("Server: Entrypoint binding with props") {
4567 TestServer test(R"((
4568 services = [
4569 ( name = "hello",
4570 worker = (
4571 compatibilityDate = "2024-02-23",
4572 compatibilityFlags = ["experimental"],
4573 modules = [
4574 ( name = "main.js",
4575 esModule =
4576 `import {WorkerEntrypoint} from "cloudflare:workers";
4577 `export default {
4578 ` async fetch(request, env) {
4579 ` return new Response("got: " + await env.MyRpc.getProps());
4580 ` }
4581 `}
4582 `export class MyRpc extends WorkerEntrypoint {
4583 ` getProps() { return this.ctx.props.foo; }
4584 `}
4585 )
4586 ],
4587 bindings = [
4588 ( name = "MyRpc",
4589 service = (
4590 name = "hello",
4591 entrypoint = "MyRpc",
4592 props = (
4593 json = `{"foo": 123}
4594 )
4595 )
4596 )
4597 ]
4598 )
4599 ),
4600 ],
4601 sockets = [
4602 ( name = "main", address = "test-addr", service = "hello" ),
4603 ]
4604 ))"_kj);
4605 
4606 test.server.allowExperimental();
4607 test.start();
4608 
4609 auto conn = test.connect("test-addr");
4610 conn.httpGet200("/", "got: 123");
4611}
4612 
4613KJ_TEST("Server: ctx.exports self-referential bindings") {
4614 TestServer test(R"((
4615 services = [
4616 ( name = "hello",
4617 worker = (
4618 compatibilityDate = "2025-02-23",
4619 compatibilityFlags = ["enable_ctx_exports"],
4620 modules = [
4621 ( name = "main.js",
4622 esModule =
4623 `import { WorkerEntrypoint, DurableObject, WorkflowEntrypoint } from "cloudflare:workers";
4624 `export default {
4625 ` async fetch(request, env, ctx) {
4626 ` // First set the actor state the old fashion way, to make sure we get
4627 ` // reconnected to the same actor when using self-referential bindings.
4628 ` {
4629 ` let bindingActor = env.NS.get(env.NS.idFromName("qux"));
4630 ` await bindingActor.setValue(234);
4631 ` }
4632 `
4633 ` let actor = ctx.exports.MyActor.get(ctx.exports.MyActor.idFromName("qux"));
4634 ` return new Response([
4635 ` await ctx.exports.MyEntrypoint.foo(123),
4636 ` await ctx.exports.AnotherEntrypoint.bar(321),
4637 ` await actor.baz(),
4638 ` await ctx.exports.default.corge(555),
4639 ` await actor.grault(456),
4640 ` ctx.exports.UnconfiguredActor.constructor.name,
4641 ` await ctx.exports.MyEntrypoint.myProps(),
4642 ` await ctx.exports.MyEntrypoint({props: {foo: 123, bar: "abc"}}).myProps(),
4643 ` MyWorkflow in ctx.exports,
4644 ` ].join(", "));
4645 ` },
4646 ` corge(i) { return `corge: ${i}` }
4647 `}
4648 `export class MyEntrypoint extends WorkerEntrypoint {
4649 ` foo(i) { return `foo: ${i}` }
4650 ` grault(i) { return `grault: ${i}` }
4651 ` myProps() { return JSON.stringify(this.ctx.props) }
4652 `}
4653 `export class AnotherEntrypoint extends WorkerEntrypoint {
4654 ` bar(i) { return `bar: ${i}` }
4655 `}
4656 `export class MyActor extends DurableObject {
4657 ` setValue(i) { this.value = i; }
4658 ` baz() { return `baz: ${this.value}` }
4659 ` grault(i) { return this.ctx.exports.MyEntrypoint.grault(i); }
4660 `}
4661 `export class UnconfiguredActor extends DurableObject {
4662 ` qux(i) { return `qux: ${i}` }
4663 `}
4664 `export class MyWorkflow extends WorkflowEntrypoint {}
4665 )
4666 ],
4667 bindings = [
4668 # A regular binding, just here to make sure it doesn't mess up self-referential
4669 # channel numbers.
4670 ( name = "INTERNET", service = "internet" ),
4671 
4672 # Similarly, an actor namespace binding.
4673 (name = "NS", durableObjectNamespace = "MyActor")
4674 ],
4675 durableObjectNamespaces = [
4676 ( className = "MyActor",
4677 uniqueKey = "mykey",
4678 )
4679 ],
4680 durableObjectStorage = (inMemory = void)
4681 )
4682 ),
4683 ],
4684 sockets = [
4685 ( name = "main", address = "test-addr", service = "hello" ),
4686 ]
4687 ))"_kj);
4688 
4689 test.server.allowExperimental();
4690 test.start();
4691 
4692 auto conn = test.connect("test-addr");
4693 conn.httpGet200("/",
4694 "foo: 123, bar: 321, baz: 234, corge: 555, grault: 456, LoopbackDurableObjectClass, "
4695 "{}, {\"foo\":123,\"bar\":\"abc\"}, false");
4696}
4697 
4698KJ_TEST("Server: loopback binding calls accept version property") {
4699 TestServer test(R"((
4700 services = [
4701 ( name = "hello",
4702 worker = (
4703 compatibilityDate = "2025-08-01",
4704 compatibilityFlags = ["enable_ctx_exports", "enable_version_api"],
4705 modules = [
4706 ( name = "main.js",
4707 esModule =
4708 `export default {
4709 ` async fetch(request, env, ctx) {
4710 ` const serviceVersions = await Promise.all([
4711 ` ctx.exports.default({ version: {} }),
4712 ` ctx.exports.default({ version: { cohort: null } }),
4713 ` ctx.exports.default({ version: { cohort: "test" } }),
4714 ` ctx.exports.default({ props: {}, version: { cohort: "test" } }),
4715 ` ].map(service => service.version));
4716 ` if (serviceVersions.every(version => version === this.version)) {
4717 ` return new Response(serviceVersions[0]);
4718 ` }
4719 ` return new Response(null, { status: 500 });
4720 ` },
4721 ` get version() { return "constant"; },
4722 `}
4723 )
4724 ],
4725 )
4726 ),
4727 ],
4728 sockets = [
4729 ( name = "main", address = "test-addr", service = "hello" ),
4730 ]
4731 ))"_kj);
4732 
4733 test.server.allowExperimental();
4734 test.start();
4735 
4736 auto conn = test.connect("test-addr");
4737 conn.httpGet200("/", "constant");
4738}
4739 
4740// =======================================================================================
4741 
4742// TODO(beta): Test TLS (send and receive)
4743// TODO(beta): Test CLI overrides
4744 
4745KJ_TEST("Server: encodeResponseBody: manual option") {
4746 TestServer test(R"((
4747 services = [
4748 ( name = "hello",
4749 worker = (
4750 compatibilityDate = "2022-08-17",
4751 modules = [
4752 ( name = "main.js",
4753 esModule =
4754 `export default {
4755 ` async fetch(request, env) {
4756 ` // Make a subrequest with encodeResponseBody: "manual"
4757 ` let response = await fetch("http://subhost/foo", {
4758 ` encodeResponseBody: "manual"
4759 ` });
4760 `
4761 ` // Get the raw bytes, which should not be decompressed
4762 ` let rawBytes = await response.arrayBuffer();
4763 ` let decoder = new TextDecoder();
4764 ` let rawText = decoder.decode(rawBytes);
4765 `
4766 ` return new Response(
4767 ` "Content-Encoding: " + response.headers.get("Content-Encoding") + "\n" +
4768 ` "Raw content: " + rawText
4769 ` );
4770 ` }
4771 `}
4772 )
4773 ]
4774 )
4775 )
4776 ],
4777 sockets = [
4778 ( name = "main",
4779 address = "test-addr",
4780 service = "hello"
4781 )
4782 ]
4783 ))"_kj);
4784 
4785 test.start();
4786 auto conn = test.connect("test-addr");
4787 conn.sendHttpGet("/");
4788 
4789 auto subreq = test.receiveInternetSubrequest("subhost");
4790 subreq.recv(R"(
4791 GET /foo HTTP/1.1
4792 Host: subhost
4793 
4794 )"_blockquote);
4795 
4796 // Send a response with Content-Encoding: gzip, but the body is not actually
4797 // compressed - it's just "fake-gzipped-content" as plain text
4798 subreq.send(R"(
4799 HTTP/1.1 200 OK
4800 Content-Length: 20
4801 Content-Encoding: gzip
4802 
4803 fake-gzipped-content
4804 )"_blockquote);
4805 
4806 // Verify that:
4807 // 1. The Content-Encoding header was preserved
4808 // 2. The body was not decompressed (we get the raw "fake-gzipped-content")
4809 conn.recvHttp200(R"(
4810 Content-Encoding: gzip
4811 Raw content: fake-gzipped-content)"_blockquote);
4812}
4813 
4814KJ_TEST("Server: encodeResponseBody: manual pass-through") {
4815 TestServer test(R"((
4816 services = [
4817 ( name = "hello",
4818 worker = (
4819 compatibilityDate = "2022-08-17",
4820 modules = [
4821 ( name = "main.js",
4822 esModule =
4823 `export default {
4824 ` async fetch(request, env) {
4825 ` // Make a subrequest with encodeResponseBody: "manual" and pass through the response
4826 ` return fetch("http://subhost/foo", {
4827 ` encodeResponseBody: "manual"
4828 ` });
4829 ` }
4830 `}
4831 )
4832 ]
4833 )
4834 )
4835 ],
4836 sockets = [
4837 ( name = "main",
4838 address = "test-addr",
4839 service = "hello"
4840 )
4841 ]
4842 ))"_kj);
4843 
4844 test.start();
4845 auto conn = test.connect("test-addr");
4846 conn.sendHttpGet("/");
4847 
4848 auto subreq = test.receiveInternetSubrequest("subhost");
4849 subreq.recv(R"(
4850 GET /foo HTTP/1.1
4851 Host: subhost
4852 
4853 )"_blockquote);
4854 
4855 // Send a response with Content-Encoding: gzip, but the body is not actually
4856 // compressed - it's just "fake-gzipped-content" as plain text
4857 subreq.send(R"(
4858 HTTP/1.1 200 OK
4859 Content-Length: 20
4860 Content-Encoding: gzip
4861 
4862 fake-gzipped-content
4863 )"_blockquote);
4864 
4865 // Verify that the response is passed through verbatim, with:
4866 // 1. The Content-Encoding header preserved
4867 // 2. The body not decompressed
4868 // 3. The body not re-encoded
4869 conn.recv(R"(
4870 HTTP/1.1 200 OK
4871 Content-Length: 20
4872 Content-Encoding: gzip
4873 
4874 fake-gzipped-content)"_blockquote);
4875}
4876 
4877KJ_TEST("Server: Catch websocket server errors") {
4878 TestServer test(R"((
4879 services = [
4880 ( name = "hello",
4881 worker = (
4882 compatibilityDate = "2025-04-01",
4883 modules = [
4884 ( name = "main.js",
4885 esModule =
4886 ` export default {
4887 ` async fetch(request) {
4888 ` try {
4889 ` return await handleRequest(request)
4890 ` } catch (e) {
4891 ` console.log("eerrrrr", e)
4892 ` return new Response("ok")
4893 ` }
4894 ` }
4895 ` }
4896 `
4897 ` let lastError = "none";
4898 `
4899 ` async function handleRequest(request) {
4900 ` const upgradeHeader = request.headers.get('Upgrade');
4901 ` if (!upgradeHeader || upgradeHeader !== 'websocket') {
4902 ` return new Response('Expected Upgrade: websocket' , { status: 426 });
4903 ` }
4904 `
4905 ` const webSocketPair = new WebSocketPair();
4906 ` const [client, server] = Object.values(webSocketPair);
4907 `
4908 ` server.accept();
4909 ` server.addEventListener('message', event => {
4910 ` if (event.data === "getLastError") {
4911 ` server.send(lastError)
4912 ` } else {
4913 ` let msg = event.data
4914 ` server.send(msg)
4915 ` }
4916 ` });
4917 `
4918 ` server.addEventListener('error', event => {
4919 ` lastError = event.message;
4920 ` });
4921 `
4922 ` return new Response(null, {
4923 ` status: 101,
4924 ` webSocket: client,
4925 ` });
4926 ` }
4927 )
4928 ]
4929 )
4930 ),
4931 ],
4932 sockets = [
4933 ( name = "main",
4934 address = "test-addr",
4935 service = "hello"
4936 )
4937 ]
4938 ))"_kj);
4939 
4940 class NotVeryGoodEntropySource: public kj::EntropySource {
4941 public:
4942 void generate(kj::ArrayPtr<byte> buffer) override {
4943 buffer.fill('4');
4944 }
4945 };
4946 
4947 KJ_EXPECT_LOG(ERROR,
4948 "jsg.Error: WebSocket protocol error; protocolError.statusCode = 1009; protocolError.description = Message is too large: 34603008 > 33554432");
4949 test.start();
4950 auto& waitScope = test.getWaitScope();
4951 
4952 kj::HttpHeaderTable headerTable;
4953 NotVeryGoodEntropySource entropySource;
4954 kj::HttpHeaders headers(headerTable);
4955 headers.setPtr(kj::HttpHeaderId::HOST, "foo");
4956 headers.setPtr(kj::HttpHeaderId::UPGRADE, "websocket");
4957 {
4958 auto wsConn = test.connect("test-addr");
4959 auto client = kj::newHttpClient(
4960 headerTable, wsConn.getStream(), kj::HttpClientSettings{.entropySource = entropySource});
4961 auto res = client->openWebSocket("/", headers).wait(waitScope);
4962 KJ_ASSERT(res.statusCode == 101, res.statusCode, res.statusText);
4963 auto ws = kj::mv(res.webSocketOrBody.get<kj::Own<kj::WebSocket>>());
4964 const auto smallMessage = kj::str("hello");
4965 ws->send(smallMessage).wait(waitScope);
4966 auto smallResponse = ws->receive().wait(waitScope);
4967 KJ_EXPECT(smallResponse.get<kj::String>() == smallMessage);
4968 const auto bigMessage = kj::heapArray<kj::byte>(33 * 1024 * 1024);
4969 auto sendProm =
4970 kj::evalNow([&]() { return ws->send(bigMessage); }).then([]() {}, [](kj::Exception ex) {});
4971 // Message is too big; we should close the connection.
4972 auto msg = ws->receive().wait(waitScope);
4973 sendProm.wait(waitScope);
4974 auto& resp = msg.get<kj::WebSocket::Close>();
4975 KJ_EXPECT(resp.code == 1009); // WebSocket-ese for "message too large"
4976 }
4977 {
4978 auto wsConn = test.connect("test-addr");
4979 headers.setPtr(kj::HttpHeaderId::HOST, "foo");
4980 headers.setPtr(kj::HttpHeaderId::UPGRADE, "websocket");
4981 auto client = kj::newHttpClient(
4982 headerTable, wsConn.getStream(), kj::HttpClientSettings{.entropySource = entropySource});
4983 auto res = client->openWebSocket("/", headers).wait(waitScope);
4984 KJ_ASSERT(res.statusCode == 101, res.statusCode, res.statusText);
4985 auto ws = kj::mv(res.webSocketOrBody.get<kj::Own<kj::WebSocket>>());
4986 const auto query = kj::str("getLastError");
4987 ws->send(query).wait(waitScope);
4988 auto response = ws->receive().wait(waitScope);
4989 
4990 kj::StringPtr responseString = response.get<kj::String>();
4991 KJ_EXPECT(responseString.find("1009"_kjc) != kj::none, responseString); // Error code
4992 KJ_EXPECT(responseString.find("Message is too large"_kjc) != kj::none, responseString);
4993 ws->close(1000, "").wait(waitScope);
4994 }
4995}
4996 
4997KJ_TEST("Server: Durable Object facets") {
4998 kj::StringPtr config = R"((
4999 services = [
5000 ( name = "hello",
5001 worker = (
5002 compatibilityDate = "2026-04-01",
5003 modules = [
5004 ( name = "main.js",
5005 esModule =
5006 `import { DurableObject } from "cloudflare:workers";
5007 `export default {
5008 ` async fetch(request, env, ctx) {
5009 ` let id = ctx.exports.MyActorClass.idFromName("name");
5010 ` let actor = ctx.exports.MyActorClass.get(id);
5011 ` return await actor.fetch(request);
5012 ` }
5013 `}
5014 `export class MyActorClass extends DurableObject {
5015 ` async fetch(request) {
5016 ` let results = [];
5017 `
5018 ` if (request.url.endsWith("/part1")) {
5019 ` let foo = this.ctx.facets.get("foo",
5020 ` () => ({class: this.ctx.exports.CounterFacet, id: "abc"}));
5021 ` results.push(await foo.increment(true)); // increments foo
5022 ` results.push(await foo.increment()); // increments foo
5023 ` results.push(await foo.increment()); // increments foo
5024 ` await foo.assertId("abc");
5025 `
5026 ` let bar = this.ctx.facets.get("bar", () => ({class: this.env.NESTED}));
5027 ` results.push(await bar.increment("foo", true)); // increments bar.foo
5028 ` results.push(await bar.increment("bar", true)); // increments bar.bar
5029 ` results.push(await bar.increment("foo")); // increments bar.foo
5030 ` await bar.assertId(this.ctx.id.toString());
5031 `
5032 ` // Get foo again to make sure we get the same object.
5033 ` let foo2 = this.ctx.facets.get("foo", () => {
5034 ` throw new Error("callback should not be called when already running");
5035 ` });
5036 ` results.push(await foo2.increment()); // increments foo
5037 ` results.push(await foo.increment()); // increments foo
5038 ` await foo.assertId("abc");
5039 ` } else if (request.url.endsWith("/part2")) {
5040 ` let callbackCount = 0;
5041 `
5042 ` // Get in a different order from before to make sure ID assignment is
5043 ` // consistent.
5044 ` let bar = this.ctx.facets.get("bar", () => {
5045 ` ++callbackCount;
5046 ` return {class: this.env.NESTED};
5047 ` });
5048 ` results.push(await bar.increment("bar", true)); // increments bar.bar
5049 ` results.push(await bar.increment("foo", true)); // increments bar.foo
5050 ` let foo = this.ctx.facets.get("foo", async () => {
5051 ` await Promise.resolve(); // prove that callback can be async
5052 ` ++callbackCount;
5053 ` return {class: this.env.COUNTER, id: "abc"};
5054 ` });
5055 ` results.push(await foo.increment(true)); // increments foo
5056 `
5057 ` if (callbackCount !== 2) {
5058 ` throw new Error(`callbackCount = ${callbackCount} (expected 2)`);
5059 ` }
5060 `
5061 ` // Force "foo" to abort, so we can start it up with a different class.
5062 ` this.ctx.facets.abort("foo", new Error("test abort facet"));
5063 `
5064 ` let foo2 = this.ctx.facets.get(
5065 ` "foo", () => ({class: this.env.EXFILTRATOR, id: "abc"}));
5066 ` results.push(await foo2.exfiltrate());
5067 `
5068 ` try {
5069 ` await foo.increment();
5070 ` throw new Error("broken stub didn't throw?");
5071 ` } catch (err) {
5072 ` if (err.message != "test abort facet") {
5073 ` throw err;
5074 ` }
5075 ` }
5076 `
5077 ` // Delete bar, which recursively deletes its children.
5078 ` this.ctx.facets.delete("bar");
5079 ` } else if (request.url.endsWith("/props")) {
5080 ` results.push(JSON.stringify(this.ctx.props));
5081 `
5082 ` let prop1 = this.ctx.facets.get("prop1",
5083 ` () => ({class: this.env.COUNTER, id: "abc"}));
5084 ` results.push(await prop1.myProps());
5085 `
5086 ` let prop2 = this.ctx.facets.get("prop2",
5087 ` () => ({class: this.ctx.exports.CounterFacet, id: "abc"}));
5088 ` results.push(await prop2.myProps());
5089 `
5090 ` let prop3 = this.ctx.facets.get("prop3",
5091 ` () => ({class: this.ctx.exports.CounterFacet({props: {bProp: 321}}),
5092 ` id: "abc"}));
5093 ` results.push(await prop3.myProps());
5094 `
5095 ` let prop4 = this.ctx.facets.get("prop4",
5096 ` () => ({class: this.ctx.exports.MyActorClass, id: "abc"}));
5097 ` results.push(await prop4.mainClassProps());
5098 `
5099 ` let prop5 = this.ctx.facets.get("prop5",
5100 ` () => ({class: this.ctx.exports.MyActorClass({props: {cProp: 555}}),
5101 ` id: "abc"}));
5102 ` results.push(await prop5.mainClassProps());
5103 ` } else {
5104 ` throw new Error(`bad url: ${request.url}`);
5105 ` }
5106 `
5107 ` return new Response(results.join(" "));
5108 ` }
5109 ` mainClassProps() { return JSON.stringify(this.ctx.props) }
5110 `}
5111 `export class CounterFacet extends DurableObject {
5112 ` async increment(first) {
5113 ` let storedI = (await this.ctx.storage.get("value")) || 0;
5114 ` if (first) {
5115 ` this.i = storedI;
5116 ` } else if (this.i != storedI) {
5117 ` throw new Error("inconsistent stored value ${storedI} != ${this.i}");
5118 ` }
5119 ` this.ctx.storage.put("value", this.i + 1);
5120 ` return this.i++;
5121 ` }
5122 ` assertId(id) {
5123 ` if (this.ctx.id.toString() != id) {
5124 ` throw new Error(`Wrong ID, expected ${id}, got ${this.ctx.id}`);
5125 ` }
5126 ` }
5127 ` myProps() { return JSON.stringify(this.ctx.props) }
5128 `}
5129 `export class NestedFacet extends DurableObject {
5130 ` increment(name, first) {
5131 ` let facet = this.ctx.facets.get(name, () => ({class: this.env.COUNTER}));
5132 ` return facet.increment(first);
5133 ` }
5134 ` assertId(id) {
5135 ` if (this.ctx.id.toString() != id) {
5136 ` throw new Error(`Wrong ID, expected ${id}, got ${this.ctx.id}`);
5137 ` }
5138 ` }
5139 `}
5140 `export class ExfiltrationFacet extends DurableObject {
5141 ` exfiltrate() {
5142 ` return this.ctx.storage.get("value");
5143 ` }
5144 `}
5145 )
5146 ],
5147 bindings = [
5148 ( name = "COUNTER",
5149 durableObjectClass = (
5150 name = "hello",
5151 entrypoint = "CounterFacet",
5152 props = (
5153 json = `{"aProp": 123}
5154 )
5155 )
5156 ),
5157 (name = "NESTED", durableObjectClass = (name = "hello", entrypoint = "NestedFacet")),
5158 ( name = "EXFILTRATOR",
5159 durableObjectClass = (name = "hello", entrypoint = "ExfiltrationFacet") )
5160 ],
5161 durableObjectNamespaces = [
5162 ( className = "MyActorClass",
5163 uniqueKey = "mykey",
5164 )
5165 ],
5166 durableObjectStorage = (localDisk = "my-disk")
5167 )
5168 ),
5169 ( name = "my-disk",
5170 disk = (
5171 path = "../../do-storage",
5172 writable = true,
5173 )
5174 ),
5175 ],
5176 sockets = [
5177 ( name = "main",
5178 address = "test-addr",
5179 service = "hello"
5180 )
5181 ]
5182 ))"_kj;
5183 
5184 // Create a directory outside of the test scope which we can use across multiple TestServers.
5185 auto dir = kj::newInMemoryDirectory(kj::nullClock());
5186 
5187 {
5188 TestServer test(config);
5189 
5190 // Link our directory into the test filesystem.
5191 test.root->transfer(
5192 kj::Path({"do-storage"_kj}), kj::WriteMode::CREATE, *dir, nullptr, kj::TransferMode::LINK);
5193 
5194 test.server.allowExperimental();
5195 test.start();
5196 auto conn = test.connect("test-addr");
5197 conn.httpGet200("/part1", "0 1 2 0 0 1 3 4");
5198 }
5199 
5200 // Verify the expected files exist.
5201 auto nsDir = dir->openSubdir(kj::Path({"mykey"}));
5202 KJ_EXPECT(nsDir->exists(
5203 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.sqlite"})));
5204 KJ_EXPECT(nsDir->exists(
5205 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.1.sqlite"})));
5206 KJ_EXPECT(nsDir->exists(
5207 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.2.sqlite"})));
5208 KJ_EXPECT(nsDir->exists(
5209 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.3.sqlite"})));
5210 KJ_EXPECT(nsDir->exists(
5211 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.4.sqlite"})));
5212 KJ_EXPECT(nsDir->exists(
5213 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.facets"})));
5214 
5215 // We should only have created four child facets (foo, bar, bar.foo, bar.bar). No ID 5 should
5216 // exist.
5217 KJ_EXPECT(!nsDir->exists(
5218 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.5.sqlite"})));
5219 
5220 // We didn't create any other durable objects in the namespace. All files in the namespace should
5221 // be prefixed with our one DO ID, except for metadata.sqlite (the per-namespace alarm scheduler).
5222 for (auto& name: nsDir->listNames()) {
5223 KJ_EXPECT(
5224 name.startsWith("3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.") ||
5225 name.startsWith("metadata.sqlite"),
5226 "unexpected file found in namespace storage", name);
5227 }
5228 
5229 // Start a new server, make sure it's able to load the files again.
5230 {
5231 TestServer test(config);
5232 
5233 // Link our directory into the test filesystem.
5234 test.root->transfer(
5235 kj::Path({"do-storage"_kj}), kj::WriteMode::CREATE, *dir, nullptr, kj::TransferMode::LINK);
5236 
5237 test.server.allowExperimental();
5238 test.start();
5239 auto conn = test.connect("test-addr");
5240 conn.httpGet200("/part2", "1 2 5 6");
5241 }
5242 
5243 // Root and foo still exist, bar does not.
5244 KJ_EXPECT(nsDir->exists(
5245 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.sqlite"})));
5246 KJ_EXPECT(nsDir->exists(
5247 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.1.sqlite"})));
5248 KJ_EXPECT(!nsDir->exists(
5249 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.2.sqlite"})));
5250 KJ_EXPECT(!nsDir->exists(
5251 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.3.sqlite"})));
5252 KJ_EXPECT(!nsDir->exists(
5253 kj::Path({"3652ef6221834806dc8df802d1d216e27b7d07e0a6b7adf6cfdaeec90f06459a.4.sqlite"})));
5254 
5255 // Test facets can have custom ctx.props.
5256 {
5257 TestServer test(config);
5258 
5259 // We don't need the existing storage but the path does have to exist for the test to work.
5260 test.root->openSubdir(kj::Path({"do-storage"_kj}), kj::WriteMode::CREATE);
5261 
5262 test.server.allowExperimental();
5263 test.start();
5264 auto conn = test.connect("test-addr");
5265 conn.httpGet200("/props", "{} {\"aProp\":123} {} {\"bProp\":321} {} {\"cProp\":555}");
5266 }
5267}
5268 
5269KJ_TEST("Server: Durable Object facet limits") {
5270 kj::StringPtr config = R"((
5271 services = [
5272 ( name = "hello",
5273 worker = (
5274 compatibilityDate = "2025-04-01",
5275 compatibilityFlags = ["experimental"],
5276 modules = [
5277 ( name = "main.js",
5278 esModule =
5279 `import { DurableObject } from "cloudflare:workers";
5280 `export default {
5281 ` async fetch(request, env, ctx) {
5282 ` let id = env.MY_ACTOR.idFromName("limits");
5283 ` let actor = env.MY_ACTOR.get(id);
5284 ` return await actor.fetch(request);
5285 ` }
5286 `}
5287 `export class MyActorClass extends DurableObject {
5288 ` async fetch(request) {
5289 ` let url = new URL(request.url);
5290 ` switch (url.pathname) {
5291 ` case "/name-too-long": {
5292 ` try {
5293 ` this.ctx.facets.get("x".repeat(257),
5294 ` () => ({class: this.env.RECURSIVE}));
5295 ` return new Response("no error");
5296 ` } catch (e) {
5297 ` return new Response(e.constructor.name + ": " + e.message);
5298 ` }
5299 ` }
5300 ` case "/name-256-ok": {
5301 ` this.ctx.facets.get("x".repeat(256),
5302 ` () => ({class: this.env.RECURSIVE}));
5303 ` return new Response("ok");
5304 ` }
5305 ` case "/abort-name-too-long": {
5306 ` try {
5307 ` this.ctx.facets.abort("x".repeat(257), new Error("test"));
5308 ` return new Response("no error");
5309 ` } catch (e) {
5310 ` return new Response(e.constructor.name + ": " + e.message);
5311 ` }
5312 ` }
5313 ` case "/delete-name-too-long": {
5314 ` try {
5315 ` this.ctx.facets.delete("x".repeat(257));
5316 ` return new Response("no error");
5317 ` } catch (e) {
5318 ` return new Response(e.constructor.name + ": " + e.message);
5319 ` }
5320 ` }
5321 ` case "/depth-ok": {
5322 ` // Create 3 levels of facets below root = 4 total (the max).
5323 ` let facet = this.ctx.facets.get("a",
5324 ` () => ({class: this.env.RECURSIVE}));
5325 ` return new Response(await facet.nestOk(2));
5326 ` }
5327 ` case "/depth-exceeded": {
5328 ` // Create 3 levels below root, then try one more.
5329 ` let facet = this.ctx.facets.get("b",
5330 ` () => ({class: this.env.RECURSIVE}));
5331 ` return new Response(await facet.nestDeeper(2));
5332 ` }
5333 ` }
5334 ` }
5335 `}
5336 `export class RecursiveFacet extends DurableObject {
5337 ` async nestOk(remaining) {
5338 ` if (remaining <= 0) return "ok";
5339 ` let facet = this.ctx.facets.get("child",
5340 ` () => ({class: this.env.RECURSIVE}));
5341 ` return await facet.nestOk(remaining - 1);
5342 ` }
5343 ` async nestDeeper(remaining) {
5344 ` if (remaining <= 0) {
5345 ` try {
5346 ` this.ctx.facets.get("too-deep",
5347 ` () => ({class: this.env.RECURSIVE}));
5348 ` return "no error, unexpected";
5349 ` } catch (e) {
5350 ` return e.constructor.name + ": " + e.message;
5351 ` }
5352 ` }
5353 ` let facet = this.ctx.facets.get("child",
5354 ` () => ({class: this.env.RECURSIVE}));
5355 ` return await facet.nestDeeper(remaining - 1);
5356 ` }
5357 `}
5358 )
5359 ],
5360 bindings = [
5361 (name = "MY_ACTOR", durableObjectNamespace = "MyActorClass"),
5362 (name = "RECURSIVE",
5363 durableObjectClass = (name = "hello", entrypoint = "RecursiveFacet"))
5364 ],
5365 durableObjectNamespaces = [
5366 ( className = "MyActorClass",
5367 uniqueKey = "mykey",
5368 )
5369 ],
5370 durableObjectStorage = (localDisk = "my-disk")
5371 )
5372 ),
5373 ( name = "my-disk",
5374 disk = (
5375 path = "../../do-storage",
5376 writable = true,
5377 )
5378 ),
5379 ],
5380 sockets = [
5381 ( name = "main",
5382 address = "test-addr",
5383 service = "hello"
5384 )
5385 ]
5386 ))"_kj;
5387 
5388 TestServer test(config);
5389 test.root->openSubdir(kj::Path({"do-storage"_kj}), kj::WriteMode::CREATE);
5390 test.server.allowExperimental();
5391 test.start();
5392 auto conn = test.connect("test-addr");
5393 
5394 // Name length limit.
5395 conn.httpGet200("/name-too-long", "TypeError: Facet name is too long (max 256 characters).");
5396 conn.httpGet200("/name-256-ok", "ok");
5397 conn.httpGet200(
5398 "/abort-name-too-long", "TypeError: Facet name is too long (max 256 characters).");
5399 conn.httpGet200(
5400 "/delete-name-too-long", "TypeError: Facet name is too long (max 256 characters).");
5401 
5402 // Depth limit.
5403 conn.httpGet200("/depth-ok", "ok");
5404 conn.httpGet200("/depth-exceeded",
5405 "Error: Facet nesting depth limit exceeded. "
5406 "The maximum depth including the root Durable Object is 4.");
5407}
5408 
5409KJ_TEST("Server: Pass service stubs in ctx.props.") {
5410 TestServer test(R"((
5411 services = [
5412 ( name = "hello",
5413 worker = (
5414 compatibilityDate = "2025-08-01",
5415 compatibilityFlags = ["enable_ctx_exports"],
5416 modules = [
5417 ( name = "main.js",
5418 esModule =
5419 `import { WorkerEntrypoint } from "cloudflare:workers";
5420 `export default {
5421 ` async fetch(request, env, ctx) {
5422 ` let props = {
5423 ` foo: ctx.exports.FooEntry({props: {greeting: "Hello"}}),
5424 ` foo2: ctx.exports.FooEntry({props: {greeting: "Welcome"}}),
5425 ` }
5426 ` let result = await ctx.exports.BarEntry({props}).run();
5427 ` return new Response(result);
5428 ` },
5429 `}
5430 `export class FooEntry extends WorkerEntrypoint {
5431 ` greet(name) { return `${this.ctx.props.greeting}, ${name}!` }
5432 `}
5433 `export class BarEntry extends WorkerEntrypoint {
5434 ` async run() {
5435 ` let greet1 = await this.ctx.props.foo.greet("Alice");
5436 ` let greet2 = await this.ctx.props.foo2.greet("Bob");
5437 ` return [greet1, greet2].join("\n");
5438 ` }
5439 `}
5440 )
5441 ],
5442 )
5443 ),
5444 ],
5445 sockets = [
5446 ( name = "main", address = "test-addr", service = "hello" ),
5447 ]
5448 ))"_kj);
5449 
5450 test.server.allowExperimental();
5451 test.start();
5452 
5453 auto conn = test.connect("test-addr");
5454 conn.httpGet200("/", "Hello, Alice!\nWelcome, Bob!");
5455}
5456 
5457#if __linux__
5458// This test uses pipe2 and dup2 to capture stdout which is far easier on linux.
5459 
5460struct FdPair {
5461 kj::AutoCloseFd output;
5462 kj::AutoCloseFd input;
5463};
5464 
5465auto makePipeFds() {
5466 int pipeFds[2];
5467 KJ_SYSCALL(pipe2(pipeFds, 0));
5468 
5469 return FdPair{
5470 .output = kj::AutoCloseFd(pipeFds[0]),
5471 .input = kj::AutoCloseFd(pipeFds[1]),
5472 };
5473}
5474 
5475template <typename Func>
5476auto expectLogLine(int fd, Func&& f) {
5477 char buffer[4096];
5478 int pos = 0;
5479 char c;
5480 while (read(fd, &c, 1) == 1) {
5481 if (c == '\n') {
5482 break;
5483 }
5484 if (pos < sizeof(buffer) - 1) {
5485 buffer[pos++] = c;
5486 }
5487 }
5488 buffer[pos] = '\0'; // null-terminate
5489 
5490 kj::StringPtr logline(buffer);
5491 f(logline);
5492}
5493 
5494KJ_TEST("Server: structured logging with console methods") {
5495 TestServer test(R"((
5496 services = [
5497 ( name = "hello",
5498 worker = (
5499 compatibilityDate = "2024-11-01",
5500 compatibilityFlags = [
5501 "nodejs_compat",
5502 "experimental",
5503 "enable_nodejs_process_v2"
5504 ],
5505 modules = [
5506 ( name = "main.js",
5507 esModule =
5508 `export default {
5509 ` async fetch(request, env, ctx) {
5510 ` console.log("This is a log message", { key: "value" });
5511 ` console.info("This is an info message");
5512 ` console.warn("This is a warning message");
5513 ` console.error("This is an error message");
5514 ` console.debug("This is a debug message");
5515 ` console.debug({a: 1});
5516 `
5517 ` process.stdout.write("stdout");
5518 ` process.stdout.write("stdout with\nmultiple\nnewlines\nlog");
5519 ` process.stdout.write("ged");
5520 ` process.stderr.write("stderr");
5521 ` await 0;
5522 ` process.stderr.write("after await");
5523 `
5524 ` try {
5525 ` throw new Error("Test exception for structured logging");
5526 ` } catch (e) {
5527 ` console.error(e);
5528 ` }
5529 `
5530 ` return new Response("Structured logging test completed");
5531 ` }
5532 `}
5533 )
5534 ]
5535 )
5536 )
5537 ],
5538 sockets = [
5539 ( name = "main",
5540 address = "test-addr",
5541 service = "hello"
5542 )
5543 ],
5544 # Enable structured logging for this test
5545 structuredLogging = true
5546 ))"_kj,
5547 Worker::ConsoleMode::STDOUT);
5548 auto interceptorPipe = makePipeFds();
5549 int originalStdout = dup(STDOUT_FILENO);
5550 int originalStderr = dup(STDERR_FILENO);
5551 KJ_SYSCALL(dup2(interceptorPipe.input.get(), STDOUT_FILENO));
5552 KJ_SYSCALL(dup2(interceptorPipe.input.get(), STDERR_FILENO));
5553 interceptorPipe.input = nullptr;
5554 KJ_DEFER({
5555 // Restore stdout/stderr
5556 KJ_SYSCALL(dup2(originalStdout, STDOUT_FILENO));
5557 close(originalStdout);
5558 KJ_SYSCALL(dup2(originalStderr, STDERR_FILENO));
5559 close(originalStderr);
5560 });
5561 
5562 test.server.allowExperimental();
5563 test.start();
5564 auto conn = test.connect("test-addr");
5565 
5566 conn.sendHttpGet("/");
5567 conn.recvHttp200("Structured logging test completed");
5568 
5569 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5570 KJ_ASSERT(logline.contains(R"({"timestamp")"), logline);
5571 KJ_ASSERT(logline.contains(R"("level":"log")"), logline);
5572 KJ_ASSERT(logline.contains(R"("message":"This is a log message { key: 'value' }")"), logline);
5573 });
5574 
5575 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5576 KJ_ASSERT(logline.contains(R"("level":"info")"), logline);
5577 KJ_ASSERT(logline.contains(R"("message":"This is an info message")"), logline);
5578 });
5579 
5580 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5581 KJ_ASSERT(logline.contains(R"("level":"warn")"), logline);
5582 KJ_ASSERT(logline.contains(R"("message":"This is a warning message")"), logline);
5583 });
5584 
5585 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5586 KJ_ASSERT(logline.contains(R"("level":"error")"), logline);
5587 KJ_ASSERT(logline.contains(R"("message":"This is an error message")"), logline);
5588 });
5589 
5590 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5591 KJ_ASSERT(logline.contains(R"("level":"debug")"), logline);
5592 KJ_ASSERT(logline.contains(R"("message":"This is a debug message")"), logline);
5593 });
5594 
5595 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5596 KJ_ASSERT(logline.contains(R"("level":"debug")"), logline);
5597 KJ_ASSERT(logline.contains(R"("message":"{ a: 1 }")"), logline);
5598 });
5599 
5600 // process.stdout should be logs split by newline
5601 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5602 KJ_ASSERT(logline.contains(R"("level":"log")"), logline);
5603 KJ_ASSERT(logline.contains(R"("message":"stdout: stdoutstdout with")"), logline);
5604 });
5605 
5606 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5607 KJ_ASSERT(logline.contains(R"("level":"log")"), logline);
5608 KJ_ASSERT(logline.contains(R"("message":"stdout: multiple")"), logline);
5609 });
5610 
5611 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5612 KJ_ASSERT(logline.contains(R"("level":"log")"), logline);
5613 KJ_ASSERT(logline.contains(R"("message":"stdout: newlines")"), logline);
5614 });
5615 
5616 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5617 KJ_ASSERT(logline.contains(R"("level":"log")"), logline);
5618 KJ_ASSERT(logline.contains(R"("message":"stdout: logged")"), logline);
5619 });
5620 
5621 // process.stderr should be info
5622 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5623 KJ_ASSERT(logline.contains(R"("level":"log")"), logline);
5624 KJ_ASSERT(logline.contains(R"("message":"stderr: stderr")"), logline);
5625 });
5626 
5627 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5628 KJ_ASSERT(logline.contains(R"("level":"error")"), logline);
5629 KJ_ASSERT(
5630 logline.contains(
5631 R"_("message":"Error: Test exception for structured logging\n at Object.fetch (main.js:18:13)")_"),
5632 logline);
5633 });
5634 
5635 expectLogLine(interceptorPipe.output.get(), [](kj::StringPtr logline) {
5636 KJ_ASSERT(logline.contains(R"("level":"log")"), logline);
5637 KJ_ASSERT(logline.contains(R"("message":"stderr: after await")"), logline);
5638 });
5639}
5640 
5641KJ_TEST("Server: transpiled typescript") {
5642 TestServer test(singleWorker(R"((
5643 compatibilityDate = "2025-08-01",
5644 compatibilityFlags = ["typescript_strip_types"],
5645 modules = [
5646 ( name = "main.ts",
5647 esModule =
5648 `export default {
5649 ` async fetch(request): Promise<Response> {
5650 ` return new Response("Hello from typescript");
5651 ` }
5652 `} satisfies ExportedHandler<Env>;
5653 )
5654 ]
5655 ))"_kj));
5656 test.server.allowExperimental();
5657 test.start();
5658 auto conn = test.connect("test-addr");
5659 conn.httpGet200("/", "Hello from typescript");
5660}
5661 
5662KJ_TEST("Server: transpiled typescript failure") {
5663 TestServer test(singleWorker(R"((
5664 compatibilityDate = "2025-08-01",
5665 compatibilityFlags = ["typescript_strip_types"],
5666 modules = [
5667 ( name = "main.ts",
5668 esModule =
5669 `enum Foo { A, B }
5670 `export default {
5671 ` async fetch(request): Promise<Response> {
5672 ` return new Response("Hello from typescript");
5673 ` }
5674 `} satisfies ExportedHandler<Env>;
5675 )
5676 ]
5677 ))"_kj));
5678 test.server.allowExperimental();
5679 
5680 test.expectErrors(R"(service hello: Error transpiling main.ts : Unsupported syntax
5681 TypeScript enum is not supported in strip-only mode
5682service hello: Uncaught TypeError: Main module must be an ES module.
5683)");
5684}
5685 
5686#endif // __linux__
5687 
5688// Helper types for V8 serialization in tests
5689class SerializationContextGlobalObject: public jsg::Object, public jsg::ContextGlobal {};
5690struct SerializationTestContext: public SerializationContextGlobalObject {
5691 JSG_RESOURCE_TYPE(SerializationTestContext) {}
5692};
5693JSG_DECLARE_ISOLATE_TYPE(SerializationTestIsolate, SerializationTestContext);
5694 
5695// Helper function to serialize JavaScript values using V8
5696kj::Array<kj::byte> serializeJsArguments(
5697 std::initializer_list<std::function<jsg::JsValue(jsg::Lock&)>> argBuilders) {
5698 // Create an evaluator to get access to a V8 isolate
5699 jsg::test::Evaluator<SerializationTestContext, SerializationTestIsolate> evaluator(v8System);
5700 
5701 kj::Array<kj::byte> result;
5702 evaluator.run([&](auto& lock) {
5703 jsg::Lock& js = lock;
5704 
5705 // Create an array with the arguments
5706 auto argsArray = js.arr();
5707 for (auto& builder: argBuilders) {
5708 argsArray.add(js, builder(js));
5709 }
5710 
5711 // Serialize the array using jsg::Serializer
5712 jsg::Serializer serializer(js,
5713 jsg::Serializer::Options{
5714 .version = 15,
5715 .omitHeader = false,
5716 });
5717 serializer.write(js, jsg::JsValue(argsArray));
5718 result = serializer.release().data;
5719 });
5720 
5721 return result;
5722}
5723 
5724// Helper function to deserialize V8 data and convert to JSON string
5725kj::String deserializeV8ToJson(kj::ArrayPtr<const kj::byte> data) {
5726 jsg::test::Evaluator<SerializationTestContext, SerializationTestIsolate> evaluator(v8System);
5727 
5728 kj::String result;
5729 evaluator.run([&](auto& lock) {
5730 jsg::Lock& js = lock;
5731 
5732 // Deserialize the V8 data
5733 jsg::Deserializer deserializer(js, data, kj::none, kj::none, jsg::Deserializer::Options{});
5734 auto value = deserializer.readValue(js);
5735 
5736 // Convert to JSON string
5737 result = js.serializeJson(value);
5738 });
5739 
5740 return result;
5741}
5742 
5743KJ_TEST("Server: debug port RPC calls") {
5744 // This test connects to the debug port via Cap'n Proto RPC and makes actual RPC calls.
5745 TestServer test(R"((
5746 services = [
5747 ( name = "hello",
5748 worker = (
5749 compatibilityDate = "2024-01-01",
5750 modules = [
5751 ( name = "worker.js",
5752 esModule =
5753 `export default {
5754 ` async fetch(request) {
5755 ` return new Response("Hello from hello service");
5756 ` }
5757 `}
5758 )
5759 ]
5760 )
5761 ),
5762 ( name = "world",
5763 worker = (
5764 compatibilityDate = "2024-01-01",
5765 modules = [
5766 ( name = "worker.js",
5767 esModule =
5768 `export default {
5769 ` async fetch(request) {
5770 ` return new Response("Hello from world service");
5771 ` }
5772 `}
5773 )
5774 ]
5775 )
5776 ),
5777 ( name = "named-entrypoint",
5778 worker = (
5779 compatibilityDate = "2024-01-01",
5780 modules = [
5781 ( name = "worker.js",
5782 esModule =
5783 `export let customHandler = {
5784 ` async fetch(request) {
5785 ` return new Response("Hello from custom entrypoint");
5786 ` }
5787 `}
5788 `export default {
5789 ` async fetch(request) {
5790 ` return new Response("Default handler");
5791 ` }
5792 `}
5793 )
5794 ]
5795 )
5796 ),
5797 ( name = "props-service",
5798 worker = (
5799 compatibilityDate = "2024-01-01",
5800 modules = [
5801 ( name = "worker.js",
5802 esModule =
5803 `export default {
5804 ` async fetch(request, env, ctx) {
5805 ` const greeting = ctx?.props?.greeting || "no greeting";
5806 ` const name = ctx?.props?.name || "no name";
5807 ` return new Response("Props: " + greeting + " " + name);
5808 ` }
5809 `}
5810 )
5811 ]
5812 )
5813 ),
5814 ( name = "actor-service",
5815 worker = (
5816 compatibilityDate = "2024-01-01",
5817 modules = [
5818 ( name = "worker.js",
5819 esModule =
5820 `export class MyActor {
5821 ` constructor(state, env) {
5822 ` this.state = state;
5823 ` }
5824 ` async fetch(request) {
5825 ` const url = new URL(request.url);
5826 ` if (url.pathname === "/increment") {
5827 ` let count = (await this.state.storage.get("count")) || 0;
5828 ` count++;
5829 ` await this.state.storage.put("count", count);
5830 ` return new Response("Count: " + count);
5831 ` }
5832 ` return new Response("Actor: " + this.state.id.toString());
5833 ` }
5834 `}
5835 )
5836 ],
5837 durableObjectNamespaces = [
5838 ( className = "MyActor", uniqueKey = "test-actor" )
5839 ],
5840 durableObjectStorage = ( inMemory = void )
5841 )
5842 ),
5843 ( name = "rpc-service",
5844 worker = (
5845 compatibilityDate = "2024-09-02",
5846 compatibilityFlags = ["experimental"],
5847 modules = [
5848 ( name = "worker.js",
5849 esModule =
5850 `import {WorkerEntrypoint} from "cloudflare:workers";
5851 `export default class extends WorkerEntrypoint {
5852 ` async add(a, b) {
5853 ` return a + b;
5854 ` }
5855 ` async multiply(x, y) {
5856 ` return x * y;
5857 ` }
5858 ` async greet(name) {
5859 ` return "Hello, " + name + "!";
5860 ` }
5861 `}
5862 )
5863 ]
5864 )
5865 )
5866 ],
5867 sockets = [
5868 ( name = "main", address = "test-addr", service = "hello" )
5869 ]
5870 ))"_kj);
5871 
5872 // Enable the debug port on a unique address
5873 test.server.enableDebugPort(kj::str("debug-addr"));
5874 
5875 // Allow experimental features for RPC service
5876 test.server.allowExperimental();
5877 
5878 test.start();
5879 
5880 // Connect to the debug port
5881 auto debugConn = test.connect("debug-addr");
5882 
5883 // Create a TwoPartyClient for Cap'n Proto RPC
5884 capnp::TwoPartyClient client(debugConn.getStream());
5885 
5886 // Get the debug port capability
5887 auto debugPort = client.bootstrap().castAs<rpc::WorkerdDebugPort>();
5888 
5889 // Set up HTTP-over-Cap'n-Proto factory to convert Cap'n Proto HttpService to KJ HttpService
5890 capnp::ByteStreamFactory byteStreamFactory;
5891 kj::HttpHeaderTable::Builder headerTableBuilder;
5892 capnp::HttpOverCapnpFactory httpOverCapnpFactory(
5893 byteStreamFactory, headerTableBuilder, capnp::HttpOverCapnpFactory::LEVEL_2);
5894 auto headerTable = headerTableBuilder.build();
5895 
5896 // Helper to get bootstrap from service and entrypoint
5897 auto getBootstrap = [&](kj::StringPtr service, kj::Maybe<kj::StringPtr> entrypoint,
5898 auto&& propsBuilder) {
5899 auto req = debugPort.getEntrypointRequest();
5900 req.setService(service);
5901 KJ_IF_SOME(e, entrypoint) {
5902 req.setEntrypoint(e);
5903 }
5904 auto props = req.initProps();
5905 propsBuilder(props);
5906 auto resp = req.send().wait(test.ws);
5907 return resp.getEntrypoint();
5908 };
5909 
5910 // Helper to get a dispatcher from a bootstrap client
5911 auto getDispatcherFromBootstrap = [&](rpc::WorkerdBootstrap::Client bootstrap) {
5912 auto eventResp = bootstrap.startEventRequest().send().wait(test.ws);
5913 return eventResp.getDispatcher();
5914 };
5915 
5916 // Helper to get dispatcher from service and entrypoint (composes the two above)
5917 auto getDispatcher = [&](kj::StringPtr service, kj::Maybe<kj::StringPtr> entrypoint,
5918 auto&& propsBuilder) {
5919 return getDispatcherFromBootstrap(getBootstrap(service, entrypoint, propsBuilder));
5920 };
5921 
5922 // Helper to make HTTP request from a dispatcher
5923 auto makeHttpRequestFromDispatcher = [&](rpc::EventDispatcher::Client dispatcher,
5924 kj::StringPtr path) -> kj::String {
5925 auto capnpHttpService = dispatcher.getHttpServiceRequest().send().wait(test.ws).getHttp();
5926 
5927 // Convert to KJ HttpService and make request
5928 auto kjHttpService = httpOverCapnpFactory.capnpToKj(kj::mv(capnpHttpService));
5929 auto httpClient = kj::newHttpClient(*kjHttpService);
5930 auto url = kj::str("http://test", path);
5931 auto httpResponse = httpClient->request(kj::HttpMethod::GET, url, kj::HttpHeaders(*headerTable))
5932 .response.wait(test.ws);
5933 
5934 KJ_EXPECT(httpResponse.statusCode == 200);
5935 return httpResponse.body->readAllText().wait(test.ws);
5936 };
5937 
5938 // Helper to make HTTP request from a bootstrap client (works for both entrypoints and actors)
5939 auto makeHttpRequestFromBootstrap = [&](rpc::WorkerdBootstrap::Client bootstrap,
5940 kj::StringPtr path) -> kj::String {
5941 return makeHttpRequestFromDispatcher(getDispatcherFromBootstrap(kj::mv(bootstrap)), path);
5942 };
5943 
5944 // Helper to make HTTP request through an entrypoint with custom props
5945 auto makeHttpRequestImpl = [&](kj::StringPtr service, kj::Maybe<kj::StringPtr> entrypoint,
5946 auto&& propsBuilder) {
5947 return makeHttpRequestFromDispatcher(getDispatcher(service, entrypoint, propsBuilder), "/");
5948 };
5949 
5950 // Convenience wrapper with default empty props
5951 auto makeHttpRequest = [&](kj::StringPtr service, kj::Maybe<kj::StringPtr> entrypoint) {
5952 return makeHttpRequestImpl(service, entrypoint, [](auto& props) { props.setEmptyObject(); });
5953 };
5954 
5955 // Test 1: Request a non-existent service should fail
5956 KJ_EXPECT_THROW_MESSAGE("Worker \"nonexistent\" not found",
5957 getBootstrap("nonexistent", kj::none, [](auto& props) { props.setEmptyObject(); }));
5958 
5959 // Test 2: Get entrypoint for different services
5960 KJ_EXPECT(makeHttpRequest("hello", kj::none) == "Hello from hello service");
5961 KJ_EXPECT(makeHttpRequest("world", kj::none) == "Hello from world service");
5962 
5963 // Test 3: Named entrypoint works
5964 KJ_EXPECT(
5965 makeHttpRequest("named-entrypoint", "customHandler"_kjc) == "Hello from custom entrypoint");
5966 
5967 // Test 4: Passing props object works
5968 KJ_EXPECT(makeHttpRequestImpl("props-service", kj::none, [](auto& props) {
5969 props.setEmptyObject();
5970 auto properties = props.initProperties(2);
5971 properties[0].setName("greeting");
5972 properties[0].setJson("\"Hello\"");
5973 properties[1].setName("name");
5974 properties[1].setJson("\"World\"");
5975 }) == "Props: Hello World");
5976 
5977 // Test 5: Getting an actor works and we can call methods on it
5978 {
5979 // Create a deterministic actor ID
5980 kj::byte actorIdBytes[32] = {};
5981 for (size_t i = 0; i < sizeof(actorIdBytes); i++) {
5982 actorIdBytes[i] = static_cast<kj::byte>(i);
5983 }
5984 
5985 // Helper to make an HTTP request to the actor
5986 auto makeActorRequest = [&](kj::StringPtr path) -> kj::String {
5987 auto req = debugPort.getActorRequest();
5988 req.setService("actor-service");
5989 req.setEntrypoint("MyActor");
5990 // Convert actor ID bytes to hex string
5991 req.setActorId(kj::encodeHex(kj::arrayPtr(actorIdBytes, sizeof(actorIdBytes))));
5992 auto resp = req.send().wait(test.ws);
5993 return makeHttpRequestFromBootstrap(resp.getActor(), path);
5994 };
5995 
5996 // Make a first request to increment the counter
5997 {
5998 auto bodyText = makeActorRequest("/increment");
5999 KJ_EXPECT(bodyText == "Count: 1");
6000 }
6001 
6002 // Make a second request to increment again - verifies state persistence
6003 {
6004 auto bodyText = makeActorRequest("/increment");
6005 KJ_EXPECT(bodyText == "Count: 2");
6006 }
6007 
6008 // Make a request to verify the actor ID is correct
6009 {
6010 auto bodyText = makeActorRequest("/");
6011 
6012 // The actor should return its ID as a hex string
6013 // Convert our actor ID bytes to hex string to compare
6014 kj::String expectedId = kj::encodeHex(kj::arrayPtr(actorIdBytes, sizeof(actorIdBytes)));
6015 kj::String expectedResponse = kj::str("Actor: ", expectedId);
6016 KJ_EXPECT(bodyText == expectedResponse, bodyText, expectedResponse);
6017 }
6018 }
6019 
6020 // Test 6: Call RPC methods using jsRpcSession with V8-serialized arguments
6021 {
6022 // Get dispatcher and JS RPC session - use pipelining because jsRpcSession() doesn't return until session closes
6023 auto dispatcher =
6024 getDispatcher("rpc-service", kj::none, [](auto& props) { props.setEmptyObject(); });
6025 auto rpcSessionReq = dispatcher.jsRpcSessionRequest();
6026 auto sessionPromise = rpcSessionReq.send();
6027 auto rpcTarget = sessionPromise.getTopLevel();
6028 
6029 // Test calling add(5, 3) -> 8
6030 auto v8SerializedArgs = serializeJsArguments({[](jsg::Lock& js) {
6031 return jsg::JsValue(js.num(5));
6032 }, [](jsg::Lock& js) { return jsg::JsValue(js.num(3)); }});
6033 
6034 auto callReq = rpcTarget.callRequest();
6035 callReq.setMethodName("add");
6036 auto operation = callReq.initOperation();
6037 auto jsValue = operation.initCallWithArgs();
6038 jsValue.setV8Serialized(v8SerializedArgs);
6039 
6040 auto callResp = callReq.send().wait(test.ws);
6041 auto result = callResp.getResult();
6042 
6043 auto resultData = result.getV8Serialized();
6044 KJ_EXPECT(resultData.size() > 0, "Result should be non-empty");
6045 
6046 auto jsonResult = deserializeV8ToJson(resultData);
6047 KJ_EXPECT(jsonResult == "8", jsonResult, "Expected result to be 8");
6048 }
6049}
6050 
6051KJ_TEST("Server: workerdDebugPort binding loopback test") {
6052 // This test verifies that a worker can use the workerdDebugPort binding to connect
6053 // back to the same workerd instance's debug port and access other services.
6054 TestServer test(R"((
6055 services = [
6056 ( name = "target-service",
6057 worker = (
6058 compatibilityDate = "2024-01-01",
6059 modules = [
6060 ( name = "worker.js",
6061 esModule =
6062 `export default {
6063 ` async fetch(request) {
6064 ` return new Response("Hello from target!");
6065 ` }
6066 `}
6067 `export let namedHandler = {
6068 ` async fetch(request) {
6069 ` return new Response("Hello from named entrypoint!");
6070 ` }
6071 `}
6072 )
6073 ]
6074 )
6075 ),
6076 ( name = "test-service",
6077 worker = (
6078 compatibilityDate = "2024-01-01",
6079 compatibilityFlags = ["experimental"],
6080 modules = [
6081 ( name = "worker.js",
6082 esModule =
6083 `export default {
6084 ` async fetch(request, env, ctx) {
6085 ` // Connect to the debug port
6086 ` const client = await env.debugPort.connect("debug-addr");
6087 `
6088 ` // Test 1: Access the default entrypoint
6089 ` const defaultFetcher = client.getEntrypoint("target-service");
6090 ` const defaultResp = await defaultFetcher.fetch("http://fake-host/");
6091 ` const defaultText = await defaultResp.text();
6092 ` if (defaultText !== "Hello from target!") {
6093 ` throw new Error("Expected 'Hello from target!' but got: " + defaultText);
6094 ` }
6095 `
6096 ` // Test 2: Access a named entrypoint
6097 ` const namedFetcher = client.getEntrypoint("target-service", "namedHandler");
6098 ` const namedResp = await namedFetcher.fetch("http://fake-host/");
6099 ` const namedText = await namedResp.text();
6100 ` if (namedText !== "Hello from named entrypoint!") {
6101 ` throw new Error("Expected 'Hello from named entrypoint!' but got: " + namedText);
6102 ` }
6103 `
6104 ` return new Response("All tests passed!");
6105 ` }
6106 `}
6107 )
6108 ],
6109 bindings = [
6110 ( name = "debugPort",
6111 workerdDebugPort = void
6112 )
6113 ]
6114 )
6115 )
6116 ],
6117 sockets = [
6118 ( name = "main", address = "test-addr", service = "test-service" )
6119 ]
6120 ))"_kj);
6121 
6122 // Enable the debug port on a known address
6123 test.server.enableDebugPort(kj::str("debug-addr"));
6124 test.server.allowExperimental();
6125 
6126 test.start();
6127 
6128 // Run the test by invoking the fetch handler
6129 auto conn = test.connect("test-addr");
6130 conn.httpGet200("/", "All tests passed!");
6131}
6132 
6133KJ_TEST("Server: workerdDebugPort binding with props") {
6134 // This test verifies that props can be passed through the workerdDebugPort binding.
6135 TestServer test(R"((
6136 services = [
6137 ( name = "target-service",
6138 worker = (
6139 compatibilityDate = "2024-01-01",
6140 compatibilityFlags = ["experimental"],
6141 modules = [
6142 ( name = "worker.js",
6143 esModule =
6144 `import {WorkerEntrypoint} from "cloudflare:workers";
6145 `export class PropsHandler extends WorkerEntrypoint {
6146 ` async fetch(request) {
6147 ` const props = this.ctx.props;
6148 ` return new Response(JSON.stringify(props));
6149 ` }
6150 `}
6151 )
6152 ]
6153 )
6154 ),
6155 ( name = "test-service",
6156 worker = (
6157 compatibilityDate = "2024-01-01",
6158 compatibilityFlags = ["experimental"],
6159 modules = [
6160 ( name = "worker.js",
6161 esModule =
6162 `export default {
6163 ` async fetch(request, env, ctx) {
6164 ` // Connect to the debug port
6165 ` const client = await env.debugPort.connect("debug-addr");
6166 `
6167 ` // Test passing props to the entrypoint
6168 ` const fetcher = client.getEntrypoint(
6169 ` "target-service", "PropsHandler", {foo: "bar", num: 42});
6170 ` const resp = await fetcher.fetch("http://fake-host/");
6171 ` const props = await resp.json();
6172 `
6173 ` if (props.foo !== "bar") {
6174 ` throw new Error("Expected props.foo to be 'bar' but got: " + props.foo);
6175 ` }
6176 ` if (props.num !== 42) {
6177 ` throw new Error("Expected props.num to be 42 but got: " + props.num);
6178 ` }
6179 `
6180 ` return new Response("Props test passed!");
6181 ` }
6182 `}
6183 )
6184 ],
6185 bindings = [
6186 ( name = "debugPort",
6187 workerdDebugPort = void
6188 )
6189 ]
6190 )
6191 )
6192 ],
6193 sockets = [
6194 ( name = "main", address = "test-addr", service = "test-service" )
6195 ]
6196 ))"_kj);
6197 
6198 test.server.enableDebugPort(kj::str("debug-addr"));
6199 test.server.allowExperimental();
6200 
6201 test.start();
6202 
6203 auto conn = test.connect("test-addr");
6204 conn.httpGet200("/", "Props test passed!");
6205}
6206 
6207KJ_TEST("Server: workerdDebugPort binding getActor") {
6208 // This test verifies that getActor can be used to access Durable Objects via the debug port.
6209 TestServer test(R"((
6210 services = [
6211 ( name = "do-service",
6212 worker = (
6213 compatibilityDate = "2024-01-01",
6214 compatibilityFlags = ["experimental"],
6215 modules = [
6216 ( name = "worker.js",
6217 esModule =
6218 `import {DurableObject} from "cloudflare:workers";
6219 `export default {
6220 ` async fetch(request) {
6221 ` return new Response("DO service default handler");
6222 ` }
6223 `}
6224 `export class Counter extends DurableObject {
6225 ` counter = 0;
6226 ` async fetch(request) {
6227 ` this.counter++;
6228 ` return new Response("Counter: " + this.counter);
6229 ` }
6230 `}
6231 )
6232 ],
6233 durableObjectNamespaces = [
6234 ( className = "Counter",
6235 uniqueKey = "test-do-key"
6236 )
6237 ],
6238 durableObjectStorage = (inMemory = void)
6239 )
6240 ),
6241 ( name = "test-service",
6242 worker = (
6243 compatibilityDate = "2024-01-01",
6244 compatibilityFlags = ["experimental"],
6245 modules = [
6246 ( name = "worker.js",
6247 esModule =
6248 `export default {
6249 ` async fetch(request, env, ctx) {
6250 ` // Connect to the debug port
6251 ` const client = await env.debugPort.connect("debug-addr");
6252 `
6253 ` // Get the same actor twice using a fixed ID
6254 ` const actorId = "0".repeat(64);
6255 `
6256 ` const actor1 = client.getActor("do-service", "Counter", actorId);
6257 ` const resp1 = await actor1.fetch("http://fake-host/");
6258 ` const text1 = await resp1.text();
6259 ` if (text1 !== "Counter: 1") {
6260 ` throw new Error("Expected 'Counter: 1' but got: " + text1);
6261 ` }
6262 `
6263 ` // Second request to same actor should increment counter
6264 ` const actor2 = client.getActor("do-service", "Counter", actorId);
6265 ` const resp2 = await actor2.fetch("http://fake-host/");
6266 ` const text2 = await resp2.text();
6267 ` if (text2 !== "Counter: 2") {
6268 ` throw new Error("Expected 'Counter: 2' but got: " + text2);
6269 ` }
6270 `
6271 ` // Different actor ID should have independent state (counter starts at 1)
6272 ` const differentActorId = "1".repeat(64);
6273 ` const actor3 = client.getActor("do-service", "Counter", differentActorId);
6274 ` const resp3 = await actor3.fetch("http://fake-host/");
6275 ` const text3 = await resp3.text();
6276 ` if (text3 !== "Counter: 1") {
6277 ` throw new Error("Expected 'Counter: 1' for different actor but got: " + text3);
6278 ` }
6279 `
6280 ` return new Response("DO actor test passed!");
6281 ` }
6282 `}
6283 )
6284 ],
6285 bindings = [
6286 ( name = "debugPort",
6287 workerdDebugPort = void
6288 )
6289 ]
6290 )
6291 )
6292 ],
6293 sockets = [
6294 ( name = "main", address = "test-addr", service = "test-service" )
6295 ]
6296 ))"_kj);
6297 
6298 test.server.enableDebugPort(kj::str("debug-addr"));
6299 test.server.allowExperimental();
6300 
6301 test.start();
6302 
6303 auto conn = test.connect("test-addr");
6304 conn.httpGet200("/", "DO actor test passed!");
6305}
6306 
6307KJ_TEST("Server: workerdDebugPort WebSocket passthrough via WorkerEntrypoint") {
6308 // This test verifies that a WebSocket obtained via the debug port can be passed through
6309 // a service binding response (from a WorkerEntrypoint). This was previously broken because
6310 // the debug port connection was destroyed when the intermediate IoContext finished.
6311 TestServer test(R"((
6312 services = [
6313 ( name = "target-service",
6314 worker = (
6315 compatibilityDate = "2024-01-01",
6316 modules = [
6317 ( name = "worker.js",
6318 esModule =
6319 `export default {
6320 ` async fetch(request) {
6321 ` // Accept WebSocket upgrade and echo messages with a prefix
6322 ` const upgradeHeader = request.headers.get("Upgrade");
6323 ` if (upgradeHeader === "websocket") {
6324 ` const pair = new WebSocketPair();
6325 ` pair[1].accept();
6326 ` pair[1].addEventListener("message", (e) => {
6327 ` pair[1].send("echo:" + e.data);
6328 ` });
6329 ` return new Response(null, { status: 101, webSocket: pair[0] });
6330 ` }
6331 ` return new Response("Not a WebSocket request");
6332 ` }
6333 `}
6334 )
6335 ]
6336 )
6337 ),
6338 ( name = "proxy-service",
6339 worker = (
6340 compatibilityDate = "2024-01-01",
6341 compatibilityFlags = ["experimental"],
6342 modules = [
6343 ( name = "worker.js",
6344 esModule =
6345 `import {WorkerEntrypoint} from "cloudflare:workers";
6346 `
6347 `// This WorkerEntrypoint gets a WebSocket via debug port and passes it through
6348 `export class Proxy extends WorkerEntrypoint {
6349 ` async fetch(request) {
6350 ` const client = await this.env.debugPort.connect("debug-addr");
6351 ` const fetcher = client.getEntrypoint("target-service");
6352 ` const response = await fetcher.fetch(request);
6353 ` if (response.webSocket) {
6354 ` // Pass through the WebSocket from the debug port
6355 ` return new Response(null, { status: 101, webSocket: response.webSocket });
6356 ` }
6357 ` return response;
6358 ` }
6359 `}
6360 `
6361 `export default {
6362 ` async fetch(request, env) {
6363 ` // Route through the Proxy entrypoint to test WebSocket passthrough
6364 ` return env.proxy.fetch(request);
6365 ` }
6366 `}
6367 )
6368 ],
6369 bindings = [
6370 ( name = "debugPort", workerdDebugPort = void ),
6371 ( name = "proxy", service = (name = "proxy-service", entrypoint = "Proxy") )
6372 ]
6373 )
6374 )
6375 ],
6376 sockets = [
6377 ( name = "main", address = "test-addr", service = "proxy-service" )
6378 ]
6379 ))"_kj);
6380 
6381 test.server.enableDebugPort(kj::str("debug-addr"));
6382 test.server.allowExperimental();
6383 
6384 test.start();
6385 
6386 // Connect and upgrade to WebSocket
6387 auto wsConn = test.connect("test-addr");
6388 wsConn.upgradeToWebSocket();
6389 
6390 // Send a message and verify we get the echoed response
6391 // WebSocket frame: 0x81 = final frame + text, 0x05 = payload length 5
6392 constexpr kj::StringPtr testMessage = "hello"_kj;
6393 wsConn.send(kj::str("\x81\x05", testMessage));
6394 wsConn.recvWebSocket("echo:hello");
6395 
6396 // Send another message to verify the connection stays alive
6397 constexpr kj::StringPtr testMessage2 = "world"_kj;
6398 wsConn.send(kj::str("\x81\x05", testMessage2));
6399 wsConn.recvWebSocket("echo:world");
6400}
6401} // namespace
6402} // namespace workerd::server