NOTE
2.91 3. Register the Channel
Select a NioEventLoop, create a ChannelPromise/Future, register the channel through Unsafe with the JDK NIO selector, and invoke handlerAdded/initChannel callbacks.
This is a historical learning note and may contain outdated or incomplete understanding.
1. Code to Analyze
ChannelFuture regFuture = config().group().register(channel)
2. Select a NioEventLoop
config()returnsServerBootstrap.ServerBootstrap.group()returnsEventLoopGroup.EventLoopGroupextendsMultithreadEventLoopGroup, so theregistermethod ofMultithreadEventLoopGroupis called.MultithreadEventLoopGroup#register
public ChannelFuture register(Channel channel) {
return next().register(channel);
}
next() calls MultithreadEventExecutorGroup#next:
public EventExecutor next() {
// chooser is either PowerOfTwoEventExecutorChooser
// or GenericEventExecutorChooser.
return chooser.next();
}
The chooser’s next method selects a NioEventLoop, then calls its register method.
The call eventually reaches SingleThreadEventLoop#register.
3. Create a Future
SingleThreadEventLoop#register
public ChannelFuture register(Channel channel) {
return register(new DefaultChannelPromise(channel, this));
}
Pass the Channel and EventLoopGroup into the DefaultChannelPromise constructor. First look at DefaultChannelPromise.
3.1. DefaultChannelFuture Class Hierarchy

