NOTE
Pipeline Initialization
1. When pipeline is created 2. What pipeline looks like 3. Tail and Head
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