MCPcopy Create free account
hub / github.com/Pagghiu/SaneCppLibraries / asyncWriteWritable

Method asyncWriteWritable

Libraries/AsyncStreams/AsyncStreams.cpp:1368–1397  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

1366}
1367
1368void 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
1399bool AsyncPipeline::retryPendingWrites()
1400{

Callers

nothing calls this directly

Calls 5

refBufferMethod · 0.80
pauseMethod · 0.80
unrefBufferMethod · 0.80
isValidMethod · 0.45
writeMethod · 0.45

Tested by

no test coverage detected