Skip to content
File

Blob: src/workerd/api/tests/js-rpc-test.js

javascript2120 lines
1// Copyright (c) 2024 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
4import assert from 'node:assert';
5import { waitUntil } from 'cloudflare:workers';
6import {
7 WorkerEntrypoint,
8 DurableObject,
9 RpcPromise,
10 RpcProperty,
11 RpcStub,
12 RpcTarget,
13 ServiceStub,
14} from 'cloudflare:workers';
15 
16try {
17 waitUntil(null);
18 throw new Error('This should have thrown');
19} catch (error) {
20 assert.match(
21 error.message,
22 /Disallowed operation called within global scope./
23 );
24}
25 
26class MyCounter extends RpcTarget {
27 constructor(i = 0) {
28 super();
29 this.i = i;
30 }
31 
32 async increment(j) {
33 this.i += j;
34 return this.i;
35 }
36 
37 disposed = false;
38 
39 #disposedResolver;
40 #onDisposedPromise = new Promise(
41 (resolve) => (this.#disposedResolver = resolve)
42 );
43 
44 [Symbol.dispose]() {
45 this.disposed = true;
46 this.#disposedResolver();
47 }
48 
49 onDisposed() {
50 return Promise.race([
51 this.#onDisposedPromise,
52 scheduler.wait(1000).then(() => {
53 throw new Error('timed out waiting for disposal');
54 }),
55 ]);
56 }
57 
58 // Tests that `fetch()` is not special for RpcTargets.
59 async fetch(a, b, c) {
60 return `${this.i} ${a} ${b} ${c}`;
61 }
62}
63 
64class RpcBox extends RpcTarget {
65 #value;
66 
67 constructor(value) {
68 super();
69 this.#value = value;
70 }
71 
72 get value() {
73 return this.#value;
74 }
75}
76 
77class NonRpcClass {
78 foo() {
79 return 123;
80 }
81 get bar() {
82 return {
83 baz() {
84 return 456;
85 },
86 };
87 }
88 
89 i = 0;
90 increment(i) {
91 this.i += i;
92 }
93}
94 
95export let nonClass = {
96 async noArgs(x, env, ctx) {
97 assert.strictEqual(typeof ctx.waitUntil, 'function');
98 return x === undefined ? env.twelve + 1 : 'param not undefined?';
99 },
100 
101 async oneArg(i, env, ctx) {
102 assert.strictEqual(typeof ctx.waitUntil, 'function');
103 return env.twelve * i;
104 },
105 
106 async oneArgOmitCtx(i, env) {
107 return env.twelve * i + 1;
108 },
109 
110 async oneArgOmitEnvCtx(i) {
111 return 2 * i;
112 },
113 
114 async twoArgs(i, j, env, ctx) {
115 assert.strictEqual(typeof ctx.waitUntil, 'function');
116 return i * j + env.twelve;
117 },
118 
119 async fetch(req, env, ctx) {
120 // This is used in the stream test to fetch some gziped data.
121 if (req.url.endsWith('/gzip')) {
122 return new Response('this text was gzipped', {
123 headers: {
124 'Content-Encoding': 'gzip',
125 },
126 });
127 } else if (req.url.endsWith('/stream-from-rpc')) {
128 let stream = await env.MyService.returnReadableStream();
129 return new Response(stream);
130 } else {
131 throw new Error('unknown route');
132 }
133 },
134};
135 
136// Globals used to test passing RPC promises or properties across I/O contexts (which is expected
137// to fail).
138let globalRpcPromise;
139 
140// Promise initialized by testWaitUntil() and then resolved shortly later, in a waitUntil task.
141let globalWaitUntilPromise;
142 
143export class MyService extends WorkerEntrypoint {
144 constructor(ctx, env) {
145 super(ctx, env);
146 
147 assert.strictEqual(this.ctx, ctx);
148 assert.strictEqual(this.env, env);
149 
150 // This shouldn't be callable!
151 this.instanceMethod = () => {
152 return 'nope';
153 };
154 
155 this.instanceObject = {
156 func: (a) => a * 5,
157 };
158 }
159 
160 async noArgsMethod(x) {
161 return x === undefined ? this.env.twelve + 1 : 'param not undefined?';
162 }
163 
164 async oneArgMethod(i) {
165 return this.env.twelve * i;
166 }
167 
168 async twoArgsMethod(i, j) {
169 return i * j + this.env.twelve;
170 }
171 
172 async makeCounter(i) {
173 return new MyCounter(i);
174 }
175 
176 async incrementCounter(counter, i) {
177 return await counter.increment(i);
178 }
179 
180 async getAnObject(i) {
181 return { foo: 123 + i, counter: new MyCounter(i) };
182 }
183 
184 async getADeeperObject(i) {
185 return { foo: 123 + i, box: new RpcBox(new RpcStub(new MyCounter(i))) };
186 }
187 
188 async getMap() {
189 let map = new Map();
190 map.set('foo', 123);
191 map.set('bar', 456);
192 return map;
193 }
194 
195 async fetch(req, x) {
196 assert.strictEqual(x, undefined);
197 return new Response('method = ' + req.method + ', url = ' + req.url);
198 }
199 
200 async connect(socket) {
201 const enc = new TextEncoder();
202 let writer = socket.writable.getWriter();
203 await writer.write(enc.encode('hello'));
204 await writer.close();
205 }
206 
207 // Define a property to test behavior of property accessors.
208 get nonFunctionProperty() {
209 return { foo: 123 };
210 }
211 
212 get functionProperty() {
213 return (a, b) => a - b;
214 }
215 
216 get objectProperty() {
217 let nullPrototype = { foo: 123 };
218 nullPrototype.__proto__ = null;
219 
220 return {
221 func: (a, b) => a * b,
222 deeper: {
223 deepFunc: (a, b) => a / b,
224 },
225 counter5: new MyCounter(5),
226 nonRpc: new NonRpcClass(),
227 nullPrototype,
228 someText: 'hello',
229 };
230 }
231 
232 get promiseProperty() {
233 return scheduler.wait(10).then(() => 123);
234 }
235 
236 get rejectingPromiseProperty() {
237 return Promise.reject(new Error('REJECTED'));
238 }
239 
240 get throwingProperty() {
241 throw new Error('PROPERTY THREW');
242 }
243 
244 throwingMethod() {
245 const err = new Error('METHOD THREW');
246 err.abc = 123;
247 throw err;
248 }
249 
250 async neverReturn() {
251 await new Promise((resolve) => {});
252 }
253 
254 async tryUseGlobalRpcPromise() {
255 return await globalRpcPromise;
256 }
257 async tryUseGlobalRpcPromisePipeline() {
258 return await globalRpcPromise.increment(1);
259 }
260 
261 async getNonRpcClass() {
262 return { obj: new NonRpcClass() };
263 }
264 
265 async getNullPrototypeObject() {
266 let obj = { foo: 123 };
267 obj.__proto__ = null;
268 return obj;
269 }
270 
271 async getFunction() {
272 let func = (a, b) => a ^ b;
273 func.someProperty = 123;
274 return func;
275 }
276 
277 getRpcPromise(callback) {
278 return callback();
279 }
280 getNestedRpcPromise(callback) {
281 return { value: callback() };
282 }
283 getRemoteNestedRpcPromise(callback) {
284 // Use a function as a cheap way to return a JsRpcStub that has a remote property `value` which
285 // itself is initialized as a JsRpcPromise.
286 let result = () => {};
287 result.value = callback();
288 return result;
289 }
290 getRpcProperty(callback) {
291 return callback.foo;
292 }
293 getNestedRpcProperty(callback) {
294 return { value: callback.foo };
295 }
296 getRemoteNestedRpcProperty(callback) {
297 // Use a function as a cheap way to return a JsRpcStub that has a remote property `value` which
298 // itself is initialized as a JsRpcProperty.
299 let result = () => {};
300 
301 // We have to dup() the callback because otherwise it will be implicitly disposed when this
302 // function returns, but we want it to remain accessible as a property of the returned
303 // function.
304 let dupCallback = callback.dup();
305 result.value = dupCallback.foo;
306 
307 // Not really
308 result[Symbol.dispose] = () => dupCallback[Symbol.dispose]();
309 
310 return result;
311 }
312 
313 get(a) {
314 return a + 1;
315 }
316 put(a, b) {
317 return a + b;
318 }
319 delete(a) {
320 return a - 1;
321 }
322 
323 async testDispose(counter) {
324 // Prove that the counter works at this point.
325 let count = await counter.increment(5);
326 
327 let counterDup = counter.dup();
328 
329 return {
330 count,
331 counter: counter.dup(), // need to dup() for return
332 async incrementOriginal(n) {
333 // This will fail because after the call ends, the counter stub is disposed.
334 return await counter.increment(n);
335 },
336 async incrementDup(n) {
337 // This will succeed since we kept a dup.
338 return await counterDup.increment(n);
339 },
340 [Symbol.dispose]() {
341 counterDup[Symbol.dispose]();
342 },
343 };
344 }
345 
346 async leak(stub) {
347 // Leak the input stub
348 stub.dup();
349 
350 let ctx = this.ctx;
351 
352 // Return something that contains stubs, holding the context open.
353 return {
354 noop() {},
355 abort() {
356 ctx.abort(new RangeError('foo bar abort reason'));
357 },
358 };
359 }
360 
361 async leakButReturnPlainObject(stub) {
362 // Leak the input stub, so it will be disposed when the context is torn down.
363 stub.dup();
364 
365 // Return a plain object (one with no stubs and no disposer). This should NOT
366 // hold the context open, so the stub should be dropped promptly.
367 return { foo: 123 };
368 }
369 
370 async writeToStream(stream) {
371 let writer = stream.getWriter();
372 let enc = new TextEncoder();
373 await writer.write(enc.encode('foo, '));
374 await writer.write(enc.encode('bar, '));
375 await writer.write(enc.encode('baz!'));
376 await writer.close();
377 }
378 
379 async writeToStreamExpectingError(stream) {
380 // We expect the writes to fail on the remote end. However, that won't happen
381 // right away. Due to backpressure and flow control, write does not wait until
382 // the remote end has received the data. It resolves as soon as there is space
383 // in the buffer to perform another write. Here we'll just keep writing until
384 // we get the expected error.
385 let writer = stream.getWriter();
386 let enc = new TextEncoder();
387 for (;;) {
388 await writer.write(enc.encode('foo, '));
389 }
390 }
391 
392 async writeToStreamAbort(stream) {
393 // In this test, aborting the stream should propagate back to the remote
394 // side, causing the stream to be errored and the abort algorithm to be
395 // called with the provided error. Unfortunately the current implementation
396 // does not allow for that. The reason passed to abort is cached locally and
397 // is never communicated to the remote. Instead, the remote side will end up
398 // with a generic disconnect error. Sad face.
399 let writer = stream.getWriter();
400 writer.abort(new Error('boom'));
401 }
402 
403 async readFromStream(stream) {
404 return await new Response(stream).text();
405 }
406 
407 async returnReadableStream() {
408 let { readable, writable } = new IdentityTransformStream();
409 this.ctx.waitUntil(this.writeToStream(writable));
410 return readable;
411 }
412 
413 async returnMultipleReadableStreams() {
414 let { readable, writable } = new IdentityTransformStream();
415 this.ctx.waitUntil(this.writeToStream(writable));
416 
417 let pair2 = new IdentityTransformStream();
418 this.ctx.waitUntil(this.writeToStream(pair2.writable));
419 let readable2 = pair2.readable;
420 
421 return [readable, readable2];
422 }
423 
424 async roundTrip(value) {
425 return value;
426 }
427 
428 async returnEmptyHeaders() {
429 return new Headers();
430 }
431 
432 async returnHeaders() {
433 let result = new Headers();
434 result.append('foo', 'bar');
435 result.append('Set-Cookie', 'abc');
436 result.append('set-cookie', 'def');
437 result.append('corge', '!@#');
438 result.append('Content-Length', '123');
439 return result;
440 }
441 
442 async returnRequest() {
443 return new Request('http://my-url.com', {
444 method: 'PUT',
445 headers: {
446 'Accept-Encoding': 'bazzip',
447 Foo: 'Bar',
448 },
449 redirect: 'manual',
450 body: 'Hello every body!',
451 cf: {
452 abc: 123,
453 hello: 'goodbye',
454 },
455 });
456 }
457 
458 async returnResponse() {
459 return new Response('Response body!', {
460 status: 404,
461 headers: {
462 'Content-Type': 'abc',
463 },
464 cf: { foo: 123, bar: 'def' },
465 });
466 }
467 
468 testWaitUntil() {
469 // Initialize globalWaitUntilPromise to a promise that will be resolved in a waitUntil task
470 // later on. We'll perform a cross-context wait to verify that the waitUntil task actually
471 // completes and resolves the promise.
472 let resolve;
473 globalWaitUntilPromise = new Promise((r) => {
474 resolve = r;
475 });
476 
477 this.ctx.waitUntil(scheduler.wait(100).then(resolve));
478 }
479 
480 testImportedWaitUntil() {
481 // Initialize globalWaitUntilPromise to a promise that will be resolved in a waitUntil task
482 // later on. We'll perform a cross-context wait to verify that the waitUntil task actually
483 // completes and resolves the promise.
484 let resolve;
485 globalWaitUntilPromise = new Promise((r) => {
486 resolve = r;
487 });
488 
489 waitUntil(scheduler.wait(100).then(resolve));
490 }
491 
492 async call(func, arg) {
493 return await func(arg);
494 }
495 
496 async getProp(obj, prop) {
497 return await obj[prop];
498 }
499 
500 // Useful to test pipelining.
501 async identity(x) {
502 return x;
503 }
504}
505 
506// An entrypoint which forwards methods calls to MyService, thus acting as a proxy.
507export class MyServiceProxy extends WorkerEntrypoint {
508 makeCounter(i) {
509 return this.env.MyService.makeCounter(i);
510 }
511 
512 getAnObject(i) {
513 return this.env.MyService.getAnObject(i);
514 }
515 
516 getADeeperObject(i) {
517 return this.env.MyService.getADeeperObject(i);
518 }
519}
520 
521class PostAbortCallTester extends RpcTarget {
522 constructor(ctx) {
523 super();
524 this.ctx = ctx;
525 }
526 
527 ping() {
528 return 'pong';
529 }
530 
531 hang() {
532 return new Promise((resolve) => {});
533 }
534 
535 abort() {
536 this.ctx.abort('test aborted by abort()');
537 }
538 
539 async failCriticalSection() {
540 await this.ctx.blockConcurrencyWhile(() => {
541 throw new Error('test broken critical section');
542 });
543 }
544}
545 
546export class MyActor extends DurableObject {
547 #counter = 0;
548 
549 constructor(ctx, env) {
550 super(ctx, env);
551 
552 assert.strictEqual(this.ctx, ctx);
553 assert.strictEqual(this.env, env);
554 }
555 
556 async increment(amount) {
557 this.#counter += amount;
558 return this.#counter;
559 }
560 
561 async doCallbackBlockingConcurrency() {
562 // Check that we can receive RPC callbacks during blockConcurrencyWhile(), if they are from
563 // an RPC running inside the block. This verifies that the critical section is captured
564 // correctly in IoContext::makeReentryCallback().
565 return this.ctx.blockConcurrencyWhile(async () => {
566 let func = () => {
567 return 12345;
568 };
569 return await this.env.MyService.getRpcPromise(func);
570 });
571 }
572 
573 makePostAbortCallTester() {
574 return new PostAbortCallTester(this.ctx);
575 }
576}
577 
578export class ActorNoExtends {
579 async fetch(req) {
580 return new Response('from ActorNoExtends');
581 }
582 
583 // This can't be called!
584 async foo() {
585 return 123;
586 }
587}
588 
589export default class DefaultService extends WorkerEntrypoint {
590 async fetch(req) {
591 // Test this.env here just to prove omitting the constructor entirely works.
592 return new Response('default service ' + this.env.twelve);
593 }
594}
595 
596export let basicServiceBinding = {
597 async test(controller, env, ctx) {
598 // Test service binding RPC.
599 assert.strictEqual(await env.self.oneArg(3), 36);
600 assert.strictEqual(await env.self.oneArgOmitCtx(3), 37);
601 assert.strictEqual(await env.self.oneArgOmitEnvCtx(3), 6);
602 await assert.rejects(() => env.self.twoArgs(123, 2), {
603 name: 'TypeError',
604 message:
605 'Cannot call handler function "twoArgs" over RPC because it has the wrong ' +
606 'number of arguments. A simple function handler can only be called over RPC if it has ' +
607 'exactly the arguments (arg, env, ctx), where only the first argument comes from the ' +
608 'client. To support multi-argument RPC functions, use class-based syntax (extending ' +
609 'WorkerEntrypoint) instead.',
610 });
611 await assert.rejects(() => env.self.noArgs(), {
612 name: 'TypeError',
613 message:
614 'Attempted to call RPC function "noArgs" with the wrong number of arguments. ' +
615 'When calling a top-level handler function that is not declared as part of a class, you ' +
616 'must always send exactly one argument. In order to support variable numbers of ' +
617 'arguments, the server must use class-based syntax (extending WorkerEntrypoint) ' +
618 'instead.',
619 });
620 await assert.rejects(() => env.self.oneArg(1, 2), {
621 name: 'TypeError',
622 message:
623 'Attempted to call RPC function "oneArg" with the wrong number of arguments. ' +
624 'When calling a top-level handler function that is not declared as part of a class, you ' +
625 'must always send exactly one argument. In order to support variable numbers of ' +
626 'arguments, the server must use class-based syntax (extending WorkerEntrypoint) ' +
627 'instead.',
628 });
629 
630 // If we restore multi-arg support, remove the `rejects` checks above and un-comment these:
631 // assert.strictEqual(await env.self.noArgs(), 13);
632 // assert.strictEqual(await env.self.twoArgs(123, 2), 258);
633 // assert.strictEqual(await env.self.twoArgs(123, 2, "foo", "bar", "baz"), 258);
634 },
635};
636 
637export let extendingEntrypointClasses = {
638 async test(controller, env, ctx) {
639 // Verify that we can instantiate classes that inherit built-in classes.
640 let svc = new MyService(ctx, env);
641 assert.equal(svc instanceof WorkerEntrypoint, true);
642 },
643};
644export let connectBinding = {
645 async test(controller, env, ctx) {
646 let socket = await env.MyService.connect('localhost:8081');
647 await socket.opened;
648 const dec = new TextDecoder();
649 let result = '';
650 for await (const chunk of socket.readable) {
651 result += dec.decode(chunk, { stream: true });
652 }
653 result += dec.decode();
654 assert.strictEqual(result, 'hello');
655 await socket.closed;
656 },
657};
658 
659export let namedServiceBinding = {
660 async test(controller, env, ctx) {
661 assert.strictEqual(await env.MyService.noArgsMethod(), 13);
662 assert.strictEqual(await env.MyService.oneArgMethod(3), 36);
663 assert.strictEqual(await env.MyService.twoArgsMethod(123, 2), 258);
664 
665 // Properties that return a function can actually be called.
666 assert.strictEqual(await env.MyService.functionProperty(19, 6), 13);
667 
668 // Members of an object-typed property can be invoked.
669 assert.strictEqual(await env.MyService.objectProperty.func(6, 4), 24);
670 assert.strictEqual(
671 await env.MyService.objectProperty.deeper.deepFunc(6, 3),
672 2
673 );
674 assert.strictEqual(
675 await env.MyService.objectProperty.counter5.increment(3),
676 8
677 );
678 
679 // Awaiting a property itself gets the value.
680 assert.strictEqual(
681 JSON.stringify(await env.MyService.nonFunctionProperty),
682 '{"foo":123}'
683 );
684 assert.strictEqual(await env.MyService.objectProperty.someText, 'hello');
685 {
686 let counter = await env.MyService.objectProperty.counter5;
687 assert.strictEqual(await counter.increment(3), 8);
688 assert.strictEqual(await counter.increment(7), 15);
689 }
690 
691 {
692 let func = await env.MyService.objectProperty.func;
693 assert.strictEqual(await func(3, 7), 21);
694 }
695 {
696 let func = await env.MyService.getFunction();
697 assert.strictEqual(await func(3, 6), 5);
698 assert.strictEqual(await func.someProperty, 123);
699 }
700 {
701 // Pipeline the function call.
702 let func = env.MyService.getFunction();
703 assert.strictEqual(await func(3, 6), 5);
704 assert.strictEqual(await func.someProperty, 123);
705 }
706 
707 // A property that returns a Promise will wait for the Promise.
708 assert.strictEqual(await env.MyService.promiseProperty, 123);
709 
710 let sawFinally = false;
711 assert.strictEqual(
712 await env.MyService.promiseProperty.finally(() => {
713 sawFinally = true;
714 }),
715 123
716 );
717 assert.strictEqual(sawFinally, true);
718 
719 // `fetch()` is special, the params get converted into a Request.
720 let resp = await env.MyService.fetch('http://foo/', { method: 'POST' });
721 assert.strictEqual(await resp.text(), 'method = POST, url = http://foo/');
722 
723 await assert.rejects(() => env.MyService.instanceMethod(), {
724 name: 'TypeError',
725 message:
726 'The RPC receiver does not implement the method "instanceMethod".',
727 });
728 
729 await assert.rejects(() => env.MyService.instanceObject.func(3), {
730 name: 'TypeError',
731 message:
732 'The RPC receiver does not implement the method "instanceObject".',
733 });
734 
735 await assert.rejects(() => env.MyService.instanceObject, {
736 name: 'TypeError',
737 message:
738 'The RPC receiver does not implement the method "instanceObject".',
739 });
740 
741 await assert.rejects(() => env.MyService.throwingProperty, {
742 name: 'Error',
743 message: 'PROPERTY THREW',
744 });
745 await assert.rejects(() => env.MyService.throwingMethod(), {
746 name: 'Error',
747 message: 'METHOD THREW',
748 });
749 
750 await assert.rejects(() => env.MyService.rejectingPromiseProperty, {
751 name: 'Error',
752 message: 'REJECTED',
753 });
754 assert.strictEqual(
755 await env.MyService.rejectingPromiseProperty.catch((err) => {
756 assert.strictEqual(err.message, 'REJECTED');
757 return 234;
758 }),
759 234
760 );
761 assert.strictEqual(
762 await env.MyService.rejectingPromiseProperty.then(
763 () => 432,
764 (err) => {
765 assert.strictEqual(err.message, 'REJECTED');
766 return 234;
767 }
768 ),
769 234
770 );
771 
772 let getByName = (name) => {
773 return env.MyService.getRpcMethodForTestOnly(name);
774 };
775 
776 // Check getRpcMethodForTestOnly() actually works.
777 assert.strictEqual(await getByName('twoArgsMethod')(2, 3), 18);
778 
779 // Check we cannot call reserved methods.
780 await assert.rejects(() => getByName('constructor')(), {
781 name: 'TypeError',
782 message:
783 "'constructor' is a reserved method and cannot be called over RPC.",
784 });
785 await assert.rejects(() => getByName('fetch')(), {
786 name: 'TypeError',
787 message: "'fetch' is a reserved method and cannot be called over RPC.",
788 });
789 
790 // Check we cannot call methods of Object.
791 await assert.rejects(() => getByName('toString')(), {
792 name: 'TypeError',
793 message: 'The RPC receiver does not implement the method "toString".',
794 });
795 await assert.rejects(() => getByName('hasOwnProperty')(), {
796 name: 'TypeError',
797 message:
798 'The RPC receiver does not implement the method "hasOwnProperty".',
799 });
800 
801 // Check we cannot access `env` or `ctx`.
802 await assert.rejects(() => getByName('env')(), {
803 name: 'TypeError',
804 message: 'The RPC receiver does not implement the method "env".',
805 });
806 await assert.rejects(() => getByName('ctx')(), {
807 name: 'TypeError',
808 message: 'The RPC receiver does not implement the method "ctx".',
809 });
810 
811 // Check what happens if we access something that's actually declared as a property on the
812 // class. The difference in error message proves to us that `env` and `ctx` weren't visible at
813 // all, which is what we want.
814 await assert.rejects(() => getByName('nonFunctionProperty')(), {
815 name: 'TypeError',
816 message: '"nonFunctionProperty" is not a function.',
817 });
818 await assert.rejects(() => getByName('nonFunctionProperty').foo(), {
819 name: 'TypeError',
820 message: '"nonFunctionProperty.foo" is not a function.',
821 });
822 
823 // Check that we can't access contents of a property that is a class but not derived from
824 // RpcTarget.
825 await assert.rejects(() => env.MyService.objectProperty.nonRpc.foo(), {
826 name: 'TypeError',
827 message: 'The RPC receiver does not implement the method "nonRpc".',
828 });
829 await assert.rejects(() => env.MyService.objectProperty.nonRpc.bar.baz(), {
830 name: 'TypeError',
831 message: 'The RPC receiver does not implement the method "nonRpc".',
832 });
833 await assert.rejects(() => env.MyService.objectProperty.nullPrototype.foo, {
834 name: 'TypeError',
835 message:
836 'The RPC receiver does not implement the method "nullPrototype".',
837 });
838 
839 // Extra-paranoid check that we can't access methods on env or ctx.
840 await assert.rejects(
841 () => env.MyService.objectProperty.env.MyService.noArgsMethod(),
842 {
843 name: 'TypeError',
844 message: 'The RPC receiver does not implement the method "env".',
845 }
846 );
847 await assert.rejects(() => env.MyService.objectProperty.ctx.waitUntil(), {
848 name: 'TypeError',
849 message: 'The RPC receiver does not implement the method "ctx".',
850 });
851 
852 // Can't serialize instances of classes that aren't derived from RpcTarget.
853 await assert.rejects(() => env.MyService.getNonRpcClass(), {
854 name: 'DataCloneError',
855 message:
856 'Could not serialize object of type "NonRpcClass". This type does not support ' +
857 'serialization.',
858 });
859 await assert.rejects(() => env.MyService.getNullPrototypeObject(), {
860 name: 'DataCloneError',
861 message:
862 'Could not serialize object of type "Object". This type does not support ' +
863 'serialization.',
864 });
865 
866 // A stateless entryponit method that never returns should fail due to PendingEvent tracking.
867 await assert.rejects(() => env.MyService.neverReturn(), {
868 name: 'Error',
869 message:
870 "The Workers runtime canceled this request because it detected that your Worker's code " +
871 'had hung and would never generate a response. Refer to: ' +
872 'https://developers.cloudflare.com/workers/observability/errors/',
873 });
874 
875 {
876 let map = await env.MyService.getMap();
877 assert.strictEqual(map.get('foo'), 123);
878 assert.strictEqual(map.get('bar'), 456);
879 assert.strictEqual(map.get('baz'), undefined);
880 }
881 },
882};
883 
884export let namedActorBinding = {
885 async test(controller, env, ctx) {
886 let id = env.MyActor.idFromName('foo');
887 let stub = env.MyActor.get(id);
888 
889 assert.strictEqual(await stub.increment(5), 5);
890 assert.strictEqual(await stub.increment(2), 7);
891 assert.strictEqual(await stub.increment(8), 15);
892 
893 assert.strictEqual(await stub.doCallbackBlockingConcurrency(), 12345);
894 },
895};
896 
897// Test that if the actor class doesn't extend `DurableObject`, we don't allow RPC.
898export let actorWithoutExtendsRejectsRpc = {
899 async test(controller, env, ctx) {
900 let id = env.ActorNoExtends.idFromName('foo');
901 let stub = env.ActorNoExtends.get(id);
902 
903 // fetch() works.
904 let resp = await stub.fetch('http://foo');
905 assert.strictEqual(await resp.text(), 'from ActorNoExtends');
906 
907 // RPC does not.
908 await assert.rejects(() => stub.foo(), {
909 name: 'TypeError',
910 message:
911 'The receiving Durable Object does not support RPC, because its class was not declared ' +
912 'with `extends DurableObject`. In order to enable RPC, make sure your class ' +
913 'extends the special class `DurableObject`, which can be imported from the module ' +
914 '"cloudflare:workers".',
915 });
916 },
917};
918 
919// Test calling the default export when it is a class.
920export let defaultExportClass = {
921 async test(controller, env, ctx) {
922 let resp = await env.defaultExport.fetch('http://foo');
923 assert.strictEqual(await resp.text(), 'default service 12');
924 },
925};
926 
927export let loopbackJsRpcTarget = {
928 async test(controller, env, ctx) {
929 {
930 let counter = new MyCounter(4);
931 let stub = new RpcStub(counter);
932 assert.strictEqual(await stub.increment(5), 9);
933 assert.strictEqual(await stub.increment(7), 16);
934 
935 assert.strictEqual(await stub.fetch(true, 123, 'baz'), '16 true 123 baz');
936 
937 assert.strictEqual(counter.disposed, false);
938 stub[Symbol.dispose]();
939 
940 await assert.rejects(stub.increment(2), {
941 name: 'Error',
942 message: 'RPC stub used after being disposed.',
943 });
944 
945 await counter.onDisposed();
946 assert.strictEqual(counter.disposed, true);
947 
948 assert.strictEqual(stub instanceof RpcStub, true);
949 assert.strictEqual(stub.increment instanceof RpcProperty, true);
950 assert.strictEqual(stub.increment(1) instanceof RpcPromise, true);
951 assert.strictEqual(env.MyService instanceof ServiceStub, true);
952 }
953 
954 // In fact, RpcStubs can be created from any old object.
955 {
956 let stub = new RpcStub({
957 sum(a, b) {
958 return a + b;
959 },
960 });
961 
962 assert.strictEqual(await stub.sum(12, 34), 46);
963 }
964 
965 // Or function.
966 {
967 let func = (a, b) => {
968 return a + b;
969 };
970 func.ownProperty = 'hello';
971 let stub = new RpcStub(func);
972 
973 assert.strictEqual(await stub(12, 34), 46);
974 assert.strictEqual(await stub.ownProperty, 'hello');
975 }
976 
977 // Or Proxy of an RpcTarget.
978 {
979 let counter = new MyCounter(4);
980 let stub = new RpcStub(new Proxy(counter, {}));
981 assert.strictEqual(await stub.increment(5), 9);
982 }
983 },
984};
985 
986export let sendStubOverRpc = {
987 async test(controller, env, ctx) {
988 let stub = new RpcStub(new MyCounter(4));
989 let stubDup = stub.dup();
990 
991 assert.strictEqual(await env.MyService.incrementCounter(stub, 5), 9);
992 
993 await assert.rejects(() => stub.increment(7), {
994 name: 'Error',
995 message: 'RPC stub used after being disposed.',
996 });
997 
998 assert.strictEqual(await stubDup.increment(7), 16);
999 },
1000};
1001 
1002export let receiveStubOverRpc = {
1003 async test(controller, env, ctx) {
1004 let stub = await env.MyService.makeCounter(17);
1005 assert.strictEqual(await stub.increment(2), 19);
1006 assert.strictEqual(await stub.increment(-10), 9);
1007 
1008 // Do multiple concurrent calls, they should be delivered in the order in which they were made.
1009 let promise1 = stub.increment(6);
1010 let promise2 = stub.increment(4);
1011 let promise3 = stub.increment(3);
1012 assert.deepEqual(
1013 await Promise.all([promise1, promise2, promise3]),
1014 [15, 19, 22]
1015 );
1016 },
1017};
1018 
1019export let promisePipelining = {
1020 async test(controller, env, ctx) {
1021 assert.strictEqual(await env.MyService.makeCounter(12).increment(3), 15);
1022 
1023 assert.strictEqual(await env.MyService.getAnObject(5).foo, 128);
1024 assert.strictEqual(
1025 await env.MyService.getAnObject(5).counter.increment(7),
1026 12
1027 );
1028 
1029 assert.rejects(() => env.MyService.oneArgMethod(5).foo(), {
1030 name: 'TypeError',
1031 message: 'The RPC receiver does not implement the method "foo".',
1032 });
1033 
1034 assert.rejects(() => env.MyService.getMap().foo(), {
1035 name: 'TypeError',
1036 message: 'The RPC receiver does not implement the method "foo".',
1037 });
1038 },
1039};
1040 
1041// Test promise pipelining through a proxy.
1042export let promisePipeliningProxy = {
1043 async test(controller, env, ctx) {
1044 // Pipeline on a proxied call that just returns a stub.
1045 {
1046 let counter = env.MyServiceProxy.makeCounter(12);
1047 let promise1 = counter.increment(3);
1048 let promise2 = counter.increment(5);
1049 assert.strictEqual(await promise1, 15);
1050 assert.strictEqual(await promise2, 20);
1051 }
1052 
1053 // Pipeline on a proxied call that returns an object containing a stub.
1054 {
1055 let counter = env.MyServiceProxy.getAnObject(12).counter;
1056 let promise1 = counter.increment(3);
1057 let promise2 = counter.increment(5);
1058 assert.strictEqual(await promise1, 15);
1059 assert.strictEqual(await promise2, 20);
1060 }
1061 
1062 // Pipeline on a proxied call that returns an object containing an object that contains a
1063 // stub. (This ensures that pipelining can traverse JsRpcProperty values.)
1064 {
1065 let counter = env.MyServiceProxy.getADeeperObject(12).box.value;
1066 let promise1 = counter.increment(3);
1067 let promise2 = counter.increment(5);
1068 assert.strictEqual(await promise1, 15);
1069 assert.strictEqual(await promise2, 20);
1070 }
1071 },
1072};
1073 
1074export let disposal = {
1075 async test(controller, env, ctx) {
1076 // Call function that returns plain stub. Dispose it.
1077 {
1078 let counter = await env.MyService.makeCounter(12);
1079 assert.strictEqual(await counter.increment(3), 15);
1080 counter[Symbol.dispose]();
1081 await assert.rejects(counter.increment(2), {
1082 name: 'Error',
1083 message: 'RPC stub used after being disposed.',
1084 });
1085 }
1086 
1087 // Call function that returns an object containing a stub. The object has a dispose() method
1088 // added that disposes everything inside.
1089 {
1090 let obj = await env.MyService.getAnObject(5);
1091 assert.strictEqual(await obj.counter.increment(7), 12);
1092 
1093 obj[Symbol.dispose]();
1094 await assert.rejects(obj.counter.increment(2), {
1095 name: 'Error',
1096 message: 'RPC stub used after being disposed.',
1097 });
1098 }
1099 
1100 // Verify a capability passed as an RPC param receives the disposal callback.
1101 {
1102 let counter = new MyCounter(3);
1103 assert.strictEqual(await env.MyService.incrementCounter(counter, 5), 8);
1104 await counter.onDisposed();
1105 assert.strictEqual(counter.disposed, true);
1106 }
1107 
1108 // A more complex case with testDispose().
1109 {
1110 let counter = new MyCounter(3);
1111 let obj = await env.MyService.testDispose(counter);
1112 
1113 // The counter was able to be incremented during the call.
1114 assert.strictEqual(obj.count, 8);
1115 
1116 // But after the call, the stub is disposed.
1117 await assert.rejects(obj.incrementOriginal(), {
1118 name: 'Error',
1119 message: 'RPC stub used after being disposed.',
1120 });
1121 
1122 // The duplicate we created in the call still works though.
1123 assert.strictEqual(await obj.incrementDup(7), 15);
1124 
1125 // Also, the returned copy works.
1126 assert.strictEqual(await obj.counter.increment(4), 19);
1127 
1128 // But of course, disposing the return value overall breaks everything.
1129 obj[Symbol.dispose]();
1130 await assert.rejects(obj.counter.increment(2), {
1131 name: 'Error',
1132 message: 'RPC stub used after being disposed.',
1133 });
1134 await assert.rejects(obj.incrementDup(7), {
1135 name: 'Error',
1136 message: 'RPC stub used after being disposed.',
1137 });
1138 
1139 await counter.onDisposed();
1140 assert.strictEqual(counter.disposed, true);
1141 }
1142 
1143 // Test a leak situation.
1144 {
1145 let counter = new MyCounter(3);
1146 let obj = await env.MyService.leak(counter);
1147 
1148 // Give a chance for disposal to happen.
1149 await obj.noop();
1150 await scheduler.wait(10);
1151 
1152 // It should not have happened.
1153 assert.strictEqual(counter.disposed, false);
1154 
1155 // Even if we GC the leaked stub, still no disposal!
1156 gc();
1157 await scheduler.wait(10);
1158 assert.strictEqual(counter.disposed, false);
1159 
1160 // If we abort the server's I/O context, though, then the counter is disposed.
1161 await assert.rejects(obj.abort(), {
1162 name: 'RangeError',
1163 message: 'foo bar abort reason',
1164 });
1165 
1166 await counter.onDisposed();
1167 assert.strictEqual(counter.disposed, true);
1168 }
1169 
1170 // Test that a call which returns a plain object does not need to be disposed.
1171 // Historically, the callee context would not be torn down promptly.
1172 {
1173 let counter = new MyCounter(3);
1174 await env.MyService.leakButReturnPlainObject(counter);
1175 
1176 // Give a chance for disposal to happen. When running with a real socket, this involves
1177 // non-deterministic scheduling, and it can be relatively slower when running under ASAN,
1178 // QEMU, etc. so we'll try to dynamically adjust how much time we wait, up to 10s.
1179 for (let i = 0; i < 1000; i++) {
1180 if (counter.disposed) break;
1181 await scheduler.wait(10);
1182 }
1183 
1184 // It should have happened! A call that returns a plain object should NOT
1185 // require disposal to clean up its context!
1186 assert.strictEqual(counter.disposed, true);
1187 }
1188 },
1189};
1190 
1191export let crossContextSharingDoesntWork = {
1192 async test(controller, env, ctx) {
1193 // Test what happens if a JsRpcPromise or JsRpcProperty is shared cross-context. This is not
1194 // intended to work in general, but there are specific cases where it does work, and we should
1195 // avoid breaking those with future changes.
1196 
1197 // Sharing an RPC promise between contexts works as long as the promise returns a simple value
1198 // (with no I/O objects), since JsRpcPromise wraps a simple JS promise and we support sharing
1199 // JS promises.
1200 globalRpcPromise = env.MyService.oneArgMethod(2);
1201 assert.strictEqual(await env.MyService.tryUseGlobalRpcPromise(), 24);
1202 
1203 // Sharing a property of a service binding works, because the service binding itself is not
1204 // tied to an I/O context. Awaiting the property actually initiates a new RPC session from
1205 // whatever context performed the await.
1206 globalRpcPromise = env.MyService.nonFunctionProperty;
1207 assert.strictEqual(
1208 JSON.stringify(await env.MyService.tryUseGlobalRpcPromise()),
1209 '{"foo":123}'
1210 );
1211 
1212 // OK, now let's look at cases that do NOT work. These all produce the same error.
1213 let _expectedError = {
1214 name: 'Error',
1215 message:
1216 'Cannot perform I/O on behalf of a different request. I/O objects (such as streams, ' +
1217 'request/response bodies, and others) created in the context of one request handler ' +
1218 "cannot be accessed from a different request's handler. This is a limitation of " +
1219 'Cloudflare Workers which allows us to improve overall performance.',
1220 };
1221 
1222 // A promise which resolves to a value that contains a stub. The stub cannot be used from a
1223 // different context.
1224 //
1225 // Note that the part that actually fails here is not awaiting the promise, but rather when
1226 // tryUseGlobalRpcPromise() tries to return the result, it tries to serialize the stub, but
1227 // it can't do that from the wrong context.
1228 globalRpcPromise = env.MyService.makeCounter(12);
1229 
1230 await assert.rejects(() => env.MyService.tryUseGlobalRpcPromise(), {
1231 name: 'Error',
1232 message:
1233 'Cannot perform I/O on behalf of a different request. I/O objects (such as streams, ' +
1234 'request/response bodies, and others) created in the context of one request handler ' +
1235 "cannot be accessed from a different request's handler. This is a limitation of " +
1236 'Cloudflare Workers which allows us to improve overall performance. (I/O type: Client)',
1237 });
1238 
1239 // Pipelining on someone else's promise straight-up doesn't work.
1240 await assert.rejects(() => env.MyService.tryUseGlobalRpcPromisePipeline(), {
1241 name: 'Error',
1242 message:
1243 'Cannot perform I/O on behalf of a different request. I/O objects (such as streams, ' +
1244 'request/response bodies, and others) created in the context of one request handler ' +
1245 "cannot be accessed from a different request's handler. This is a limitation of " +
1246 'Cloudflare Workers which allows us to improve overall performance. ' +
1247 '(I/O type: JsRpcPromise)',
1248 });
1249 
1250 // Now let's try accessing a JsRpcProperty, where the property is NOT a direct property of a
1251 // top-level service binding. This works even less than a JsRpcPromise, since there's no inner
1252 // JS promise, it tries to create one on-demand, which fails because the parent object is
1253 // tied to the original I/O context.
1254 globalRpcPromise = env.MyService.getAnObject(5).counter;
1255 
1256 await assert.rejects(() => env.MyService.tryUseGlobalRpcPromise(), {
1257 name: 'Error',
1258 message:
1259 'Cannot perform I/O on behalf of a different request. I/O objects (such as streams, ' +
1260 'request/response bodies, and others) created in the context of one request handler ' +
1261 "cannot be accessed from a different request's handler. This is a limitation of " +
1262 'Cloudflare Workers which allows us to improve overall performance. (I/O type: Pipeline)',
1263 });
1264 
1265 await assert.rejects(() => env.MyService.tryUseGlobalRpcPromisePipeline(), {
1266 name: 'Error',
1267 message:
1268 'Cannot perform I/O on behalf of a different request. I/O objects (such as streams, ' +
1269 'request/response bodies, and others) created in the context of one request handler ' +
1270 "cannot be accessed from a different request's handler. This is a limitation of " +
1271 'Cloudflare Workers which allows us to improve overall performance. ' +
1272 '(I/O type: JsRpcPromise)',
1273 });
1274 },
1275};
1276 
1277export let waitUntilWorks = {
1278 async test(controller, env, ctx) {
1279 // Tests ctx.waitUntil
1280 {
1281 globalWaitUntilPromise = null;
1282 await env.MyService.testWaitUntil();
1283 
1284 assert.ok(globalWaitUntilPromise instanceof Promise);
1285 await globalWaitUntilPromise;
1286 }
1287 
1288 // Tests `import { waitUntil } from 'cloudflare:workers` on WorkerEntrypoint
1289 {
1290 globalWaitUntilPromise = null;
1291 await env.MyService.testImportedWaitUntil();
1292 
1293 assert.ok(globalWaitUntilPromise instanceof Promise);
1294 await globalWaitUntilPromise;
1295 }
1296 },
1297};
1298 
1299export let serializeRpcPromiseOrProprety = {
1300 async test(controller, env, ctx) {
1301 // What happens if we actually try to serialize a JsRpcPromise or JsRpcProperty? Let's make
1302 // sure these aren't, for instance, treated as functions because they are callable.
1303 
1304 let func = () => {
1305 return { x: 123 };
1306 };
1307 func.foo = { x: 456 };
1308 
1309 // If we directly return returning a JsRpcPromise, the system automatically awaits it on the
1310 // server side because it's a thenable.
1311 assert.deepEqual(await env.MyService.getRpcPromise(func), {
1312 x: 123,
1313 });
1314 
1315 // Pipelining also works.
1316 assert.strictEqual(await env.MyService.getRpcPromise(func).x, 123);
1317 
1318 // If a JsRpcPromise appears somewhere in the serialization tree, it'll just fail serialization.
1319 // NOTE: We could choose to make this work later.
1320 await assert.rejects(() => env.MyService.getNestedRpcPromise(func), {
1321 name: 'DataCloneError',
1322 message:
1323 'Could not serialize object of type "RpcPromise". This type does not support ' +
1324 'serialization.',
1325 });
1326 await assert.rejects(() => env.MyService.getNestedRpcPromise(func).value, {
1327 name: 'DataCloneError',
1328 message:
1329 'Could not serialize object of type "RpcPromise". This type does not support ' +
1330 'serialization.',
1331 });
1332 await assert.rejects(
1333 () => env.MyService.getNestedRpcPromise(func).value.x,
1334 {
1335 name: 'DataCloneError',
1336 message:
1337 'Could not serialize object of type "RpcPromise". This type does not support ' +
1338 'serialization.',
1339 }
1340 );
1341 
1342 // Things get a little weird when we return a stub which itself has properties that reflect
1343 // our RPC promise. If we await fetch the JsRpcPromise itself, this works, again because
1344 // somewhere along the line V8 says "oh look a thenable" and awaits it, before it can be
1345 // subject to serialization. That's fine.
1346 assert.deepEqual(
1347 await env.MyService.getRemoteNestedRpcPromise(func).value,
1348 { x: 123 }
1349 );
1350 await assert.rejects(
1351 () => env.MyService.getRemoteNestedRpcPromise(func).value.x,
1352 {
1353 name: 'TypeError',
1354 message: 'The RPC receiver does not implement the method "value".',
1355 }
1356 );
1357 
1358 // The story is similar for a JsRpcProperty -- though the implementation details differ.
1359 assert.deepEqual(await env.MyService.getRpcProperty(func), {
1360 x: 456,
1361 });
1362 assert.strictEqual(await env.MyService.getRpcProperty(func).x, 456);
1363 await assert.rejects(() => env.MyService.getNestedRpcProperty(func), {
1364 name: 'DataCloneError',
1365 message:
1366 'Could not serialize object of type "RpcProperty". This type does not support ' +
1367 'serialization.',
1368 });
1369 await assert.rejects(() => env.MyService.getNestedRpcProperty(func).value, {
1370 name: 'DataCloneError',
1371 message:
1372 'Could not serialize object of type "RpcProperty". This type does not support ' +
1373 'serialization.',
1374 });
1375 await assert.rejects(
1376 () => env.MyService.getNestedRpcProperty(func).value.x,
1377 {
1378 name: 'DataCloneError',
1379 message:
1380 'Could not serialize object of type "RpcProperty". This type does not support ' +
1381 'serialization.',
1382 }
1383 );
1384 
1385 assert.deepEqual(
1386 await env.MyService.getRemoteNestedRpcProperty(func).value,
1387 { x: 456 }
1388 );
1389 await assert.rejects(
1390 () => env.MyService.getRemoteNestedRpcProperty(func).value.x,
1391 {
1392 name: 'TypeError',
1393 message: 'The RPC receiver does not implement the method "value".',
1394 }
1395 );
1396 await assert.rejects(
1397 () => env.MyService.getRemoteNestedRpcProperty(func).value(),
1398 {
1399 name: 'TypeError',
1400 message: '"foo" is not a function.',
1401 }
1402 );
1403 },
1404};
1405 
1406export let streams = {
1407 async test(controller, env, ctx) {
1408 // Send WritableStream.
1409 {
1410 let { readable, writable } = new IdentityTransformStream();
1411 let promise = env.MyService.writeToStream(writable);
1412 let text = await new Response(readable).text();
1413 assert.strictEqual(text, 'foo, bar, baz!');
1414 await promise;
1415 }
1416 
1417 {
1418 const dec = new TextDecoder();
1419 let result = '';
1420 const { promise, resolve } = Promise.withResolvers();
1421 const writable = new WritableStream({
1422 write(chunk) {
1423 result += dec.decode(chunk, { stream: true });
1424 },
1425 close() {
1426 result += dec.decode();
1427 resolve();
1428 },
1429 });
1430 const p1 = env.MyService.writeToStream(writable);
1431 await promise;
1432 assert.strictEqual(result, 'foo, bar, baz!');
1433 await p1;
1434 }
1435 
1436 {
1437 // In this test, the remote side writes a chunk to the stream below, which throws
1438 // an error. Ideally the error would propagate back to the calling side so that the
1439 // remote knows the stream failed and can no longer be written to. The call to
1440 // writeToStreamExpectingError should throw because the error should be propagated
1441 // through the round trip.
1442 let _result = '';
1443 let writeCalled = 0;
1444 const writable = new WritableStream({
1445 write(chunk) {
1446 writeCalled++;
1447 throw new Error('boom');
1448 },
1449 });
1450 
1451 try {
1452 await env.MyService.writeToStreamExpectingError(writable);
1453 throw new Error('should have thrown');
1454 } catch (err) {
1455 assert.strictEqual(err.message, 'boom');
1456 // The write method should have been called once.
1457 assert.strictEqual(writeCalled, 1);
1458 }
1459 }
1460 
1461 {
1462 // In this test, the remote side aborts the writable stream it receives.
1463 // The abort should propagate such that the abort algorithm is called and the
1464 // writeToStreamAbort call should succeed. The error passed on to abort(reason)
1465 // should be the error that was given on the remote side when abort is called,
1466 // but we currently do not propagate the abort reason through. What ends up
1467 // happening is that the local stream is dropped with a generic cancelation
1468 // error.
1469 const { promise, resolve } = Promise.withResolvers();
1470 const writable = new WritableStream({
1471 write(chunk) {},
1472 abort(reason) {
1473 resolve(reason);
1474 },
1475 });
1476 await env.MyService.writeToStreamAbort(writable);
1477 const reason = await promise;
1478 // TODO(someday): The reason should be the error that was passed to abort on the
1479 // remote side, but we currently do not propagate this. We end up with a generic
1480 // disconnection error instead, which certainly not ideal.
1481 assert.strictEqual(
1482 reason.message,
1483 'WritableStream received over RPC was disconnected because the remote execution ' +
1484 'context has endeded.'
1485 );
1486 }
1487 
1488 // TODO(someday): Is there any way to construct an encoded WritableStream? Only system
1489 // streams can be encoded, but there's no API that returns an encoded WritableStream I think.
1490 
1491 // Send ReadableStream.
1492 {
1493 let { readable, writable } = new IdentityTransformStream();
1494 let promise = env.MyService.readFromStream(readable);
1495 
1496 let writer = writable.getWriter();
1497 let enc = new TextEncoder();
1498 await writer.write(enc.encode('foo, '));
1499 await writer.write(enc.encode('bar, '));
1500 await writer.write(enc.encode('baz!'));
1501 await writer.close();
1502 
1503 assert.strictEqual(await promise, 'foo, bar, baz!');
1504 }
1505 
1506 // Send a JS-backed ReadableStream.
1507 {
1508 let controller;
1509 let readable = new ReadableStream({
1510 start(c) {
1511 controller = c;
1512 },
1513 });
1514 let promise = env.MyService.readFromStream(readable);
1515 
1516 let enc = new TextEncoder();
1517 controller.enqueue(enc.encode('foo, '));
1518 controller.enqueue(enc.encode('bar, '));
1519 controller.enqueue(enc.encode('baz!'));
1520 controller.close();
1521 
1522 assert.strictEqual(await promise, 'foo, bar, baz!');
1523 }
1524 
1525 // Send streams that are locked.
1526 {
1527 let { readable, writable } = new IdentityTransformStream();
1528 
1529 let writer = writable.getWriter();
1530 let enc = new TextEncoder();
1531 writer.write(enc.encode('foo'));
1532 
1533 let reader = readable.getReader();
1534 
1535 assert.rejects(env.MyService.writeToStream(writable), {
1536 name: 'TypeError',
1537 message: 'The WritableStream has been locked to a writer.',
1538 });
1539 assert.rejects(env.MyService.readFromStream(readable), {
1540 name: 'TypeError',
1541 message: 'The ReadableStream has been locked to a reader.',
1542 });
1543 
1544 // Verify the streams still work.
1545 let dec = new TextDecoder();
1546 assert.strictEqual(dec.decode((await reader.read()).value), 'foo');
1547 
1548 writer.write(enc.encode('bar'));
1549 assert.strictEqual(dec.decode((await reader.read()).value), 'bar');
1550 
1551 writer.close();
1552 assert.strictEqual((await reader.read()).done, true);
1553 }
1554 
1555 // Receive ReadableStream.
1556 {
1557 let readable = await env.MyService.returnReadableStream();
1558 let text = await new Response(readable).text();
1559 assert.strictEqual(text, 'foo, bar, baz!');
1560 }
1561 
1562 // Receive multiple ReadableStreams.
1563 {
1564 let readables = await env.MyService.returnMultipleReadableStreams();
1565 assert.strictEqual(
1566 await new Response(readables[0]).text(),
1567 'foo, bar, baz!'
1568 );
1569 assert.strictEqual(
1570 await new Response(readables[1]).text(),
1571 'foo, bar, baz!'
1572 );
1573 }
1574 
1575 // Send ReadableStream, but fail to fully write it.
1576 {
1577 let { readable, writable } = new IdentityTransformStream();
1578 let promise = env.MyService.readFromStream(readable);
1579 
1580 let writer = writable.getWriter();
1581 let enc = new TextEncoder();
1582 await writer.write(enc.encode('foo, '));
1583 await writer.write(enc.encode('bar, '));
1584 await writer.write(enc.encode('baz!'));
1585 await writer.abort('foo');
1586 
1587 await assert.rejects(promise, {
1588 name: 'Error',
1589 // TODO(someday): Propagate the actual error.
1590 message: 'ReadableStream received over RPC disconnected prematurely.',
1591 });
1592 }
1593 
1594 // Send fixed-length ReadableStream.
1595 {
1596 let { readable, writable } = new FixedLengthStream(
1597 'foo, bar, baz!'.length
1598 );
1599 let promise = env.MyService.readFromStream(readable);
1600 
1601 let writer = writable.getWriter();
1602 let enc = new TextEncoder();
1603 await writer.write(enc.encode('foo, '));
1604 await writer.write(enc.encode('bar, '));
1605 await writer.write(enc.encode('baz!'));
1606 await writer.close();
1607 
1608 assert.strictEqual(await promise, 'foo, bar, baz!');
1609 }
1610 
1611 // Send an encoded ReadableStream
1612 {
1613 let gzippedResp = await env.self.fetch('http://foo/gzip');
1614 
1615 let text = await env.MyService.readFromStream(gzippedResp.body);
1616 
1617 assert.strictEqual(text, 'this text was gzipped');
1618 }
1619 
1620 // Round trip streams.
1621 {
1622 let { readable, writable } = new IdentityTransformStream();
1623 
1624 readable = await env.MyService.roundTrip(readable);
1625 writable = await env.MyService.roundTrip(writable);
1626 
1627 let readPromise = new Response(readable).text();
1628 
1629 let writer = writable.getWriter();
1630 let enc = new TextEncoder();
1631 await writer.write(enc.encode('foo, '));
1632 await writer.write(enc.encode('bar, '));
1633 await writer.write(enc.encode('baz!'));
1634 await writer.close();
1635 
1636 assert.strictEqual(await readPromise, 'foo, bar, baz!');
1637 }
1638 
1639 // Perform an HTTP request whose response uses a ReadableStream obtained over RPC.
1640 {
1641 let resp = await env.self.fetch('http://foo/stream-from-rpc');
1642 
1643 assert.strictEqual(await resp.text(), 'foo, bar, baz!');
1644 }
1645 },
1646};
1647 
1648export let serializeHttpTypes = {
1649 async test(controller, env, ctx) {
1650 {
1651 let headers = await env.MyService.returnEmptyHeaders();
1652 assert.deepEqual([...headers], []);
1653 }
1654 
1655 {
1656 let headers = await env.MyService.returnHeaders();
1657 assert.strictEqual(headers instanceof Headers, true);
1658 
1659 // Awkwardly, there's actually no API to get the non-lowercased header names.
1660 assert.deepEqual(
1661 [...headers],
1662 [
1663 ['content-length', '123'],
1664 ['corge', '!@#'],
1665 ['foo', 'bar'],
1666 ['set-cookie', 'abc'],
1667 ['set-cookie', 'def'],
1668 ]
1669 );
1670 
1671 assert.deepEqual(headers.getSetCookie(), ['abc', 'def']);
1672 }
1673 
1674 {
1675 let req = await env.MyService.returnRequest();
1676 
1677 assert.strictEqual(req.url, 'http://my-url.com');
1678 assert.strictEqual(req.method, 'PUT');
1679 assert.strictEqual(req.headers.get('Accept-Encoding'), 'bazzip');
1680 assert.strictEqual(req.headers.get('Foo'), 'Bar');
1681 assert.strictEqual(req.redirect, 'manual');
1682 
1683 assert.strictEqual(await req.text(), 'Hello every body!');
1684 
1685 assert.deepEqual(req.cf, {
1686 abc: 123,
1687 hello: 'goodbye',
1688 });
1689 }
1690 
1691 {
1692 let req = await env.MyService.roundTrip(
1693 new Request('http://foo', { signal: AbortSignal.timeout(100) })
1694 );
1695 assert.strictEqual(req.url, 'http://foo');
1696 assert.ok(req.signal instanceof AbortSignal);
1697 }
1698 
1699 {
1700 let req = await env.MyService.returnResponse();
1701 
1702 assert.strictEqual(req.status, 404);
1703 assert.strictEqual(req.statusText, 'Not Found');
1704 assert.strictEqual(req.headers.get('Content-Type'), 'abc');
1705 assert.deepEqual(req.cf, { foo: 123, bar: 'def' });
1706 
1707 assert.strictEqual(await req.text(), 'Response body!');
1708 }
1709 },
1710};
1711 
1712// Test that exceptions thrown from async native functions have a proper stack trace. (This is
1713// not specific to RPC but RPC is a convenient place to test it since we can easily define the
1714// callee to throw an exception.)
1715//
1716// Note that it's only a *local* stack trace of the client-side stack leading up to the call. The
1717// stack on the server side is not, at present, transmitted to the client.
1718export let testAsyncStackTrace = {
1719 async test(controller, env, ctx) {
1720 try {
1721 await env.MyService.throwingMethod();
1722 } catch (e) {
1723 // check that the custom property made it through
1724 assert.strictEqual(e.abc, 123);
1725 // verify a local stack trace was produced
1726 assert.strictEqual(e.stack.includes('at async Object.test'), true);
1727 }
1728 },
1729};
1730 
1731// Test that exceptions thrown over RPC have the .remote property.
1732export let testExceptionProperties = {
1733 async test(controller, env, ctx) {
1734 try {
1735 await env.MyService.throwingMethod();
1736 } catch (e) {
1737 assert.strictEqual(e.abc, 123);
1738 assert.strictEqual(e.remote, true);
1739 assert.strictEqual(e.message, 'METHOD THREW');
1740 }
1741 },
1742};
1743 
1744// Test that get(), put(), and delete() are valid RPC method names, not hijacked by Fetcher.
1745export let canUseGetPutDelete = {
1746 async test(controller, env, ctx) {
1747 assert.strictEqual(await env.MyService.get(12), 13);
1748 assert.strictEqual(await env.MyService.put(5, 7), 12);
1749 assert.strictEqual(await env.MyService.delete(3), 2);
1750 },
1751};
1752 
1753// Test that stubs can still be used after logging them.
1754export let logging = {
1755 async test(controller, env, ctx) {
1756 let counter = new MyCounter(0);
1757 let stub = new RpcStub(counter);
1758 assert.strictEqual(await stub.increment(1), 1);
1759 assert.strictEqual(await stub.increment(1), 2);
1760 console.log(stub);
1761 assert.strictEqual(await stub.increment(1), 3);
1762 },
1763};
1764 
1765export let proxiedRpcTarget = {
1766 async test(controller, env, ctx) {
1767 // Proxy RPC target.
1768 {
1769 let counter = new MyCounter(0);
1770 let proxy = new Proxy(counter, {
1771 get(target, prop, receiver) {
1772 if (prop == 'increment') {
1773 return (i) => target.increment(i + 123);
1774 } else {
1775 let result = target[prop];
1776 if (result instanceof Function) {
1777 result = result.bind(target);
1778 }
1779 return result;
1780 }
1781 },
1782 });
1783 
1784 await env.MyService.incrementCounter(proxy, 1);
1785 
1786 assert.strictEqual(counter.i, 124);
1787 }
1788 
1789 // Proxy plain object.
1790 {
1791 let counter = {
1792 i: 0,
1793 increment(i) {
1794 this.i += i;
1795 },
1796 };
1797 
1798 let proxy = new Proxy(counter, {
1799 get(target, prop, receiver) {
1800 if (prop == 'increment') {
1801 return (i) => target.increment(i + 123);
1802 } else {
1803 let result = target[prop];
1804 if (result instanceof Function) {
1805 result = result.bind(target);
1806 }
1807 return result;
1808 }
1809 },
1810 });
1811 
1812 await env.MyService.incrementCounter(proxy, 1);
1813 
1814 assert.strictEqual(counter.i, 124);
1815 }
1816 
1817 // Proxy function.
1818 {
1819 let func = (i) => {
1820 return i * 3;
1821 };
1822 func.ownProp = 123;
1823 let proxy = new Proxy(func, {
1824 apply(target, thisArg, argumentsList) {
1825 return target(...argumentsList) + 2;
1826 },
1827 });
1828 
1829 assert.strictEqual(await env.MyService.call(proxy, 2), 8);
1830 assert.strictEqual(await env.MyService.getProp(proxy, 'ownProp'), 123);
1831 
1832 // Try pipelining.
1833 assert.strictEqual(await env.MyService.identity(() => proxy)()(4), 14);
1834 assert.strictEqual(
1835 await env.MyService.identity(() => proxy)().ownProp,
1836 123
1837 );
1838 
1839 assert.strictEqual(
1840 await env.MyService.identity(() => ({ x: proxy }))().x(4),
1841 14
1842 );
1843 assert.strictEqual(
1844 await env.MyService.identity(() => ({ x: proxy }))().x.ownProp,
1845 123
1846 );
1847 }
1848 
1849 // Proxy RPC target that is callable.
1850 {
1851 let counter = new MyCounter(0);
1852 
1853 // We make the proxy target be a function so that it is callable, but we implement
1854 // getPrototypeOf() to make it appear to implement RpcTarget.
1855 let func = (i) => i * 11;
1856 func.ownProp = 123;
1857 let proxy = new Proxy(func, {
1858 get(target, prop, receiver) {
1859 if (prop == 'increment') {
1860 return (i) => counter.increment(i + 123);
1861 } else if (prop == 'ownProp') {
1862 return target.ownProp;
1863 } else {
1864 let result = counter[prop];
1865 if (result instanceof Function) {
1866 result = result.bind(counter);
1867 }
1868 return result;
1869 }
1870 },
1871 has(target, prop) {
1872 return prop in counter || prop === 'ownProp';
1873 },
1874 getPrototypeOf(target) {
1875 return Object.getPrototypeOf(counter);
1876 },
1877 });
1878 
1879 // We can call it.
1880 assert.strictEqual(await env.MyService.call(proxy, 3), 33);
1881 
1882 // We *cannot* access own properties of the function.
1883 assert.rejects(() => env.MyService.getProp(proxy, 'ownProp'), {
1884 name: 'TypeError',
1885 message: 'The RPC receiver does not implement the method "ownProp".',
1886 });
1887 
1888 // We *can* access prototype properties, becaues it's an RpcTarget.
1889 await env.MyService.incrementCounter(proxy, 1);
1890 assert.strictEqual(counter.i, 124);
1891 
1892 // Try pipelined calls.
1893 assert.strictEqual(
1894 await env.MyService.identity(() => proxy)().increment(3),
1895 250
1896 );
1897 assert.strictEqual(
1898 await env.MyService.identity(() => ({ p: proxy }))().p.increment(4),
1899 377
1900 );
1901 }
1902 
1903 // Can't proxy a class that doesn't extend `RpcTarget`.
1904 {
1905 let nonRpc = new NonRpcClass();
1906 let proxy = new Proxy(nonRpc, {
1907 get(target, prop, receiver) {
1908 if (prop == 'increment') {
1909 return (i) => target.increment(i + 123);
1910 } else {
1911 let result = target[prop];
1912 if (result instanceof Function) {
1913 result = result.bind(target);
1914 }
1915 return result;
1916 }
1917 },
1918 });
1919 
1920 await assert.rejects(() => env.MyService.incrementCounter(proxy, 1), {
1921 name: 'DataCloneError',
1922 message:
1923 'Proxy could not be serialized because it is not a valid RPC receiver type. The ' +
1924 'Proxy must emulate either a plain object or an RpcTarget, as indicated by the ' +
1925 "Proxy's prototype chain.",
1926 });
1927 }
1928 
1929 // Can't proxy an RpcTarget if we've overridden the prototype to say it's something else.
1930 {
1931 let counter = new MyCounter(0);
1932 let proxy = new Proxy(counter, {
1933 getPrototypeOf(target) {
1934 return NonRpcClass.prototype;
1935 },
1936 });
1937 
1938 await assert.rejects(() => env.MyService.incrementCounter(proxy, 1), {
1939 name: 'DataCloneError',
1940 message:
1941 'Proxy could not be serialized because it is not a valid RPC receiver type. The ' +
1942 'Proxy must emulate either a plain object or an RpcTarget, as indicated by the ' +
1943 "Proxy's prototype chain.",
1944 });
1945 }
1946 
1947 // CAN proxy a class that doesn't extend `RpcTarget` if we fake the prototype.
1948 {
1949 let nonRpc = new NonRpcClass();
1950 let proxy = new Proxy(nonRpc, {
1951 getPrototypeOf(target) {
1952 return RpcTarget.prototype;
1953 },
1954 });
1955 
1956 await env.MyService.incrementCounter(proxy, 321);
1957 assert.strictEqual(nonRpc.i, 321);
1958 }
1959 },
1960};
1961 
1962// Test that we can construct a WorkerEntrypoint
1963export class MyEntrypoint extends WorkerEntrypoint {
1964 rpcFunc() {
1965 return 'hello from entrypoint';
1966 }
1967}
1968function constructEntrypoint(cls, env) {
1969 return new cls({ waitUntil: () => {} }, env);
1970}
1971export let testConstructEntrypoint = {
1972 async test(controller, env, ctx) {
1973 const constructed = constructEntrypoint(MyEntrypoint, env);
1974 assert.strictEqual(await constructed.rpcFunc(), 'hello from entrypoint');
1975 },
1976};
1977 
1978// Test that calls to an RpcTarget made after the context is aborted don't get delivered.
1979export let portAbortCall = {
1980 async test(controller, env, ctx) {
1981 {
1982 let id = env.MyActor.newUniqueId();
1983 let actor = env.MyActor.get(id);
1984 let stub = await actor.makePostAbortCallTester();
1985 
1986 let hangPromise = stub.hang();
1987 assert.strictEqual(await stub.ping(), 'pong');
1988 let abortPromise = stub.abort();
1989 let pingPromise = stub.ping();
1990 
1991 await assert.rejects(abortPromise, {
1992 name: 'Error',
1993 message: 'test aborted by abort()',
1994 });
1995 await assert.rejects(pingPromise, {
1996 name: 'Error',
1997 message: 'test aborted by abort()',
1998 });
1999 await assert.rejects(hangPromise, {
2000 name: 'Error',
2001 message: 'test aborted by abort()',
2002 });
2003 // TODO(bug): This should propagate the abort reason.
2004 await assert.rejects(stub.ping(), {
2005 name: 'Error',
2006 message:
2007 'The execution context which hosts this callback is no longer running.',
2008 });
2009 await assert.rejects(actor.increment(2), {
2010 name: 'Error',
2011 message: 'test aborted by abort()',
2012 });
2013 }
2014 
2015 // Start over with a new stub, this time use failCriticalSection() to break the actor. As of
2016 // this writing, this differs significantly from plain `abort()` in that
2017 // `IoContext::abortException` never gets set, since `IoContext::abort()` is not directly
2018 // called, but instead the exception is joined into the on-abort promise.
2019 {
2020 let id = env.MyActor.newUniqueId();
2021 let actor = env.MyActor.get(id);
2022 let stub = await actor.makePostAbortCallTester();
2023 
2024 let hangPromise = stub.hang();
2025 assert.strictEqual(await stub.ping(), 'pong');
2026 let failPromise = stub.failCriticalSection();
2027 let pingPromise = stub.ping();
2028 
2029 await assert.rejects(failPromise, {
2030 name: 'Error',
2031 message: 'test broken critical section',
2032 });
2033 await assert.rejects(pingPromise, {
2034 name: 'Error',
2035 message: 'test broken critical section',
2036 });
2037 await assert.rejects(hangPromise, {
2038 name: 'Error',
2039 message: 'test broken critical section',
2040 });
2041 // TODO(bug): This should propagate the abort reason.
2042 await assert.rejects(stub.ping(), {
2043 name: 'Error',
2044 message:
2045 'The execution context which hosts this callback is no longer running.',
2046 });
2047 await assert.rejects(actor.increment(2), {
2048 name: 'Error',
2049 message: 'test broken critical section',
2050 });
2051 }
2052 },
2053};
2054 
2055export class Greeter extends WorkerEntrypoint {
2056 async greet(name) {
2057 return `${this.ctx.props.greeting}, ${name}!`;
2058 }
2059}
2060 
2061export class GreeterFactory extends WorkerEntrypoint {
2062 async makeGreeter(greeting) {
2063 return this.ctx.exports.Greeter({ props: { greeting } });
2064 }
2065 async makeGreeterWrapped(greeting) {
2066 return { greeter: this.ctx.exports.Greeter({ props: { greeting } }) };
2067 }
2068}
2069 
2070export let sendServiceStubOverRpc = {
2071 async test(controller, env, ctx) {
2072 {
2073 let greeter = await env.GreeterFactory.makeGreeter('Yo');
2074 assert.strictEqual(await greeter.greet('Alice'), 'Yo, Alice!');
2075 }
2076 
2077 // Test that we can pipeline on service stubs.
2078 {
2079 let greeter = env.GreeterFactory.makeGreeter('Yo');
2080 assert.strictEqual(await greeter.greet('Alice'), 'Yo, Alice!');
2081 }
2082 
2083 // Pipelining works a little differently when the service stub is returned as the top-level
2084 // value vs. an inner value, so test an inner value too.
2085 {
2086 let greeter = env.GreeterFactory.makeGreeterWrapped('Yo').greeter;
2087 assert.strictEqual(await greeter.greet('Alice'), 'Yo, Alice!');
2088 }
2089 },
2090};
2091 
2092// Make sure that calls are delivered in e-order, even in the presence of pushed externals.
2093export let eOrderTest = {
2094 async test(controller, env, ctx) {
2095 let abortController = new AbortController();
2096 let abortSignal = abortController.signal;
2097 
2098 let _readableController;
2099 let readableStream = new ReadableStream({
2100 start(c) {
2101 _readableController = c;
2102 },
2103 });
2104 
2105 let stub = await env.MyService.makeCounter(0);
2106 
2107 let promises = [];
2108 promises.push(stub.increment(1));
2109 promises.push(stub.increment(1));
2110 promises.push(stub.increment(1, abortSignal));
2111 promises.push(stub.increment(1));
2112 promises.push(stub.increment(1, readableStream));
2113 promises.push(stub.increment(1));
2114 
2115 let results = await Promise.all(promises);
2116 
2117 assert.deepEqual(results, [1, 2, 3, 4, 5, 6]);
2118 },
2119};