NOTE

1. Detect a New Connection

1. Set a breakpoint 2. New connection enters

JavaCreated Updated 1 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. Set a Breakpoint

When a new connection arrives, the OP_ACCEPT logic in processSelectedKeys inside NioEventLoop.run is called.

private void processSelectedKey(SelectionKey k, AbstractNioChannel ch) {
    //...
    if ((readyOps & (SelectionKey.OP_READ | SelectionKey.OP_ACCEPT)) != 0 || readyOps == 0) {
        unsafe.read();
    }
    //...
}

1.1. Start a Client Connection

Set a breakpoint at unsafe.read, start the server, and connect with telnet or nc.

nc localhost 8000

2. New Connection Enters

  • Enter AbstractNioMessageChannel.NioMessageUnsafe#read
public void read() {
        assert eventLoop().inEventLoop();
        final ChannelConfig config = config();
        final ChannelPipeline pipeline = pipeline();
        // This allocHandler controls the server-side connection-accept rate.
        final RecvByteBufAllocator.Handle allocHandle = unsafe().recvBufAllocHandle();
        allocHandle.reset(config);


        boolean closed = false;
        Throwable exception = null;
        try {
            try {
                // Keep accepting client connections until there are no more
                // incoming connections or the rate limit is reached.
                do {
                    // NioServerSocketChannel#doReadMessages
                    // Obtain the JDK low-level channel.
                    int localRead = doReadMessages(readBuf);
                    if (localRead == 0) {
                        break;
                    }
                    if (localRead < 0) {
                        closed = true;
                        break;
                    }
                    // Count accepted connections.
                    allocHandle.incMessagesRead(localRead);
                    // Rate control.
                } while (allocHandle.continueReading());
            } catch (Throwable t) {
                exception = t;
            }

            int size = readBuf.size();
            for (int i = 0; i < size; i ++) {
                readPending = false;
                // Allocate a thread and register with the selector.
                pipeline.fireChannelRead(readBuf.get(i));
            }
            readBuf.clear();
            allocHandle.readComplete();
            pipeline.fireChannelReadComplete();

            if (exception != null) {
                closed = closeOnReadError(exception);

                pipeline.fireExceptionCaught(exception);
            }

            if (closed) {
                inputShutdown = true;
                if (isOpen()) {
                    close(voidPromise());
                }
            }
        } finally {
            // Check if there is a readPending which was not processed yet.
            // This could be for two reasons:
            // * The user called Channel.read() or ChannelHandlerContext.read() in channelRead(...) method
            // * The user called Channel.read() or ChannelHandlerContext.read() in channelReadComplete(...) method
            //
            // See https://github.com/netty/netty/issues/2254
            if (!readPending && !config.isAutoRead()) {
                removeReadOp();
            }
        }
    }
}

After a client connects, three main things happen: accept the new connection and obtain the underlying JDK channel, control the rate at which connections are accepted, and allocate a thread and register the selector.

2.1. Accept the New Connection and Obtain the Underlying JDK Channel

protected int doReadMessages(List<Object> buf) throws Exception {
    // Obtain the underlying JDK channel after accept.
    SocketChannel ch = SocketUtils.accept(javaChannel());

    try {
        if (ch != null) {
            // Create NioSocketChannel and put it into the collection.
            buf.add(new NioSocketChannel(this, ch));
            return 1;// Returning 1 means there is one client channel.
        }
    } catch (Throwable t) {
        logger.warn("Failed to create a new channel from an accepted socket.", t);

        try {
            ch.close();
        } catch (Throwable t2) {
            logger.warn("Failed to close a socket.", t2);
        }
    }

    return 0;
}

2.1.1. Creating NioSocketChannel

2. Create NioSocketChannel

2.2. Control the Connection-Accept Rate

DefaultMaxMessagesRecvByteBufAllocator.MaxMessageHandle#continueReading(io.netty.util.UncheckedBooleanSupplier)
public boolean continueReading(UncheckedBooleanSupplier maybeMoreDataSupplier) {
    return config.isAutoRead() &&
           (!respectMaybeMoreData || maybeMoreDataSupplier.get()) &&
           // Current connection count < 16.
           totalMessages < maxMessagePerRead &&
           totalBytesRead > 0;
}

2.3. Allocate a Thread and Register the Selector

3. Allocate a Thread and Register the Selector

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub