Skip to content

Commit 8f21ea8

Browse files
authored
Merge pull request #7500 from cloudflare/jasnell/ts-streams-pipe-abort-algorithm
Observe the pipe's AbortSignal through an abort algorithm
2 parents 2854260 + fdfeb04 commit 8f21ea8

12 files changed

Lines changed: 240 additions & 10 deletions

File tree

‎src/per_isolate/per_isolate-env.d.ts‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,8 +93,22 @@ declare const utils: {
9393
fileHandle: object,
9494
keepExistingData: boolean
9595
): FileSystemWriteContext;
96+
// The DOM's "add an abort algorithm" (api/abort-bootstrap.h): `algorithm`
97+
// runs when `signal` aborts, before the 'abort' event, and never for a
98+
// synthetic 'abort' event. Consumed by the webstreams pipe.
99+
addAbortAlgorithm(
100+
signal: AbortSignal,
101+
algorithm: () => void
102+
): AbortAlgorithmHandle;
96103
};
97104

105+
// An abort algorithm registration. Obtained only from
106+
// utils.addAbortAlgorithm(); it has no JS-reachable constructor.
107+
declare interface AbortAlgorithmHandle {
108+
// Unregisters the algorithm. Idempotent.
109+
remove(): void;
110+
}
111+
98112
// Native incremental digest, backing the TypeScript DigestStream. Obtained only
99113
// from utils.createDigestContext(); it has no JS-reachable constructor.
100114
//

‎src/per_isolate/webstreams/readable.ts‎

