NOTE

Outbound Events

1. Add handlers for the experiment 2. Start from the pipeline 3. Enter TailContext 4-7. Traverse outbound handlers back to HeadContext

JavaCreated Updated 1 min readhistorical

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

1. Add Handlers for the Experiment

  • OutboundHandlerA and OutboundHandlerC
package com.example.server.handler;

import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelOutboundHandlerAdapter;
import io.netty.channel.ChannelPromise;

/**
 * @description:
 * @author: <AUTHOR>
 * @create: 2019-12-10 23:08
 **/
public class OutboundHandlerA extends ChannelOutboundHandlerAdapter
{
    @Override
    public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception
    {
        System.out.println(this.getClass().getSimpleName() + ":" + msg);
        ctx.write(msg, promise);
    }
}
  • OutboundHandlerB
package com.example.server.handler;

import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelOutboundHandlerAdapter;
import io.netty.channel.ChannelPromise;

import java.util.concurrent.TimeUnit;

/**
 * @description:
 * @author: <AUTHOR>
 * @create: 2019-12-10 23:08
 **/
public class OutboundHandlerB extends ChannelOutboundHandlerAdapter
{
    @Override
    public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception
    {
        System.out.println(this.getClass().getSimpleName() + ":" + msg);
        ctx.write(msg, promise);
    }

    @Override
    public void handlerAdded(ChannelHandlerContext ctx) throws Exception
    {
        ctx.executor().schedule(()->{
            ctx.channel().writshoue("Hello World");
        },3, TimeUnit.SECONDS);
    }
}
  • NettyServer
.childHandler(new ChannelInitializer<SocketChannel>()
{
    @Override
    protected void initChannel(SocketChannel channel) throws Exception
    {
        ChannelPipeline pipeline = channel.pipeline();
        pipeline.addLast(new OutboundHandlerA())
                .addLast(new OutboundHandlerC())
                .addLast(new OutboundHandlerB());
    }
});

1.1. Connect Using nc

nc localhost 8000

1.2. Output

OutboundHandlerB:Hello World
OutboundHandlerC:Hello World
OutboundHandlerA:Hello World

1.3. Explanation

We can see that our handlers are invoked from the tail toward the head.

2. Start from the Pipeline

Set a breakpoint in OutboundHandlerB.handlerAdded, connect using nc, and start debugging. Execution first enters AbstractChannel#write.

public ChannelFuture write(Object msg) {
    return pipeline.write(msg);
}

The data enters the pipeline for propagation.

3. Enter TailContext First

DefaultChannelPipeline#write

public final ChannelFuture write(Object msg) {
    return tail.write(msg);
}

It is then forwarded through AbstractChannelHandlerContext.

3.1. Forward Through AbstractChannelHandlerContext

AbstractChannelHandlerContext#write

public ChannelFuture write(Object msg) {
    return write(msg, newPromise());
}

public ChannelFuture write(final Object msg, final ChannelPromise promise) {
    if (msg == null) {
        throw new NullPointerException("msg");
    }

    try {
        if (isNotValidPromise(promise, true)) {
            ReferenceCountUtil.release(msg);
            // cancelled
            return promise;
        }
    } catch (RuntimeException e) {
        ReferenceCountUtil.release(msg);
        throw e;
    }
    // Here.
    write(msg, false, promise);

    return promise;
}

private void write(Object msg, boolean flush, ChannelPromise promise) {
    // Find the next node.
    AbstractChannelHandlerContext next = findContextOutbound();
    final Object m = pipeline.touch(msg, next);
    EventExecutor executor = next.executor();
    if (executor.inEventLoop()) {
        if (flush) {
            next.invokeWriteAndFlush(m, promise);
        } else {
            // Here.
            next.invokeWrite(m, promise);
        }
    } else {
        AbstractWriteTask task;
        if (flush) {
            task = WriteAndFlushTask.newInstance(next, m, promise);
        }  else {
            task = WriteTask.newInstance(next, m, promise);
        }
        safeExecute(executor, task, promise, m);
    }
}

We can see that it first finds the next node, then calls that node’s write method.

3.1.1. How the Next Node Is Found

private AbstractChannelHandlerContext findContextOutbound() {
    AbstractChannelHandlerContext ctx = this;
    do {
        ctx = ctx.prev;
    } while (!ctx.outbound);
    return ctx;
}

It simply walks backward from the current node until it finds the previous outbound node.

What is the outbound node before tail? It is our OutboundHandlerB.

4. Enter OutboundHandlerB

io.netty.channel.AbstractChannelHandlerContext#invokeWrite

private void invokeWrite(Object msg, ChannelPromise promise) {
    if (invokeHandler()) {
        invokeWrite0(msg, promise);
    } else {
        write(msg, promise);
    }
}


 private void invokeWrite0(Object msg, ChannelPromise promise) {
    try {
        // This enters our OutboundHandlerB.
        ((ChannelOutboundHandler) handler()).write(this, msg, promise);
    } catch (Throwable t) {
        notifyOutboundHandlerException(t, promise);
    }
}
  • OutboundHandlerB
public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception
{
    System.out.println(this.getClass().getSimpleName() + ":" + msg);
    ctx.write(msg, promise);
}

Execution enters our OutboundHandlerB. After printing, it continues through AbstractChannelHandlerContext#write.

4.1. Forward Through AbstractChannelHandlerContext

5. Enter OutboundHandlerC

  • OutboundHandlerC
public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception
{
    System.out.println(this.getClass().getSimpleName() + ":" + msg);
    ctx.write(msg, promise);
}

Execution enters our OutboundHandlerC. After printing, it continues through AbstractChannelHandlerContext#write.

5.1. Forward Through AbstractChannelHandlerContext

6. Enter OutboundHandlerA

  • OutboundHandlerA
public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception
{
    System.out.println(this.getClass().getSimpleName() + ":" + msg);
    ctx.write(msg, promise);
}

Execution enters our OutboundHandlerA. After printing, it continues through AbstractChannelHandlerContext#write.

6.1. Forward Through AbstractChannelHandlerContext

7. Finally Enter HeadContext

  • io.netty.channel.DefaultChannelPipeline.HeadContext#write
public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception {
    unsafe.write(msg, promise);
}

public final void write(Object msg, ChannelPromise promise) {
    assertEventLoop();

    ChannelOutboundBuffer outboundBuffer = this.outboundBuffer;
    if (outboundBuffer == null) {
        // If the outboundBuffer is null we know the channel was closed and so
        // need to fail the future right away. If it is not null the handling of the rest
        // will be done in flush0()
        // See https://github.com/netty/netty/issues/2362
        safeSetFailure(promise, WRITE_CLOSED_CHANNEL_EXCEPTION);
        // release message now to prevent resource-leak
        ReferenceCountUtil.release(msg);
        return;
    }

    int size;
    try {
        msg = filterOutboundMessage(msg);
        size = pipeline.estimatorHandle().size(msg);
        if (size < 0) {
            size = 0;
        }
    } catch (Throwable t) {
        safeSetFailure(promise, t);
        ReferenceCountUtil.release(msg);
        return;
    }

    outboundBuffer.addMessage(msg, size, promise);
}

Execution enters HeadContext, which continues into AbstractUnsafe and performs the final work.

Discussion

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