File
Blob: src/workerd/api/streams/writable-sink-adapter-test.c++
| 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 | |
| 12 | namespace workerd::api::streams { |
| 13 | |
| 14 | namespace { |
| 15 | struct 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 | |
| 39 | struct 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 | |
| 52 | struct 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 | |
| 68 | KJ_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 | |
| 91 | KJ_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 | |
| 105 | KJ_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 | |
| 120 | KJ_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 | |
| 137 | KJ_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 | |
| 193 | KJ_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 | |
| 233 | KJ_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 | |
| 280 | KJ_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 | |
| 301 | KJ_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 | |
| 332 | KJ_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 | |
| 345 | KJ_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 | |
| 377 | KJ_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 | |
| 404 | KJ_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 | |
| 431 | KJ_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 | |
| 456 | KJ_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 | |
| 477 | KJ_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 | |
| 501 | KJ_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 | |
| 515 | KJ_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 | |
| 533 | KJ_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 | |
| 562 | KJ_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 | |
| 587 | KJ_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 | |
| 606 | KJ_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 | |
| 629 | KJ_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 | |
| 659 | KJ_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 | |
| 689 | KJ_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 | |
| 719 | KJ_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 | |
| 749 | KJ_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 | |
| 775 | KJ_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 | |
| 805 | KJ_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 | |
| 830 | KJ_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 | |
| 855 | KJ_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 | |
| 871 | KJ_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 | |
| 902 | namespace { |
| 903 | struct WritableStreamContext { |
| 904 | kj::Vector<kj::Array<const kj::byte>> chunks; |
| 905 | bool closed = false; |
| 906 | kj::Maybe<jsg::JsRef<jsg::JsValue>> maybeAbort; |
| 907 | }; |
| 908 | |
| 909 | jsg::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 | |
| 933 | jsg::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 | |
| 944 | struct 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 | |
| 957 | struct 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 | |
| 964 | KJ_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 | |
| 977 | KJ_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 | |
| 998 | KJ_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 | |
| 1013 | KJ_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 | |
| 1028 | KJ_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 | |
| 1042 | KJ_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 | |
| 1056 | KJ_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 | |
| 1081 | KJ_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 | |
| 1103 | KJ_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 | |
| 1129 | KJ_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 | |
| 1155 | KJ_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 |