NOTE
1. Detect a New Connection
1. Set a breakpoint 2. New connection enters
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.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;
}
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub