暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Netty源码分析之注册服务端通道(下)

徘徊笔记 2019-06-08
141

来源: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() {
          @Override
          public 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() {
            @Override
            public 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进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

                        评论