1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16 package io.netty.handler.codec.redis;
17
18 import io.netty.channel.ChannelHandlerContext;
19 import io.netty.handler.codec.CodecException;
20 import io.netty.handler.codec.MessageToMessageDecoder;
21 import io.netty.handler.codec.PrematureChannelClosureException;
22 import io.netty.util.ReferenceCountUtil;
23 import io.netty.util.internal.ObjectUtil;
24 import io.netty.util.internal.UnstableApi;
25
26 import java.util.ArrayDeque;
27 import java.util.ArrayList;
28 import java.util.Deque;
29 import java.util.List;
30
31
32
33
34
35 @UnstableApi
36 public final class RedisArrayAggregator extends MessageToMessageDecoder<RedisMessage> {
37
38 private static final int DEFAULT_MAX_ARRAY_LENGTH = RedisConstants.REDIS_MAX_ARRAY_LENGTH;
39
40
41
42
43 private static final int INITIAL_CHILDREN_CAPACITY = 32;
44
45 private final int maxNestedArrayDepth;
46 private final Deque<AggregateState> depths = new ArrayDeque<AggregateState>(4);
47 private final int maxElements;
48
49
50
51
52
53 private long pendingElements;
54
55
56
57
58
59
60
61
62
63
64 @Deprecated
65 public RedisArrayAggregator() {
66
67 this(DEFAULT_MAX_ARRAY_LENGTH, 1024);
68 }
69
70
71
72
73
74
75
76
77
78
79 public RedisArrayAggregator(int maxElements, int maxNestedArrayDepth) {
80 super(RedisMessage.class);
81 this.maxElements = ObjectUtil.checkPositive(maxElements, "maxElements");
82 this.maxNestedArrayDepth = ObjectUtil.checkPositive(maxNestedArrayDepth, "maxNestedArrayDepth");
83 }
84
85 @Override
86 protected void decode(ChannelHandlerContext ctx, RedisMessage msg, List<Object> out) throws Exception {
87 if (msg instanceof ArrayHeaderRedisMessage) {
88 msg = decodeRedisArrayHeader((ArrayHeaderRedisMessage) msg);
89 if (msg == null) {
90 return;
91 }
92 } else {
93 ReferenceCountUtil.retain(msg);
94 }
95
96 while (!depths.isEmpty()) {
97 AggregateState current = depths.peek();
98 current.children.add(msg);
99
100
101 if (current.children.size() == current.length) {
102 msg = new ArrayRedisMessage(current.children);
103 depths.pop();
104 pendingElements -= current.length;
105 } else {
106
107 return;
108 }
109 }
110
111 out.add(msg);
112 }
113
114 private CodecException clearAndCreateException(String msg) {
115 releaseAndClearDepths();
116 return new CodecException(msg);
117 }
118
119 private RedisMessage decodeRedisArrayHeader(ArrayHeaderRedisMessage header) {
120 if (header.isNull()) {
121 return ArrayRedisMessage.NULL_INSTANCE;
122 } else if (header.length() == 0L) {
123 return ArrayRedisMessage.EMPTY_INSTANCE;
124 } else if (header.length() > 0L) {
125
126 if (header.length() > maxElements) {
127 throw clearAndCreateException("this codec doesn't support longer length than " + maxElements);
128 }
129
130 if (depths.size() >= maxNestedArrayDepth) {
131 throw clearAndCreateException("max nested array depth exceeded: " + maxNestedArrayDepth);
132 }
133
134
135
136
137
138
139 long newPendingElements = pendingElements + header.length();
140 if (newPendingElements > maxElements) {
141 throw clearAndCreateException(
142 "total outstanding array elements exceeds " + maxElements);
143 }
144 pendingElements = newPendingElements;
145
146
147 depths.push(new AggregateState((int) header.length()));
148 return null;
149 } else {
150 throw clearAndCreateException("bad length: " + header.length());
151 }
152 }
153
154 private static final class AggregateState {
155 private final int length;
156 private final List<RedisMessage> children;
157 AggregateState(int length) {
158 this.length = length;
159 this.children = new ArrayList<RedisMessage>(Math.min(length, INITIAL_CHILDREN_CAPACITY));
160 }
161 }
162
163 @Override
164 public void handlerRemoved(ChannelHandlerContext ctx) throws Exception {
165 super.handlerRemoved(ctx);
166 releaseAndClearDepths();
167 }
168
169 private void releaseAndClearDepths() {
170 for (AggregateState state : depths) {
171 for (RedisMessage message : state.children) {
172 ReferenceCountUtil.safeRelease(message);
173 }
174 }
175 depths.clear();
176 pendingElements = 0;
177 }
178
179 @Override
180 public void channelInactive(ChannelHandlerContext ctx) throws Exception {
181 super.channelInactive(ctx);
182
183 if (!depths.isEmpty()) {
184 ctx.fireExceptionCaught(new PrematureChannelClosureException(
185 "channel gone inactive with " + depths.size() +
186 " messages still incomplete"));
187 }
188 }
189 }