package io.netty.handler.traffic;

import g.a.a.a.a;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.ByteBufHolder;
import io.netty.channel.Channel;
import io.netty.channel.ChannelConfig;
import io.netty.channel.ChannelDuplexHandler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelOutboundBuffer;
import io.netty.channel.ChannelPromise;
import io.netty.util.Attribute;
import io.netty.util.AttributeKey;
import io.netty.util.internal.logging.InternalLogger;
import io.netty.util.internal.logging.InternalLoggerFactory;
import java.util.concurrent.TimeUnit;
import kotlinx.serialization.json.internal.AbstractJsonLexerKt;

/* JADX INFO: loaded from: classes.dex */
public abstract class AbstractTrafficShapingHandler extends ChannelDuplexHandler {
    public static final int CHANNEL_DEFAULT_USER_DEFINED_WRITABILITY_INDEX = 1;
    public static final long DEFAULT_CHECK_INTERVAL = 1000;
    public static final long DEFAULT_MAX_SIZE = 4194304;
    public static final long DEFAULT_MAX_TIME = 15000;
    public static final int GLOBALCHANNEL_DEFAULT_USER_DEFINED_WRITABILITY_INDEX = 3;
    public static final int GLOBAL_DEFAULT_USER_DEFINED_WRITABILITY_INDEX = 2;
    public static final long MINIMAL_WAIT = 10;
    public volatile long checkInterval;
    public volatile long maxTime;
    public volatile long maxWriteDelay;
    public volatile long maxWriteSize;
    private volatile long readLimit;
    public TrafficCounter trafficCounter;
    public final int userDefinedWritabilityIndex;
    private volatile long writeLimit;
    private static final InternalLogger logger = InternalLoggerFactory.getInstance((Class<?>) AbstractTrafficShapingHandler.class);
    public static final AttributeKey<Boolean> READ_SUSPENDED = AttributeKey.valueOf(AbstractTrafficShapingHandler.class.getName() + ".READ_SUSPENDED");
    public static final AttributeKey<Runnable> REOPEN_TASK = AttributeKey.valueOf(AbstractTrafficShapingHandler.class.getName() + ".REOPEN_TASK");

    public static final class ReopenReadTimerTask implements Runnable {
        public final ChannelHandlerContext ctx;

        public ReopenReadTimerTask(ChannelHandlerContext channelHandlerContext) {
            this.ctx = channelHandlerContext;
        }

        @Override // java.lang.Runnable
        public void run() {
            InternalLogger internalLogger;
            StringBuilder sb;
            String str;
            Channel channel = this.ctx.channel();
            ChannelConfig channelConfigConfig = channel.config();
            if (channelConfigConfig.isAutoRead() || !AbstractTrafficShapingHandler.isHandlerActive(this.ctx)) {
                if (AbstractTrafficShapingHandler.logger.isDebugEnabled()) {
                    if (!channelConfigConfig.isAutoRead() || AbstractTrafficShapingHandler.isHandlerActive(this.ctx)) {
                        internalLogger = AbstractTrafficShapingHandler.logger;
                        sb = new StringBuilder();
                        str = "Normal unsuspend: ";
                    } else {
                        internalLogger = AbstractTrafficShapingHandler.logger;
                        sb = new StringBuilder();
                        str = "Unsuspend: ";
                    }
                    sb.append(str);
                    sb.append(channelConfigConfig.isAutoRead());
                    sb.append(AbstractJsonLexerKt.COLON);
                    sb.append(AbstractTrafficShapingHandler.isHandlerActive(this.ctx));
                    internalLogger.debug(sb.toString());
                }
                channel.attr(AbstractTrafficShapingHandler.READ_SUSPENDED).set(Boolean.FALSE);
                channelConfigConfig.setAutoRead(true);
                channel.read();
            } else {
                if (AbstractTrafficShapingHandler.logger.isDebugEnabled()) {
                    InternalLogger internalLogger2 = AbstractTrafficShapingHandler.logger;
                    StringBuilder sbF = a.F("Not unsuspend: ");
                    sbF.append(channelConfigConfig.isAutoRead());
                    sbF.append(AbstractJsonLexerKt.COLON);
                    sbF.append(AbstractTrafficShapingHandler.isHandlerActive(this.ctx));
                    internalLogger2.debug(sbF.toString());
                }
                channel.attr(AbstractTrafficShapingHandler.READ_SUSPENDED).set(Boolean.FALSE);
            }
            if (AbstractTrafficShapingHandler.logger.isDebugEnabled()) {
                InternalLogger internalLogger3 = AbstractTrafficShapingHandler.logger;
                StringBuilder sbF2 = a.F("Unsuspend final status => ");
                sbF2.append(channelConfigConfig.isAutoRead());
                sbF2.append(AbstractJsonLexerKt.COLON);
                sbF2.append(AbstractTrafficShapingHandler.isHandlerActive(this.ctx));
                internalLogger3.debug(sbF2.toString());
            }
        }
    }

