package io.netty.util.internal.shaded.org.jctools.queues;

import io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue;
import io.netty.util.internal.shaded.org.jctools.util.Pow2;
import io.netty.util.internal.shaded.org.jctools.util.UnsafeAccess;
import io.netty.util.internal.shaded.org.jctools.util.UnsafeRefArrayAccess;
import java.util.Iterator;

/* JADX INFO: loaded from: classes.dex */
public class MpscChunkedArrayQueue<E> extends MpscChunkedArrayQueueConsumerFields<E> implements MessagePassingQueue<E>, QueueProgressIndicators {
    private static final long C_INDEX_OFFSET;
    private static final Object JUMP;
    private static final long P_INDEX_OFFSET;
    private static final long P_LIMIT_OFFSET;
    long p0;
    long p1;
    long p10;
    long p11;
    long p12;
    long p13;
    long p14;
    long p15;
    long p16;
    long p17;
    long p2;
    long p3;
    long p4;
    long p5;
    long p6;
    long p7;

    static {
        try {
            P_INDEX_OFFSET = UnsafeAccess.UNSAFE.objectFieldOffset(MpscChunkedArrayQueueProducerFields.class.getDeclaredField("producerIndex"));
            try {
                C_INDEX_OFFSET = UnsafeAccess.UNSAFE.objectFieldOffset(MpscChunkedArrayQueueConsumerFields.class.getDeclaredField("consumerIndex"));
                try {
                    P_LIMIT_OFFSET = UnsafeAccess.UNSAFE.objectFieldOffset(MpscChunkedArrayQueueColdProducerFields.class.getDeclaredField("producerLimit"));
                    JUMP = new Object();
                } catch (NoSuchFieldException e2) {
                    throw new RuntimeException(e2);
                }
            } catch (NoSuchFieldException e3) {
                throw new RuntimeException(e3);
            }
        } catch (NoSuchFieldException e4) {
            throw new RuntimeException(e4);
        }
    }

    public MpscChunkedArrayQueue(int i2) {
        this(Math.max(2, Pow2.roundToPowerOfTwo(i2 / 8)), i2, false);
    }

    private boolean casProducerIndex(long j2, long j3) {
        return UnsafeAccess.UNSAFE.compareAndSwapLong(this, P_INDEX_OFFSET, j2, j3);
    }

    private boolean casProducerLimit(long j2, long j3) {
        return UnsafeAccess.UNSAFE.compareAndSwapLong(this, P_LIMIT_OFFSET, j2, j3);
    }

    private E[] getNextBuffer(E[] eArr, long j2) {
        long jNextArrayOffset = nextArrayOffset(j2);
        E[] eArr2 = (E[]) ((Object[]) UnsafeRefArrayAccess.lvElement(eArr, jNextArrayOffset));
        UnsafeRefArrayAccess.soElement(eArr, jNextArrayOffset, null);
        return eArr2;
    }

    private int getNextBufferCapacity(E[] eArr, long j2) {
        int length = eArr.length;
        if (this.isFixedChunkSize) {
            return eArr.length;
        }
        if (eArr.length - 1 != j2) {
            return (eArr.length * 2) - 1;
        }
        throw new IllegalStateException();
    }

    private long lvConsumerIndex() {
        return UnsafeAccess.UNSAFE.getLongVolatile(this, C_INDEX_OFFSET);
    }

    private long lvProducerIndex() {
        return UnsafeAccess.UNSAFE.getLongVolatile(this, P_INDEX_OFFSET);
    }

    private long lvProducerLimit() {
        return this.producerLimit;
    }

    private static long modifiedCalcElementOffset(long j2, long j3) {
        return UnsafeRefArrayAccess.REF_ARRAY_BASE + ((j2 & j3) << (UnsafeRefArrayAccess.REF_ELEMENT_SHIFT - 1));
    }

    private long newBufferAndOffset(E[] eArr, long j2) {
        this.consumerBuffer = eArr;
        long length = (eArr.length - 2) << 1;
        this.consumerMask = length;
        return modifiedCalcElementOffset(j2, length);
    }

    private E newBufferPeek(E[] eArr, long j2) {
        E e2 = (E) UnsafeRefArrayAccess.lvElement(eArr, newBufferAndOffset(eArr, j2));
        if (e2 != null) {
            return e2;
        }
        throw new IllegalStateException("new buffer must have at least one element");
    }

