Skip to content
File

Blob: src/workerd/api/streams/identitytransformstream-backpressure-test.js

javascript143 lines
1// Copyright (c) 2023 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 
5import { notStrictEqual, strictEqual } from 'node:assert';
6 
7export const identityTransformStream = {
8 async test(ctrl, env, ctx) {
9 const ts = new IdentityTransformStream({ highWaterMark: 10 });
10 const writer = ts.writable.getWriter();
11 const reader = ts.readable.getReader();
12 
13 strictEqual(writer.desiredSize, 10);
14 
15 // We shouldn't have to wait here.
16 const firstReady = writer.ready;
17 await writer.ready;
18 
19 writer.write(new Uint8Array(1));
20 strictEqual(writer.desiredSize, 9);
21 
22 // Let's write a second chunk that will be buffered. This one
23 // should impact the desiredSize and the backpressure signal.
24 writer.write(new Uint8Array(9));
25 strictEqual(writer.desiredSize, 0);
26 
27 // The ready promise should have been replaced
28 notStrictEqual(firstReady, writer.ready);
29 
30 async function waitForReady() {
31 strictEqual(writer.desiredSize, 0);
32 await writer.ready;
33 // The backpressure should have been relieved a bit,
34 // but only by the amount of what we've currently read.
35 strictEqual(writer.desiredSize, 1);
36 }
37 
38 await Promise.all([
39 // We call the waitForReady first to ensure that we set up waiting on
40 // the ready promise before we relieve the backpressure using the read.
41 // If the backpressure signal is not working correctly, the test will
42 // fail with an error indicating that a hanging promise was canceled.
43 waitForReady(),
44 reader.read(),
45 ]);
46 
47 // If we read again, the backpressure should be fully resolved.
48 await reader.read();
49 strictEqual(writer.desiredSize, 10);
50 },
51};
52 
53export const identityTransformStreamNoHWM = {
54 async test(ctrl, env, ctx) {
55 // Test that the original default behavior still works as expected.
56 
57 const ts = new IdentityTransformStream();
58 const writer = ts.writable.getWriter();
59 const reader = ts.readable.getReader();
60 
61 strictEqual(writer.desiredSize, 1);
62 
63 // We shouldn't have to wait here.
64 const firstReady = writer.ready;
65 await writer.ready;
66 
67 writer.write(new Uint8Array(1));
68 strictEqual(writer.desiredSize, 1);
69 
70 // Let's write a second chunk that will be buffered. There should
71 // be no indication that the desired size has changed.
72 writer.write(new Uint8Array(9));
73 strictEqual(writer.desiredSize, 1);
74 
75 // The ready promise should be exactly the same...
76 strictEqual(firstReady, writer.ready);
77 
78 async function waitForReady() {
79 strictEqual(writer.desiredSize, 1);
80 await writer.ready;
81 strictEqual(writer.desiredSize, 1);
82 }
83 
84 await Promise.all([
85 // We call the waitForReady first to ensure that we set up waiting on
86 // the ready promise before we relieve the backpressure using the read.
87 // If the backpressure signal is not working correctly, the test will
88 // fail with an error indicating that a hanging promise was canceled.
89 waitForReady(),
90 reader.read(),
91 ]);
92 
93 // If we read again, the backpressure should be fully resolved.
94 await reader.read();
95 strictEqual(writer.desiredSize, 1);
96 },
97};
98 
99export const fixedLengthStream = {
100 async test(ctrl, env, ctx) {
101 const ts = new FixedLengthStream(10, { highWaterMark: 100 });
102 const writer = ts.writable.getWriter();
103 const reader = ts.readable.getReader();
104 
105 // Even tho we specified 100 as our highWaterMark, we only expect 10
106 // bytes total, so we'll make that our highWaterMark instead.
107 strictEqual(writer.desiredSize, 10);
108 
109 // We shouldn't have to wait here.
110 const firstReady = writer.ready;
111 await writer.ready;
112 
113 writer.write(new Uint8Array(1));
114 strictEqual(writer.desiredSize, 9);
115 
116 // Let's write a second chunk that will be buffered. This one
117 // should impact the desiredSize and the backpressure signal.
118 writer.write(new Uint8Array(9));
119 strictEqual(writer.desiredSize, 0);
120 
121 // The ready promise should have been replaced
122 notStrictEqual(firstReady, writer.ready);
123 
124 async function waitForReady() {
125 await writer.ready;
126 // The backpressure should have been relieved a bit,
127 // but only by the amount of what we've currently read.
128 strictEqual(writer.desiredSize, 1);
129 }
130 
131 await Promise.all([
132 // We call the waitForReady first to ensure that we set up waiting on
133 // the ready promise before we relieve the backpressure using the read.
134 waitForReady(),
135 reader.read(),
136 ]);
137 
138 // If we read again, the backpressure should be fully resolved.
139 await reader.read();
140 strictEqual(writer.desiredSize, 10);
141 },
142};