Class ReactiveDataProducer
java.lang.Object
org.apache.hc.core5.reactive.ReactiveDataProducer
- All Implemented Interfaces:
AsyncDataProducer, ResourceHolder, org.reactivestreams.Subscriber<ByteBuffer>
@Contract(threading=SAFE)
final class ReactiveDataProducer
extends Object
implements AsyncDataProducer, org.reactivestreams.Subscriber<ByteBuffer>
An asynchronous data producer that supports Reactive Streams.
- Since:
- 5.0
-
Field Summary
FieldsModifier and TypeFieldDescriptionprivate static final intprivate final ArrayDeque<ByteBuffer> private final AtomicBooleanprivate final AtomicReference<Throwable> private final ReentrantLockprivate final org.reactivestreams.Publisher<ByteBuffer> private final AtomicReference<DataStreamChannel> private final AtomicReference<org.reactivestreams.Subscription> -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionintReturns the number of bytes immediately available for output.voidvoidvoidonNext(ByteBuffer byteBuffer) voidonSubscribe(org.reactivestreams.Subscription subscription) voidproduce(DataStreamChannel channel) Triggered to signal the ability of the underlying data channel to accept more data.void(package private) voidsetChannel(DataStreamChannel channel) private void
-
Field Details
-
BUFFER_WINDOW_SIZE
private static final int BUFFER_WINDOW_SIZE- See Also:
-
requestChannel
-
exception
-
complete
-
publisher
-
subscription
-
buffers
-
lock
-
-
Constructor Details
-
ReactiveDataProducer
-
-
Method Details
-
setChannel
-
onSubscribe
public void onSubscribe(org.reactivestreams.Subscription subscription) - Specified by:
onSubscribein interfaceorg.reactivestreams.Subscriber<ByteBuffer>
-
onNext
- Specified by:
onNextin interfaceorg.reactivestreams.Subscriber<ByteBuffer>
-
onError
- Specified by:
onErrorin interfaceorg.reactivestreams.Subscriber<ByteBuffer>
-
onComplete
public void onComplete()- Specified by:
onCompletein interfaceorg.reactivestreams.Subscriber<ByteBuffer>
-
signalReadiness
private void signalReadiness() -
available
public int available()Description copied from interface:AsyncDataProducerReturns the number of bytes immediately available for output. This method can be used as a hint to control output events of the underlying I/O session.Please note this method should return zero if the data producer is unable to produce any more data, in which case
AsyncDataProducer.produce(DataStreamChannel)method will not get triggered. The producer can resume writing out data asynchronously once more data becomes available or request output readiness events withDataStreamChannel.requestOutput().- Specified by:
availablein interfaceAsyncDataProducer- Returns:
- the number of bytes immediately available for output
- See Also:
-
produce
Description copied from interface:AsyncDataProducerTriggered to signal the ability of the underlying data channel to accept more data. The data producer can choose to write data immediately inside the call or asynchronously at some later point.Please note this method gets triggered only if
AsyncDataProducer.available()returns a positive value.- Specified by:
producein interfaceAsyncDataProducer- Parameters:
channel- the data channel capable of accepting more data.- Throws:
IOException- in case of an I/O error.- See Also:
-
releaseResources
public void releaseResources()- Specified by:
releaseResourcesin interfaceResourceHolder
-