| 1366 | } |
| 1367 | |
| 1368 | void AsyncPipeline::asyncWriteWritable(AsyncBufferView::ID bufferID, AsyncReadableStream& readable, |
| 1369 | AsyncWritableStream& writable) |
| 1370 | { |
| 1371 | source->getBuffersPool().refBuffer(bufferID); |
| 1372 | if (not writable.canAcceptWrite()) |
| 1373 | { |
| 1374 | PendingWrite* pendingWrite = findPendingWrite(readable, writable); |
| 1375 | SC_ASYNC_STREAMS_ASSERT_RELEASE(pendingWrite != nullptr); |
| 1376 | SC_ASYNC_STREAMS_ASSERT_RELEASE(not pendingWrite->bufferID.isValid()); |
| 1377 | const bool wasAlreadyBackpressured = hasPendingWritesForReadable(readable); |
| 1378 | pendingWrite->readable = &readable; |
| 1379 | pendingWrite->writable = &writable; |
| 1380 | pendingWrite->bufferID = bufferID; |
| 1381 | if (not wasAlreadyBackpressured) |
| 1382 | { |
| 1383 | readable.pause(); |
| 1384 | } |
| 1385 | return; |
| 1386 | } |
| 1387 | |
| 1388 | Function<void(AsyncBufferView::ID)> func; |
| 1389 | func.template bind<AsyncPipeline, &AsyncPipeline::afterWrite>(*this); |
| 1390 | // TODO: We should probably block when closing for in-flight writes |
| 1391 | Result res = writable.write(bufferID, func); |
| 1392 | if (not res) |
| 1393 | { |
| 1394 | source->getBuffersPool().unrefBuffer(bufferID); |
| 1395 | eventError.emit(res); |
| 1396 | } |
| 1397 | } |
| 1398 | |
| 1399 | bool AsyncPipeline::retryPendingWrites() |
| 1400 | { |
nothing calls this directly
no test coverage detected