File
Blob: src/workerd/api/tests/request-signal-enabled.js
| 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 |
| 4 | import { DurableObject, WorkerEntrypoint } from 'cloudflare:workers'; |
| 5 | import assert from 'node:assert'; |
| 6 | |
| 7 | export 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 | |
| 16 | let reqSignalReason = null; |
| 17 | let rpcSignalReason = null; |
| 18 | export 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 | |
| 33 | export 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 | |
| 122 | export 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 | |
| 136 | export 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 | |
| 149 | export 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 | |
| 168 | export 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 | |
| 196 | export 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 | |
| 223 | export 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 | }; |