    private E newBufferPoll(E[] eArr, long j2) {
        long jNewBufferAndOffset = newBufferAndOffset(eArr, j2);
        E e2 = (E) UnsafeRefArrayAccess.lvElement(eArr, jNewBufferAndOffset);
        if (e2 == null) {
            throw new IllegalStateException("new buffer must have at least one element");
        }
        UnsafeRefArrayAccess.soElement(eArr, jNewBufferAndOffset, null);
        soConsumerIndex(j2 + 2);
        return e2;
    }

    private long nextArrayOffset(long j2) {
        return modifiedCalcElementOffset(j2 + 2, Long.MAX_VALUE);
    }

    private int offerSlowPath(long j2, E[] eArr, long j3, long j4) {
        long jLvConsumerIndex = lvConsumerIndex();
        long j5 = this.maxQueueCapacity;
        long currentBufferCapacity = getCurrentBufferCapacity(j2, j5) + jLvConsumerIndex;
        if (currentBufferCapacity > j3) {
            return !casProducerLimit(j4, currentBufferCapacity) ? 1 : 0;
        }
        if (jLvConsumerIndex == j3 - j5) {
            return 2;
        }
        return casProducerIndex(j3, 1 + j3) ? 3 : 1;
    }

    private void resize(long j2, E[] eArr, long j3, long j4, long j5, E e2) {
        E[] eArr2 = (E[]) CircularArrayOffsetCalculator.allocate(getNextBufferCapacity(eArr, j5));
        this.producerBuffer = eArr2;
        this.producerMask = (r8 - 2) << 1;
        long jModifiedCalcElementOffset = modifiedCalcElementOffset(j3, j2);
        UnsafeRefArrayAccess.soElement(eArr2, modifiedCalcElementOffset(j3, this.producerMask), e2);
        UnsafeRefArrayAccess.soElement(eArr, nextArrayOffset(j2), eArr2);
        long j6 = j5 - (j3 - j4);
        if (j6 <= 0) {
            throw new IllegalStateException();
        }
        soProducerLimit(Math.min(j2, j6) + j3);
        UnsafeRefArrayAccess.soElement(eArr, jModifiedCalcElementOffset, JUMP);
        soProducerIndex(2 + j3);
    }

    private void soConsumerIndex(long j2) {
        UnsafeAccess.UNSAFE.putOrderedLong(this, C_INDEX_OFFSET, j2);
    }

    private void soProducerIndex(long j2) {
        UnsafeAccess.UNSAFE.putOrderedLong(this, P_INDEX_OFFSET, j2);
    }

