Skip to content
File

Blob: src/workerd/api/tests/cross-context-promise-test.js

javascript459 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 { match, rejects, strictEqual, throws } from 'assert';
5import { AsyncLocalStorage } from 'async_hooks';
6import { inspect } from 'util';
7import { mock } from 'node:test';
8 
9export const crossContextResolveWorks = {
10 async test(_, env) {
11 // We're going to send two simultaneous requests to the same endpoint.
12 const results = await Promise.allSettled([
13 env.subrequest.fetch('http://example.org/resolve'),
14 env.subrequest.fetch('http://example.org/resolve'),
15 ]);
16 strictEqual(results[0].status, 'fulfilled');
17 strictEqual(results[1].status, 'fulfilled');
18 strictEqual(results[0].value.status, 200);
19 strictEqual(results[1].value.status, 200);
20 strictEqual(await results[0].value.text(), 'ok');
21 strictEqual(await results[1].value.text(), 'ok');
22 },
23};
24 
25export const crossContextRejectWorks = {
26 async test(_, env) {
27 // We're going to send two simultaneous requests to the same endpoint.
28 const results = await Promise.allSettled([
29 env.subrequest.fetch('http://example.org/reject'),
30 env.subrequest.fetch('http://example.org/reject'),
31 ]);
32 strictEqual(results[0].status, 'fulfilled');
33 strictEqual(results[1].status, 'fulfilled');
34 strictEqual(results[0].value.status, 200);
35 strictEqual(results[1].value.status, 200);
36 strictEqual(await results[0].value.text(), 'ok');
37 strictEqual(await results[1].value.text(), 'ok');
38 },
39};
40 
41export const crossContextStreamWorks = {
42 async test(_, env) {
43 // We're going to send two simultaneous requests to the same endpoint.
44 const results = await Promise.allSettled([
45 env.subrequest.fetch('http://example.org/stream'),
46 env.subrequest.fetch('http://example.org/stream'),
47 ]);
48 strictEqual(results[0].status, 'fulfilled');
49 strictEqual(results[1].status, 'fulfilled');
50 strictEqual(results[0].value.status, 200);
51 strictEqual(results[1].value.status, 200);
52 strictEqual(await results[0].value.text(), 'ok');
53 strictEqual(await results[1].value.text(), 'ok');
54 },
55};
56 
57export const customThenableWorks = {
58 async test(_, env) {
59 // We're going to send two simultaneous requests to the same endpoint.
60 const results = await Promise.allSettled([
61 env.subrequest.fetch('http://example.org/thenable'),
62 env.subrequest.fetch('http://example.org/thenable'),
63 ]);
64 strictEqual(results[0].status, 'fulfilled');
65 strictEqual(results[1].status, 'fulfilled');
66 strictEqual(results[0].value.status, 200);
67 strictEqual(results[1].value.status, 200);
68 strictEqual(await results[0].value.text(), 'ok');
69 strictEqual(await results[1].value.text(), 'ok');
70 },
71};
72 
73export const unhandledRejectionWorks = {
74 async test(_, env) {
75 // We're going to send two simultaneous requests to the same endpoint.
76 const results = await Promise.allSettled([
77 env.subrequest.fetch('http://example.org/unhandled'),
78 env.subrequest.fetch('http://example.org/unhandled'),
79 ]);
80 strictEqual(results[0].status, 'fulfilled');
81 strictEqual(results[1].status, 'fulfilled');
82 strictEqual(results[0].value.status, 200);
83 strictEqual(results[1].value.status, 200);
84 strictEqual(await results[0].value.text(), 'ok');
85 strictEqual(await results[1].value.text(), 'ok');
86 },
87};
88 
89export const expiredContextWorks = {
90 async test(_, env) {
91 // We're going to send two simultaneous requests to the same endpoint.
92 const results = await Promise.allSettled([
93 env.subrequest.fetch('http://example.org/expired'),
94 env.subrequest.fetch('http://example.org/expired'),
95 ]);
96 strictEqual(results[0].status, 'rejected');
97 strictEqual(results[1].status, 'fulfilled');
98 strictEqual(
99 results[0].reason.message,
100 "The Workers runtime canceled this request because it detected that your Worker's code " +
101 'had hung and would never generate a response. Refer to: ' +
102 'https://developers.cloudflare.com/workers/observability/errors/'
103 );
104 strictEqual(results[1].value.status, 200);
105 strictEqual(await results[1].value.text(), 'ok');
106 // Wait a tick for things to settle out before checking the global.
107 // We're just making sure here that the promise in the first request
108 // was canceled correctly.
109 await scheduler.wait(100);
110 strictEqual(globalThis.expiredRan, undefined);
111 },
112};
113 
114export const asyncIterWorks = {
115 async test(_, env) {
116 // We're going to send two simultaneous requests to the same endpoint.
117 const results = await Promise.allSettled([
118 env.subrequest.fetch('http://example.org/asynciter'),
119 env.subrequest.fetch('http://example.org/asynciter'),
120 ]);
121 strictEqual(results[0].status, 'fulfilled');
122 strictEqual(results[1].status, 'fulfilled');
123 strictEqual(results[0].value.status, 200);
124 strictEqual(results[1].value.status, 200);
125 strictEqual(await results[0].value.text(), 'ok');
126 strictEqual(await results[1].value.text(), 'ok');
127 },
128};
129 
130export const cyclicAwaitsWorks = {
131 async test(_, env) {
132 // We're going to send two simultaneous requests to the same endpoint.
133 const results = await Promise.allSettled([
134 env.subrequest.fetch('http://example.org/cyclic'),
135 env.subrequest.fetch('http://example.org/cyclic'),
136 ]);
137 strictEqual(results[0].status, 'rejected');
138 strictEqual(results[1].status, 'fulfilled');
139 },
140};
141 
142export default {
143 async fetch(req, env, ctx) {
144 if (req.url.endsWith('/resolve')) {
145 return resolveTest(req, env, ctx);
146 } else if (req.url.endsWith('/reject')) {
147 return rejectTest(req, env, ctx);
148 } else if (req.url.endsWith('/stream')) {
149 return crossRequestStream(req, env, ctx);
150 } else if (req.url.endsWith('/thenable')) {
151 return customThenable(req, env, ctx);
152 } else if (req.url.endsWith('/unhandled')) {
153 return unhandledRejection(req, env, ctx);
154 } else if (req.url.endsWith('/expired')) {
155 return expiredContext(req, env, ctx);
156 } else if (req.url.endsWith('/asynciter')) {
157 return asyncIterator(req, env, ctx);
158 } else if (req.url.endsWith('/cyclic')) {
159 return cyclicPromise(req, env, ctx);
160 }
161 throw new Error('Invalid URL');
162 },
163};
164 
165function setupWaiter(ctx) {
166 const { promise, resolve } = Promise.withResolvers();
167 setTimeout(resolve, 1000);
168 ctx.waitUntil(promise);
169}
170 
171async function resolveTest(req, env, ctx) {
172 const als = new AsyncLocalStorage();
173 // This will be called twice. The first time in one request where we will
174 // create the promise and resolver. The second time for the second request
175 // where the promise will be resolved.
176 if (globalThis.request1 === undefined) {
177 setupWaiter(ctx);
178 const { promise, resolve } = Promise.withResolvers();
179 globalThis.request1 = { promise, resolve };
180 const ab = AbortSignal.abort();
181 strictEqual(ab.aborted, true);
182 await als.run(123, async () => {
183 await promise;
184 strictEqual(als.getStore(), 123);
185 });
186 // This part is the main test. It will not run until after the promise
187 // is resolved in the second request.
188 // We use an AbortSignal because it is bound to the IoContext and will
189 // throw an error if ab.aborted is checked from the wrong IoContext.
190 // If this line runes, it is proof that the promise continuation is
191 // running in the correct IoContext.
192 strictEqual(ab.aborted, true);
193 return new Response('ok');
194 }
195 
196 // This is our second request. Here, all we do is resolve the promise.
197 
198 // While we are deferring the continuations from the promise, the promise state
199 // change should happen immediately. Before calling resolve, the state should
200 // be pending. After calling resolve, the state should be resolved showing
201 // an undefined value. Updating the state of the promise immediately and
202 // synchronously is required by the language specification.
203 // See: https://tc39.es/ecma262/#sec-promise-resolve-functions
204 strictEqual(inspect(globalThis.request1.promise), 'Promise { <pending> }');
205 als.run('abc', () => globalThis.request1.resolve());
206 strictEqual(inspect(globalThis.request1.promise), 'Promise { undefined }');
207 
208 const p = globalThis.request1.promise;
209 
210 globalThis.request1 = undefined;
211 
212 // We ought to be able to do a cross-request wait on the promise still.
213 await p;
214 
215 return new Response('ok');
216}
217 
218const reason = new Error('boom');
219 
220async function rejectTest(req, env, ctx) {
221 if (globalThis.request2 === undefined) {
222 setupWaiter(ctx);
223 const { promise, reject } = Promise.withResolvers();
224 globalThis.request2 = { reject };
225 const ab = AbortSignal.abort();
226 strictEqual(ab.aborted, true);
227 try {
228 // The promise will be rejected from the other request.
229 await promise;
230 throw new Error('should not get here');
231 } catch (err) {
232 // The reason provided by the other request should be carried
233 // through here. If the ab.aborted check throws, then the continuation
234 // is running in the wrong IoContext, which is the main thing we are
235 // testing for here.
236 strictEqual(err, reason);
237 strictEqual(ab.aborted, true);
238 }
239 return new Response('ok');
240 }
241 
242 // This is our second request. Here, all we do is reject the promise.
243 globalThis.request2.reject(reason);
244 globalThis.request2 = undefined;
245 return new Response('ok');
246}
247 
248async function crossRequestStream(req, env, ctx) {
249 // Here, we are going to create a stream that will be used across
250 // requests. The first request will create the stream and queue up
251 // a pending read. The second request will fulfill the read by providing
252 // data to the stream.
253 if (globalThis.stream === undefined) {
254 setupWaiter(ctx);
255 let controller;
256 const readable = new ReadableStream({
257 start(c) {
258 controller = c;
259 },
260 });
261 globalThis.stream = { controller };
262 const reader = readable.getReader();
263 const ab = AbortSignal.abort();
264 strictEqual(ab.aborted, true);
265 const _read = await reader.read();
266 strictEqual(ab.aborted, true);
267 return new Response('ok');
268 }
269 
270 // This is our second request. Here, all we do is provide data to the stream.
271 // This will cause our pending read to be fulfilled.
272 const enc = new TextEncoder();
273 globalThis.stream.controller.enqueue(enc.encode('hello'));
274 globalThis.stream = undefined;
275 return new Response('ok');
276}
277 
278async function customThenable(req, env, ctx) {
279 // This will be called twice. The first time in one request where we will
280 // create the promise and resolver. The second time for the second request
281 // where the promise will be resolved.
282 if (globalThis.thenable === undefined) {
283 setupWaiter(ctx);
284 const { promise, resolve } = Promise.withResolvers();
285 globalThis.thenable = { resolve };
286 const ab = AbortSignal.abort();
287 strictEqual(ab.aborted, true);
288 
289 // We check to make sure the value provided by the custom thenable is
290 // property passed through to the promise resolution.
291 
292 strictEqual(await promise, 1);
293 // This part is the main test. It will not run until after the promise
294 // is resolved in the second request.
295 // We use an AbortSignal because it is bound to the IoContext and will
296 // throw an error if ab.aborted is checked from the wrong IoContext.
297 // If this line runes, it is proof that the promise continuation is
298 // running in the correct IoContext.
299 strictEqual(ab.aborted, true);
300 return new Response('ok');
301 }
302 
303 // This is our second request. Here, all we do is resolve the promise.
304 const ab = AbortSignal.abort();
305 strictEqual(ab.aborted, true);
306 
307 const then = mock.fn((resolve) => {
308 // The thenable should be invoked in the second request's IoContext.
309 // If it is not, then the ab.aborted check below will fail.
310 strictEqual(ab.aborted, true);
311 resolve(1);
312 });
313 
314 globalThis.thenable.resolve({ then });
315 globalThis.thenable = undefined;
316 
317 // Confirm that the thenable was called exactly once in this request.
318 await Promise.resolve();
319 strictEqual(then.mock.calls.length, 1);
320 
321 return new Response('ok');
322}
323 
324async function unhandledRejection(req, env, ctx) {
325 // This will be called twice. The first time in one request where we will
326 // create the promise and resolver. The second time for the second request
327 // where the promise will be rejected.
328 if (globalThis.unhandled === undefined) {
329 setupWaiter(ctx);
330 const { reject } = Promise.withResolvers();
331 globalThis.unhandled = { reject };
332 const ab = AbortSignal.abort();
333 strictEqual(ab.aborted, true);
334 
335 const rejectPromise = Promise.withResolvers();
336 globalThis.addEventListener(
337 'unhandledrejection',
338 (event) => {
339 // Here we have a gotcha! The unhandledrejection event is dispatched
340 // synchronously when the promise is rejected. It does not get deferred.
341 // so the IoContext here will be the second request's IoContext! This
342 // means that our ab.aborted check will fail!
343 throws(() => ab.aborted, {
344 message: /I\/O type: RefcountedCanceler/,
345 });
346 strictEqual(event.reason, reason);
347 rejectPromise.resolve();
348 },
349 { once: true }
350 );
351 await rejectPromise.promise;
352 
353 return new Response('ok');
354 }
355 
356 // This is our second request. Here, all we do is reject the promise.
357 globalThis.unhandled.reject(reason);
358 globalThis.unhandled = undefined;
359 return new Response('ok');
360}
361 
362async function expiredContext(req, env, ctx) {
363 if (globalThis.expired === undefined) {
364 // We do not arrange for a waiter here. We want the request context
365 // to be canceled before the promise is resolved. In the error log
366 // (which we unfortunately cannot check here) we should see a message
367 // about the hanging promise being canceled and another message about
368 // a promise resolved across request contexts but the target request
369 // is not longer active. On the calling side, the first request
370 // should reject and the second request should resolve.
371 const { promise, resolve } = Promise.withResolvers();
372 globalThis.expired = { resolve };
373 promise.then(() => {
374 // This should never run...
375 globalThis.expiredRan = true;
376 });
377 await new Promise(() => {});
378 return new Response('ok');
379 }
380 
381 // This is our second request. Here, all we do is resolve the promise.
382 // Let's wait a bit to make sure the other request has had time to
383 // be canceled and destroyed.
384 await scheduler.wait(100);
385 globalThis.expired.resolve();
386 globalThis.expired = undefined;
387 return new Response('ok');
388}
389 
390async function* gen(ab) {
391 let c = 0;
392 for (;;) {
393 await scheduler.wait(10);
394 strictEqual(ab.aborted, true);
395 yield c++;
396 }
397}
398 
399async function asyncIterator(req, env, ctx) {
400 if (globalThis.asynciter === undefined) {
401 const ab = AbortSignal.abort();
402 globalThis.asynciter = gen(ab);
403 globalThis.asyncIter2 = gen(ab);
404 return new Response('ok');
405 }
406 
407 // Advancing the iterator should throw an I/O error because the generator
408 // is bound to the wrong IoContext here. This test just confirms that our
409 // promise context patch does not change the expected behavior of the async
410 // generator and async iterators. Those should still throw I/O errors when
411 // the generator is bound to the wrong IoContext via something that it captures.
412 await rejects(globalThis.asynciter.next(), {
413 message: /Cannot perform I\/O/,
414 });
415 
416 const foo = {
417 [Symbol.asyncIterator]() {
418 return globalThis.asyncIter2;
419 },
420 };
421 
422 try {
423 for await (const _ of foo) {
424 // intentionally empty
425 }
426 throw new Error('should not get here');
427 } catch (err) {
428 match(err.message, /Cannot perform I\/O/);
429 }
430 
431 return new Response('ok');
432}
433 
434async function cyclicPromise(req, env, ctx) {
435 if (globalThis.cyclic === undefined) {
436 setupWaiter(ctx);
437 const { promise, resolve } = Promise.withResolvers();
438 globalThis.cyclic = { promise, resolve };
439 const ab = AbortSignal.abort();
440 strictEqual(ab.aborted, true);
441 await promise;
442 throw new Error('should never get here');
443 }
444 
445 async function foo() {
446 await globalThis.cyclic.promise;
447 throw new Error('shnould never get here');
448 }
449 
450 // The first request will hang because the promise is resolved
451 // with a promise that waits on the first, creating a hang
452 // condition that we can detect. The test log will contain
453 // a message about a hanging promise being canceled.
454 globalThis.cyclic.resolve(foo());
455 globalThis.cyclic = undefined;
456 
457 return new Response('ok');
458}