NOTE

Inbound Events

1. Add handlers for the experiment 2. Start from the pipeline 3. Forward through AbstractChannelHandlerContext 4. Enter HeadContext 5-8. Traverse inbound handlers to TailContext

JavaCreated Updated 2 min readhistorical

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

1. Add Handlers for the Experiment

Add three InboundHandler instances for the experiment.

  • InboundHandlerA
package com.example.server.handler;

import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;

/**
 * @description:
 * @author: <AUTHOR>
 * @create: 2019-12-10 20:41
 **/
public class InboundHandlerA extends ChannelInboundHandlerAdapter
{
    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception
    {
        // Each handler simply prints its class name and the message.
        System.out.println(this.getClass().getSimpleName() + ":" + msg);
        ctx.fireChannelRead(msg);
    }

    @Override
    public void channelActive(ChannelHandlerContext ctx) throws Exception
    {
        ctx.pipeline().fireChannelRead("Hello World");
    }
}
  • InboundHandlerB and InboundHandlerC are the same.
package com.example.server.handler;

import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;

/**
 * @description:
 * @author: <AUTHOR>
 * @create: 2019-12-10 20:41
 **/
public class InboundHandlerA extends ChannelInboundHandlerAdapter
{
    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception
    {
        // Each handler simply prints its class name and the message.
        System.out.println(this.getClass().getSimpleName() + ":" + msg);
        ctx.fireChannelRead(msg);
    }
}
  • NettyServer
.childHandler(new ChannelInitializer<SocketChannel>()
{
    @Override
    protected void initChannel(SocketChannel channel) throws Exception
    {
        ChannelPipeline pipeline = channel.pipeline();
        // Add them in order.
        pipeline.addLast(new InboundHandlerA())
                .addLast(new InboundHandlerB())
                .addLast(new InboundHandlerC());
    }
});

2. Start from the Pipeline

When we use ctx.pipeline().fireChannelRead("Hello World"), the message is first passed to the pipeline.

  • DefaultChannelPipeline#fireChannelRead
public final ChannelPipeline fireChannelRead(Object msg) {
    AbstractChannelHandlerContext.invokeChannelRead(head, msg);
    return this;
}

The message is forwarded through AbstractChannelHandlerContext.

3. Forward Through AbstractChannelHandlerContext

static void invokeChannelRead(final AbstractChannelHandlerContext next, Object msg) {
    final Object m = next.pipeline.touch(ObjectUtil.checkNotNull(msg, "msg"), next);
    EventExecutor executor = next.executor();
    if (executor.inEventLoop()) {
        next.invokeChannelRead(m);
    } else {
        executor.execute(new Runnable() {
            @Override
            public void run() {
                next.invokeChannelRead(m);
            }
        });
    }
}

We can debug to see what next is: DefaultChannelPipeline$HeadContext. This shows that after a message comes from the pipeline, the first node it enters is Head.

  • AbstractChannelHandlerContext#invokeChannelRead
private void invokeChannelRead(Object msg) {
    if (invokeHandler()) {
        try {
            // io.netty.channel.DefaultChannelPipeline.HeadContext#channelRead
            ((ChannelInboundHandler) handler()).channelRead(this, msg);
        } catch (Throwable t) {
            notifyHandlerException(t);
        }
    } else {
        fireChannelRead(msg);
    }
}

4. Enter HeadContext First

DefaultChannelPipeline.HeadContext#channelRead

public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
    // io.netty.channel.AbstractChannelHandlerContext#fireChannelRead
    ctx.fireChannelRead(msg);
}

HeadContext itself does not do much; it simply calls AbstractChannelHandlerContext#fireChannelRead again.

public ChannelHandlerContext fireChannelRead(final Object msg) {
    invokeChannelRead(findContextInbound(), msg);
    return this;
}

This is the interesting part: it first uses findContextInbound to find the next node, then uses invokeChannelRead to call that node’s channelRead method.

4.1. How the Next Node Is Found

private AbstractChannelHandlerContext findContextInbound() {
    AbstractChannelHandlerContext ctx = this;
    do {
        ctx = ctx.next;
    } while (!ctx.inbound);
    return ctx;
}

It simply starts from the current node and follows the linked list forward until it finds an inbound node. The next node is our InboundHandlerA.

The logic in invokeChannelRead is the same as described in the section on forwarding through AbstractChannelHandlerContext, so execution enters InboundHandlerA.channelRead.

5. Enter InboundHandlerA

public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception
{
    System.out.println(this.getClass().getSimpleName() + ":" + msg);
    ctx.fireChannelRead(msg);
}

After printing the message, this time it calls AbstractChannelHandlerContext#fireChannelRead rather than the pipeline. The logic is the same as above.

public ChannelHandlerContext fireChannelRead(final Object msg) {
    invokeChannelRead(findContextInbound(), msg);
    return this;
}

Again, it first uses findContextInbound to find the next node, then uses invokeChannelRead to call that node’s channelRead method.

5.1. How the Next Node Is Found

private AbstractChannelHandlerContext findContextInbound() {
    AbstractChannelHandlerContext ctx = this;
    do {
        ctx = ctx.next;
    } while (!ctx.inbound);
    return ctx;
}

It simply walks forward through the linked list from the current node until it finds an inbound node. The next node is our InboundHandlerB.

6. Enter InboundHandlerB

public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception
{
    System.out.println(this.getClass().getSimpleName() + ":" + msg);
    ctx.fireChannelRead(msg);
}

After printing the message, the logic is the same as entering InboundHandlerA; this time it continues to InboundHandlerC.

7. Enter InboundHandlerC

public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception
{
    System.out.println(this.getClass().getSimpleName() + ":" + msg);
    ctx.fireChannelRead(msg);
}

After C finishes printing, the next node is tail.

8. Enter TailContext

public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
    onUnhandledInboundMessage(msg);
}

protected void onUnhandledInboundMessage(Object msg) {
    try {
        logger.debug(
                "Discarded inbound message {} that reached at the tail of the pipeline. " +
                        "Please check your pipeline configuration.", msg);
    } finally {
        ReferenceCountUtil.release(msg);
    }
}

If debug logging is enabled, it prints information and finally releases the resource. This shows that Netty handles cleanup carefully at the end of the pipeline.

Discussion

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