    private void soProducerLimit(long j2) {
        UnsafeAccess.UNSAFE.putOrderedLong(this, P_LIMIT_OFFSET, j2);
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public int capacity() {
        return (int) (this.maxQueueCapacity / 2);
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.QueueProgressIndicators
    public long currentConsumerIndex() {
        return lvConsumerIndex();
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.QueueProgressIndicators
    public long currentProducerIndex() {
        return lvProducerIndex();
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public void drain(MessagePassingQueue.Consumer<E> consumer, MessagePassingQueue.WaitStrategy waitStrategy, MessagePassingQueue.ExitCondition exitCondition) {
        E eRelaxedPoll;
        while (true) {
            while (exitCondition.keepRunning()) {
                eRelaxedPoll = relaxedPoll();
                int iIdle = eRelaxedPoll == null ? waitStrategy.idle(iIdle) : 0;
            }
            return;
            consumer.accept(eRelaxedPoll);
        }
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public int fill(MessagePassingQueue.Supplier<E> supplier, int i2) {
        long j2;
        while (true) {
            long jLvProducerLimit = lvProducerLimit();
            long jLvProducerIndex = lvProducerIndex();
            if ((jLvProducerIndex & 1) != 1) {
                long j3 = this.producerMask;
                E[] eArr = this.producerBuffer;
                long jMin = Math.min(jLvProducerLimit, ((long) (i2 * 2)) + jLvProducerIndex);
                if (jLvProducerIndex == jLvProducerLimit || jLvProducerLimit < jMin) {
                    int iOfferSlowPath = offerSlowPath(j3, eArr, jLvProducerIndex, jLvProducerLimit);
                    if (iOfferSlowPath == 1) {
                        continue;
                    } else {
                        if (iOfferSlowPath == 2) {
                            return 0;
                        }
                        if (iOfferSlowPath == 3) {
                            resize(j3, eArr, jLvProducerIndex, this.consumerIndex, this.maxQueueCapacity, supplier.get());
                            return 1;
                        }
                        j2 = jMin;
                    }
                } else {
                    j2 = jMin;
                }
                if (casProducerIndex(jLvProducerIndex, j2)) {
                    int i3 = (int) ((j2 - jLvProducerIndex) / 2);
                    for (int i4 = 0; i4 < i3; i4++) {
                        UnsafeRefArrayAccess.soElement(eArr, modifiedCalcElementOffset(((long) (i4 * 2)) + jLvProducerIndex, j3), supplier.get());
                    }
                    return i3;
                }
            }
        }
    }

    protected long getCurrentBufferCapacity(long j2, long j3) {
        return (this.isFixedChunkSize || 2 + j2 != j3) ? j2 : j3;
    }

    @Override // java.util.AbstractCollection, java.util.Collection, java.lang.Iterable
    public final Iterator<E> iterator() {
        throw new UnsupportedOperationException();
    }

    @Override // java.util.Queue, io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public boolean offer(E e2) {
        if (e2 == null) {
            throw null;
        }
        while (true) {
            long jLvProducerLimit = lvProducerLimit();
            long jLvProducerIndex = lvProducerIndex();
            if ((jLvProducerIndex & 1) != 1) {
                long j2 = this.producerMask;
                E[] eArr = this.producerBuffer;
                if (jLvProducerLimit <= jLvProducerIndex) {
                    int iOfferSlowPath = offerSlowPath(j2, eArr, jLvProducerIndex, jLvProducerLimit);
                    if (iOfferSlowPath == 1) {
                        continue;
                    } else {
                        if (iOfferSlowPath == 2) {
                            return false;
                        }
                        if (iOfferSlowPath == 3) {
                            resize(j2, eArr, jLvProducerIndex, this.consumerIndex, this.maxQueueCapacity, e2);
                            return true;
                        }
                    }
                }
                if (casProducerIndex(jLvProducerIndex, 2 + jLvProducerIndex)) {
                    UnsafeRefArrayAccess.soElement(eArr, modifiedCalcElementOffset(jLvProducerIndex, j2), e2);
                    return true;
                }
            }
        }
    }

    @Override // java.util.Queue, io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public E peek() {
        E[] eArr = this.consumerBuffer;
        long j2 = this.consumerIndex;
        long j3 = this.consumerMask;
        long jModifiedCalcElementOffset = modifiedCalcElementOffset(j2, j3);
        E e2 = (E) UnsafeRefArrayAccess.lvElement(eArr, jModifiedCalcElementOffset);
        if (e2 == null && j2 != lvProducerIndex()) {
            do {
                e2 = (E) UnsafeRefArrayAccess.lvElement(eArr, jModifiedCalcElementOffset);
            } while (e2 == null);
        }
        return e2 == JUMP ? newBufferPeek(getNextBuffer(eArr, j3), j2) : e2;
    }

    @Override // java.util.Queue, io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public E poll() {
        E[] eArr = this.consumerBuffer;
        long j2 = this.consumerIndex;
        long j3 = this.consumerMask;
        long jModifiedCalcElementOffset = modifiedCalcElementOffset(j2, j3);
        E e2 = (E) UnsafeRefArrayAccess.lvElement(eArr, jModifiedCalcElementOffset);
        if (e2 == null) {
            if (j2 == lvProducerIndex()) {
                return null;
            }
            do {
                e2 = (E) UnsafeRefArrayAccess.lvElement(eArr, jModifiedCalcElementOffset);
            } while (e2 == null);
        }
        if (e2 == JUMP) {
            return newBufferPoll(getNextBuffer(eArr, j3), j2);
        }
        UnsafeRefArrayAccess.soElement(eArr, jModifiedCalcElementOffset, null);
        soConsumerIndex(j2 + 2);
        return e2;
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public boolean relaxedOffer(E e2) {
        return offer(e2);
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public E relaxedPeek() {
        E[] eArr = this.consumerBuffer;
        long j2 = this.consumerIndex;
        long j3 = this.consumerMask;
        E e2 = (E) UnsafeRefArrayAccess.lvElement(eArr, modifiedCalcElementOffset(j2, j3));
        return e2 == JUMP ? newBufferPeek(getNextBuffer(eArr, j3), j2) : e2;
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public E relaxedPoll() {
        E[] eArr = this.consumerBuffer;
        long j2 = this.consumerIndex;
        long j3 = this.consumerMask;
        long jModifiedCalcElementOffset = modifiedCalcElementOffset(j2, j3);
        E e2 = (E) UnsafeRefArrayAccess.lvElement(eArr, jModifiedCalcElementOffset);
        if (e2 == null) {
            return null;
        }
        if (e2 == JUMP) {
            return newBufferPoll(getNextBuffer(eArr, j3), j2);
        }
        UnsafeRefArrayAccess.soElement(eArr, jModifiedCalcElementOffset, null);
        soConsumerIndex(j2 + 2);
        return e2;
    }

    @Override // java.util.AbstractCollection, java.util.Collection, io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public final int size() {
        long jLvConsumerIndex = lvConsumerIndex();
        while (true) {
            long jLvProducerIndex = lvProducerIndex();
            long jLvConsumerIndex2 = lvConsumerIndex();
            if (jLvConsumerIndex == jLvConsumerIndex2) {
                return ((int) (jLvProducerIndex - jLvConsumerIndex2)) >> 1;
            }
            jLvConsumerIndex = jLvConsumerIndex2;
        }
    }

    public MpscChunkedArrayQueue(int i2, int i3, boolean z) {
        if (i2 < 2) {
            throw new IllegalArgumentException("Initial capacity must be 2 or more");
        }
        if (i3 < 4) {
            throw new IllegalArgumentException("Max capacity must be 4 or more");
        }
        if (Pow2.roundToPowerOfTwo(i2) >= Pow2.roundToPowerOfTwo(i3)) {
            throw new IllegalArgumentException("Initial capacity cannot exceed maximum capacity(both rounded up to a power of 2)");
        }
        int iRoundToPowerOfTwo = Pow2.roundToPowerOfTwo(i2);
        long j2 = (iRoundToPowerOfTwo - 1) << 1;
        E[] eArr = (E[]) CircularArrayOffsetCalculator.allocate(iRoundToPowerOfTwo + 1);
        this.producerBuffer = eArr;
        this.producerMask = j2;
        this.consumerBuffer = eArr;
        this.consumerMask = j2;
        this.maxQueueCapacity = ((long) Pow2.roundToPowerOfTwo(i3)) << 1;
        this.isFixedChunkSize = z;
        soProducerLimit(j2);
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public int drain(MessagePassingQueue.Consumer<E> consumer) {
        return drain(consumer, capacity());
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public int drain(MessagePassingQueue.Consumer<E> consumer, int i2) {
        int i3 = 0;
        while (i3 < i2) {
            E eRelaxedPoll = relaxedPoll();
            if (eRelaxedPoll == null) {
                break;
            }
            consumer.accept(eRelaxedPoll);
            i3++;
        }
        return i3;
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public int fill(MessagePassingQueue.Supplier<E> supplier) {
        int iCapacity = capacity();
        long j2 = 0;
        do {
            int iFill = fill(supplier, MpmcArrayQueue.RECOMENDED_OFFER_BATCH);
            if (iFill == 0) {
                return (int) j2;
            }
            j2 += (long) iFill;
        } while (j2 <= iCapacity);
        return (int) j2;
    }

    @Override // io.netty.util.internal.shaded.org.jctools.queues.MessagePassingQueue
    public void fill(MessagePassingQueue.Supplier<E> supplier, MessagePassingQueue.WaitStrategy waitStrategy, MessagePassingQueue.ExitCondition exitCondition) {
        while (exitCondition.keepRunning()) {
            while (fill(supplier, MpmcArrayQueue.RECOMENDED_OFFER_BATCH) != 0) {
            }
            int iIdle = 0;
            while (fill(supplier, MpmcArrayQueue.RECOMENDED_OFFER_BATCH) == 0 && exitCondition.keepRunning()) {
                iIdle = waitStrategy.idle(iIdle);
            }
        }
    }
}
