Skip to content
File

Blob: src/workerd/api/tests/request-signal-enabled.js

javascript244 lines
1// Copyright (c) 2025 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 { DurableObject, WorkerEntrypoint } from 'cloudflare:workers';
5import assert from 'node:assert';
6 
7export class AbortTracker extends DurableObject {
8 async getAborted(key) {
9 return this.ctx.storage.get(key);
10 }
11 async setAborted(key, value) {
12 await this.ctx.storage.put(key, value);
13 }
14}
15 
16let reqSignalReason = null;
17let rpcSignalReason = null;
18export class OtherServer extends WorkerEntrypoint {
19 async fetch(req) {
20 await scheduler.wait(300);
21 return new Response('completed');
22 }
23 
24 async echo(val) {
25 return val;
26 }
27 
28 async saveReqReason(req) {
29 rpcSignalReason = req.signal.reason?.message;
30 }
31}
32 
33export class Server extends WorkerEntrypoint {
34 async fetch(req) {
35 const key = new URL(req.url).pathname.slice(1);
36 let abortTracker = this.env.AbortTracker.get(
37 this.env.AbortTracker.idFromName('AbortTracker')
38 );
39 await abortTracker.setAborted(key, false);
40 
41 req.signal.onabort = () => {
42 this.ctx.waitUntil(abortTracker.setAborted(key, true));
43 };
44 
45 return this[key](req);
46 }
47 
48 async valid() {
49 return new Response('hello world');
50 }
51 
52 async error() {
53 throw new Error('boom');
54 }
55 
56 async hang() {
57 for (;;) {
58 await scheduler.wait(86400);
59 }
60 }
61 
62 async hangAfterSendingSomeData() {
63 const { readable, writable } = new IdentityTransformStream();
64 this.ctx.waitUntil(this.sendSomeData(writable));
65 
66 return new Response(readable);
67 }
68 
69 async sendSomeData(writable) {
70 const writer = writable.getWriter();
71 const enc = new TextEncoder();
72 await writer.write(enc.encode('hello world'));
73 await this.hang();
74 }
75 
76 async triggerSubrequest(req) {
77 this.ctx.waitUntil(this.callOtherServer(req));
78 await this.hang();
79 }
80 
81 async callOtherServer(req) {
82 const key = 'subrequest';
83 
84 let abortTracker = this.env.AbortTracker.get(
85 this.env.AbortTracker.idFromName('AbortTracker')
86 );
87 
88 const passedThroughReq = new Request(req);
89 passedThroughReq.onabort = () => {
90 this.ctx.waitUntil(abortTracker.setAborted(key, true));
91 };
92 
93 const res = await this.env.OtherServer.fetch(passedThroughReq);
94 const text = await res.text();
95 
96 if (text == 'completed') {
97 await abortTracker.setAborted(key, false);
98 }
99 }
100 
101 async passIncomingRequestOverRpc(req) {
102 // req.signal is in the "this signal" slot. It's an incoming request signal so it has the
103 // IGNORE_FOR_SUBREQUESTS flag set.
104 req.signal.onabort = () => {
105 // If we're here we know that the reason is filled in on req.signal
106 reqSignalReason = req.signal.reason.message;
107 
108 // Ensure that the incoming request signal is in the "signal" slot, which is the one we
109 // actually serialize currently.
110 const newReq = req.clone();
111 
112 // Since the signal is not included in serialization, OtherServer won't be able to read the
113 // reason.
114 this.ctx.waitUntil(this.env.OtherServer.saveReqReason(newReq));
115 };
116 
117 // Just hang and wait for the client to give up
118 await scheduler.wait(86400);
119 }
120}
121 
122export const noAbortOnSimpleResponse = {
123 async test(ctrl, env, ctx) {
124 let abortTracker = env.AbortTracker.get(
125 env.AbortTracker.idFromName('AbortTracker')
126 );
127 
128 const req = env.Server.fetch('http://example.com/valid');
129 
130 const res = await req;
131 assert.strictEqual(await res.text(), 'hello world');
132 assert.strictEqual(await abortTracker.getAborted('valid'), false);
133 },
134};
135 
136export const noAbortIfServerThrows = {
137 async test(ctrl, env, ctx) {
138 let abortTracker = env.AbortTracker.get(
139 env.AbortTracker.idFromName('AbortTracker')
140 );
141 
142 const req = env.Server.fetch('http://example.com/error');
143 
144 await assert.rejects(() => req, { name: 'Error', message: 'boom' });
145 assert.strictEqual(await abortTracker.getAborted('error'), false);
146 },
147};
148 
149export const abortIfClientAbandonsRequest = {
150 async test(ctrl, env, ctx) {
151 let abortTracker = env.AbortTracker.get(
152 env.AbortTracker.idFromName('AbortTracker')
153 );
154 
155 // This endpoint never generates a response, so we can timeout after an arbitrary time.
156 const req = env.Server.fetch('http://example.com/hang', {
157 signal: AbortSignal.timeout(500),
158 });
159 
160 await assert.rejects(() => req, {
161 name: 'TimeoutError',
162 message: 'The operation was aborted due to timeout',
163 });
164 assert.strictEqual(await abortTracker.getAborted('hang'), true);
165 },
166};
167 
168export const abortIfClientCancelsReadingResponse = {
169 async test(ctrl, env, ctx) {
170 let abortTracker = env.AbortTracker.get(
171 env.AbortTracker.idFromName('AbortTracker')
172 );
173 
174 // This endpoint begins generating a response but then hangs
175 const req = env.Server.fetch('http://example.com/hangAfterSendingSomeData');
176 const res = await req;
177 const reader = res.body.getReader();
178 
179 const { value, done } = await reader.read();
180 assert.strictEqual(new TextDecoder().decode(value), 'hello world');
181 assert.ok(!done);
182 
183 // Give up reading
184 await reader.cancel();
185 
186 // Waste a bit of time so the server cleans up
187 await scheduler.wait(0);
188 
189 assert.strictEqual(
190 await abortTracker.getAborted('hangAfterSendingSomeData'),
191 true
192 );
193 },
194};
195 
196export const abortedRequestDoesNotAbortSubrequest = {
197 async test(ctrl, env, ctx) {
198 let abortTracker = env.AbortTracker.get(
199 env.AbortTracker.idFromName('AbortTracker')
200 );
201 
202 // This endpoint calls another endpoint that eventually completes after wasting 300 ms
203 // So, we abort the initial request quickly...
204 const req = env.Server.fetch('http://example.com/triggerSubrequest', {
205 signal: AbortSignal.timeout(100),
206 });
207 
208 await assert.rejects(() => req, {
209 name: 'TimeoutError',
210 message: 'The operation was aborted due to timeout',
211 });
212 assert.strictEqual(
213 await abortTracker.getAborted('triggerSubrequest'),
214 true
215 );
216 
217 // Then make sure that the subrequest wasn't also aborted
218 await scheduler.wait(500);
219 assert.strictEqual(await abortTracker.getAborted('subrequest'), false);
220 },
221};
222 
223export const requestPassedOverRpcDoesNotIncludeSignal = {
224 async test(ctrl, env, ctx) {
225 // This test covers the case when the `request_signal_passthrough` compat flag is not set.
226 // The incoming request's signal is not serialized.
227 // See request-signal-psssthrough.js for the corresponding test when the flag is set.
228 
229 await assert.rejects(
230 () =>
231 env.Server.fetch('http://example.com/passIncomingRequestOverRpc', {
232 signal: AbortSignal.timeout(100),
233 }),
234 { message: 'The operation was aborted due to timeout' }
235 );
236 
237 // Yield so the workers can write the abort reasons
238 await scheduler.wait(0);
239 
240 assert.strictEqual(reqSignalReason, 'The client has disconnected');
241 assert.strictEqual(rpcSignalReason, undefined);
242 },
243};