    public AbstractTrafficShapingHandler() {
        this(0L, 0L, 1000L, DEFAULT_MAX_TIME);
    }

    public AbstractTrafficShapingHandler(long j2) {
        this(0L, 0L, j2, DEFAULT_MAX_TIME);
    }

    public AbstractTrafficShapingHandler(long j2, long j3) {
        this(j2, j3, 1000L, DEFAULT_MAX_TIME);
    }

    public AbstractTrafficShapingHandler(long j2, long j3, long j4) {
        this(j2, j3, j4, DEFAULT_MAX_TIME);
    }

    public AbstractTrafficShapingHandler(long j2, long j3, long j4, long j5) {
        this.maxTime = DEFAULT_MAX_TIME;
        this.checkInterval = 1000L;
        this.maxWriteDelay = 4000L;
        this.maxWriteSize = DEFAULT_MAX_SIZE;
        if (j5 <= 0) {
            throw new IllegalArgumentException("maxTime must be positive");
        }
        this.userDefinedWritabilityIndex = userDefinedWritabilityIndex();
        this.writeLimit = j2;
        this.readLimit = j3;
        this.checkInterval = j4;
        this.maxTime = j5;
    }

    public static boolean isHandlerActive(ChannelHandlerContext channelHandlerContext) {
        Boolean bool = (Boolean) channelHandlerContext.channel().attr(READ_SUSPENDED).get();
        return bool == null || Boolean.FALSE.equals(bool);
    }

    public long calculateSize(Object obj) {
        ByteBuf byteBufContent;
        if (obj instanceof ByteBuf) {
            byteBufContent = (ByteBuf) obj;
        } else {
            if (!(obj instanceof ByteBufHolder)) {
                return -1L;
            }
            byteBufContent = ((ByteBufHolder) obj).content();
        }
        return byteBufContent.readableBytes();
    }

    @Override // io.netty.channel.ChannelInboundHandlerAdapter, io.netty.channel.ChannelInboundHandler
    public void channelRead(ChannelHandlerContext channelHandlerContext, Object obj) {
        long jCalculateSize = calculateSize(obj);
        long jMilliSecondFromNano = TrafficCounter.milliSecondFromNano();
        if (jCalculateSize > 0) {
            long jCheckWaitReadTime = checkWaitReadTime(channelHandlerContext, this.trafficCounter.readTimeToWait(jCalculateSize, this.readLimit, this.maxTime, jMilliSecondFromNano), jMilliSecondFromNano);
            if (jCheckWaitReadTime >= 10) {
                Channel channel = channelHandlerContext.channel();
                ChannelConfig channelConfigConfig = channel.config();
                InternalLogger internalLogger = logger;
                if (internalLogger.isDebugEnabled()) {
                    internalLogger.debug("Read suspend: " + jCheckWaitReadTime + AbstractJsonLexerKt.COLON + channelConfigConfig.isAutoRead() + AbstractJsonLexerKt.COLON + isHandlerActive(channelHandlerContext));
                }
                if (channelConfigConfig.isAutoRead() && isHandlerActive(channelHandlerContext)) {
                    channelConfigConfig.setAutoRead(false);
                    channel.attr(READ_SUSPENDED).set(Boolean.TRUE);
                    Attribute attributeAttr = channel.attr(REOPEN_TASK);
                    Runnable reopenReadTimerTask = (Runnable) attributeAttr.get();
                    if (reopenReadTimerTask == null) {
                        reopenReadTimerTask = new ReopenReadTimerTask(channelHandlerContext);
                        attributeAttr.set(reopenReadTimerTask);
                    }
                    channelHandlerContext.executor().schedule(reopenReadTimerTask, jCheckWaitReadTime, TimeUnit.MILLISECONDS);
                    if (internalLogger.isDebugEnabled()) {
                        StringBuilder sbF = a.F("Suspend final status => ");
                        sbF.append(channelConfigConfig.isAutoRead());
                        sbF.append(AbstractJsonLexerKt.COLON);
                        sbF.append(isHandlerActive(channelHandlerContext));
                        sbF.append(" will reopened at: ");
                        sbF.append(jCheckWaitReadTime);
                        internalLogger.debug(sbF.toString());
                    }
                }
            }
        }
        informReadOperation(channelHandlerContext, jMilliSecondFromNano);
        channelHandlerContext.fireChannelRead(obj);
    }

