NOTE

Pipeline Initialization

1. When pipeline is created 2. What pipeline looks like 3. Tail and Head

JavaCreated Updated 1 min readhistorical

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

1. When Is the Pipeline Created?

Whether it is a client-side or server-side channel, the pipeline is created in the AbstractChannel constructor. Every channel has its own pipeline.

protected AbstractChannel(Channel parent) {
    this.parent = parent;
    id = newId();
    unsafe = newUnsafe();
    pipeline = newChannelPipeline();// Create pipeline.
}
  • newPipeline
protected DefaultChannelPipeline newChannelPipeline() {
    return new DefaultChannelPipeline(this);
}

2. What Does the Pipeline Look Like?

  • DefaultChannelPipeline
protected DefaultChannelPipeline(Channel channel) {
    this.channel = ObjectUtil.checkNotNull(channel, "channel");
    succeededFuture = new SucceededChannelFuture(channel, null);
    voidPromise =  new VoidChannelPromise(channel, true);

    // Head and tail nodes.
    tail = new TailContext(this);
    head = new HeadContext(this);

    // Link head and tail as a doubly linked list.
    head.next = tail;
    tail.prev = head;
}

The pipeline is a doubly linked list.

2.1. Pipeline Node Data Structure: ChannelHandlerContext

// The default implementation of this interface is AbstractChannelHandlerContext.
public interface ChannelHandlerContext extends AttributeMap, ChannelInboundInvoker, ChannelOutboundInvoker {
    // Channel associated with this node.
    Channel channel();

    // Thread pool used to execute tasks.
    EventExecutor executor();

    // Name of this handler.
    String name();

    // Handler that executes the business logic.
    ChannelHandler handler();

    // Methods from ChannelInboundInvoker that propagate read events.

    // Methods from ChannelOutboundInvoker that propagate write events.
}

3. Tail and Head Analysis

3.1. TailContext Is Inbound

final class TailContext extends AbstractChannelHandlerContext implements ChannelInboundHandler {// Implements Inbound.

    TailContext(DefaultChannelPipeline pipeline) {
        // AbstractChannelHandlerContext(DefaultChannelPipeline pipeline, EventExecutor executor, String name,
        //                          boolean inbound, boolean outbound)
        // TAIL is inbound.
        super(pipeline, null, TAIL_NAME, true, false);
        // Marks this handlerContext as fully added.
        setAddComplete();
    }

    @Override
    // Shows that context itself is the handler.
    public ChannelHandler handler() {
        return this;
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
        onUnhandledInboundException(cause);
        /* Ultimately the following method is called. This shows that tail mainly
           performs final cleanup; if earlier handlers did not process the event,
           it logs a warning.
        protected void onUnhandledInboundException(Throwable cause) {
            try {
                    logger.warn(
                            "An exceptionCaught() event was fired, and it reached at the tail of the pipeline. " +
                                    "It usually means the last handler in the pipeline did not handle the exception.",
                            cause);
                } finally {
                    ReferenceCountUtil.release(cause);
                }
            }*/
    }

    @Override
    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);
            }
        }
        */
    }
}

3.2. HeadContext Is Outbound

final class HeadContext extends AbstractChannelHandlerContext
            implements ChannelOutboundHandler, ChannelInboundHandler {// Implements both inbound and outbound.
    // Has an unsafe instance used for read/write operations.
    private final Unsafe unsafe;

    HeadContext(DefaultChannelPipeline pipeline) {
        // AbstractChannelHandlerContext(DefaultChannelPipeline pipeline, EventExecutor executor, String name,
        //                          boolean inbound, boolean outbound)
        // head is outbound.
        super(pipeline, null, HEAD_NAME, false, true);
        unsafe = pipeline.channel().unsafe();
        setAddComplete();
    }

    // Data-processing operations are delegated to unsafe.
    @Override
    public void bind(
            ChannelHandlerContext ctx, SocketAddress localAddress, ChannelPromise promise)
            throws Exception {
        unsafe.bind(localAddress, promise);
    }

    @Override
    public void connect(
            ChannelHandlerContext ctx,
            SocketAddress remoteAddress, SocketAddress localAddress,
            ChannelPromise promise) throws Exception {
        unsafe.connect(remoteAddress, localAddress, promise);
    }

    // Propagate events.
    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
        ctx.fireExceptionCaught(cause);
    }

    @Override
    public void channelRegistered(ChannelHandlerContext ctx) throws Exception {
        invokeHandlerAddedIfNeeded();
        ctx.fireChannelRegistered();
    }
}

Discussion

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