JDK Future provides interfaces for obtaining the result of an asynchronous task:
Future
cancel
isCancelled
isDone
get
get
In addition to these, Netty’s Future provides sync() and await() for blocking waits, and Listeners for callbacks when a task finishes. Compared with actively polling a JDK Future, Netty can use asynchronous callbacks.
ChannelFuture extends Netty’s Future and binds a Channel [I/O operation] to a Future.
Promise extends Netty’s Future and provides setSuccess and setFailure methods used to wake threads waiting in sync and await.
ChannelPromise combines ChannelFuture and Promise.
At the bottom, DefaultChannelPromise is the default implementation of ChannelPromise.
As described above, there are two programming styles:
- Synchronous blocking
- Use
sync/awaitto block until another thread callssetSuccess/setFailure.
- Use
- Asynchronous callback
- Bind a listener and receive asynchronous notification when a result is available.
4. Call Unsafe to register
Pass the created DefaultChannelPromise into register.
public ChannelFuture register(final ChannelPromise promise) {
ObjectUtil.checkNotNull(promise, "promise");
// AbstractUnsafe#register associates the current channel with the eventLoop.
promise.channel().unsafe().register(this, promise);
return promise;
}
AbstractUnsafe
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;
}
// Finally associate the NioEventLoop with the channel.
AbstractChannel.this.eventLoop = eventLoop;
// If the current thread is already in the EventLoop, register directly.
if (eventLoop.inEventLoop()) {
register0(promise);
} else {
// Otherwise wrap it as a Runnable and submit it to SingleThreadEventExecutor#execute.
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.1. Associate the channel with NioEventLoop
AbstractChannel.this.eventLoop = eventLoop;
4.2. Actual Registration
private void register0(ChannelPromise promise) {
try {
// check if the channel is still open as it could be closed in the mean time when the register
// call was outside of the eventLoop
if (!promise.setUncancellable() || !ensureOpen(promise)) {
return;
}
boolean firstRegistration = neverRegistered;
// Actually register.
doRegister();
neverRegistered = false;
registered = true;
// Invoke handlerAdded on handlers.
pipeline.invokeHandlerAddedIfNeeded();
safeSetSuccess(promise);
// Invoke the Handler's channelRegistered method.
pipeline.fireChannelRegistered();
// Only fire a channelActive if the channel has never been registered. This prevents firing
// multiple channel actives if the channel is deregistered and re-registered.
// This returns false here; it returns true only during bind.
if (isActive()) {
if (firstRegistration) {
// So handler.channelActive is not invoked here.
pipeline.fireChannelActive();
} else if (config().isAutoRead()) {
// This channel was registered before and autoRead() is set. This means we need to begin read
// again so that we process inbound data.
//
// See https://github.com/netty/netty/issues/4805
beginRead();
}
}
} catch (Throwable t) {
// Close the channel directly to avoid FD leak.
closeForcibly();
closeFuture.setClosed();
safeSetFailure(promise, t);
}
}
There are two key steps: registration itself and invoking the handlerAdded method.
4.2.1. Register the channel with the JDK NIO selector
protected void doRegister() throws Exception {
boolean selected = false;
for (;;) {
try {
// Call the underlying JDK channel registration.
// eventLoop().unwrappedSelector() is the underlying JDK selector.
// 0 means no events are of interest yet.
// this is attached to the JDK selector as the current AbstractNioChannel.
selectionKey = javaChannel().register(eventLoop().unwrappedSelector(), 0, this);
return;
} catch (CancelledKeyException e) {
if (!selected) {
// Force the Selector to select now as the "canceled" SelectionKey may still be
// cached and not removed because no Select.select(..) operation was called yet.
eventLoop().selectNow();
selected = true;
} else {
// We forced a select operation on the selector before but the SelectionKey is still cached
// for whatever reason. JDK bug ?
throw e;
}
}
}
}
4.2.2. Invoke handlerAdded
final void invokeHandlerAddedIfNeeded() {
assert channel.eventLoop().inEventLoop();
if (firstRegistration) {
firstRegistration = false;
// We are now registered to the EventLoop. It's time to call the callbacks for the ChannelHandlers,
// that were added before the registration was done.
callHandlerAddedForAllHandlers();
}
}
private void callHandlerAddedForAllHandlers() {
final PendingHandlerCallback pendingHandlerCallbackHead;
synchronized (this) {
assert !registered;
// This Channel itself was registered.
registered = true;
pendingHandlerCallbackHead = this.pendingHandlerCallbackHead;
// Null out so it can be GC'ed.
this.pendingHandlerCallbackHead = null;
}
// This must happen outside of the synchronized(...) block as otherwise handlerAdded(...) may be called while
// holding the lock and so produce a deadlock if handlerAdded(...) will try to add another handler from outside
// the EventLoop.
PendingHandlerCallback task = pendingHandlerCallbackHead;
// One of the tasks here is
// io.netty.channel.DefaultChannelPipeline.PendingHandlerAddedTask#execute.
while (task != null) {
task.execute();
task = task.next;
}
}
DefaultChannelPipeline.PendingHandlerAddedTask#execute
void execute() {
EventExecutor executor = ctx.executor();
if (executor.inEventLoop()) {
// Calls io.netty.channel.DefaultChannelPipeline#callHandlerAdded0.
callHandlerAdded0(ctx);
} else {
try {
executor.execute(this);
} catch (RejectedExecutionException e) {
if (logger.isWarnEnabled()) {
logger.warn(
"Can't invoke handlerAdded() as the EventExecutor {} rejected it, removing handler {}.",
executor, ctx.name(), e);
}
remove0(ctx);
ctx.setRemoved();
}
}
}
DefaultChannelPipeline#callHandlerAdded0
private void callHandlerAdded0(final AbstractChannelHandlerContext ctx) {
try {
// We must call setAddComplete before calling handlerAdded. Otherwise if the handlerAdded method generates
// any pipeline events ctx.handler() will miss them because the state will not allow it.
ctx.setAddComplete();
// Invoke io.netty.channel.ChannelInitializer#handlerAdded.
// Propagation starts from the pipeline.
ctx.handler().handlerAdded(ctx);
} catch (Throwable t) {
boolean removed = false;
try {
remove0(ctx);
try {
ctx.handler().handlerRemoved(ctx);
} finally {
ctx.setRemoved();
}
removed = true;
} catch (Throwable t2) {
if (logger.isWarnEnabled()) {
logger.warn("Failed to remove a handler: " + ctx.name(), t2);
}
}
if (removed) {
fireExceptionCaught(new ChannelPipelineException(
ctx.handler().getClass().getName() +
".handlerAdded() has thrown an exception; removed.", t));
} else {
fireExceptionCaught(new ChannelPipelineException(
ctx.handler().getClass().getName() +
".handlerAdded() has thrown an exception; also failed to remove.", t));
}
}
}
ChannelInitializer#handlerAdded
public void handlerAdded(ChannelHandlerContext ctx) throws Exception {
if (ctx.channel().isRegistered()) {
// Call io.netty.channel.ChannelInitializer#initChannel(ChannelHandlerContext).
initChannel(ctx);
}
}
4.2.2.1. Callback ChannelInitializer#initChannel
ChannelInitializer#initChannel
private boolean initChannel(ChannelHandlerContext ctx) throws Exception {
if (initMap.putIfAbsent(ctx, Boolean.TRUE) == null) { // Guard against re-entrance.
try {
// Finally call back to io.netty.channel.ChannelInitializer#initChannel.
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;
}
When was ChannelInitializer added?
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub