Class AsyncIterablePublisher.SubscriptionImpl
- java.lang.Object
-
- org.reactivestreams.example.unicast.AsyncIterablePublisher.SubscriptionImpl
-
- All Implemented Interfaces:
java.lang.Runnable,Subscription
- Enclosing class:
- AsyncIterablePublisher<T>
final class AsyncIterablePublisher.SubscriptionImpl extends java.lang.Object implements Subscription, java.lang.Runnable
-
-
Field Summary
Fields Modifier and Type Field Description private booleancancelledprivate longdemandprivate java.util.concurrent.ConcurrentLinkedQueue<AsyncIterablePublisher.Signal>inboundSignalsprivate java.util.Iterator<T>iteratorprivate java.util.concurrent.atomic.AtomicBooleanon(package private) Subscriber<? super T>subscriber
-
Constructor Summary
Constructors Constructor Description SubscriptionImpl(Subscriber<? super T> subscriber)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description voidcancel()Request thePublisherto stop sending data and clean up resources.private voiddoCancel()private voiddoRequest(long n)private voiddoSend()private voiddoSubscribe()(package private) voidinit()voidrequest(long n)No events will be sent by aPublisheruntil demand is signaled via this method.voidrun()private voidsignal(AsyncIterablePublisher.Signal signal)private voidterminateDueTo(java.lang.Throwable t)private voidtryScheduleToExecute()
-
-
-
Field Detail
-
subscriber
final Subscriber<? super T> subscriber
-
cancelled
private boolean cancelled
-
demand
private long demand
-
iterator
private java.util.Iterator<T> iterator
-
inboundSignals
private final java.util.concurrent.ConcurrentLinkedQueue<AsyncIterablePublisher.Signal> inboundSignals
-
on
private final java.util.concurrent.atomic.AtomicBoolean on
-
-
Constructor Detail
-
SubscriptionImpl
SubscriptionImpl(Subscriber<? super T> subscriber)
-
-
Method Detail
-
doRequest
private void doRequest(long n)
-
doCancel
private void doCancel()
-
doSubscribe
private void doSubscribe()
-
doSend
private void doSend()
-
terminateDueTo
private void terminateDueTo(java.lang.Throwable t)
-
signal
private void signal(AsyncIterablePublisher.Signal signal)
-
run
public final void run()
- Specified by:
runin interfacejava.lang.Runnable
-
tryScheduleToExecute
private final void tryScheduleToExecute()
-
request
public void request(long n)
Description copied from interface:SubscriptionNo events will be sent by aPublisheruntil demand is signaled via this method.It can be called however often and whenever needed—but if the outstanding cumulative demand ever becomes Long.MAX_VALUE or more, it may be treated by the
Publisheras "effectively unbounded".Whatever has been requested can be sent by the
Publisherso only signal demand for what can be safely handled.A
Publishercan send less than is requested if the stream ends but then must emit eitherSubscriber.onError(Throwable)orSubscriber.onComplete().- Specified by:
requestin interfaceSubscription- Parameters:
n- the strictly positive number of elements to requests to the upstreamPublisher
-
cancel
public void cancel()
Description copied from interface:SubscriptionRequest thePublisherto stop sending data and clean up resources.Data may still be sent to meet previously signalled demand after calling cancel.
- Specified by:
cancelin interfaceSubscription
-
init
void init()
-
-