Skip to content
File

Blob: src/workerd/api/streams/writable-sink-adapter-test.c++

46.4 KB
1#include "standard.h"
2#include "writable-sink-adapter.h"
3#include "writable.h"
4 
5#include <workerd/api/system-streams.h>
6#include <workerd/jsg/jsg-test.h>
7#include <workerd/jsg/jsg.h>
8#include <workerd/tests/test-fixture.h>
9#include <workerd/util/own-util.h>
10#include <workerd/util/stream-utils.h>
11 
12namespace workerd::api::streams {
13 
14namespace {
15struct SimpleEventRecordingSink final: public kj::AsyncOutputStream {
16 struct State {
17 size_t writeCalled = 0;
18 };
19 State state;
20 
21 State& getState() {
22 return state;
23 }
24 
25 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
26 state.writeCalled++;
27 return kj::READY_NOW;
28 }
29 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
30 state.writeCalled++;
31 return kj::READY_NOW;
32 }
33 
34 kj::Promise<void> whenWriteDisconnected() override {
35 return kj::NEVER_DONE;
36 }
37};
38 
39struct NeverReadySink final: public kj::AsyncOutputStream {
40 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
41 return kj::NEVER_DONE;
42 }
43 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
44 return kj::NEVER_DONE;
45 }
46 
47 kj::Promise<void> whenWriteDisconnected() override {
48 return kj::NEVER_DONE;
49 }
50};
51 
52struct ThrowingSink final: public kj::AsyncOutputStream {
53 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
54 KJ_FAIL_REQUIRE("worker_do_not_log; write() always throws");
55 co_return;
56 }
57 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
58 KJ_FAIL_REQUIRE("worker_do_not_log; write() always throws");
59 co_return;
60 }
61 
62 kj::Promise<void> whenWriteDisconnected() override {
63 return kj::NEVER_DONE;
64 }
65};
66} // namespace
67 
68KJ_TEST("Basic construction with default options") {
69 TestFixture fixture;
70 
71 fixture.runInIoContext([&](const TestFixture::Environment& env) {
72 auto sink =
73 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
74 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink));
75 
76 KJ_ASSERT(!adapter->isClosed(), "Adapter should not be closed upon construction");
77 KJ_ASSERT(!adapter->isClosing(), "Adapter should not be closing upon construction");
78 KJ_ASSERT(adapter->isErrored() == kj::none, "Adapter should not be errored upon construction");
79 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter->getDesiredSize()) == 16384,
80 "Adapter should have default highWaterMark of 16384");
81 auto& options = KJ_ASSERT_NONNULL(adapter->getOptions());
82 KJ_ASSERT(options.highWaterMark == 16384);
83 KJ_ASSERT(options.detachOnWrite == false);
84 
85 auto readyPromise = adapter->getReady(env.js);
86 KJ_ASSERT(readyPromise.getState(env.js) == jsg::Promise<void>::State::FULFILLED,
87 "Initial ready promise should be fulfilled");
88 });
89}
90 
91KJ_TEST("Construction with custom highWaterMark option") {
92 TestFixture fixture;
93 
94 fixture.runInIoContext([&](const TestFixture::Environment& env) {
95 auto sink =
96 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
97 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink),
98 WritableStreamSinkJsAdapter::Options{.highWaterMark = 100});
99 
100 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter->getDesiredSize()) == 100,
101 "Adapter should have custom highWaterMark of 100");
102 });
103}
104 
105KJ_TEST("Construction with detachOnWrite=true option") {
106 TestFixture fixture;
107 
108 fixture.runInIoContext([&](const TestFixture::Environment& env) {
109 auto sink =
110 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
111 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink),
112 WritableStreamSinkJsAdapter::Options{
113 .detachOnWrite = true,
114 });
115 auto& options = KJ_ASSERT_NONNULL(adapter->getOptions());
116 KJ_ASSERT(options.detachOnWrite == true);
117 });
118}
119 
120KJ_TEST("Construction with all custom options combined") {
121 TestFixture fixture;
122 
123 fixture.runInIoContext([&](const TestFixture::Environment& env) {
124 auto sink =
125 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
126 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink),
127 WritableStreamSinkJsAdapter::Options{
128 .highWaterMark = 100,
129 .detachOnWrite = true,
130 });
131 auto& options = KJ_ASSERT_NONNULL(adapter->getOptions());
132 KJ_ASSERT(options.highWaterMark == 100);
133 KJ_ASSERT(options.detachOnWrite == true);
134 });
135}
136 
137KJ_TEST("Basic end() operation completes successfully") {
138 TestFixture fixture;
139 
140 fixture.runInIoContext([&](const TestFixture::Environment& env) {
141 auto sink =
142 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
143 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink));
144 
145 auto endPromise = adapter->end(env.js);
146 
147 KJ_ASSERT(endPromise.getState(env.js) == jsg::Promise<void>::State::PENDING,
148 "End promise should be pending immediately after end() call");
149 KJ_ASSERT(
150 adapter->isClosed() == false, "Adapter should not be closed immediately after end() call");
151 KJ_ASSERT(adapter->isClosing() == true,
152 "Adapter should be in closing state immediately after end() call");
153 KJ_ASSERT(adapter->isErrored() == kj::none, "Adapter should not be errored after end() call");
154 
155 auto rejectedEnd = adapter->end(env.js);
156 KJ_ASSERT(rejectedEnd.getState(env.js) == jsg::Promise<void>::State::REJECTED,
157 "Second end() call should be rejected");
158 
159 auto rejectedWrite = adapter->write(env.js, env.js.str("data"_kj));
160 KJ_ASSERT(rejectedWrite.getState(env.js) == jsg::Promise<void>::State::REJECTED,
161 "Write after end() call should be rejected");
162 
163 auto rejectedFlush = adapter->flush(env.js);
164 KJ_ASSERT(rejectedFlush.getState(env.js) == jsg::Promise<void>::State::REJECTED,
165 "Flush after end() call should be rejected");
166 
167 return env.context
168 .awaitJs(env.js, endPromise.then(env.js, [&adapter = *adapter](jsg::Lock& js) {
169 KJ_ASSERT(
170 adapter.isClosed() == true, "Adapter should be closed after end() promise resolves");
171 KJ_ASSERT(adapter.isClosing() == false,
172 "Adapter should not be in closing state after end() promise resolves");
173 KJ_ASSERT(
174 adapter.isErrored() == kj::none, "Adapter should not be errored after successful end()");
175 KJ_ASSERT(adapter.getDesiredSize() == kj::none,
176 "Desired size should be none after adapter is closed");
177 
178 auto rejectedEnd = adapter.end(js);
179 KJ_ASSERT(rejectedEnd.getState(js) == jsg::Promise<void>::State::FULFILLED,
180 "Second end() call should be fulfilled");
181 
182 auto rejectedWrite = adapter.write(js, js.str("data"_kj));
183 KJ_ASSERT(rejectedWrite.getState(js) == jsg::Promise<void>::State::REJECTED,
184 "Write after end() call should be rejected");
185 
186 auto rejectedFlush = adapter.flush(js);
187 KJ_ASSERT(rejectedFlush.getState(js) == jsg::Promise<void>::State::REJECTED,
188 "Flush after end() call should be rejected");
189 })).attach(kj::mv(adapter));
190 });
191}
192 
193KJ_TEST("Basic abort() operation") {
194 TestFixture fixture;
195 
196 fixture.runInIoContext([&](const TestFixture::Environment& env) {
197 auto sink =
198 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
199 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink));
200 
201 adapter->abort(env.js, env.js.str("Abort reason"_kj));
202 
203 KJ_ASSERT(adapter->isClosed() == false, "Adapter should not be closed after abort()");
204 KJ_ASSERT(adapter->isClosing() == false, "Adapter should not be closing after abort()");
205 auto& exception =
206 KJ_ASSERT_NONNULL(adapter->isErrored(), "Adapter should be in errored state after abort()");
207 KJ_ASSERT(exception.getDescription().contains("Abort reason"),
208 "Adapter should be in errored state after abort()");
209 
210 KJ_ASSERT(adapter->getDesiredSize() == kj::none,
211 "Desired size should be none after adapter is errored");
212 
213 auto rejectedWrite = adapter->write(env.js, env.js.str("data"_kj));
214 KJ_ASSERT(rejectedWrite.getState(env.js) == jsg::Promise<void>::State::REJECTED,
215 "Write after abort() call should be rejected");
216 
217 auto rejectedFlush = adapter->flush(env.js);
218 KJ_ASSERT(rejectedFlush.getState(env.js) == jsg::Promise<void>::State::REJECTED,
219 "Flush after abort() call should be rejected");
220 
221 auto rejectedEnd = adapter->end(env.js);
222 KJ_ASSERT(rejectedEnd.getState(env.js) == jsg::Promise<void>::State::REJECTED,
223 "End after abort() call should be rejected");
224 
225 adapter->abort(env.js, env.js.str("Abort reason 2"_kj));
226 auto& exception2 = KJ_ASSERT_NONNULL(
227 adapter->isErrored(), "Adapter should still be in errored state after second abort()");
228 KJ_ASSERT(exception2.getDescription().contains("Abort reason 2"),
229 "Adapter should reflect reason from second abort() call");
230 });
231}
232 
233KJ_TEST("Abort from closing state supersedes close") {
234 TestFixture fixture;
235 
236 fixture.runInIoContext([&](const TestFixture::Environment& env) {
237 auto sink =
238 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
239 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink));
240 
241 auto endPromise = adapter->end(env.js);
242 KJ_ASSERT(endPromise.getState(env.js) == jsg::Promise<void>::State::PENDING,
243 "End promise should be pending immediately after end() call");
244 KJ_ASSERT(adapter->isClosing() == true,
245 "Adapter should be in closing state immediately after end() call");
246 
247 adapter->abort(env.js, env.js.str("Abort reason"_kj));
248 
249 KJ_ASSERT(adapter->isClosed() == false, "Adapter should not be closed after abort()");
250 KJ_ASSERT(adapter->isClosing() == false, "Adapter should not be closing after abort()");
251 auto& exception =
252 KJ_ASSERT_NONNULL(adapter->isErrored(), "Adapter should be in errored state after abort()");
253 KJ_ASSERT(exception.getDescription().contains("Abort reason"),
254 "Adapter should be in errored state after abort()");
255 
256 KJ_ASSERT(adapter->getDesiredSize() == kj::none,
257 "Desired size should be none after adapter is errored");
258 
259 auto rejectedWrite = adapter->write(env.js, env.js.str("data"_kj));
260 KJ_ASSERT(rejectedWrite.getState(env.js) == jsg::Promise<void>::State::REJECTED,
261 "Write after abort() call should be rejected");
262 
263 auto rejectedFlush = adapter->flush(env.js);
264 KJ_ASSERT(rejectedFlush.getState(env.js) == jsg::Promise<void>::State::REJECTED,
265 "Flush after abort() call should be rejected");
266 
267 auto rejectedEnd = adapter->end(env.js);
268 KJ_ASSERT(rejectedEnd.getState(env.js) == jsg::Promise<void>::State::REJECTED,
269 "End after abort() call should be rejected");
270 
271 return env.context
272 .awaitJs(env.js, endPromise.then(env.js, [](jsg::Lock& js) {
273 return js.rejectedPromise<void>(js.error("End promise should not resolve after abort()"));
274 }, [](jsg::Lock& js, jsg::Value exception) {
275 return js.resolvedPromise();
276 })).attach(kj::mv(adapter));
277 });
278}
279 
280KJ_TEST("Abort from closed state") {
281 TestFixture fixture;
282 
283 fixture.runInIoContext([&](const TestFixture::Environment& env) {
284 auto sink =
285 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
286 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink));
287 
288 auto endPromise = adapter->end(env.js);
289 return env.context.awaitJs(env.js, kj::mv(endPromise))
290 .then([adapter = kj::mv(adapter)]() mutable {
291 KJ_ASSERT(adapter->isClosed() == true, "Adapter should be closed after end()");
292 adapter->abort(KJ_EXCEPTION(FAILED, "Abort after closed should be no-op"));
293 KJ_ASSERT(adapter->isClosed() == false,
294 "Adapter switches to errored state after abort() from closed state");
295 KJ_ASSERT_NONNULL(adapter->isErrored(),
296 "Adapter should be in errored state after abort() from closed state");
297 });
298 });
299}
300 
301KJ_TEST("Abort rejects ready promise with abort reason") {
302 TestFixture fixture;
303 
304 fixture.runInIoContext([&](const TestFixture::Environment& env) {
305 auto sink =
306 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
307 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink),
308 WritableStreamSinkJsAdapter::Options{
309 .highWaterMark = 1,
310 });
311 adapter->write(env.js, env.js.str("data"_kj));
312 adapter->write(env.js, env.js.str("data2"_kj));
313 
314 auto readyPromise = adapter->getReady(env.js);
315 KJ_ASSERT(readyPromise.getState(env.js) == jsg::Promise<void>::State::PENDING,
316 "Teady promise should be pending");
317 
318 adapter->abort(env.js, env.js.str("Abort reason"_kj));
319 
320 return env.context
321 .awaitJs(env.js, readyPromise.then(env.js, [](jsg::Lock& js) {
322 return js.rejectedPromise<void>(js.error("Ready promise should not resolve after abort()"));
323 }, [](jsg::Lock& js, jsg::Value exception) {
324 auto ex = jsg::JsValue(exception.getHandle(js));
325 KJ_ASSERT(ex.toString(js).contains("Abort reason"),
326 "Ready promise should be rejected with abort reason");
327 return js.resolvedPromise();
328 })).attach(kj::mv(adapter));
329 });
330}
331 
332KJ_TEST("Abort aborts underlying sink") {
333 TestFixture fixture;
334 
335 fixture.runInIoContext([&](const TestFixture::Environment& env) {
336 SimpleEventRecordingSink sink;
337 kj::Own<SimpleEventRecordingSink> fake(&sink, kj::NullDisposer::instance);
338 auto adapter =
339 kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, newWritableSink(kj::mv(fake)));
340 adapter->abort(env.js, env.js.str("Abort reason"_kj));
341 KJ_ASSERT_NONNULL(adapter->isErrored(), "Underlying sink's abort() should have been called");
342 });
343}
344 
345KJ_TEST("Abort rejects in-flight operations") {
346 TestFixture fixture;
347 
348 fixture.runInIoContext([&](const TestFixture::Environment& env) {
349 auto neverDoneSink = newWritableSink(kj::heap<NeverReadySink>());
350 auto adapter =
351 kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(neverDoneSink));
352 
353 auto writePromise = adapter->write(env.js, env.js.str("data"_kj));
354 auto flushPromise = adapter->flush(env.js);
355 auto endPromise = adapter->end(env.js);
356 
357 adapter->abort(env.js, env.js.str("Abort reason"_kj));
358 
359 return env.context
360 .awaitJs(env.js,
361 endPromise.then(env.js,
362 [](jsg::Lock& js) {
363 return js.rejectedPromise<void>(js.error("End promise should not resolve after abort()"));
364 },
365 [writePromise = kj::mv(writePromise), flushPromise = kj::mv(flushPromise)](
366 jsg::Lock& js, jsg::Value exception) mutable {
367 KJ_ASSERT(writePromise.getState(js) == jsg::Promise<void>::State::REJECTED,
368 "Write promise should be rejected after abort()");
369 KJ_ASSERT(flushPromise.getState(js) == jsg::Promise<void>::State::REJECTED,
370 "Flush promise should be rejected after abort()");
371 
372 return js.resolvedPromise();
373 })).attach(kj::mv(adapter));
374 });
375}
376 
377KJ_TEST("end() waits for all pending writes to complete") {
378 TestFixture fixture;
379 SimpleEventRecordingSink sink;
380 
381 fixture.runInIoContext([&](const TestFixture::Environment& env) {
382 kj::Own<SimpleEventRecordingSink> fake(&sink, kj::NullDisposer::instance);
383 auto adapter =
384 kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, newWritableSink(kj::mv(fake)));
385 
386 adapter->write(env.js, env.js.str("data1"_kj));
387 adapter->write(env.js, env.js.str("data2"_kj));
388 adapter->write(env.js, env.js.str("data3"_kj));
389 adapter->write(env.js, env.js.str("data4"_kj));
390 KJ_ASSERT(sink.getState().writeCalled == 1,
391 "Underlying sink's write() should have been called four times");
392 
393 auto endPromise = adapter->end(env.js);
394 
395 return env.context
396 .awaitJs(env.js, endPromise.then(env.js, [&state = sink.getState()](jsg::Lock& js) {
397 KJ_ASSERT(state.writeCalled == 4,
398 "Underlying sink's write() should have been called four times before end() resolves");
399 return js.resolvedPromise();
400 })).attach(kj::mv(adapter));
401 });
402}
403 
404KJ_TEST("end() waits for all pending flushes to complete") {
405 TestFixture fixture;
406 SimpleEventRecordingSink sink;
407 fixture.runInIoContext([&](const TestFixture::Environment& env) {
408 kj::Own<SimpleEventRecordingSink> fake(&sink, kj::NullDisposer::instance);
409 auto adapter =
410 kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, newWritableSink(kj::mv(fake)));
411 
412 auto flush1 = adapter->flush(env.js);
413 auto flush2 = adapter->flush(env.js);
414 
415 auto endPromise = adapter->end(env.js);
416 
417 return env.context
418 .awaitJs(env.js,
419 endPromise.then(env.js,
420 [&state = sink.getState(), flush1 = kj::mv(flush1), flush2 = kj::mv(flush2)](
421 jsg::Lock& js) mutable {
422 KJ_ASSERT(flush1.getState(js) == jsg::Promise<void>::State::FULFILLED,
423 "First flush() promise should be fulfilled before end() resolves");
424 KJ_ASSERT(flush2.getState(js) == jsg::Promise<void>::State::FULFILLED,
425 "Second flush() promise should be fulfilled before end() resolves");
426 return js.resolvedPromise();
427 })).attach(kj::mv(adapter));
428 });
429}
430 
431KJ_TEST("end() with large queue of pending operations") {
432 TestFixture fixture;
433 SimpleEventRecordingSink sink;
434 
435 fixture.runInIoContext([&](const TestFixture::Environment& env) {
436 kj::Own<SimpleEventRecordingSink> fake(&sink, kj::NullDisposer::instance);
437 auto adapter =
438 kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, newWritableSink(kj::mv(fake)));
439 
440 for (int i = 0; i < 1024; i++) {
441 adapter->write(env.js, env.js.str("data"_kj));
442 adapter->flush(env.js);
443 }
444 
445 auto endPromise = adapter->end(env.js);
446 
447 return env.context
448 .awaitJs(env.js, endPromise.then(env.js, [&state = sink.getState()](jsg::Lock& js) {
449 KJ_ASSERT(state.writeCalled == 1024,
450 "Underlying sink's write() should have been called four times before end() resolves");
451 return js.resolvedPromise();
452 })).attach(kj::mv(adapter));
453 });
454}
455 
456KJ_TEST("end() when underlyink sink.end() fails should error adapter") {
457 TestFixture fixture;
458 
459 fixture.runInIoContext([&](const TestFixture::Environment& env) {
460 auto throwingSink = newWritableSink(kj::heap<ThrowingSink>());
461 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(throwingSink));
462 
463 auto writePromise = adapter->write(env.js, env.js.str("hello"_kj));
464 
465 return env.context
466 .awaitJs(env.js, writePromise.then(env.js, [](jsg::Lock& js) {
467 return js.rejectedPromise<void>(js.error("Write promise should not resolve when sink fails"));
468 }, [&adapter = *adapter](jsg::Lock& js, jsg::Value exception) {
469 auto err = jsg::JsValue(exception.getHandle(js));
470 KJ_ASSERT(err.toString(js).contains("internal error"));
471 KJ_ASSERT(adapter.isErrored() != kj::none, "Adapter should be in errored state");
472 return js.resolvedPromise();
473 })).attach(kj::mv(adapter));
474 });
475}
476 
477KJ_TEST("flush() completes after all prior writes") {
478 TestFixture fixture;
479 
480 fixture.runInIoContext([&](const TestFixture::Environment& env) {
481 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
482 auto& state = recordingSink->getState();
483 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(
484 env.js, env.context, newWritableSink(kj::mv(recordingSink)));
485 
486 adapter->write(env.js, env.js.str("data1"_kj));
487 adapter->write(env.js, env.js.str("data2"_kj));
488 KJ_ASSERT(state.writeCalled == 1, "Underlying sink's write() should have been called twice");
489 
490 auto flushPromise = adapter->flush(env.js);
491 
492 return env.context
493 .awaitJs(env.js, flushPromise.then(env.js, [&state](jsg::Lock& js) {
494 KJ_ASSERT(state.writeCalled == 2,
495 "Underlying sink's write() should have been called twice before flush() resolves");
496 return js.resolvedPromise();
497 })).attach(kj::mv(state), kj::mv(adapter));
498 });
499}
500 
501KJ_TEST("flush() with no writes completes immediately") {
502 TestFixture fixture;
503 
504 fixture.runInIoContext([&](const TestFixture::Environment& env) {
505 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
506 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(
507 env.js, env.context, newWritableSink(kj::mv(recordingSink)));
508 
509 auto flushPromise = adapter->flush(env.js);
510 
511 return env.context.awaitJs(env.js, kj::mv(flushPromise)).attach(kj::mv(adapter));
512 });
513}
514 
515KJ_TEST("multiple sequential flush() calls") {
516 TestFixture fixture;
517 
518 fixture.runInIoContext([&](const TestFixture::Environment& env) {
519 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
520 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(
521 env.js, env.context, newWritableSink(kj::mv(recordingSink)));
522 
523 auto flush1 = adapter->flush(env.js);
524 auto flush2 = adapter->flush(env.js);
525 
526 return env.context
527 .awaitJs(env.js, flush1.then(env.js, [flush2 = kj::mv(flush2)](jsg::Lock& js) mutable {
528 return kj::mv(flush2);
529 })).attach(kj::mv(adapter));
530 });
531}
532 
533KJ_TEST("write() when underlyink sink.write() fails should error adapter") {
534 TestFixture fixture;
535 
536 fixture.runInIoContext([&](const TestFixture::Environment& env) {
537 auto throwingSink = newWritableSink(kj::heap<ThrowingSink>());
538 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(throwingSink));
539 
540 auto writePromise = adapter->write(env.js, env.js.str("data"_kj));
541 auto flushPromise = adapter->flush(env.js);
542 
543 return env.context
544 .awaitJs(env.js,
545 writePromise.then(env.js,
546 [](jsg::Lock& js) {
547 return js.rejectedPromise<void>(js.error("Write promise should not resolve when sink fails"));
548 },
549 [&adapter = *adapter, flushPromise = kj::mv(flushPromise)](
550 jsg::Lock& js, jsg::Value exception) mutable {
551 auto err = jsg::JsValue(exception.getHandle(js));
552 KJ_ASSERT(err.toString(js).contains("internal error"));
553 KJ_ASSERT(adapter.isErrored() != kj::none, "Adapter should be in errored state");
554 return flushPromise.then(js, [](jsg::Lock& js) {
555 return js.rejectedPromise<void>(
556 js.error("Flush promise should not resolve when sink fails"));
557 }, [](jsg::Lock& js, jsg::Value exception) { return js.resolvedPromise(); });
558 })).attach(kj::mv(adapter));
559 });
560}
561 
562KJ_TEST("multiple writes() should only error adapter once") {
563 TestFixture fixture;
564 
565 fixture.runInIoContext([&](const TestFixture::Environment& env) {
566 auto throwingSink = newWritableSink(kj::heap<ThrowingSink>());
567 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(throwingSink));
568 
569 auto write1 = adapter->write(env.js, env.js.str("data"_kj));
570 auto write2 = adapter->write(env.js, env.js.str("data"_kj));
571 
572 return env.context
573 .awaitJs(env.js, write1.then(env.js, [](jsg::Lock& js) {
574 return js.rejectedPromise<void>(js.error("Write promise should not resolve when sink fails"));
575 }, [&adapter = *adapter, write2 = kj::mv(write2)](jsg::Lock& js, jsg::Value exception) mutable {
576 auto err = jsg::JsValue(exception.getHandle(js));
577 KJ_ASSERT(err.toString(js).contains("internal error"));
578 KJ_ASSERT(adapter.isErrored() != kj::none, "Adapter should be in errored state");
579 return write2.then(js, [](jsg::Lock& js) {
580 return js.rejectedPromise<void>(
581 js.error("Write promise should not resolve when sink fails"));
582 }, [](jsg::Lock& js, jsg::Value exception) { return js.resolvedPromise(); });
583 })).attach(kj::mv(adapter));
584 });
585}
586 
587KJ_TEST("zero-length writes are a non-op (string)") {
588 TestFixture fixture;
589 
590 fixture.runInIoContext([&](const TestFixture::Environment& env) {
591 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
592 auto& state = recordingSink->getState();
593 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(
594 env.js, env.context, newWritableSink(kj::mv(recordingSink)));
595 
596 auto writePromise = adapter->write(env.js, env.js.str(""_kj));
597 KJ_ASSERT(state.writeCalled == 0, "Underlying sink's write() should not have been called");
598 
599 return env.context
600 .awaitJs(env.js, writePromise.then(env.js, [&state](jsg::Lock& js) {
601 KJ_ASSERT(state.writeCalled == 0, "Underlying sink's write() should not have been called");
602 })).attach(kj::mv(adapter));
603 });
604}
605 
606KJ_TEST("zero-length writes are a non-op (ArrayBuffer)") {
607 TestFixture fixture;
608 
609 fixture.runInIoContext([&](const TestFixture::Environment& env) {
610 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
611 auto& state = recordingSink->getState();
612 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(
613 env.js, env.context, newWritableSink(kj::mv(recordingSink)));
614 
615 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, 0);
616 jsg::BufferSource source(env.js, kj::mv(backing));
617 jsg::JsValue handle(source.getHandle(env.js));
618 
619 auto writePromise = adapter->write(env.js, handle);
620 KJ_ASSERT(state.writeCalled == 0, "Underlying sink's write() should not have been called");
621 
622 return env.context
623 .awaitJs(env.js, writePromise.then(env.js, [&state](jsg::Lock& js) {
624 KJ_ASSERT(state.writeCalled == 0, "Underlying sink's write() should not have been called");
625 })).attach(kj::mv(adapter));
626 });
627}
628 
629KJ_TEST("writing small ArrayBuffer") {
630 TestFixture fixture;
631 
632 fixture.runInIoContext([&](const TestFixture::Environment& env) {
633 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
634 auto& state = recordingSink->getState();
635 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context,
636 newWritableSink(kj::mv(recordingSink)),
637 WritableStreamSinkJsAdapter::Options{
638 .highWaterMark = 10,
639 });
640 
641 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, 10);
642 jsg::BufferSource source(env.js, kj::mv(backing));
643 jsg::JsValue handle(source.getHandle(env.js));
644 
645 auto writePromise = adapter->write(env.js, handle);
646 KJ_ASSERT(state.writeCalled == 1, "Underlying sink's write() should not have been called");
647 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter->getDesiredSize()) == 0,
648 "Adapter's desired size should be 0 after writing highWaterMark bytes");
649 
650 return env.context
651 .awaitJs(env.js, writePromise.then(env.js, [&state, &adapter = *adapter](jsg::Lock& js) {
652 KJ_ASSERT(state.writeCalled == 1, "Underlying sink's write() should not have been called");
653 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter.getDesiredSize()) == 10,
654 "Back to initial desired size after write completes");
655 })).attach(kj::mv(adapter));
656 });
657}
658 
659KJ_TEST("writing medium ArrayBuffer") {
660 TestFixture fixture;
661 
662 fixture.runInIoContext([&](const TestFixture::Environment& env) {
663 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
664 auto& state = recordingSink->getState();
665 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context,
666 newWritableSink(kj::mv(recordingSink)),
667 WritableStreamSinkJsAdapter::Options{
668 .highWaterMark = 5 * 1024,
669 });
670 
671 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, 4 * 1024);
672 jsg::BufferSource source(env.js, kj::mv(backing));
673 jsg::JsValue handle(source.getHandle(env.js));
674 
675 auto writePromise = adapter->write(env.js, handle);
676 KJ_ASSERT(state.writeCalled == 1, "Underlying sink's write() should not have been called");
677 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter->getDesiredSize()) == 1024,
678 "Adapter's desired size should be 1024 after writing 4 * 1024 bytes");
679 
680 return env.context
681 .awaitJs(env.js, writePromise.then(env.js, [&state, &adapter = *adapter](jsg::Lock& js) {
682 KJ_ASSERT(state.writeCalled == 1, "Underlying sink's write() should not have been called");
683 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter.getDesiredSize()) == 5 * 1024,
684 "Back to initial desired size after write completes");
685 })).attach(kj::mv(adapter));
686 });
687}
688 
689KJ_TEST("writing large ArrayBuffer") {
690 TestFixture fixture;
691 
692 fixture.runInIoContext([&](const TestFixture::Environment& env) {
693 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
694 auto& state = recordingSink->getState();
695 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context,
696 newWritableSink(kj::mv(recordingSink)),
697 WritableStreamSinkJsAdapter::Options{
698 .highWaterMark = 8 * 1024,
699 });
700 
701 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, 16 * 1024);
702 jsg::BufferSource source(env.js, kj::mv(backing));
703 jsg::JsValue handle(source.getHandle(env.js));
704 
705 auto writePromise = adapter->write(env.js, handle);
706 KJ_ASSERT(state.writeCalled == 1, "Underlying sink's write() should not have been called");
707 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter->getDesiredSize()) == -(8 * 1024),
708 "Adapter's desired size should be negative after writing 16 * 1024 bytes");
709 
710 return env.context
711 .awaitJs(env.js, writePromise.then(env.js, [&state, &adapter = *adapter](jsg::Lock& js) {
712 KJ_ASSERT(state.writeCalled == 1, "Underlying sink's write() should not have been called");
713 KJ_ASSERT(KJ_ASSERT_NONNULL(adapter.getDesiredSize()) == 8 * 1024,
714 "Back to initial desired size after write completes");
715 })).attach(kj::mv(adapter));
716 });
717}
718 
719KJ_TEST("writing the wrong types reject") {
720 TestFixture fixture;
721 
722 fixture.runInIoContext([&](const TestFixture::Environment& env) {
723 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
724 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(
725 env.js, env.context, newWritableSink(kj::mv(recordingSink)));
726 
727 auto writeNull = adapter->write(env.js, env.js.null());
728 KJ_ASSERT(writeNull.getState(env.js) == jsg::Promise<void>::State::REJECTED,
729 "Write of null should be rejected");
730 
731 auto writeUndefined = adapter->write(env.js, env.js.undefined());
732 KJ_ASSERT(writeUndefined.getState(env.js) == jsg::Promise<void>::State::REJECTED,
733 "Write of undefined should be rejected");
734 
735 auto writeNumber = adapter->write(env.js, env.js.num(42));
736 KJ_ASSERT(writeNumber.getState(env.js) == jsg::Promise<void>::State::REJECTED,
737 "Write of number should be rejected");
738 
739 auto writeBoolean = adapter->write(env.js, env.js.boolean(true));
740 KJ_ASSERT(writeBoolean.getState(env.js) == jsg::Promise<void>::State::REJECTED,
741 "Write of boolean should be rejected");
742 
743 auto writeObject = adapter->write(env.js, env.js.obj());
744 KJ_ASSERT(writeObject.getState(env.js) == jsg::Promise<void>::State::REJECTED,
745 "Write of plain object should be rejected");
746 });
747}
748 
749KJ_TEST("large number of large writes") {
750 TestFixture fixture;
751 SimpleEventRecordingSink sink;
752 
753 fixture.runInIoContext([&](const TestFixture::Environment& env) {
754 kj::Own<SimpleEventRecordingSink> fake(&sink, kj::NullDisposer::instance);
755 auto adapter =
756 kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, newWritableSink(kj::mv(fake)));
757 
758 for (int i = 0; i < 1000; i++) {
759 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, 16 * 1024);
760 jsg::BufferSource source(env.js, kj::mv(backing));
761 jsg::JsValue handle(source.getHandle(env.js));
762 
763 adapter->write(env.js, handle);
764 }
765 auto endPromise = adapter->end(env.js);
766 
767 return env.context
768 .awaitJs(env.js,
769 endPromise.then(env.js, [&state = sink.getState(), &adapter = *adapter](jsg::Lock& js) {
770 KJ_ASSERT(state.writeCalled == 1000, "Underlying sink's write() should have been called");
771 })).attach(kj::mv(adapter));
772 });
773}
774 
775KJ_TEST("ready promise signals backpressure correctly") {
776 TestFixture fixture;
777 
778 fixture.runInIoContext([&](const TestFixture::Environment& env) {
779 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
780 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context,
781 newWritableSink(kj::mv(recordingSink)),
782 WritableStreamSinkJsAdapter::Options{
783 .highWaterMark = 10,
784 });
785 
786 auto readyPromise = adapter->getReady(env.js);
787 KJ_ASSERT(readyPromise.getState(env.js) == jsg::Promise<void>::State::FULFILLED,
788 "Ready promise should be fulfilled when no backpressure");
789 
790 auto writePromise = adapter->write(env.js, env.js.str("12345678909876543210"_kj));
791 
792 readyPromise = adapter->getReady(env.js);
793 KJ_ASSERT(readyPromise.getState(env.js) == jsg::Promise<void>::State::PENDING,
794 "Ready promise should be fulfilled when no backpressure");
795 
796 return env.context
797 .awaitJs(env.js, writePromise.then(env.js, [&adapter = *adapter](jsg::Lock& js) {
798 auto readyPromise = adapter.getReady(js);
799 KJ_ASSERT(readyPromise.getState(js) == jsg::Promise<void>::State::FULFILLED,
800 "Ready promise should be fulfilled when no backpressure");
801 })).attach(kj::mv(adapter));
802 });
803}
804 
805KJ_TEST("detachOnWrite option detaches ArrayBuffer before write") {
806 TestFixture fixture;
807 
808 fixture.runInIoContext([&](const TestFixture::Environment& env) {
809 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
810 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context,
811 newWritableSink(kj::mv(recordingSink)),
812 WritableStreamSinkJsAdapter::Options{
813 .detachOnWrite = true,
814 });
815 
816 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(env.js, 10);
817 jsg::BufferSource source(env.js, kj::mv(backing));
818 KJ_ASSERT(!source.isDetached());
819 jsg::JsValue handle(source.getHandle(env.js));
820 
821 auto writePromise = adapter->write(env.js, handle);
822 
823 jsg::BufferSource source2(env.js, handle);
824 KJ_ASSERT(source2.size() == 0);
825 
826 return env.context.awaitJs(env.js, kj::mv(writePromise)).attach(kj::mv(adapter));
827 });
828}
829 
830KJ_TEST("detachOnWrite option detaches Uint8Array before write") {
831 TestFixture fixture;
832 
833 fixture.runInIoContext([&](const TestFixture::Environment& env) {
834 auto recordingSink = kj::heap<SimpleEventRecordingSink>();
835 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context,
836 newWritableSink(kj::mv(recordingSink)),
837 WritableStreamSinkJsAdapter::Options{
838 .detachOnWrite = true,
839 });
840 
841 auto backing = jsg::BackingStore::alloc<v8::Uint8Array>(env.js, 10);
842 jsg::BufferSource source(env.js, kj::mv(backing));
843 KJ_ASSERT(!source.isDetached());
844 jsg::JsValue handle(source.getHandle(env.js));
845 
846 auto writePromise = adapter->write(env.js, handle);
847 
848 jsg::BufferSource source2(env.js, handle);
849 KJ_ASSERT(source2.size() == 0);
850 
851 return env.context.awaitJs(env.js, kj::mv(writePromise)).attach(kj::mv(adapter));
852 });
853}
854 
855KJ_TEST("Creating adapter and dropping it with pending operations") {
856 TestFixture fixture;
857 
858 fixture.runInIoContext([&](const TestFixture::Environment& env) {
859 auto sink =
860 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
861 auto adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink));
862 
863 adapter->write(env.js, env.js.str("data"_kj));
864 adapter->flush(env.js);
865 adapter->end(env.js);
866 
867 // Dropping the adapter here should not crash or leak memory.
868 });
869}
870 
871KJ_TEST("Dropping the IoContext with pending operations and using the adapter in another context") {
872 TestFixture fixture;
873 kj::Maybe<kj::Own<WritableStreamSinkJsAdapter>> adapter;
874 
875 fixture.runInIoContext([&](const TestFixture::Environment& env) {
876 auto sink =
877 newIoContextWrappedWritableSink(env.context, newWritableSink(newNullOutputStream()));
878 adapter = kj::heap<WritableStreamSinkJsAdapter>(env.js, env.context, kj::mv(sink));
879 auto& adapterRef = *KJ_ASSERT_NONNULL(adapter);
880 
881 adapterRef.write(env.js, env.js.str("data"_kj));
882 adapterRef.flush(env.js);
883 adapterRef.end(env.js);
884 
885 // Dropping the IoContext here should not crash or leak memory.
886 });
887 
888 fixture.runInIoContext([&](const TestFixture::Environment& env) {
889 auto& otherContext = KJ_ASSERT_NONNULL(adapter);
890 try {
891 otherContext->write(env.js, env.js.str("data2"_kj));
892 } catch (...) {
893 auto ex = kj::getCaughtExceptionAsKj();
894 KJ_ASSERT(ex.getDescription().startsWith(
895 "jsg.Error: Cannot perform I/O on behalf of a different request."));
896 }
897 });
898}
899 
900// ================================================================================================
901 
902namespace {
903struct WritableStreamContext {
904 kj::Vector<kj::Array<const kj::byte>> chunks;
905 bool closed = false;
906 kj::Maybe<jsg::JsRef<jsg::JsValue>> maybeAbort;
907};
908 
909jsg::Ref<WritableStream> createSimpleWritableStream(jsg::Lock& js, WritableStreamContext& context) {
910 return WritableStream::constructor(js,
911 UnderlyingSink{
912 .write =
913 [&context](jsg::Lock& js, auto chunk, auto) {
914 jsg::BufferSource source(js, chunk);
915 auto data = kj::heapArray<kj::byte>(source.asArrayPtr());
916 context.chunks.add(kj::mv(data));
917 return js.resolvedPromise();
918 },
919 .abort =
920 [&context](jsg::Lock& js, auto reason) {
921 context.maybeAbort = jsg::JsRef<jsg::JsValue>(js, jsg::JsValue(reason));
922 return js.resolvedPromise();
923 },
924 .close =
925 [&context](jsg::Lock& js) {
926 context.closed = true;
927 return js.resolvedPromise();
928 },
929 },
930 StreamQueuingStrategy{});
931}
932 
933jsg::Ref<WritableStream> createErroredStream(jsg::Lock& js) {
934 return WritableStream::constructor(js,
935 UnderlyingSink{
936 .write = [](jsg::Lock& js, auto chunk,
937 auto) { return js.rejectedPromise<void>(js.error("Write error")); },
938 .abort = [](jsg::Lock& js, auto reason) { return js.resolvedPromise(); },
939 .close = [](jsg::Lock& js) { return js.resolvedPromise(); },
940 },
941 StreamQueuingStrategy{});
942}
943 
944struct FiniteReadableStreamSource final: public ReadableStreamSource {
945 int counter = 0;
946 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override {
947 if (++counter < 5) {
948 kj::ArrayPtr<kj::byte> buf(static_cast<kj::byte*>(buffer), maxBytes);
949 buf.fill('a');
950 return maxBytes;
951 }
952 static constexpr size_t eof = 0;
953 return eof;
954 }
955};
956 
957struct ErroringStreamSource final: public ReadableStreamSource {
958 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override {
959 return KJ_EXCEPTION(FAILED, "worker_do_not_log: Read error");
960 }
961};
962} // namespace
963 
964KJ_TEST("WritableStreamSinkKjAdapter construction") {
965 capnp::MallocMessageBuilder message;
966 auto flags = message.initRoot<CompatibilityFlags>();
967 flags.setStreamsJavaScriptControllers(true);
968 TestFixture fixture({.featureFlags = flags.asReader()});
969 WritableStreamContext context;
970 
971 fixture.runInIoContext([&](const TestFixture::Environment& env) {
972 auto stream = createSimpleWritableStream(env.js, context);
973 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
974 });
975}
976 
977KJ_TEST("WritableStreamSinkKjAdapter construction with locked stream") {
978 capnp::MallocMessageBuilder message;
979 auto flags = message.initRoot<CompatibilityFlags>();
980 flags.setStreamsJavaScriptControllers(true);
981 TestFixture fixture({.featureFlags = flags.asReader()});
982 WritableStreamContext context;
983 
984 fixture.runInIoContext([&](const TestFixture::Environment& env) {
985 auto stream = createSimpleWritableStream(env.js, context);
986 auto writer = stream->getWriter(env.js);
987 
988 try {
989 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
990 KJ_FAIL_ASSERT("Construction with locked stream should have thrown");
991 } catch (...) {
992 auto ex = kj::getCaughtExceptionAsKj();
993 KJ_ASSERT(ex.getDescription().contains("WritableStream is locked"));
994 }
995 });
996}
997 
998KJ_TEST("WritableStreamSinkKjAdapter construction with closed stream") {
999 capnp::MallocMessageBuilder message;
1000 auto flags = message.initRoot<CompatibilityFlags>();
1001 flags.setStreamsJavaScriptControllers(true);
1002 TestFixture fixture({.featureFlags = flags.asReader()});
1003 WritableStreamContext context;
1004 
1005 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1006 auto stream = createSimpleWritableStream(env.js, context);
1007 stream->close(env.js);
1008 
1009 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
1010 });
1011}
1012 
1013KJ_TEST("WritableStreamSinkKjAdapter construction with errored stream") {
1014 capnp::MallocMessageBuilder message;
1015 auto flags = message.initRoot<CompatibilityFlags>();
1016 flags.setStreamsJavaScriptControllers(true);
1017 TestFixture fixture({.featureFlags = flags.asReader()});
1018 WritableStreamContext context;
1019 
1020 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1021 auto stream = createSimpleWritableStream(env.js, context);
1022 stream->abort(env.js, env.js.str("Abort reason"_kj));
1023 
1024 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
1025 });
1026}
1027 
1028KJ_TEST("WritableStreamSinkKjAdapter construction with immediate end") {
1029 capnp::MallocMessageBuilder message;
1030 auto flags = message.initRoot<CompatibilityFlags>();
1031 flags.setStreamsJavaScriptControllers(true);
1032 TestFixture fixture({.featureFlags = flags.asReader()});
1033 WritableStreamContext context;
1034 
1035 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1036 auto stream = createSimpleWritableStream(env.js, context);
1037 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
1038 return adapter->end().attach(kj::mv(adapter));
1039 });
1040}
1041 
1042KJ_TEST("WritableStreamSinkKjAdapter construction with immediate abort") {
1043 capnp::MallocMessageBuilder message;
1044 auto flags = message.initRoot<CompatibilityFlags>();
1045 flags.setStreamsJavaScriptControllers(true);
1046 TestFixture fixture({.featureFlags = flags.asReader()});
1047 WritableStreamContext context;
1048 
1049 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1050 auto stream = createSimpleWritableStream(env.js, context);
1051 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
1052 adapter->abort(KJ_EXCEPTION(DISCONNECTED, "Abort reason"));
1053 });
1054}
1055 
1056KJ_TEST("WritableStreamSinkKjAdapter single write") {
1057 capnp::MallocMessageBuilder message;
1058 auto flags = message.initRoot<CompatibilityFlags>();
1059 flags.setStreamsJavaScriptControllers(true);
1060 TestFixture fixture({.featureFlags = flags.asReader()});
1061 WritableStreamContext context;
1062 kj::FixedArray<kj::byte, 1024> buffer;
1063 
1064 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1065 auto stream = createSimpleWritableStream(env.js, context);
1066 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
1067 
1068 buffer.fill('a');
1069 return adapter->write(buffer.asPtr()).then([&adapter = *adapter]() {
1070 return adapter.end();
1071 }).attach(kj::mv(adapter));
1072 });
1073 
1074 KJ_ASSERT(context.chunks.size() == 1, "Underlying stream should have received one chunk");
1075 KJ_ASSERT(context.chunks[0].size() == 1024, "Underlying stream chunk should be 1024 bytes");
1076 KJ_ASSERT(context.chunks[0] == buffer, "Underlying stream chunk should match written data");
1077 KJ_ASSERT(context.closed, "Underlying stream should be closed");
1078 KJ_ASSERT(context.maybeAbort == kj::none, "Underlying stream should not be aborted");
1079}
1080 
1081KJ_TEST("WritableStreamSinkKjAdapter zero-length write") {
1082 capnp::MallocMessageBuilder message;
1083 auto flags = message.initRoot<CompatibilityFlags>();
1084 flags.setStreamsJavaScriptControllers(true);
1085 TestFixture fixture({.featureFlags = flags.asReader()});
1086 WritableStreamContext context;
1087 kj::ArrayPtr<kj::byte> buffer = nullptr;
1088 
1089 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1090 auto stream = createSimpleWritableStream(env.js, context);
1091 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
1092 
1093 return adapter->write(buffer).then([&adapter = *adapter]() {
1094 return adapter.end();
1095 }).attach(kj::mv(adapter));
1096 });
1097 
1098 KJ_ASSERT(context.chunks.size() == 0, "Underlying stream should not have chunks");
1099 KJ_ASSERT(context.closed, "Underlying stream should be closed");
1100 KJ_ASSERT(context.maybeAbort == kj::none, "Underlying stream should not be aborted");
1101}
1102 
1103KJ_TEST("WritableStreamSinkKjAdapter concurrent writes forbidden") {
1104 capnp::MallocMessageBuilder message;
1105 auto flags = message.initRoot<CompatibilityFlags>();
1106 flags.setStreamsJavaScriptControllers(true);
1107 TestFixture fixture({.featureFlags = flags.asReader()});
1108 WritableStreamContext context;
1109 kj::FixedArray<kj::byte, 100> buffer;
1110 
1111 try {
1112 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1113 auto stream = createSimpleWritableStream(env.js, context);
1114 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
1115 
1116 buffer.asPtr().fill('a');
1117 
1118 auto p1 = adapter->write(buffer.asPtr());
1119 // The second one should fail.
1120 
1121 return adapter->write(buffer.asPtr()).attach(kj::mv(adapter));
1122 });
1123 } catch (...) {
1124 auto ex = kj::getCaughtExceptionAsKj();
1125 KJ_ASSERT(ex.getDescription().contains("Cannot have multiple concurrent writes"));
1126 }
1127}
1128 
1129KJ_TEST("WritableStreamSinkKjAdapter write after close") {
1130 capnp::MallocMessageBuilder message;
1131 auto flags = message.initRoot<CompatibilityFlags>();
1132 flags.setStreamsJavaScriptControllers(true);
1133 TestFixture fixture({.featureFlags = flags.asReader()});
1134 WritableStreamContext context;
1135 kj::FixedArray<kj::byte, 100> buffer;
1136 
1137 try {
1138 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1139 auto stream = createSimpleWritableStream(env.js, context);
1140 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
1141 
1142 buffer.asPtr().fill('a');
1143 
1144 return adapter->end()
1145 .then([&adapter = *adapter, &buffer]() {
1146 return adapter.write(buffer.asPtr());
1147 }).attach(kj::mv(adapter));
1148 });
1149 } catch (...) {
1150 auto ex = kj::getCaughtExceptionAsKj();
1151 KJ_ASSERT(ex.getDescription().contains("Cannot write after close"));
1152 }
1153}
1154 
1155KJ_TEST("WritableStreamSinkKjAdapter single errored") {
1156 capnp::MallocMessageBuilder message;
1157 auto flags = message.initRoot<CompatibilityFlags>();
1158 flags.setStreamsJavaScriptControllers(true);
1159 TestFixture fixture({.featureFlags = flags.asReader()});
1160 kj::FixedArray<kj::byte, 1024> buffer;
1161 
1162 fixture.runInIoContext([&](const TestFixture::Environment& env) {
1163 auto stream = createErroredStream(env.js);
1164 auto adapter = kj::heap<WritableStreamSinkKjAdapter>(env.js, env.context, kj::mv(stream));
1165 
1166 buffer.fill('a');
1167 return adapter->write(buffer.asPtr())
1168 .then([&adapter = *adapter]() { KJ_FAIL_ASSERT("Write should have failed"); },
1169 [](kj::Exception exception) {
1170 KJ_ASSERT(exception.getDescription().contains("Write error"),
1171 "Write should have failed with underlying stream error");
1172 }).attach(kj::mv(adapter));
1173 });
1174}
1175} // namespace workerd::api::streams