1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16 package io.netty.channel;
17
18 import io.netty.channel.Channel.Unsafe;
19 import io.netty.util.ReferenceCountUtil;
20 import io.netty.util.ResourceLeakDetector;
21 import io.netty.util.concurrent.EventExecutor;
22 import io.netty.util.concurrent.EventExecutorGroup;
23 import io.netty.util.concurrent.FastThreadLocal;
24 import io.netty.util.internal.ObjectUtil;
25 import io.netty.util.internal.StringUtil;
26 import io.netty.util.internal.logging.InternalLogger;
27 import io.netty.util.internal.logging.InternalLoggerFactory;
28
29 import java.net.SocketAddress;
30 import java.util.ArrayList;
31 import java.util.IdentityHashMap;
32 import java.util.Iterator;
33 import java.util.LinkedHashMap;
34 import java.util.List;
35 import java.util.Map;
36 import java.util.NoSuchElementException;
37 import java.util.WeakHashMap;
38 import java.util.concurrent.RejectedExecutionException;
39 import java.util.concurrent.atomic.AtomicReferenceFieldUpdater;
40
41
42
43
44
45 public class DefaultChannelPipeline implements ChannelPipeline {
46
47 static final InternalLogger logger = InternalLoggerFactory.getInstance(DefaultChannelPipeline.class);
48
49 private static final String HEAD_NAME = generateName0(HeadContext.class);
50 private static final String TAIL_NAME = generateName0(TailContext.class);
51
52 private static final FastThreadLocal<Map<Class<?>, String>> nameCaches =
53 new FastThreadLocal<Map<Class<?>, String>>() {
54 @Override
55 protected Map<Class<?>, String> initialValue() {
56 return new WeakHashMap<Class<?>, String>();
57 }
58 };
59
60 private static final AtomicReferenceFieldUpdater<DefaultChannelPipeline, MessageSizeEstimator.Handle> ESTIMATOR =
61 AtomicReferenceFieldUpdater.newUpdater(
62 DefaultChannelPipeline.class, MessageSizeEstimator.Handle.class, "estimatorHandle");
63 final HeadContext head;
64 final TailContext tail;
65
66 private final Channel channel;
67 private final ChannelFuture succeededFuture;
68 private final VoidChannelPromise voidPromise;
69 private final boolean touch = ResourceLeakDetector.isEnabled();
70
71 private Map<EventExecutorGroup, EventExecutor> childExecutors;
72 private volatile MessageSizeEstimator.Handle estimatorHandle;
73 private boolean firstRegistration = true;
74
75
76
77
78
79
80
81
82
83 private PendingHandlerCallback pendingHandlerCallbackHead;
84
85
86
87
88
89 private boolean registered;
90
91 protected DefaultChannelPipeline(Channel channel) {
92 this.channel = ObjectUtil.checkNotNull(channel, "channel");
93 succeededFuture = new SucceededChannelFuture(channel, null);
94 voidPromise = new VoidChannelPromise(channel, true);
95
96 tail = new TailContext(this);
97 head = new HeadContext(this);
98
99 head.next = tail;
100 tail.prev = head;
101 }
102
103 final MessageSizeEstimator.Handle estimatorHandle() {
104 MessageSizeEstimator.Handle handle = estimatorHandle;
105 if (handle == null) {
106 handle = channel.config().getMessageSizeEstimator().newHandle();
107 if (!ESTIMATOR.compareAndSet(this, null, handle)) {
108 handle = estimatorHandle;
109 }
110 }
111 return handle;
112 }
113
114 final Object touch(Object msg, AbstractChannelHandlerContext next) {
115 return touch ? ReferenceCountUtil.touch(msg, next) : msg;
116 }
117
118 private AbstractChannelHandlerContext newContext(EventExecutorGroup group, String name, ChannelHandler handler) {
119 return new DefaultChannelHandlerContext(this, childExecutor(group), name, handler);
120 }
121
122 private EventExecutor childExecutor(EventExecutorGroup group) {
123 if (group == null) {
124 return null;
125 }
126 Boolean pinEventExecutor = channel.config().getOption(ChannelOption.SINGLE_EVENTEXECUTOR_PER_GROUP);
127 if (pinEventExecutor != null && !pinEventExecutor) {
128 return group.next();
129 }
130 Map<EventExecutorGroup, EventExecutor> childExecutors = this.childExecutors;
131 if (childExecutors == null) {
132
133 childExecutors = this.childExecutors = new IdentityHashMap<EventExecutorGroup, EventExecutor>(4);
134 }
135
136
137 EventExecutor childExecutor = childExecutors.get(group);
138 if (childExecutor == null) {
139 childExecutor = group.next();
140 childExecutors.put(group, childExecutor);
141 }
142 return childExecutor;
143 }
144 @Override
145 public final Channel channel() {
146 return channel;
147 }
148
149 @Override
150 public final ChannelPipeline addFirst(String name, ChannelHandler handler) {
151 return addFirst(null, name, handler);
152 }
153
154 private enum AddStrategy {
155 ADD_FIRST,
156 ADD_LAST,
157 ADD_BEFORE,
158 ADD_AFTER;
159 }
160
161 private ChannelPipeline internalAdd(EventExecutorGroup group, String name,
162 ChannelHandler handler, String baseName,
163 AddStrategy addStrategy) {
164 final AbstractChannelHandlerContext newCtx;
165 synchronized (this) {
166 checkMultiplicity(handler);
167 name = filterName(name, handler);
168
169 newCtx = newContext(group, name, handler);
170
171 switch (addStrategy) {
172 case ADD_FIRST:
173 addFirst0(newCtx);
174 break;
175 case ADD_LAST:
176 addLast0(newCtx);
177 break;
178 case ADD_BEFORE:
179 addBefore0(getContextOrDie(baseName), newCtx);
180 break;
181 case ADD_AFTER:
182 addAfter0(getContextOrDie(baseName), newCtx);
183 break;
184 default:
185 throw new IllegalArgumentException("unknown add strategy: " + addStrategy);
186 }
187
188
189
190
191 if (!registered) {
192 newCtx.setAddPending();
193 callHandlerCallbackLater(newCtx, true);
194 return this;
195 }
196
197 EventExecutor executor = newCtx.executor();
198 if (!executor.inEventLoop()) {
199 callHandlerAddedInEventLoop(newCtx, executor);
200 return this;
201 }
202 }
203 callHandlerAdded0(newCtx);
204 return this;
205 }
206
207 @Override
208 public final ChannelPipeline addFirst(EventExecutorGroup group, String name, ChannelHandler handler) {
209 return internalAdd(group, name, handler, null, AddStrategy.ADD_FIRST);
210 }
211
212 private void addFirst0(AbstractChannelHandlerContext newCtx) {
213 AbstractChannelHandlerContext nextCtx = head.next;
214 newCtx.prev = head;
215 newCtx.next = nextCtx;
216 head.next = newCtx;
217 nextCtx.prev = newCtx;
218 }
219
220 @Override
221 public final ChannelPipeline addLast(String name, ChannelHandler handler) {
222 return addLast(null, name, handler);
223 }
224
225 @Override
226 public final ChannelPipeline addLast(EventExecutorGroup group, String name, ChannelHandler handler) {
227 return internalAdd(group, name, handler, null, AddStrategy.ADD_LAST);
228 }
229
230 private void addLast0(AbstractChannelHandlerContext newCtx) {
231 AbstractChannelHandlerContext prev = tail.prev;
232 newCtx.prev = prev;
233 newCtx.next = tail;
234 prev.next = newCtx;
235 tail.prev = newCtx;
236 }
237
238 @Override
239 public final ChannelPipeline addBefore(String baseName, String name, ChannelHandler handler) {
240 return addBefore(null, baseName, name, handler);
241 }
242
243 @Override
244 public final ChannelPipeline addBefore(
245 EventExecutorGroup group, String baseName, String name, ChannelHandler handler) {
246 return internalAdd(group, name, handler, baseName, AddStrategy.ADD_BEFORE);
247 }
248
249 private static void addBefore0(AbstractChannelHandlerContext ctx, AbstractChannelHandlerContext newCtx) {
250 newCtx.prev = ctx.prev;
251 newCtx.next = ctx;
252 ctx.prev.next = newCtx;
253 ctx.prev = newCtx;
254 }
255
256 private String filterName(String name, ChannelHandler handler) {
257 if (name == null) {
258 return generateName(handler);
259 }
260 checkDuplicateName(name);
261 return name;
262 }
263
264 @Override
265 public final ChannelPipeline addAfter(String baseName, String name, ChannelHandler handler) {
266 return addAfter(null, baseName, name, handler);
267 }
268
269 @Override
270 public final ChannelPipeline addAfter(
271 EventExecutorGroup group, String baseName, String name, ChannelHandler handler) {
272 return internalAdd(group, name, handler, baseName, AddStrategy.ADD_AFTER);
273 }
274
275 private static void addAfter0(AbstractChannelHandlerContext ctx, AbstractChannelHandlerContext newCtx) {
276 newCtx.prev = ctx;
277 newCtx.next = ctx.next;
278 ctx.next.prev = newCtx;
279 ctx.next = newCtx;
280 }
281
282 public final ChannelPipeline addFirst(ChannelHandler handler) {
283 return addFirst(null, handler);
284 }
285
286 @Override
287 public final ChannelPipeline addFirst(ChannelHandler... handlers) {
288 return addFirst(null, handlers);
289 }
290
291 @Override
292 public final ChannelPipeline addFirst(EventExecutorGroup executor, ChannelHandler... handlers) {
293 ObjectUtil.checkNotNull(handlers, "handlers");
294 if (handlers.length == 0 || handlers[0] == null) {
295 return this;
296 }
297
298 int size;
299 for (size = 1; size < handlers.length; size ++) {
300 if (handlers[size] == null) {
301 break;
302 }
303 }
304
305 for (int i = size - 1; i >= 0; i --) {
306 ChannelHandler h = handlers[i];
307 addFirst(executor, null, h);
308 }
309
310 return this;
311 }
312
313 public final ChannelPipeline addLast(ChannelHandler handler) {
314 return addLast(null, handler);
315 }
316
317 @Override
318 public final ChannelPipeline addLast(ChannelHandler... handlers) {
319 return addLast(null, handlers);
320 }
321
322 @Override
323 public final ChannelPipeline addLast(EventExecutorGroup executor, ChannelHandler... handlers) {
324 ObjectUtil.checkNotNull(handlers, "handlers");
325
326 for (ChannelHandler h: handlers) {
327 if (h == null) {
328 break;
329 }
330 addLast(executor, null, h);
331 }
332
333 return this;
334 }
335
336 private String generateName(ChannelHandler handler) {
337 Map<Class<?>, String> cache = nameCaches.get();
338 Class<?> handlerType = handler.getClass();
339 String name = cache.get(handlerType);
340 if (name == null) {
341 name = generateName0(handlerType);
342 cache.put(handlerType, name);
343 }
344
345
346
347 if (context0(name) != null) {
348 String baseName = name.substring(0, name.length() - 1);
349 for (int i = 1;; i ++) {
350 String newName = baseName + i;
351 if (context0(newName) == null) {
352 name = newName;
353 break;
354 }
355 }
356 }
357 return name;
358 }
359
360 private static String generateName0(Class<?> handlerType) {
361 return StringUtil.simpleClassName(handlerType) + "#0";
362 }
363
364 @Override
365 public final ChannelPipeline remove(ChannelHandler handler) {
366 remove(getContextOrDie(handler));
367 return this;
368 }
369
370 @Override
371 public final ChannelHandler remove(String name) {
372 return remove(getContextOrDie(name)).handler();
373 }
374
375 @SuppressWarnings("unchecked")
376 @Override
377 public final <T extends ChannelHandler> T remove(Class<T> handlerType) {
378 return (T) remove(getContextOrDie(handlerType)).handler();
379 }
380
381 public final <T extends ChannelHandler> T removeIfExists(String name) {
382 return removeIfExists(context(name));
383 }
384
385 public final <T extends ChannelHandler> T removeIfExists(Class<T> handlerType) {
386 return removeIfExists(context(handlerType));
387 }
388
389 public final <T extends ChannelHandler> T removeIfExists(ChannelHandler handler) {
390 return removeIfExists(context(handler));
391 }
392
393 @SuppressWarnings("unchecked")
394 private <T extends ChannelHandler> T removeIfExists(ChannelHandlerContext ctx) {
395 if (ctx == null || !tryRemove((AbstractChannelHandlerContext) ctx)) {
396 return null;
397 }
398 return (T) ctx.handler();
399 }
400
401 private AbstractChannelHandlerContext remove(final AbstractChannelHandlerContext ctx) {
402 if (!tryRemove(ctx)) {
403
404 throw new NoSuchElementException(ctx.name());
405 }
406 return ctx;
407 }
408
409
410
411
412
413
414
415 private boolean tryRemove(final AbstractChannelHandlerContext ctx) {
416 assert ctx != head && ctx != tail;
417
418 synchronized (this) {
419 if (!atomicRemoveFromHandlerList(ctx)) {
420 return false;
421 }
422
423
424
425
426 if (!registered) {
427 callHandlerCallbackLater(ctx, false);
428 return true;
429 }
430
431 EventExecutor executor = ctx.executor();
432 if (!executor.inEventLoop()) {
433 executor.execute(new Runnable() {
434 @Override
435 public void run() {
436 callHandlerRemoved0(ctx);
437 }
438 });
439 return true;
440 }
441 }
442 callHandlerRemoved0(ctx);
443 return true;
444 }
445
446
447
448
449
450
451 private boolean isLinked(AbstractChannelHandlerContext ctx) {
452 assert Thread.holdsLock(this);
453 return ctx.prev.next == ctx;
454 }
455
456
457
458
459
460
461 private synchronized boolean atomicRemoveFromHandlerList(AbstractChannelHandlerContext ctx) {
462 if (!isLinked(ctx)) {
463
464 return false;
465 }
466 AbstractChannelHandlerContext prev = ctx.prev;
467 AbstractChannelHandlerContext next = ctx.next;
468 prev.next = next;
469 next.prev = prev;
470 return true;
471 }
472
473 @Override
474 public final ChannelHandler removeFirst() {
475 if (head.next == tail) {
476 throw new NoSuchElementException();
477 }
478 return remove(head.next).handler();
479 }
480
481 @Override
482 public final ChannelHandler removeLast() {
483 if (head.next == tail) {
484 throw new NoSuchElementException();
485 }
486 return remove(tail.prev).handler();
487 }
488
489 @Override
490 public final ChannelPipeline replace(ChannelHandler oldHandler, String newName, ChannelHandler newHandler) {
491 replace(getContextOrDie(oldHandler), newName, newHandler);
492 return this;
493 }
494
495 @Override
496 public final ChannelHandler replace(String oldName, String newName, ChannelHandler newHandler) {
497 return replace(getContextOrDie(oldName), newName, newHandler);
498 }
499
500 @Override
501 @SuppressWarnings("unchecked")
502 public final <T extends ChannelHandler> T replace(
503 Class<T> oldHandlerType, String newName, ChannelHandler newHandler) {
504 return (T) replace(getContextOrDie(oldHandlerType), newName, newHandler);
505 }
506
507 private ChannelHandler replace(
508 final AbstractChannelHandlerContext ctx, String newName, ChannelHandler newHandler) {
509 assert ctx != head && ctx != tail;
510
511 final AbstractChannelHandlerContext newCtx;
512 synchronized (this) {
513 if (!isLinked(ctx)) {
514
515 throw new NoSuchElementException(ctx.name());
516 }
517 checkMultiplicity(newHandler);
518 if (newName == null) {
519 newName = generateName(newHandler);
520 } else {
521 boolean sameName = ctx.name().equals(newName);
522 if (!sameName) {
523 checkDuplicateName(newName);
524 }
525 }
526
527 newCtx = newContext(ctx.childExecutor, newName, newHandler);
528
529 replace0(ctx, newCtx);
530
531
532
533
534
535 if (!registered) {
536 callHandlerCallbackLater(newCtx, true);
537 callHandlerCallbackLater(ctx, false);
538 return ctx.handler();
539 }
540 EventExecutor executor = ctx.executor();
541 if (!executor.inEventLoop()) {
542 executor.execute(new Runnable() {
543 @Override
544 public void run() {
545
546
547
548 callHandlerAdded0(newCtx);
549 callHandlerRemoved0(ctx);
550 }
551 });
552 return ctx.handler();
553 }
554 }
555
556
557
558 callHandlerAdded0(newCtx);
559 callHandlerRemoved0(ctx);
560 return ctx.handler();
561 }
562
563 private static void replace0(AbstractChannelHandlerContext oldCtx, AbstractChannelHandlerContext newCtx) {
564 AbstractChannelHandlerContext prev = oldCtx.prev;
565 AbstractChannelHandlerContext next = oldCtx.next;
566 newCtx.prev = prev;
567 newCtx.next = next;
568
569
570
571
572
573 prev.next = newCtx;
574 next.prev = newCtx;
575
576
577 oldCtx.prev = newCtx;
578 oldCtx.next = newCtx;
579 }
580
581 private static void checkMultiplicity(ChannelHandler handler) {
582 if (handler instanceof ChannelHandlerAdapter) {
583 ChannelHandlerAdapter h = (ChannelHandlerAdapter) handler;
584 if (!h.isSharable() && h.added) {
585 throw new ChannelPipelineException(
586 h.getClass().getName() +
587 " is not a @Sharable handler, so can't be added or removed multiple times.");
588 }
589 h.added = true;
590 }
591 }
592
593 private void callHandlerAdded0(final AbstractChannelHandlerContext ctx) {
594 try {
595 ctx.callHandlerAdded();
596 } catch (Throwable t) {
597 boolean removed = false;
598 try {
599
600
601 if (atomicRemoveFromHandlerList(ctx)) {
602 ctx.callHandlerRemoved();
603 }
604 removed = true;
605 } catch (Throwable t2) {
606 if (logger.isWarnEnabled()) {
607 logger.warn("Failed to remove a handler: " + ctx.name(), t2);
608 }
609 }
610
611 if (removed) {
612 fireExceptionCaught(new ChannelPipelineException(
613 ctx.handler().getClass().getName() +
614 ".handlerAdded() has thrown an exception; removed.", t));
615 } else {
616 fireExceptionCaught(new ChannelPipelineException(
617 ctx.handler().getClass().getName() +
618 ".handlerAdded() has thrown an exception; also failed to remove.", t));
619 }
620 }
621 }
622
623 private void callHandlerRemoved0(final AbstractChannelHandlerContext ctx) {
624
625 try {
626 ctx.callHandlerRemoved();
627 } catch (Throwable t) {
628 fireExceptionCaught(new ChannelPipelineException(
629 ctx.handler().getClass().getName() + ".handlerRemoved() has thrown an exception.", t));
630 }
631 }
632
633 final void invokeHandlerAddedIfNeeded() {
634 assert channel.eventLoop().inEventLoop();
635 if (firstRegistration) {
636 firstRegistration = false;
637
638
639 callHandlerAddedForAllHandlers();
640 }
641 }
642
643 @Override
644 public final ChannelHandler first() {
645 ChannelHandlerContext first = firstContext();
646 if (first == null) {
647 return null;
648 }
649 return first.handler();
650 }
651
652 @Override
653 public final ChannelHandlerContext firstContext() {
654 AbstractChannelHandlerContext first = head.next;
655 if (first == tail) {
656 return null;
657 }
658 return head.next;
659 }
660
661 @Override
662 public final ChannelHandler last() {
663 AbstractChannelHandlerContext last = tail.prev;
664 if (last == head) {
665 return null;
666 }
667 return last.handler();
668 }
669
670 @Override
671 public final ChannelHandlerContext lastContext() {
672 AbstractChannelHandlerContext last = tail.prev;
673 if (last == head) {
674 return null;
675 }
676 return last;
677 }
678
679 @Override
680 public final ChannelHandler get(String name) {
681 ChannelHandlerContext ctx = context(name);
682 if (ctx == null) {
683 return null;
684 } else {
685 return ctx.handler();
686 }
687 }
688
689 @SuppressWarnings("unchecked")
690 @Override
691 public final <T extends ChannelHandler> T get(Class<T> handlerType) {
692 ChannelHandlerContext ctx = context(handlerType);
693 if (ctx == null) {
694 return null;
695 } else {
696 return (T) ctx.handler();
697 }
698 }
699
700 @Override
701 public final ChannelHandlerContext context(String name) {
702 return context0(ObjectUtil.checkNotNull(name, "name"));
703 }
704
705 @Override
706 public final ChannelHandlerContext context(ChannelHandler handler) {
707 ObjectUtil.checkNotNull(handler, "handler");
708
709 AbstractChannelHandlerContext ctx = head.next;
710 for (;;) {
711
712 if (ctx == null) {
713 return null;
714 }
715
716 if (ctx.handler() == handler) {
717 return ctx;
718 }
719
720 ctx = ctx.next;
721 }
722 }
723
724 @Override
725 public final ChannelHandlerContext context(Class<? extends ChannelHandler> handlerType) {
726 ObjectUtil.checkNotNull(handlerType, "handlerType");
727
728 AbstractChannelHandlerContext ctx = head.next;
729 for (;;) {
730 if (ctx == null) {
731 return null;
732 }
733 if (handlerType.isAssignableFrom(ctx.handler().getClass())) {
734 return ctx;
735 }
736 ctx = ctx.next;
737 }
738 }
739
740 @Override
741 public final List<String> names() {
742 List<String> list = new ArrayList<String>();
743 AbstractChannelHandlerContext ctx = head.next;
744 for (;;) {
745 if (ctx == null) {
746 return list;
747 }
748 list.add(ctx.name());
749 ctx = ctx.next;
750 }
751 }
752
753 @Override
754 public final Map<String, ChannelHandler> toMap() {
755 Map<String, ChannelHandler> map = new LinkedHashMap<String, ChannelHandler>();
756 AbstractChannelHandlerContext ctx = head.next;
757 for (;;) {
758 if (ctx == tail) {
759 return map;
760 }
761 map.put(ctx.name(), ctx.handler());
762 ctx = ctx.next;
763 }
764 }
765
766 @Override
767 public final Iterator<Map.Entry<String, ChannelHandler>> iterator() {
768 return toMap().entrySet().iterator();
769 }
770
771
772
773
774 @Override
775 public final String toString() {
776 StringBuilder buf = new StringBuilder()
777 .append(StringUtil.simpleClassName(this))
778 .append('{');
779 AbstractChannelHandlerContext ctx = head.next;
780 for (;;) {
781 if (ctx == tail) {
782 break;
783 }
784
785 buf.append('(')
786 .append(ctx.name())
787 .append(" = ")
788 .append(ctx.handler().getClass().getName())
789 .append(')');
790
791 ctx = ctx.next;
792 if (ctx == tail) {
793 break;
794 }
795
796 buf.append(", ");
797 }
798 buf.append('}');
799 return buf.toString();
800 }
801
802 @Override
803 public final ChannelPipeline fireChannelRegistered() {
804 if (head.executor().inEventLoop()) {
805 if (head.invokeHandler()) {
806 head.channelRegistered(head);
807 } else {
808 head.fireChannelRegistered();
809 }
810 } else {
811 head.executor().execute(this::fireChannelRegistered);
812 }
813 return this;
814 }
815
816 @Override
817 public final ChannelPipeline fireChannelUnregistered() {
818 if (head.executor().inEventLoop()) {
819 if (head.invokeHandler()) {
820 head.channelUnregistered(head);
821 } else {
822 head.fireChannelUnregistered();
823 }
824 } else {
825 head.executor().execute(this::fireChannelUnregistered);
826 }
827 return this;
828 }
829
830
831
832
833
834
835
836
837
838
839
840 private synchronized void destroy() {
841 destroyUp(head.next, false);
842 }
843
844 private void destroyUp(AbstractChannelHandlerContext ctx, boolean inEventLoop) {
845 final Thread currentThread = Thread.currentThread();
846 final AbstractChannelHandlerContext tail = this.tail;
847 for (;;) {
848 if (ctx == tail) {
849 destroyDown(currentThread, tail.prev, inEventLoop);
850 break;
851 }
852
853 final EventExecutor executor = ctx.executor();
854 if (!inEventLoop && !executor.inEventLoop(currentThread)) {
855 final AbstractChannelHandlerContext finalCtx = ctx;
856 executor.execute(new Runnable() {
857 @Override
858 public void run() {
859 destroyUp(finalCtx, true);
860 }
861 });
862 break;
863 }
864
865 ctx = ctx.next;
866 inEventLoop = false;
867 }
868 }
869
870 private void destroyDown(Thread currentThread, AbstractChannelHandlerContext ctx, boolean inEventLoop) {
871
872 final AbstractChannelHandlerContext head = this.head;
873 for (;;) {
874 if (ctx == head) {
875 break;
876 }
877
878 final EventExecutor executor = ctx.executor();
879 if (inEventLoop || executor.inEventLoop(currentThread)) {
880
881 if (atomicRemoveFromHandlerList(ctx)) {
882 callHandlerRemoved0(ctx);
883 }
884 } else {
885 final AbstractChannelHandlerContext finalCtx = ctx;
886 executor.execute(new Runnable() {
887 @Override
888 public void run() {
889 destroyDown(Thread.currentThread(), finalCtx, true);
890 }
891 });
892 break;
893 }
894
895 ctx = ctx.prev;
896 inEventLoop = false;
897 }
898 }
899
900 @Override
901 public final ChannelPipeline fireChannelActive() {
902 if (head.executor().inEventLoop()) {
903 if (head.invokeHandler()) {
904 head.channelActive(head);
905 } else {
906 head.fireChannelActive();
907 }
908 } else {
909 head.executor().execute(this::fireChannelActive);
910 }
911 return this;
912 }
913
914 @Override
915 public final ChannelPipeline fireChannelInactive() {
916 if (head.executor().inEventLoop()) {
917 if (head.invokeHandler()) {
918 head.channelInactive(head);
919 } else {
920 head.fireChannelInactive();
921 }
922 } else {
923 head.executor().execute(this::fireChannelInactive);
924 }
925 return this;
926 }
927
928 @Override
929 public final ChannelPipeline fireExceptionCaught(Throwable cause) {
930 if (head.executor().inEventLoop()) {
931 if (head.invokeHandler()) {
932 head.exceptionCaught(head, cause);
933 } else {
934 head.fireExceptionCaught(cause);
935 }
936 } else {
937 head.executor().execute(() -> fireExceptionCaught(cause));
938 }
939 return this;
940 }
941
942 @Override
943 public final ChannelPipeline fireUserEventTriggered(Object event) {
944 if (head.executor().inEventLoop()) {
945 if (head.invokeHandler()) {
946 head.userEventTriggered(head, event);
947 } else {
948 head.fireUserEventTriggered(event);
949 }
950 } else {
951 head.executor().execute(() -> fireUserEventTriggered(event));
952 }
953 return this;
954 }
955
956 @Override
957 public final ChannelPipeline fireChannelRead(Object msg) {
958 if (head.executor().inEventLoop()) {
959 if (head.invokeHandler()) {
960 head.channelRead(head, msg);
961 } else {
962 head.fireChannelRead(msg);
963 }
964 } else {
965 head.executor().execute(() -> fireChannelRead(msg));
966 }
967 return this;
968 }
969
970 @Override
971 public final ChannelPipeline fireChannelReadComplete() {
972 if (head.executor().inEventLoop()) {
973 if (head.invokeHandler()) {
974 head.channelReadComplete(head);
975 } else {
976 head.fireChannelReadComplete();
977 }
978 } else {
979 head.executor().execute(this::fireChannelReadComplete);
980 }
981 return this;
982 }
983
984 @Override
985 public final ChannelPipeline fireChannelWritabilityChanged() {
986 if (head.executor().inEventLoop()) {
987 if (head.invokeHandler()) {
988 head.channelWritabilityChanged(head);
989 } else {
990 head.fireChannelWritabilityChanged();
991 }
992 } else {
993 head.executor().execute(this::fireChannelWritabilityChanged);
994 }
995 return this;
996 }
997
998 @Override
999 public final ChannelFuture bind(SocketAddress localAddress) {
1000 return tail.bind(localAddress);
1001 }
1002
1003 @Override
1004 public final ChannelFuture connect(SocketAddress remoteAddress) {
1005 return tail.connect(remoteAddress);
1006 }
1007
1008 @Override
1009 public final ChannelFuture connect(SocketAddress remoteAddress, SocketAddress localAddress) {
1010 return tail.connect(remoteAddress, localAddress);
1011 }
1012
1013 @Override
1014 public final ChannelFuture disconnect() {
1015 return tail.disconnect();
1016 }
1017
1018 @Override
1019 public final ChannelFuture close() {
1020 return tail.close();
1021 }
1022
1023 @Override
1024 public final ChannelFuture deregister() {
1025 return tail.deregister();
1026 }
1027
1028 @Override
1029 public final ChannelPipeline flush() {
1030 tail.flush();
1031 return this;
1032 }
1033
1034 @Override
1035 public final ChannelFuture bind(SocketAddress localAddress, ChannelPromise promise) {
1036 return tail.bind(localAddress, promise);
1037 }
1038
1039 @Override
1040 public final ChannelFuture connect(SocketAddress remoteAddress, ChannelPromise promise) {
1041 return tail.connect(remoteAddress, promise);
1042 }
1043
1044 @Override
1045 public final ChannelFuture connect(
1046 SocketAddress remoteAddress, SocketAddress localAddress, ChannelPromise promise) {
1047 return tail.connect(remoteAddress, localAddress, promise);
1048 }
1049
1050 @Override
1051 public final ChannelFuture disconnect(ChannelPromise promise) {
1052 return tail.disconnect(promise);
1053 }
1054
1055 @Override
1056 public final ChannelFuture close(ChannelPromise promise) {
1057 return tail.close(promise);
1058 }
1059
1060 @Override
1061 public final ChannelFuture deregister(final ChannelPromise promise) {
1062 return tail.deregister(promise);
1063 }
1064
1065 @Override
1066 public final ChannelPipeline read() {
1067 tail.read();
1068 return this;
1069 }
1070
1071 @Override
1072 public final ChannelFuture write(Object msg) {
1073 return tail.write(msg);
1074 }
1075
1076 @Override
1077 public final ChannelFuture write(Object msg, ChannelPromise promise) {
1078 return tail.write(msg, promise);
1079 }
1080
1081 @Override
1082 public final ChannelFuture writeAndFlush(Object msg, ChannelPromise promise) {
1083 return tail.writeAndFlush(msg, promise);
1084 }
1085
1086 @Override
1087 public final ChannelFuture writeAndFlush(Object msg) {
1088 return tail.writeAndFlush(msg);
1089 }
1090
1091 @Override
1092 public final ChannelPromise newPromise() {
1093 return new DefaultChannelPromise(channel);
1094 }
1095
1096 @Override
1097 public final ChannelProgressivePromise newProgressivePromise() {
1098 return new DefaultChannelProgressivePromise(channel);
1099 }
1100
1101 @Override
1102 public final ChannelFuture newSucceededFuture() {
1103 return succeededFuture;
1104 }
1105
1106 @Override
1107 public final ChannelFuture newFailedFuture(Throwable cause) {
1108 return new FailedChannelFuture(channel, null, cause);
1109 }
1110
1111 @Override
1112 public final ChannelPromise voidPromise() {
1113 return voidPromise;
1114 }
1115
1116 private void checkDuplicateName(String name) {
1117 if (context0(name) != null) {
1118 throw new IllegalArgumentException("Duplicate handler name: " + name);
1119 }
1120 }
1121
1122 private AbstractChannelHandlerContext context0(String name) {
1123 AbstractChannelHandlerContext context = head.next;
1124 while (context != tail) {
1125 if (context.name().equals(name)) {
1126 return context;
1127 }
1128 context = context.next;
1129 }
1130 return null;
1131 }
1132
1133 private AbstractChannelHandlerContext getContextOrDie(String name) {
1134 AbstractChannelHandlerContext ctx = (AbstractChannelHandlerContext) context(name);
1135 if (ctx == null) {
1136 throw new NoSuchElementException(name);
1137 } else {
1138 return ctx;
1139 }
1140 }
1141
1142 private AbstractChannelHandlerContext getContextOrDie(ChannelHandler handler) {
1143 AbstractChannelHandlerContext ctx = (AbstractChannelHandlerContext) context(handler);
1144 if (ctx == null) {
1145 throw new NoSuchElementException(handler.getClass().getName());
1146 } else {
1147 return ctx;
1148 }
1149 }
1150
1151 private AbstractChannelHandlerContext getContextOrDie(Class<? extends ChannelHandler> handlerType) {
1152 AbstractChannelHandlerContext ctx = (AbstractChannelHandlerContext) context(handlerType);
1153 if (ctx == null) {
1154 throw new NoSuchElementException(handlerType.getName());
1155 } else {
1156 return ctx;
1157 }
1158 }
1159
1160 private void callHandlerAddedForAllHandlers() {
1161 final PendingHandlerCallback pendingHandlerCallbackHead;
1162 synchronized (this) {
1163 assert !registered;
1164
1165
1166 registered = true;
1167
1168 pendingHandlerCallbackHead = this.pendingHandlerCallbackHead;
1169
1170 this.pendingHandlerCallbackHead = null;
1171 }
1172
1173
1174
1175
1176 PendingHandlerCallback task = pendingHandlerCallbackHead;
1177 while (task != null) {
1178 task.execute();
1179 task = task.next;
1180 }
1181 }
1182
1183 private void callHandlerCallbackLater(AbstractChannelHandlerContext ctx, boolean added) {
1184 assert !registered;
1185
1186 PendingHandlerCallback task = added ? new PendingHandlerAddedTask(ctx) : new PendingHandlerRemovedTask(ctx);
1187 PendingHandlerCallback pending = pendingHandlerCallbackHead;
1188 if (pending == null) {
1189 pendingHandlerCallbackHead = task;
1190 } else {
1191
1192 while (pending.next != null) {
1193 pending = pending.next;
1194 }
1195 pending.next = task;
1196 }
1197 }
1198
1199 private void callHandlerAddedInEventLoop(final AbstractChannelHandlerContext newCtx, EventExecutor executor) {
1200 newCtx.setAddPending();
1201 executor.execute(new Runnable() {
1202 @Override
1203 public void run() {
1204 callHandlerAdded0(newCtx);
1205 }
1206 });
1207 }
1208
1209
1210
1211
1212
1213 protected void onUnhandledInboundException(Throwable cause) {
1214 try {
1215 logger.warn(
1216 "An exceptionCaught() event was fired, and it reached at the tail of the pipeline. " +
1217 "It usually means the last handler in the pipeline did not handle the exception.",
1218 cause);
1219 } finally {
1220 ReferenceCountUtil.release(cause);
1221 }
1222 }
1223
1224
1225
1226
1227
1228 protected void onUnhandledInboundChannelActive() {
1229 }
1230
1231
1232
1233
1234
1235 protected void onUnhandledInboundChannelInactive() {
1236 }
1237
1238
1239
1240
1241
1242
1243 protected void onUnhandledInboundMessage(Object msg) {
1244 try {
1245 logger.debug(
1246 "Discarded inbound message {} that reached at the tail of the pipeline. " +
1247 "Please check your pipeline configuration.", msg);
1248 } finally {
1249 ReferenceCountUtil.release(msg);
1250 }
1251 }
1252
1253
1254
1255
1256
1257
1258 protected void onUnhandledInboundMessage(ChannelHandlerContext ctx, Object msg) {
1259 onUnhandledInboundMessage(msg);
1260 if (logger.isDebugEnabled()) {
1261 logger.debug("Discarded message pipeline : {}. Channel : {}.",
1262 ctx.pipeline().names(), ctx.channel());
1263 }
1264 }
1265
1266
1267
1268
1269
1270 protected void onUnhandledInboundChannelReadComplete() {
1271 }
1272
1273
1274
1275
1276
1277
1278 protected void onUnhandledInboundUserEventTriggered(Object evt) {
1279
1280
1281 ReferenceCountUtil.release(evt);
1282 }
1283
1284
1285
1286
1287
1288 protected void onUnhandledChannelWritabilityChanged() {
1289 }
1290
1291 protected void incrementPendingOutboundBytes(long size) {
1292 ChannelOutboundBuffer buffer = channel.unsafe().outboundBuffer();
1293 if (buffer != null) {
1294 buffer.incrementPendingOutboundBytes(size);
1295 }
1296 }
1297
1298 protected void decrementPendingOutboundBytes(long size) {
1299 ChannelOutboundBuffer buffer = channel.unsafe().outboundBuffer();
1300 if (buffer != null) {
1301 buffer.decrementPendingOutboundBytes(size);
1302 }
1303 }
1304
1305
1306 final class TailContext extends AbstractChannelHandlerContext implements ChannelInboundHandler {
1307
1308 TailContext(DefaultChannelPipeline pipeline) {
1309 super(pipeline, null, TAIL_NAME, TailContext.class);
1310 setAddComplete();
1311 }
1312
1313 @Override
1314 public ChannelHandler handler() {
1315 return this;
1316 }
1317
1318 @Override
1319 public void channelRegistered(ChannelHandlerContext ctx) { }
1320
1321 @Override
1322 public void channelUnregistered(ChannelHandlerContext ctx) { }
1323
1324 @Override
1325 public void channelActive(ChannelHandlerContext ctx) {
1326 onUnhandledInboundChannelActive();
1327 }
1328
1329 @Override
1330 public void channelInactive(ChannelHandlerContext ctx) {
1331 onUnhandledInboundChannelInactive();
1332 }
1333
1334 @Override
1335 public void channelWritabilityChanged(ChannelHandlerContext ctx) {
1336 onUnhandledChannelWritabilityChanged();
1337 }
1338
1339 @Override
1340 public void handlerAdded(ChannelHandlerContext ctx) { }
1341
1342 @Override
1343 public void handlerRemoved(ChannelHandlerContext ctx) { }
1344
1345 @Override
1346 public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
1347 onUnhandledInboundUserEventTriggered(evt);
1348 }
1349
1350 @Override
1351 public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
1352 onUnhandledInboundException(cause);
1353 }
1354
1355 @Override
1356 public void channelRead(ChannelHandlerContext ctx, Object msg) {
1357 onUnhandledInboundMessage(ctx, msg);
1358 }
1359
1360 @Override
1361 public void channelReadComplete(ChannelHandlerContext ctx) {
1362 onUnhandledInboundChannelReadComplete();
1363 }
1364 }
1365
1366 final class HeadContext extends AbstractChannelHandlerContext
1367 implements ChannelOutboundHandler, ChannelInboundHandler {
1368
1369 private final Unsafe unsafe;
1370
1371 HeadContext(DefaultChannelPipeline pipeline) {
1372 super(pipeline, null, HEAD_NAME, HeadContext.class);
1373 unsafe = pipeline.channel().unsafe();
1374 setAddComplete();
1375 }
1376
1377 @Override
1378 public ChannelHandler handler() {
1379 return this;
1380 }
1381
1382 @Override
1383 public void handlerAdded(ChannelHandlerContext ctx) {
1384
1385 }
1386
1387 @Override
1388 public void handlerRemoved(ChannelHandlerContext ctx) {
1389
1390 }
1391
1392 @Override
1393 public void bind(
1394 ChannelHandlerContext ctx, SocketAddress localAddress, ChannelPromise promise) {
1395 unsafe.bind(localAddress, promise);
1396 }
1397
1398 @Override
1399 public void connect(
1400 ChannelHandlerContext ctx,
1401 SocketAddress remoteAddress, SocketAddress localAddress,
1402 ChannelPromise promise) {
1403 unsafe.connect(remoteAddress, localAddress, promise);
1404 }
1405
1406 @Override
1407 public void disconnect(ChannelHandlerContext ctx, ChannelPromise promise) {
1408 unsafe.disconnect(promise);
1409 }
1410
1411 @Override
1412 public void close(ChannelHandlerContext ctx, ChannelPromise promise) {
1413 unsafe.close(promise);
1414 }
1415
1416 @Override
1417 public void deregister(ChannelHandlerContext ctx, ChannelPromise promise) {
1418 unsafe.deregister(promise);
1419 }
1420
1421 @Override
1422 public void read(ChannelHandlerContext ctx) {
1423 unsafe.beginRead();
1424 }
1425
1426 @Override
1427 public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) {
1428 unsafe.write(msg, promise);
1429 }
1430
1431 @Override
1432 public void flush(ChannelHandlerContext ctx) {
1433 unsafe.flush();
1434 }
1435
1436 @Override
1437 public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
1438 ctx.fireExceptionCaught(cause);
1439 }
1440
1441 @Override
1442 public void channelRegistered(ChannelHandlerContext ctx) {
1443 invokeHandlerAddedIfNeeded();
1444 ctx.fireChannelRegistered();
1445 }
1446
1447 @Override
1448 public void channelUnregistered(ChannelHandlerContext ctx) {
1449 ctx.fireChannelUnregistered();
1450
1451
1452 if (!channel.isOpen()) {
1453 destroy();
1454 }
1455 }
1456
1457 @Override
1458 public void channelActive(ChannelHandlerContext ctx) {
1459 ctx.fireChannelActive();
1460
1461 readIfIsAutoRead();
1462 }
1463
1464 @Override
1465 public void channelInactive(ChannelHandlerContext ctx) {
1466 ctx.fireChannelInactive();
1467 }
1468
1469 @Override
1470 public void channelRead(ChannelHandlerContext ctx, Object msg) {
1471 ctx.fireChannelRead(msg);
1472 }
1473
1474 @Override
1475 public void channelReadComplete(ChannelHandlerContext ctx) {
1476 ctx.fireChannelReadComplete();
1477
1478 readIfIsAutoRead();
1479 }
1480
1481 private void readIfIsAutoRead() {
1482 if (channel.config().isAutoRead()) {
1483 channel.read();
1484 }
1485 }
1486
1487 @Override
1488 public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
1489 ctx.fireUserEventTriggered(evt);
1490 }
1491
1492 @Override
1493 public void channelWritabilityChanged(ChannelHandlerContext ctx) {
1494 ctx.fireChannelWritabilityChanged();
1495 }
1496 }
1497
1498 private abstract static class PendingHandlerCallback implements Runnable {
1499 final AbstractChannelHandlerContext ctx;
1500 PendingHandlerCallback next;
1501
1502 PendingHandlerCallback(AbstractChannelHandlerContext ctx) {
1503 this.ctx = ctx;
1504 }
1505
1506 abstract void execute();
1507 }
1508
1509 private final class PendingHandlerAddedTask extends PendingHandlerCallback {
1510
1511 PendingHandlerAddedTask(AbstractChannelHandlerContext ctx) {
1512 super(ctx);
1513 }
1514
1515 @Override
1516 public void run() {
1517 callHandlerAdded0(ctx);
1518 }
1519
1520 @Override
1521 void execute() {
1522 EventExecutor executor = ctx.executor();
1523 if (executor.inEventLoop()) {
1524 callHandlerAdded0(ctx);
1525 } else {
1526 try {
1527 executor.execute(this);
1528 } catch (RejectedExecutionException e) {
1529 if (logger.isWarnEnabled()) {
1530 logger.warn(
1531 "Can't invoke handlerAdded() as the EventExecutor {} rejected it, removing handler {}.",
1532 executor, ctx.name(), e);
1533 }
1534 atomicRemoveFromHandlerList(ctx);
1535 ctx.setRemoved();
1536 }
1537 }
1538 }
1539 }
1540
1541 private final class PendingHandlerRemovedTask extends PendingHandlerCallback {
1542
1543 PendingHandlerRemovedTask(AbstractChannelHandlerContext ctx) {
1544 super(ctx);
1545 }
1546
1547 @Override
1548 public void run() {
1549 callHandlerRemoved0(ctx);
1550 }
1551
1552 @Override
1553 void execute() {
1554 EventExecutor executor = ctx.executor();
1555 if (executor.inEventLoop()) {
1556 callHandlerRemoved0(ctx);
1557 } else {
1558 try {
1559 executor.execute(this);
1560 } catch (RejectedExecutionException e) {
1561 if (logger.isWarnEnabled()) {
1562 logger.warn(
1563 "Can't invoke handlerRemoved() as the EventExecutor {} rejected it," +
1564 " removing handler {}.", executor, ctx.name(), e);
1565 }
1566
1567 ctx.setRemoved();
1568 }
1569 }
1570 }
1571 }
1572 }