    @Override // io.netty.channel.ChannelInboundHandlerAdapter, io.netty.channel.ChannelInboundHandler
    public void channelRegistered(ChannelHandlerContext channelHandlerContext) {
        setUserDefinedWritability(channelHandlerContext, true);
        super.channelRegistered(channelHandlerContext);
    }

    public long checkWaitReadTime(ChannelHandlerContext channelHandlerContext, long j2, long j3) {
        return j2;
    }

    public void checkWriteSuspend(ChannelHandlerContext channelHandlerContext, long j2, long j3) {
        if (j3 > this.maxWriteSize || j2 > this.maxWriteDelay) {
            setUserDefinedWritability(channelHandlerContext, false);
        }
    }

    public void configure(long j2) {
        this.checkInterval = j2;
        TrafficCounter trafficCounter = this.trafficCounter;
        if (trafficCounter != null) {
            trafficCounter.configure(this.checkInterval);
        }
    }

    public void configure(long j2, long j3) {
        this.writeLimit = j2;
        this.readLimit = j3;
        TrafficCounter trafficCounter = this.trafficCounter;
        if (trafficCounter != null) {
            trafficCounter.resetAccounting(TrafficCounter.milliSecondFromNano());
        }
    }

    public void configure(long j2, long j3, long j4) {
        configure(j2, j3);
        configure(j4);
    }

    public void doAccounting(TrafficCounter trafficCounter) {
    }

    public long getCheckInterval() {
        return this.checkInterval;
    }

    public long getMaxTimeWait() {
        return this.maxTime;
    }

    public long getMaxWriteDelay() {
        return this.maxWriteDelay;
    }

    public long getMaxWriteSize() {
        return this.maxWriteSize;
    }

    public long getReadLimit() {
        return this.readLimit;
    }

    public long getWriteLimit() {
        return this.writeLimit;
    }

    public void informReadOperation(ChannelHandlerContext channelHandlerContext, long j2) {
    }

    @Override // io.netty.channel.ChannelDuplexHandler, io.netty.channel.ChannelOutboundHandler
    public void read(ChannelHandlerContext channelHandlerContext) {
        if (isHandlerActive(channelHandlerContext)) {
            channelHandlerContext.read();
        }
    }

    public void releaseReadSuspended(ChannelHandlerContext channelHandlerContext) {
        Channel channel = channelHandlerContext.channel();
        channel.attr(READ_SUSPENDED).set(Boolean.FALSE);
        channel.config().setAutoRead(true);
    }

    public void releaseWriteSuspended(ChannelHandlerContext channelHandlerContext) {
        setUserDefinedWritability(channelHandlerContext, true);
    }

    public void setCheckInterval(long j2) {
        this.checkInterval = j2;
        TrafficCounter trafficCounter = this.trafficCounter;
        if (trafficCounter != null) {
            trafficCounter.configure(j2);
        }
    }

    public void setMaxTimeWait(long j2) {
        if (j2 <= 0) {
            throw new IllegalArgumentException("maxTime must be positive");
        }
        this.maxTime = j2;
    }

    public void setMaxWriteDelay(long j2) {
        if (j2 <= 0) {
            throw new IllegalArgumentException("maxWriteDelay must be positive");
        }
        this.maxWriteDelay = j2;
    }

    public void setMaxWriteSize(long j2) {
        this.maxWriteSize = j2;
    }

    public void setReadLimit(long j2) {
        this.readLimit = j2;
        TrafficCounter trafficCounter = this.trafficCounter;
        if (trafficCounter != null) {
            trafficCounter.resetAccounting(TrafficCounter.milliSecondFromNano());
        }
    }

    public void setTrafficCounter(TrafficCounter trafficCounter) {
        this.trafficCounter = trafficCounter;
    }

    public void setUserDefinedWritability(ChannelHandlerContext channelHandlerContext, boolean z2) {
        ChannelOutboundBuffer channelOutboundBufferOutboundBuffer = channelHandlerContext.channel().unsafe().outboundBuffer();
        if (channelOutboundBufferOutboundBuffer != null) {
            channelOutboundBufferOutboundBuffer.setUserDefinedWritability(this.userDefinedWritabilityIndex, z2);
        }
    }