Lines changed: 6 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -49,8 +49,6 @@ const {
4949
DataViewPrototypeGetBuffer,
5050
DataViewPrototypeGetByteLength,
5151
DataViewPrototypeGetByteOffset,
52-
EventTargetAddEventListener,
53-
EventTargetRemoveEventListener,
5452
JSONParse,
5553
MathMax,
5654
MathMin,
@@ -2962,15 +2960,13 @@ function pipeToInternal<R>(
29622960
PromiseWithResolvers() as PromiseWithResolversType<void>;
29632961

29642962
let shuttingDown = false;
2965-
let abortAlgorithm: (() => void) | undefined;
2963+
let abortRegistration: AbortAlgorithmHandle | undefined;
29662964

29672965
const finalize = (error?: { reason: unknown }): void => {
29682966
writableInternals.setReadyHook(destination, undefined);
29692967
writableInternals.writerRelease(writer);
29702968
readableStreamReaderGenericRelease(reader);
2971-
if (signal !== undefined && abortAlgorithm !== undefined) {
2972-
EventTargetRemoveEventListener(signal, 'abort', abortAlgorithm);
2973-
}
2969+
abortRegistration?.remove();
29742970
if (error !== undefined) {
29752971
reject(error.reason);
29762972
} else {
@@ -3167,8 +3163,10 @@ function pipeToInternal<R>(
31673163
);
31683164
};
31693165

3166+
// Spec: an abort algorithm, not an 'abort' listener, so a synthetic event
3167+
// cannot abort the pipe and a user listener cannot prevent it.
31703168
if (signal !== undefined) {
3171-
abortAlgorithm = () => {
3169+
const abortAlgorithm = (): void => {
31723170
const abortReason = AbortSignalReasonGet(signal);
31733171
const actions: (() => Promise<unknown>)[] = [];
31743172
if (!preventAbort) {
@@ -3197,7 +3195,7 @@ function pipeToInternal<R>(
31973195
if (AbortSignalAbortedGet(signal)) {
31983196
abortAlgorithm();
31993197
} else {
3200-
EventTargetAddEventListener(signal, 'abort', abortAlgorithm);
3198+
abortRegistration = utils.addAbortAlgorithm(signal, abortAlgorithm);
32013199
}
32023200
}
32033201

‎src/tests/streams/piping/AGENTS.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,7 @@ the source FIRST, then releasing the write (`pipeStopsPullingWhenDestStalls`).
7676
| --- | --- |
7777
| `pipe-matrix.js` | migrated pipe-streams-test.js wholesale (35): pipeThrough + pipeTo across JS↔native in all directions, prevent* combos, pre-aborted and mid-read AbortSignals, tee'd pipes, queued-destination close (ledger #1-#4, #13) |
7878
| `api-surface.js` | brand checks (ledger #5), option getter order, throwing getters, invalid signal, locked pipeThrough endpoints |
79+
| `abort-signal.js` | the pipe's AbortSignal is an abort algorithm: a synthetic 'abort' event is ignored, a listener's stopImmediatePropagation() cannot block the abort, an abort after the pipe settles does nothing |
7980
| `error-propagation.js` | forward matrix (starts-errored × prevent* × truthy), hwm-0 dest (ledger #6), custom-error preservation (migrated from streams-error-edge-cases-test.js) |
8081
| `close-propagation.js` | the WPT-disabled backward territory, bounded: external close/abort on piped dest, write-throw backward propagation, idle dest-controller error and one with a write in flight (ledger #7) |
8182
| `flow-control.js` | backpressure chain (migrated from streams-backpressure-test.js), stalled-dest read-ahead bound |
Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,85 @@
1+
// Copyright (c) 2026 Cloudflare, Inc.
2+
// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3+
// https://opensource.org/licenses/Apache-2.0
4+
5+
// How a pipe observes its AbortSignal: only a real abort stops it (the
6+
// spec's abort algorithm), not the 'abort' event, and nothing a user's
7+
// 'abort' listener does can prevent it.
8+
9+
import { strictEqual, deepStrictEqual } from 'node:assert';
10+
11+
function recordingPipe(signal) {
12+
let rc;
13+
const rs = new ReadableStream({
14+
start(c) {
15+
rc = c;
16+
},
17+
});
18+
const written = [];
19+
const ws = new WritableStream({
20+
write(chunk) {
21+
written.push(chunk);
22+
},
23+
});
24+
const pipe = rs.pipeTo(ws, { signal });
25+
return { rs, ws, rc, written, pipe };
26+
}
27+
28+
// A synthetic 'abort' event on a signal that is not aborted leaves the
29+
// pipe running.
30+
export const syntheticAbortEventIgnored = {
31+
async test() {
32+
const ac = new AbortController();
33+
const { rs, ws, rc, written, pipe } = recordingPipe(ac.signal);
34+
await scheduler.wait(0);
35+
ac.signal.dispatchEvent(new Event('abort'));
36+
rc.enqueue('a');
37+
rc.close();
38+
await pipe;
39+
deepStrictEqual(written, ['a']);
40+
strictEqual(rs.locked, false);
41+
strictEqual(ws.locked, false);
42+
},
43+
};
44+
45+
// A listener registered before the pipe that stops immediate propagation
46+
// does not keep the abort from the pipe.
47+
export const stopImmediatePropagationDoesNotBlockAbort = {
48+
async test() {
49+
const ac = new AbortController();
50+
const reason = new Error('boom');
51+
ac.signal.addEventListener('abort', (e) => e.stopImmediatePropagation());
52+
const { rs, ws, rc, pipe } = recordingPipe(ac.signal);
53+
await scheduler.wait(0);
54+
ac.abort(reason);
55+
// The C++ pipe notices the abort on its next step.
56+
rc.enqueue('a');
57+
const outcome = await Promise.race([
58+
pipe.then(
59+
() => 'fulfilled',
60+
(e) => e
61+
),
62+
scheduler.wait(1000).then(() => 'pending'),
63+
]);
64+
strictEqual(outcome, reason);
65+
strictEqual(rs.locked, false);
66+
strictEqual(ws.locked, false);
67+
},
68+
};
69+
70+
// Aborting after the pipe has settled has no effect on either stream.
71+
export const abortAfterPipeSettled = {
72+
async test() {
73+
const ac = new AbortController();
74+
const { rs, ws, rc, written, pipe } = recordingPipe(ac.signal);
75+
rc.enqueue('a');
76+
rc.close();
77+
await pipe;
78+
ac.abort(new Error('late'));
79+
await scheduler.wait(0);
80+
deepStrictEqual(written, ['a']);
81+
const writer = ws.getWriter();
82+
await writer.closed;
83+
strictEqual(rs.locked, false);
84+
},
85+
};

‎src/tests/streams/piping/main.js‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,12 @@ export {
5050
pipeThroughLockedEndpoints,
5151
} from 'api-surface';
5252

53+
export {
54+
syntheticAbortEventIgnored,
55+
stopImmediatePropagationDoesNotBlockAbort,
56+
abortAfterPipeSettled,
57+
} from 'abort-signal';
58+
5359
export {
5460
sourceStartsErrored,
5561
sourceStartsErroredPreventAbort,

‎src/tests/streams/piping/piping-modules.capnp‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ const modules :List(Workerd.Worker.Module) = [
88
(name = "which-impl", esModule = embed "which-impl.js"),
99
(name = "pipe-matrix", esModule = embed "pipe-matrix.js"),
1010
(name = "api-surface", esModule = embed "api-surface.js"),
11+
(name = "abort-signal", esModule = embed "abort-signal.js"),
1112
(name = "error-propagation", esModule = embed "error-propagation.js"),
1213
(name = "close-propagation", esModule = embed "close-propagation.js"),
1314
(name = "flow-control", esModule = embed "flow-control.js"),

‎src/workerd/api/BUILD.bazel‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@ filegroup(
6464
filegroup(
6565
name = "hdrs",
6666
srcs = [
67+
"abort-bootstrap.h",
6768
"actor.h",
6869
"actor-call-retry.h",
6970
"actor-state.h",

‎src/workerd/api/abort-bootstrap.h‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
// Copyright (c) 2026 Cloudflare, Inc.
2+
// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3+
// https://opensource.org/licenses/Apache-2.0
4+
5+
#pragma once
6+
7+
#include <workerd/jsg/jsg.h>
8+
9+
namespace workerd::api {
10+
11+
// The DOM's "add an algorithm to signal's abort algorithms" for the per-isolate
12+
// bootstrap (utils.addAbortAlgorithm): `algorithm` is called with no arguments
13+
// when `signal` aborts, before the 'abort' event is dispatched, and never for a
14+
// synthetic 'abort' event. Returns an AbortAlgorithmHandle whose remove()
15+
// unregisters it. See AbortSignal::addAbortAlgorithm() in basics.h.
16+
//
17+
// Throws a TypeError if `signal` is not an AbortSignal or `algorithm` is not a
18+
// function. The caller is expected to have checked that the signal has not
19+
// aborted: algorithms are never invoked retroactively.
20+
jsg::JsValue addAbortAlgorithmForBootstrap(
21+
jsg::Lock& js, jsg::JsValue signal, jsg::JsValue algorithm);
22+
23+
} // namespace workerd::api

‎src/workerd/api/basics-test.c++‎

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
// test without pulling in the world.
77
#define WORKERD_API_BASICS_TEST 1
88

9+
#include "abort-bootstrap.h"
910
#include "actor-state.h"
1011
#include "actor.h"
1112
#include "basics.h"
@@ -75,10 +76,25 @@ struct BasicsContext: public jsg::Object, public jsg::ContextGlobal {
7576
return true;
7677
}
7778

79+
jsg::Ref<api::AbortSignal> newSignal(jsg::Lock& js) {
80+
return js.alloc<api::AbortSignal>();
81+
}
82+
83+
void abortSignal(jsg::Lock& js, jsg::Ref<api::AbortSignal> signal) {
84+
signal->triggerAbort(js, kj::none);
85+
}
86+
87+
jsg::JsValue addAbortAlgorithm(jsg::Lock& js, jsg::JsValue signal, jsg::JsValue algorithm) {
88+
return addAbortAlgorithmForBootstrap(js, signal, algorithm);
89+
}
90+
7891
JSG_RESOURCE_TYPE(BasicsContext) {
7992
JSG_METHOD(testAbortAlgorithmsRun);
8093
JSG_METHOD(testAbortAlgorithmHandleAfterSignalGone);
8194
JSG_METHOD(testAbortAlgorithmAddedWhileAborted);
95+
JSG_METHOD(newSignal);
96+
JSG_METHOD(abortSignal);
97+
JSG_METHOD(addAbortAlgorithm);
8298
}
8399
};
84100
JSG_DECLARE_ISOLATE_TYPE(BasicsIsolate,
@@ -101,5 +117,25 @@ KJ_TEST("AbortSignal abort algorithms registered after abort never run") {
101117
e.expectEval("testAbortAlgorithmAddedWhileAborted()", "boolean", "true");
102118
}
103119

120+
KJ_TEST("Bootstrap abort algorithms run on abort unless removed") {
121+
jsg::test::Evaluator<BasicsContext, BasicsIsolate, CompatibilityFlags::Reader> e(v8System);
122+
e.expectEval("const s = newSignal(); const calls = [];"
123+
"addAbortAlgorithm(s, () => calls.push(1));"
124+
"const h = addAbortAlgorithm(s, () => calls.push(2));"
125+
"h.remove(); h.remove();"
126+
"addAbortAlgorithm(s, () => calls.push(3));"
127+
"abortSignal(s);"
128+
"calls.join(',')",
129+
"string", "1,3");
130+
}
131+
132+
KJ_TEST("Bootstrap abort algorithms reject bad arguments") {
133+
jsg::test::Evaluator<BasicsContext, BasicsIsolate, CompatibilityFlags::Reader> e(v8System);
134+
e.expectEval("addAbortAlgorithm({}, () => {})", "throws",
135+
"TypeError: addAbortAlgorithm() expects an AbortSignal");
136+
e.expectEval("addAbortAlgorithm(newSignal(), {})", "throws",
137+
"TypeError: addAbortAlgorithm() expects a function");
138+
}
139+
104140
} // namespace
105141
} // namespace workerd::api

‎src/workerd/api/basics.c++‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44

55
#include "basics.h"
66

7+
#include "abort-bootstrap.h"
78
#include "actor-state.h"
89
#include "global-scope.h"
910

@@ -884,6 +885,29 @@ void AbortSignal::removeAbortAlgorithm(uint64_t token) {
884885
}
885886
}
886887

888+
jsg::JsValue addAbortAlgorithmForBootstrap(
889+
jsg::Lock& js, jsg::JsValue signal, jsg::JsValue algorithm) {
890+
// Both types are registered in EW_BASICS_ISOLATE_TYPES, so the handlers are always
891+
// present. Asserting keeps the caller (a raw bootstrap callback) free of a failure mode
892+
// it could not act on.
893+
auto& signalHandler = KJ_ASSERT_NONNULL(js.tryGetTypeHandler<jsg::Ref<AbortSignal>>(),
894+
"AbortSignal is missing from the isolate type list");
895+
auto& handleHandler = KJ_ASSERT_NONNULL(js.tryGetTypeHandler<jsg::Ref<AbortAlgorithmHandle>>(),
896+
"AbortAlgorithmHandle is missing from the isolate type list");
897+
898+
auto target = JSG_REQUIRE_NONNULL(
899+
signalHandler.tryUnwrap(js, signal), TypeError, "addAbortAlgorithm() expects an AbortSignal");
900+
auto fn = JSG_REQUIRE_NONNULL(
901+
algorithm.tryCast<jsg::JsFunction>(), TypeError, "addAbortAlgorithm() expects a function");
902+
903+
auto registration = target->addAbortAlgorithm(js,
904+
JSG_VISITABLE_LAMBDA((fn = fn.addRef(js)), (fn), (jsg::Lock& js) {
905+
v8::LocalVector<v8::Value> args(js.v8Isolate);
906+
fn.getHandle(js).call(js, js.undefined(), args);
907+
}));
908+
return jsg::JsValue(handleHandler.wrap(js, js.alloc<AbortAlgorithmHandle>(kj::mv(registration))));
909+
}
910+
887911
kj::Own<void> AbortSignal::registerPendingCancellation(jsg::Lock& js, ReleasingCanceler& canceler) {
888912
// Capturing by reference is safe: the returned handle guarantees the action never runs
889913
// once the handle has been destroyed, and holders destroy the handle before the canceler.

0 commit comments

Comments
 (0)