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.executor, 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 AbstractChannelHandlerContext.invokeChannelRegistered(head);
805 return this;
806 }
807
808 @Override
809 public final ChannelPipeline fireChannelUnregistered() {
810 AbstractChannelHandlerContext.invokeChannelUnregistered(head);
811 return this;
812 }
813
814
815
816
817
818
819
820
821
822
823
824 private synchronized void destroy() {
825 destroyUp(head.next, false);
826 }
827
828 private void destroyUp(AbstractChannelHandlerContext ctx, boolean inEventLoop) {
829 final Thread currentThread = Thread.currentThread();
830 final AbstractChannelHandlerContext tail = this.tail;
831 for (;;) {
832 if (ctx == tail) {
833 destroyDown(currentThread, tail.prev, inEventLoop);
834 break;
835 }
836
837 final EventExecutor executor = ctx.executor();
838 if (!inEventLoop && !executor.inEventLoop(currentThread)) {
839 final AbstractChannelHandlerContext finalCtx = ctx;
840 executor.execute(new Runnable() {
841 @Override
842 public void run() {
843 destroyUp(finalCtx, true);
844 }
845 });
846 break;
847 }
848
849 ctx = ctx.next;
850 inEventLoop = false;
851 }
852 }
853
854 private void destroyDown(Thread currentThread, AbstractChannelHandlerContext ctx, boolean inEventLoop) {
855
856 final AbstractChannelHandlerContext head = this.head;
857 for (;;) {
858 if (ctx == head) {
859 break;
860 }
861
862 final EventExecutor executor = ctx.executor();
863 if (inEventLoop || executor.inEventLoop(currentThread)) {
864
865 if (atomicRemoveFromHandlerList(ctx)) {
866 callHandlerRemoved0(ctx);
867 }
868 } else {
869 final AbstractChannelHandlerContext finalCtx = ctx;
870 executor.execute(new Runnable() {
871 @Override
872 public void run() {
873 destroyDown(Thread.currentThread(), finalCtx, true);
874 }
875 });
876 break;
877 }
878
879 ctx = ctx.prev;
880 inEventLoop = false;
881 }
882 }
883
884 @Override
885 public final ChannelPipeline fireChannelActive() {
886 AbstractChannelHandlerContext.invokeChannelActive(head);
887 return this;
888 }
889
890 @Override
891 public final ChannelPipeline fireChannelInactive() {
892 AbstractChannelHandlerContext.invokeChannelInactive(head);
893 return this;
894 }
895
896 @Override
897 public final ChannelPipeline fireExceptionCaught(Throwable cause) {
898 AbstractChannelHandlerContext.invokeExceptionCaught(head, cause);
899 return this;
900 }
901
902 @Override
903 public final ChannelPipeline fireUserEventTriggered(Object event) {
904 AbstractChannelHandlerContext.invokeUserEventTriggered(head, event);
905 return this;
906 }
907
908 @Override
909 public final ChannelPipeline fireChannelRead(Object msg) {
910 AbstractChannelHandlerContext.invokeChannelRead(head, msg);
911 return this;
912 }
913
914 @Override
915 public final ChannelPipeline fireChannelReadComplete() {
916 AbstractChannelHandlerContext.invokeChannelReadComplete(head);
917 return this;
918 }
919
920 @Override
921 public final ChannelPipeline fireChannelWritabilityChanged() {
922 AbstractChannelHandlerContext.invokeChannelWritabilityChanged(head);
923 return this;
924 }
925
926 @Override
927 public final ChannelFuture bind(SocketAddress localAddress) {
928 return tail.bind(localAddress);
929 }
930
931 @Override
932 public final ChannelFuture connect(SocketAddress remoteAddress) {
933 return tail.connect(remoteAddress);
934 }
935
936 @Override
937 public final ChannelFuture connect(SocketAddress remoteAddress, SocketAddress localAddress) {
938 return tail.connect(remoteAddress, localAddress);
939 }
940
941 @Override
942 public final ChannelFuture disconnect() {
943 return tail.disconnect();
944 }
945
946 @Override
947 public final ChannelFuture close() {
948 return tail.close();
949 }
950
951 @Override
952 public final ChannelFuture deregister() {
953 return tail.deregister();
954 }
955
956 @Override
957 public final ChannelPipeline flush() {
958 tail.flush();
959 return this;
960 }
961
962 @Override
963 public final ChannelFuture bind(SocketAddress localAddress, ChannelPromise promise) {
964 return tail.bind(localAddress, promise);
965 }
966
967 @Override
968 public final ChannelFuture connect(SocketAddress remoteAddress, ChannelPromise promise) {
969 return tail.connect(remoteAddress, promise);
970 }
971
972 @Override
973 public final ChannelFuture connect(
974 SocketAddress remoteAddress, SocketAddress localAddress, ChannelPromise promise) {
975 return tail.connect(remoteAddress, localAddress, promise);
976 }
977
978 @Override
979 public final ChannelFuture disconnect(ChannelPromise promise) {
980 return tail.disconnect(promise);
981 }
982
983 @Override
984 public final ChannelFuture close(ChannelPromise promise) {
985 return tail.close(promise);
986 }
987
988 @Override
989 public final ChannelFuture deregister(final ChannelPromise promise) {
990 return tail.deregister(promise);
991 }
992
993 @Override
994 public final ChannelPipeline read() {
995 tail.read();
996 return this;
997 }
998
999 @Override
1000 public final ChannelFuture write(Object msg) {
1001 return tail.write(msg);
1002 }
1003
1004 @Override
1005 public final ChannelFuture write(Object msg, ChannelPromise promise) {
1006 return tail.write(msg, promise);
1007 }
1008
1009 @Override
1010 public final ChannelFuture writeAndFlush(Object msg, ChannelPromise promise) {
1011 return tail.writeAndFlush(msg, promise);
1012 }
1013
1014 @Override
1015 public final ChannelFuture writeAndFlush(Object msg) {
1016 return tail.writeAndFlush(msg);
1017 }
1018
1019 @Override
1020 public final ChannelPromise newPromise() {
1021 return new DefaultChannelPromise(channel);
1022 }
1023
1024 @Override
1025 public final ChannelProgressivePromise newProgressivePromise() {
1026 return new DefaultChannelProgressivePromise(channel);
1027 }
1028
1029 @Override
1030 public final ChannelFuture newSucceededFuture() {
1031 return succeededFuture;
1032 }
1033
1034 @Override
1035 public final ChannelFuture newFailedFuture(Throwable cause) {
1036 return new FailedChannelFuture(channel, null, cause);
1037 }
1038
1039 @Override
1040 public final ChannelPromise voidPromise() {
1041 return voidPromise;
1042 }
1043
1044 private void checkDuplicateName(String name) {
1045 if (context0(name) != null) {
1046 throw new IllegalArgumentException("Duplicate handler name: " + name);
1047 }
1048 }
1049
1050 private AbstractChannelHandlerContext context0(String name) {
1051 AbstractChannelHandlerContext context = head.next;
1052 while (context != tail) {
1053 if (context.name().equals(name)) {
1054 return context;
1055 }
1056 context = context.next;
1057 }
1058 return null;
1059 }
1060
1061 private AbstractChannelHandlerContext getContextOrDie(String name) {
1062 AbstractChannelHandlerContext ctx = (AbstractChannelHandlerContext) context(name);
1063 if (ctx == null) {
1064 throw new NoSuchElementException(name);
1065 } else {
1066 return ctx;
1067 }
1068 }
1069
1070 private AbstractChannelHandlerContext getContextOrDie(ChannelHandler handler) {
1071 AbstractChannelHandlerContext ctx = (AbstractChannelHandlerContext) context(handler);
1072 if (ctx == null) {
1073 throw new NoSuchElementException(handler.getClass().getName());
1074 } else {
1075 return ctx;
1076 }
1077 }
1078
1079 private AbstractChannelHandlerContext getContextOrDie(Class<? extends ChannelHandler> handlerType) {
1080 AbstractChannelHandlerContext ctx = (AbstractChannelHandlerContext) context(handlerType);
1081 if (ctx == null) {
1082 throw new NoSuchElementException(handlerType.getName());
1083 } else {
1084 return ctx;
1085 }
1086 }
1087
1088 private void callHandlerAddedForAllHandlers() {
1089 final PendingHandlerCallback pendingHandlerCallbackHead;
1090 synchronized (this) {
1091 assert !registered;
1092
1093
1094 registered = true;
1095
1096 pendingHandlerCallbackHead = this.pendingHandlerCallbackHead;
1097
1098 this.pendingHandlerCallbackHead = null;
1099 }
1100
1101
1102
1103
1104 PendingHandlerCallback task = pendingHandlerCallbackHead;
1105 while (task != null) {
1106 task.execute();
1107 task = task.next;
1108 }
1109 }
1110
1111 private void callHandlerCallbackLater(AbstractChannelHandlerContext ctx, boolean added) {
1112 assert !registered;
1113
1114 PendingHandlerCallback task = added ? new PendingHandlerAddedTask(ctx) : new PendingHandlerRemovedTask(ctx);
1115 PendingHandlerCallback pending = pendingHandlerCallbackHead;
1116 if (pending == null) {
1117 pendingHandlerCallbackHead = task;
1118 } else {
1119
1120 while (pending.next != null) {
1121 pending = pending.next;
1122 }
1123 pending.next = task;
1124 }
1125 }
1126
1127 private void callHandlerAddedInEventLoop(final AbstractChannelHandlerContext newCtx, EventExecutor executor) {
1128 newCtx.setAddPending();
1129 executor.execute(new Runnable() {
1130 @Override
1131 public void run() {
1132 callHandlerAdded0(newCtx);
1133 }
1134 });
1135 }
1136
1137
1138
1139
1140
1141 protected void onUnhandledInboundException(Throwable cause) {
1142 try {
1143 logger.warn(
1144 "An exceptionCaught() event was fired, and it reached at the tail of the pipeline. " +
1145 "It usually means the last handler in the pipeline did not handle the exception.",
1146 cause);
1147 } finally {
1148 ReferenceCountUtil.release(cause);
1149 }
1150 }
1151
1152
1153
1154
1155
1156 protected void onUnhandledInboundChannelActive() {
1157 }
1158
1159
1160
1161
1162
1163 protected void onUnhandledInboundChannelInactive() {
1164 }
1165
1166
1167
1168
1169
1170
1171 protected void onUnhandledInboundMessage(Object msg) {
1172 try {
1173 logger.debug(
1174 "Discarded inbound message {} that reached at the tail of the pipeline. " +
1175 "Please check your pipeline configuration.", msg);
1176 } finally {
1177 ReferenceCountUtil.release(msg);
1178 }
1179 }
1180
1181
1182
1183
1184
1185
1186 protected void onUnhandledInboundMessage(ChannelHandlerContext ctx, Object msg) {
1187 onUnhandledInboundMessage(msg);
1188 if (logger.isDebugEnabled()) {
1189 logger.debug("Discarded message pipeline : {}. Channel : {}.",
1190 ctx.pipeline().names(), ctx.channel());
1191 }
1192 }
1193
1194
1195
1196
1197
1198 protected void onUnhandledInboundChannelReadComplete() {
1199 }
1200
1201
1202
1203
1204
1205
1206 protected void onUnhandledInboundUserEventTriggered(Object evt) {
1207
1208
1209 ReferenceCountUtil.release(evt);
1210 }
1211
1212
1213
1214
1215
1216 protected void onUnhandledChannelWritabilityChanged() {
1217 }
1218
1219 protected void incrementPendingOutboundBytes(long size) {
1220 ChannelOutboundBuffer buffer = channel.unsafe().outboundBuffer();
1221 if (buffer != null) {
1222 buffer.incrementPendingOutboundBytes(size);
1223 }
1224 }
1225
1226 protected void decrementPendingOutboundBytes(long size) {
1227 ChannelOutboundBuffer buffer = channel.unsafe().outboundBuffer();
1228 if (buffer != null) {
1229 buffer.decrementPendingOutboundBytes(size);
1230 }
1231 }
1232
1233
1234 final class TailContext extends AbstractChannelHandlerContext implements ChannelInboundHandler {
1235
1236 TailContext(DefaultChannelPipeline pipeline) {
1237 super(pipeline, null, TAIL_NAME, TailContext.class);
1238 setAddComplete();
1239 }
1240
1241 @Override
1242 public ChannelHandler handler() {
1243 return this;
1244 }
1245
1246 @Override
1247 public void channelRegistered(ChannelHandlerContext ctx) { }
1248
1249 @Override
1250 public void channelUnregistered(ChannelHandlerContext ctx) { }
1251
1252 @Override
1253 public void channelActive(ChannelHandlerContext ctx) {
1254 onUnhandledInboundChannelActive();
1255 }
1256
1257 @Override
1258 public void channelInactive(ChannelHandlerContext ctx) {
1259 onUnhandledInboundChannelInactive();
1260 }
1261
1262 @Override
1263 public void channelWritabilityChanged(ChannelHandlerContext ctx) {
1264 onUnhandledChannelWritabilityChanged();
1265 }
1266
1267 @Override
1268 public void handlerAdded(ChannelHandlerContext ctx) { }
1269
1270 @Override
1271 public void handlerRemoved(ChannelHandlerContext ctx) { }
1272
1273 @Override
1274 public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
1275 onUnhandledInboundUserEventTriggered(evt);
1276 }
1277
1278 @Override
1279 public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
1280 onUnhandledInboundException(cause);
1281 }
1282
1283 @Override
1284 public void channelRead(ChannelHandlerContext ctx, Object msg) {
1285 onUnhandledInboundMessage(ctx, msg);
1286 }
1287
1288 @Override
1289 public void channelReadComplete(ChannelHandlerContext ctx) {
1290 onUnhandledInboundChannelReadComplete();
1291 }
1292 }
1293
1294 final class HeadContext extends AbstractChannelHandlerContext
1295 implements ChannelOutboundHandler, ChannelInboundHandler {
1296
1297 private final Unsafe unsafe;
1298
1299 HeadContext(DefaultChannelPipeline pipeline) {
1300 super(pipeline, null, HEAD_NAME, HeadContext.class);
1301 unsafe = pipeline.channel().unsafe();
1302 setAddComplete();
1303 }
1304
1305 @Override
1306 public ChannelHandler handler() {
1307 return this;
1308 }
1309
1310 @Override
1311 public void handlerAdded(ChannelHandlerContext ctx) {
1312
1313 }
1314
1315 @Override
1316 public void handlerRemoved(ChannelHandlerContext ctx) {
1317
1318 }
1319
1320 @Override
1321 public void bind(
1322 ChannelHandlerContext ctx, SocketAddress localAddress, ChannelPromise promise) {
1323 unsafe.bind(localAddress, promise);
1324 }
1325
1326 @Override
1327 public void connect(
1328 ChannelHandlerContext ctx,
1329 SocketAddress remoteAddress, SocketAddress localAddress,
1330 ChannelPromise promise) {
1331 unsafe.connect(remoteAddress, localAddress, promise);
1332 }
1333
1334 @Override
1335 public void disconnect(ChannelHandlerContext ctx, ChannelPromise promise) {
1336 unsafe.disconnect(promise);
1337 }
1338
1339 @Override
1340 public void close(ChannelHandlerContext ctx, ChannelPromise promise) {
1341 unsafe.close(promise);
1342 }
1343
1344 @Override
1345 public void deregister(ChannelHandlerContext ctx, ChannelPromise promise) {
1346 unsafe.deregister(promise);
1347 }
1348
1349 @Override
1350 public void read(ChannelHandlerContext ctx) {
1351 unsafe.beginRead();
1352 }
1353
1354 @Override
1355 public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) {
1356 unsafe.write(msg, promise);
1357 }
1358
1359 @Override
1360 public void flush(ChannelHandlerContext ctx) {
1361 unsafe.flush();
1362 }
1363
1364 @Override
1365 public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
1366 ctx.fireExceptionCaught(cause);
1367 }
1368
1369 @Override
1370 public void channelRegistered(ChannelHandlerContext ctx) {
1371 invokeHandlerAddedIfNeeded();
1372 ctx.fireChannelRegistered();
1373 }
1374
1375 @Override
1376 public void channelUnregistered(ChannelHandlerContext ctx) {
1377 ctx.fireChannelUnregistered();
1378
1379
1380 if (!channel.isOpen()) {
1381 destroy();
1382 }
1383 }
1384
1385 @Override
1386 public void channelActive(ChannelHandlerContext ctx) {
1387 ctx.fireChannelActive();
1388
1389 readIfIsAutoRead();
1390 }
1391
1392 @Override
1393 public void channelInactive(ChannelHandlerContext ctx) {
1394 ctx.fireChannelInactive();
1395 }
1396
1397 @Override
1398 public void channelRead(ChannelHandlerContext ctx, Object msg) {
1399 ctx.fireChannelRead(msg);
1400 }
1401
1402 @Override
1403 public void channelReadComplete(ChannelHandlerContext ctx) {
1404 ctx.fireChannelReadComplete();
1405
1406 readIfIsAutoRead();
1407 }
1408
1409 private void readIfIsAutoRead() {
1410 if (channel.config().isAutoRead()) {
1411 channel.read();
1412 }
1413 }
1414
1415 @Override
1416 public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
1417 ctx.fireUserEventTriggered(evt);
1418 }
1419
1420 @Override
1421 public void channelWritabilityChanged(ChannelHandlerContext ctx) {
1422 ctx.fireChannelWritabilityChanged();
1423 }
1424 }
1425
1426 private abstract static class PendingHandlerCallback implements Runnable {
1427 final AbstractChannelHandlerContext ctx;
1428 PendingHandlerCallback next;
1429
1430 PendingHandlerCallback(AbstractChannelHandlerContext ctx) {
1431 this.ctx = ctx;
1432 }
1433
1434 abstract void execute();
1435 }
1436
1437 private final class PendingHandlerAddedTask extends PendingHandlerCallback {
1438
1439 PendingHandlerAddedTask(AbstractChannelHandlerContext ctx) {
1440 super(ctx);
1441 }
1442
1443 @Override
1444 public void run() {
1445 callHandlerAdded0(ctx);
1446 }
1447
1448 @Override
1449 void execute() {
1450 EventExecutor executor = ctx.executor();
1451 if (executor.inEventLoop()) {
1452 callHandlerAdded0(ctx);
1453 } else {
1454 try {
1455 executor.execute(this);
1456 } catch (RejectedExecutionException e) {
1457 if (logger.isWarnEnabled()) {
1458 logger.warn(
1459 "Can't invoke handlerAdded() as the EventExecutor {} rejected it, removing handler {}.",
1460 executor, ctx.name(), e);
1461 }
1462 atomicRemoveFromHandlerList(ctx);
1463 ctx.setRemoved();
1464 }
1465 }
1466 }
1467 }
1468
1469 private final class PendingHandlerRemovedTask extends PendingHandlerCallback {
1470
1471 PendingHandlerRemovedTask(AbstractChannelHandlerContext ctx) {
1472 super(ctx);
1473 }
1474
1475 @Override
1476 public void run() {
1477 callHandlerRemoved0(ctx);
1478 }
1479
1480 @Override
1481 void execute() {
1482 EventExecutor executor = ctx.executor();
1483 if (executor.inEventLoop()) {
1484 callHandlerRemoved0(ctx);
1485 } else {
1486 try {
1487 executor.execute(this);
1488 } catch (RejectedExecutionException e) {
1489 if (logger.isWarnEnabled()) {
1490 logger.warn(
1491 "Can't invoke handlerRemoved() as the EventExecutor {} rejected it," +
1492 " removing handler {}.", executor, ctx.name(), e);
1493 }
1494
1495 ctx.setRemoved();
1496 }
1497 }
1498 }
1499 }
1500 }