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

JavaCreated Updated 1 min readhistorical

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);
        }
    }
}

4.2.1. How Is It Registered?

4. Register the Read Event with the Selector

Discussion

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