NOTE
3. Allocate a Thread and Register the Selector
1. Code to analyze 2. Propagate through the pipeline 3. Pass through ServerBootstrapAcceptor 4. Analyze channelRead
This is a historical learning note and may contain outdated or incomplete understanding.
1. Code to Analyze
- There is a section of logic in
io.netty.channel.nio.AbstractNioMessageChannel.NioMessageUnsafe#read.
//...
for (int i = 0; i < size; i ++) {
readPending = false;
pipeline.fireChannelRead(readBuf.get(i));
}
//...
2. Propagate Through the Pipeline
pipeline.fireChannelRead(readBuf.get(i)) starts propagating backward through the pipeline.
3. Pass Through ServerBootstrapAcceptor
There is a section of logic in io.netty.bootstrap.ServerBootstrap#init.
p.addLast(new ChannelInitializer<Channel>() {
@Override
public void initChannel(final Channel ch) throws Exception {
final ChannelPipeline pipeline = ch.pipeline();
ChannelHandler handler = config.handler();
if (handler != null) {
pipeline.addLast(handler);
}
ch.eventLoop().execute(new Runnable() {
@Override
public void run() {
pipeline.addLast(new ServerBootstrapAcceptor(
ch, currentChildGroup, currentChildHandler, currentChildOptions, currentChildAttrs));
}
});
}
});
A ServerBootstrapAcceptor is added to the pipeline here, so the fireChannelRead call above also reaches ServerBootstrapAcceptor.
4. Analyze ServerBootstrapAcceptor#channelRead
public void channelRead(ChannelHandlerContext ctx, Object msg) {
final Channel child = (Channel) msg;
// Add the user-defined childHandler to the pipeline.
child.pipeline().addLast(childHandler);
// Set options and attrs.
setChannelOptions(child, childOptions, logger);
for (Entry<AttributeKey<?>, Object> e: childAttrs) {
child.attr((AttributeKey<Object>) e.getKey()).set(e.getValue());
}
try {
// Choose a NioEventLoop and register with the selector.
// io.netty.channel.AbstractChannel.AbstractUnsafe#register
childGroup.register(child).addListener(new ChannelFutureListener() {
@Override
public void operationComplete(ChannelFuture future) throws Exception {
if (!future.isSuccess()) {
forceClose(child, future.cause());
}
}
});
} catch (Throwable t) {
forceClose(child, t);
}
}
4.1. Add childHandler
.childHandler(new ChannelInitializer<SocketChannel>()
{
@Override
protected void initChannel(SocketChannel channel) throws Exception
{
ChannelPipeline pipeline = channel.pipeline();
pipeline.addLast(new LoggingHandler(LogLevel.INFO));
}
});
The added handler is a ChannelInitializer. Look at its handlerAdded method.
handlerAdded
public void handlerAdded(ChannelHandlerContext ctx) throws Exception {
if (ctx.channel().isRegistered()) {
// Calls initChannel.
initChannel(ctx);
}
}
initChannel
private boolean initChannel(ChannelHandlerContext ctx) throws Exception {
if (initMap.putIfAbsent(ctx, Boolean.TRUE) == null) { // Guard against re-entrance.
try {
// Calls our ChannelInitializer.initChannel method.
initChannel((C) ctx.channel());
} catch (Throwable cause) {
// Explicitly call exceptionCaught(...) as we removed the handler before calling initChannel(...).
// We do so to prevent multiple calls to initChannel(...).
exceptionCaught(ctx, cause);
} finally {
remove(ctx);
}
return true;
}
return false;
}
4.2. Choose a NioEventLoop and Register with the Selector
//...
// Call SingleThreadEventLoop.register.
childGroup.register(child)//...
//...
SingleThreadEventLoop#register(io.netty.channel.Channel)
public ChannelFuture register(Channel channel) {
return register(new DefaultChannelPromise(channel, this));
}
@Override
public ChannelFuture register(final ChannelPromise promise) {
ObjectUtil.checkNotNull(promise, "promise");
// Call io.netty.channel.AbstractChannel.AbstractUnsafe#register.
promise.channel().unsafe().register(this, promise);
return promise;
}
io.netty.channel.AbstractChannel.AbstractUnsafe#register
public final void register(EventLoop eventLoop, final ChannelPromise promise) {
if (eventLoop == null) {
throw new NullPointerException("eventLoop");
}
if (isRegistered()) {
promise.setFailure(new IllegalStateException("registered to an event loop already"));
return;
}
if (!isCompatible(eventLoop)) {
promise.setFailure(
new IllegalStateException("incompatible event loop type: " + eventLoop.getClass().getName()));
return;
}
// Associate the current channel with NioEventLoop.
AbstractChannel.this.eventLoop = eventLoop;
if (eventLoop.inEventLoop()) {
// Register.
register0(promise);
} else {
try {
eventLoop.execute(new Runnable() {
@Override
public void run() {
register0(promise);
}
});
} catch (Throwable t) {
logger.warn(
"Force-closing a channel whose registration task was not accepted by an event loop: {}",
AbstractChannel.this, t);
closeForcibly();
closeFuture.setClosed();
safeSetFailure(promise, t);
}
}
}
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub