public class io.smallrye.mutiny.operators.multi.processors.UnicastProcessor extends io.smallrye.mutiny.operators.AbstractMulti implements java.util.concurrent.Flow$Processor, java.util.concurrent.Flow$Subscription
{
private final java.lang.Runnable onTermination;
private final java.util.Queue queue;
private volatile boolean done;
private volatile java.lang.Throwable failure;
private volatile boolean cancelled;
private volatile java.util.concurrent.Flow$Subscriber downstream;
private static final java.util.concurrent.atomic.AtomicReferenceFieldUpdater DOWNSTREAM_UPDATER;
private final java.util.concurrent.atomic.AtomicInteger wip;
private final java.util.concurrent.atomic.AtomicLong requested;
private volatile boolean hasUpstream;
public static io.smallrye.mutiny.operators.multi.processors.UnicastProcessor create()
{
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
int v;
java.lang.Object v;
java.util.function.Supplier v;
v = new io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v = <io.smallrye.mutiny.helpers.queues.Queues: int BUFFER_S>;
v = staticinvoke <io.smallrye.mutiny.helpers.queues.Queues: java.util.function.Supplier unbounded(int)>(v);
v = interfaceinvoke v.<java.util.function.Supplier: java.lang.Object get()>();
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void <init>(java.util.Queue,java.lang.Runnable)>(v, null);
return v;
}
public static io.smallrye.mutiny.operators.multi.processors.UnicastProcessor create(java.util.Queue, java.lang.Runnable)
{
java.util.Queue v;
java.lang.Runnable v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
v := @parameter: java.util.Queue;
v := @parameter: java.lang.Runnable;
v = new io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void <init>(java.util.Queue,java.lang.Runnable)>(v, v);
return v;
}
private void <init>(java.util.Queue, java.lang.Runnable)
{
java.util.concurrent.atomic.AtomicLong v;
java.util.concurrent.atomic.AtomicInteger v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
java.lang.Object v;
java.util.Queue v;
java.lang.Runnable v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v := @parameter: java.util.Queue;
v := @parameter: java.lang.Runnable;
specialinvoke v.<io.smallrye.mutiny.operators.AbstractMulti: void <init>()>();
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean done> = 0;
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.lang.Throwable failure> = null;
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean cancelled> = 0;
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.Flow$Subscriber downstream> = null;
v = new java.util.concurrent.atomic.AtomicInteger;
specialinvoke v.<java.util.concurrent.atomic.AtomicInteger: void <init>()>();
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicInteger wip> = v;
v = new java.util.concurrent.atomic.AtomicLong;
specialinvoke v.<java.util.concurrent.atomic.AtomicLong: void <init>()>();
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicLong requested> = v;
v = staticinvoke <io.smallrye.mutiny.helpers.ParameterValidation: java.lang.Object nonNull(java.lang.Object,java.lang.String)>(v, "queue");
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.Queue queue> = v;
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.lang.Runnable onTermination> = v;
return;
}
private void onTerminate()
{
java.lang.Runnable v, v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.lang.Runnable onTermination>;
if v == null goto label;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.lang.Runnable onTermination>;
interfaceinvoke v.<java.lang.Runnable: void run()>();
label:
return;
}
void drainWithDownstream(java.util.concurrent.Flow$Subscriber)
{
long v, v, v;
byte v, v, v, v;
java.util.concurrent.atomic.AtomicInteger v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
int v, v;
boolean v, v, v, v, v, v;
java.util.concurrent.Flow$Subscriber v;
java.util.concurrent.atomic.AtomicLong v, v;
java.lang.Object v;
java.util.Queue v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v := @parameter: java.util.concurrent.Flow$Subscriber;
v = 1;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.Queue queue>;
label:
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicLong requested>;
v = virtualinvoke v.<java.util.concurrent.atomic.AtomicLong: long get()>();
v = 0L;
label:
v = v cmp v;
if v == 0 goto label;
v = interfaceinvoke v.<java.util.Queue: java.lang.Object poll()>();
if v != null goto label;
v = 1;
goto label;
label:
v = 0;
label:
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean done>;
v = specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean isCancelledOrDone(boolean,boolean)>(v, v);
if v == 0 goto label;
return;
label:
if v != 0 goto label;
interfaceinvoke v.<java.util.concurrent.Flow$Subscriber: void onNext(java.lang.Object)>(v);
v = v + 1L;
goto label;
label:
v = v cmp v;
if v != 0 goto label;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean done>;
v = interfaceinvoke v.<java.util.Queue: boolean isEmpty()>();
v = specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean isCancelledOrDone(boolean,boolean)>(v, v);
if v == 0 goto label;
return;
label:
v = v cmp 0L;
if v == 0 goto label;
v = v cmp 9223372036854775807L;
if v == 0 goto label;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicLong requested>;
v = neg v;
virtualinvoke v.<java.util.concurrent.atomic.AtomicLong: long addAndGet(long)>(v);
label:
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicInteger wip>;
v = neg v;
v = virtualinvoke v.<java.util.concurrent.atomic.AtomicInteger: int addAndGet(int)>(v);
if v != 0 goto label;
return;
}
private void drain()
{
java.util.concurrent.Flow$Subscriber v;
java.util.concurrent.atomic.AtomicInteger v, v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
int v, v, v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicInteger wip>;
v = virtualinvoke v.<java.util.concurrent.atomic.AtomicInteger: int getAndIncrement()>();
if v == 0 goto label;
return;
label:
v = 1;
label:
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.Flow$Subscriber downstream>;
if v == null goto label;
virtualinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void drainWithDownstream(java.util.concurrent.Flow$Subscriber)>(v);
return;
label:
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicInteger wip>;
v = neg v;
v = virtualinvoke v.<java.util.concurrent.atomic.AtomicInteger: int addAndGet(int)>(v);
if v != 0 goto label;
return;
}
private boolean isCancelledOrDone(boolean, boolean)
{
java.lang.Throwable v;
java.util.concurrent.Flow$Subscriber v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
java.util.Queue v;
boolean v, v, v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v := @parameter: boolean;
v := @parameter: boolean;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.Flow$Subscriber downstream>;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean cancelled>;
if v == 0 goto label;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.Queue queue>;
interfaceinvoke v.<java.util.Queue: void clear()>();
return 1;
label:
if v == 0 goto label;
if v == 0 goto label;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.lang.Throwable failure>;
if v == null goto label;
interfaceinvoke v.<java.util.concurrent.Flow$Subscriber: void onError(java.lang.Throwable)>(v);
goto label;
label:
interfaceinvoke v.<java.util.concurrent.Flow$Subscriber: void onComplete()>();
label:
return 1;
label:
return 0;
}
public void onSubscribe(java.util.concurrent.Flow$Subscription)
{
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
boolean v, v;
java.util.concurrent.Flow$Subscription v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v := @parameter: java.util.concurrent.Flow$Subscription;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean hasUpstream>;
if v == 0 goto label;
interfaceinvoke v.<java.util.concurrent.Flow$Subscription: void cancel()>();
return;
label:
v = specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean isDoneOrCancelled()>();
if v == 0 goto label;
interfaceinvoke v.<java.util.concurrent.Flow$Subscription: void cancel()>();
goto label;
label:
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean hasUpstream> = 1;
interfaceinvoke v.<java.util.concurrent.Flow$Subscription: void request(long)>(9223372036854775807L);
label:
return;
}
public void subscribe(io.smallrye.mutiny.subscription.MultiSubscriber)
{
java.lang.IllegalStateException v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
java.util.concurrent.atomic.AtomicReferenceFieldUpdater v;
io.smallrye.mutiny.subscription.MultiSubscriber v;
boolean v, v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v := @parameter: io.smallrye.mutiny.subscription.MultiSubscriber;
staticinvoke <io.smallrye.mutiny.helpers.ParameterValidation: java.lang.Object nonNull(java.lang.Object,java.lang.String)>(v, "downstream");
v = <io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicReferenceFieldUpdater DOWNSTREAM_UPDATER>;
v = virtualinvoke v.<java.util.concurrent.atomic.AtomicReferenceFieldUpdater: boolean compareAndSet(java.lang.Object,java.lang.Object,java.lang.Object)>(v, null, v);
if v == 0 goto label;
interfaceinvoke v.<io.smallrye.mutiny.subscription.MultiSubscriber: void onSubscribe(java.util.concurrent.Flow$Subscription)>(v);
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean cancelled>;
if v != 0 goto label;
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void drain()>();
goto label;
label:
v = new java.lang.IllegalStateException;
specialinvoke v.<java.lang.IllegalStateException: void <init>(java.lang.String)>("Already subscribed");
staticinvoke <io.smallrye.mutiny.helpers.Subscriptions: void fail(java.util.concurrent.Flow$Subscriber,java.lang.Throwable)>(v, v);
label:
return;
}
public synchronized void onNext(java.lang.Object)
{
io.smallrye.mutiny.subscription.BackPressureFailure v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
java.lang.Object v;
java.util.Queue v;
boolean v, v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v := @parameter: java.lang.Object;
v = specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean isDoneOrCancelled()>();
if v == 0 goto label;
return;
label:
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.Queue queue>;
v = interfaceinvoke v.<java.util.Queue: boolean offer(java.lang.Object)>(v);
if v != 0 goto label;
v = new io.smallrye.mutiny.subscription.BackPressureFailure;
specialinvoke v.<io.smallrye.mutiny.subscription.BackPressureFailure: void <init>(java.lang.String)>("the queue is full");
virtualinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void onError(java.lang.Throwable)>(v);
return;
label:
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void drain()>();
return;
}
private boolean isDoneOrCancelled()
{
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
boolean v, v, v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean done>;
if v != 0 goto label;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean cancelled>;
if v == 0 goto label;
label:
v = 1;
goto label;
label:
v = 0;
label:
return v;
}
public void onError(java.lang.Throwable)
{
java.lang.Throwable v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
boolean v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v := @parameter: java.lang.Throwable;
virtualinvoke v.<java.lang.Object: java.lang.Class getClass()>();
v = specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean isDoneOrCancelled()>();
if v == 0 goto label;
return;
label:
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void onTerminate()>();
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.lang.Throwable failure> = v;
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean done> = 1;
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void drain()>();
return;
}
public void onComplete()
{
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
boolean v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v = specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean isDoneOrCancelled()>();
if v == 0 goto label;
return;
label:
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void onTerminate()>();
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean done> = 1;
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void drain()>();
return;
}
public void request(long)
{
java.util.concurrent.atomic.AtomicLong v;
byte v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
long v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v := @parameter: long;
v = v cmp 0L;
if v <= 0 goto label;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicLong requested>;
staticinvoke <io.smallrye.mutiny.helpers.Subscriptions: long add(java.util.concurrent.atomic.AtomicLong,long)>(v, v);
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void drain()>();
label:
return;
}
public void cancel()
{
java.util.concurrent.atomic.AtomicInteger v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
java.util.concurrent.atomic.AtomicReferenceFieldUpdater v;
int v;
java.lang.Object v;
java.util.Queue v;
boolean v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean cancelled>;
if v == 0 goto label;
return;
label:
v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: boolean cancelled> = 1;
v = <io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicReferenceFieldUpdater DOWNSTREAM_UPDATER>;
v = virtualinvoke v.<java.util.concurrent.atomic.AtomicReferenceFieldUpdater: java.lang.Object getAndSet(java.lang.Object,java.lang.Object)>(v, null);
if v == null goto label;
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: void onTerminate()>();
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicInteger wip>;
v = virtualinvoke v.<java.util.concurrent.atomic.AtomicInteger: int getAndIncrement()>();
if v != 0 goto label;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.Queue queue>;
interfaceinvoke v.<java.util.Queue: void clear()>();
label:
return;
}
public boolean hasSubscriber()
{
java.util.concurrent.Flow$Subscriber v;
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
boolean v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v = v.<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.Flow$Subscriber downstream>;
if v == null goto label;
v = 1;
goto label;
label:
v = 0;
label:
return v;
}
public io.smallrye.mutiny.operators.multi.processors.SerializedProcessor serialized()
{
io.smallrye.mutiny.operators.multi.processors.UnicastProcessor v;
io.smallrye.mutiny.operators.multi.processors.SerializedProcessor v;
v := @this: io.smallrye.mutiny.operators.multi.processors.UnicastProcessor;
v = new io.smallrye.mutiny.operators.multi.processors.SerializedProcessor;
specialinvoke v.<io.smallrye.mutiny.operators.multi.processors.SerializedProcessor: void <init>(java.util.concurrent.Flow$Processor)>(v);
return v;
}
static void <clinit>()
{
java.util.concurrent.atomic.AtomicReferenceFieldUpdater v;
v = staticinvoke <java.util.concurrent.atomic.AtomicReferenceFieldUpdater: java.util.concurrent.atomic.AtomicReferenceFieldUpdater newUpdater(java.lang.Class,java.lang.Class,java.lang.String)>(class "Lio/smallrye/mutiny/operators/multi/processors/UnicastProcessor;", class "Ljava/util/concurrent/Flow$Subscriber;", "downstream");
<io.smallrye.mutiny.operators.multi.processors.UnicastProcessor: java.util.concurrent.atomic.AtomicReferenceFieldUpdater DOWNSTREAM_UPDATER> = v;
return;
}
}