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

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