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
This is a historical learning note and may contain outdated or incomplete understanding.
1. Add Handlers for the Experiment
OutboundHandlerAandOutboundHandlerC
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