
来源:https://github.com/netty/netty
设置promise为成功状态
protected final void safeSetSuccess(ChannelPromise promise) {if (!(promise instanceof VoidChannelPromise) && !promise.trySuccess()) {logger.warn("Failed to mark a promise as success because it is done already: {}", promise);}}public boolean trySuccess() {return trySuccess(null);}public boolean trySuccess(V result) {if (setSuccess0(result)) {notifyListeners();return true;}return false;}
结果的初始值为空,设置为SUCCESS
private static final AtomicReferenceFieldUpdater<DefaultPromise, Object> RESULT_UPDATER =AtomicReferenceFieldUpdater.newUpdater(DefaultPromise.class, Object.class, "result");private boolean setSuccess0(V result) {return setValue0(result == null ? SUCCESS : result);}private boolean setValue0(Object objResult) {if (RESULT_UPDATER.compareAndSet(this, null, objResult) ||RESULT_UPDATER.compareAndSet(this, UNCANCELLABLE, objResult)) {checkNotifyWaiters();return true;}return false;}
通知等待者
private synchronized void checkNotifyWaiters() {if (waiters > 0) {notifyAll();}}
执行监听器
private void notifyListeners() {EventExecutor executor = executor();if (executor.inEventLoop()) {final InternalThreadLocalMap threadLocals = InternalThreadLocalMap.get();final int stackDepth = threadLocals.futureListenerStackDepth();if (stackDepth < MAX_LISTENER_STACK_DEPTH) {threadLocals.setFutureListenerStackDepth(stackDepth + 1);try {notifyListenersNow();} finally {threadLocals.setFutureListenerStackDepth(stackDepth);}return;}}safeExecute(executor, new Runnable() {@Overridepublic void run() {notifyListenersNow();}});}
激活管道注册事件
public final ChannelPipeline fireChannelRegistered() {AbstractChannelHandlerContext.invokeChannelRegistered(head);return this;}static void invokeChannelRegistered(final AbstractChannelHandlerContext next) {EventExecutor executor = next.executor();if (executor.inEventLoop()) {next.invokeChannelRegistered();} else {executor.execute(new Runnable() {@Overridepublic void run() {next.invokeChannelRegistered();}});}}
判断监听器是否存在
private void notifyListenersNow() {Object listeners;synchronized (this) {// Only proceed if there are listeners to notify and we are not already notifying listeners.if (notifyingListeners || this.listeners == null) {return;}notifyingListeners = true;listeners = this.listeners;this.listeners = null;}for (;;) {if (listeners instanceof DefaultFutureListeners) {notifyListeners0((DefaultFutureListeners) listeners);} else {notifyListener0(this, (GenericFutureListener<? extends Future<V>>) listeners);}synchronized (this) {if (this.listeners == null) {// Nothing can throw from within this method, so setting notifyingListeners back to false does not// need to be in a finally block.notifyingListeners = false;return;}listeners = this.listeners;this.listeners = null;}}}
根据不同的监听器类型执行不同的方法,直至把监听器处理完毕后置空
private void notifyListeners0(DefaultFutureListeners listeners) {GenericFutureListener<?>[] a = listeners.listeners();int size = listeners.size();for (int i = 0; i < size; i ++) {notifyListener0(this, a[i]);}}@SuppressWarnings({ "unchecked", "rawtypes" })private static void notifyListener0(Future future, GenericFutureListener l) {try {l.operationComplete(future);} catch (Throwable t) {logger.warn("An exception was thrown by " + l.getClass().getName() + ".operationComplete()", t);}}
判断处理器的状态,是否已经添加完成
private void invokeChannelRegistered() {if (invokeHandler()) {try {((ChannelInboundHandler) handler()).channelRegistered(this);} catch (Throwable t) {notifyHandlerException(t);}} else {fireChannelRegistered();}}private boolean invokeHandler() {// Store in local variable to reduce volatile reads.int handlerState = this.handlerState;return handlerState == ADD_COMPLETE || (!ordered && handlerState == ADD_PENDING);}
依次执行处理器的注册事件处理方法,首先是HeadContext,判断管道中是都含有待执行的处理器任务
public void channelRegistered(ChannelHandlerContext ctx) throws Exception {invokeHandlerAddedIfNeeded();ctx.fireChannelRegistered();}
处理完成后开始执行下一个处理器上下文的注册事件
public ChannelHandlerContext fireChannelRegistered() {invokeChannelRegistered(findContextInbound());return this;}
查找处理器链路上的流进类型的处理器,最后一个是TailContext,没有针对该事件进行操作。
private AbstractChannelHandlerContext findContextInbound() {AbstractChannelHandlerContext ctx = this;do {ctx = ctx.next;} while (!ctx.inbound);return ctx;}
文章转载自徘徊笔记,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。