    public void setWriteLimit(long j2) {
        this.writeLimit = j2;
        TrafficCounter trafficCounter = this.trafficCounter;
        if (trafficCounter != null) {
            trafficCounter.resetAccounting(TrafficCounter.milliSecondFromNano());
        }
    }

    public abstract void submitWrite(ChannelHandlerContext channelHandlerContext, Object obj, long j2, long j3, long j4, ChannelPromise channelPromise);

    @Deprecated
    public void submitWrite(ChannelHandlerContext channelHandlerContext, Object obj, long j2, ChannelPromise channelPromise) {
        submitWrite(channelHandlerContext, obj, calculateSize(obj), j2, TrafficCounter.milliSecondFromNano(), channelPromise);
    }

    public String toString() {
        StringBuilder sb = new StringBuilder(290);
        sb.append("TrafficShaping with Write Limit: ");
        sb.append(this.writeLimit);
        sb.append(" Read Limit: ");
        sb.append(this.readLimit);
        sb.append(" CheckInterval: ");
        sb.append(this.checkInterval);
        sb.append(" maxDelay: ");
        sb.append(this.maxWriteDelay);
        sb.append(" maxSize: ");
        sb.append(this.maxWriteSize);
        sb.append(" and Counter: ");
        TrafficCounter trafficCounter = this.trafficCounter;
        if (trafficCounter != null) {
            sb.append(trafficCounter);
        } else {
            sb.append("none");
        }
        return sb.toString();
    }

    public TrafficCounter trafficCounter() {
        return this.trafficCounter;
    }

    public int userDefinedWritabilityIndex() {
        return 1;
    }

    /* JADX WARN: Removed duplicated region for block: B:11:0x006d  */
    @Override // io.netty.channel.ChannelDuplexHandler, io.netty.channel.ChannelOutboundHandler
    /*
        Code decompiled incorrectly, please refer to instructions dump.
        To view partially-correct add '--show-bad-code' argument
    */
    public void write(io.netty.channel.ChannelHandlerContext r21, java.lang.Object r22, io.netty.channel.ChannelPromise r23) {
        /*
            r20 = this;
            r10 = r20
            r2 = r22
            long r3 = r10.calculateSize(r2)
            long r7 = io.netty.handler.traffic.TrafficCounter.milliSecondFromNano()
            r0 = 0
            int r0 = (r3 > r0 ? 1 : (r3 == r0 ? 0 : -1))
            if (r0 <= 0) goto L6d
            io.netty.handler.traffic.TrafficCounter r11 = r10.trafficCounter
            long r14 = r10.writeLimit
            long r0 = r10.maxTime
            r12 = r3
            r16 = r0
            r18 = r7
            long r5 = r11.writeTimeToWait(r12, r14, r16, r18)
            r0 = 10
            int r0 = (r5 > r0 ? 1 : (r5 == r0 ? 0 : -1))
            if (r0 < 0) goto L6d
            io.netty.util.internal.logging.InternalLogger r0 = io.netty.handler.traffic.AbstractTrafficShapingHandler.logger
            boolean r1 = r0.isDebugEnabled()
            if (r1 == 0) goto L61
            java.lang.StringBuilder r1 = new java.lang.StringBuilder
            r1.<init>()
            java.lang.String r9 = "Write suspend: "
            r1.append(r9)
            r1.append(r5)
            r9 = 58
            r1.append(r9)
            io.netty.channel.Channel r11 = r21.channel()
            io.netty.channel.ChannelConfig r11 = r11.config()
            boolean r11 = r11.isAutoRead()
            r1.append(r11)
            r1.append(r9)
            boolean r9 = isHandlerActive(r21)
            r1.append(r9)
            java.lang.String r1 = r1.toString()
            r0.debug(r1)
        L61:
            r0 = r20
            r1 = r21
            r2 = r22
            r9 = r23
            r0.submitWrite(r1, r2, r3, r5, r7, r9)
            return
        L6d:
            r5 = 0
            goto L61
        */
        throw new UnsupportedOperationException("Method not decompiled: io.netty.handler.traffic.AbstractTrafficShapingHandler.write(io.netty.channel.ChannelHandlerContext, java.lang.Object, io.netty.channel.ChannelPromise):void");
    }
}
