View Javadoc
1   /*
2    * Copyright 2014 The Netty Project
3    *
4    * The Netty Project licenses this file to you under the Apache License,
5    * version 2.0 (the "License"); you may not use this file except in compliance
6    * with the License. You may obtain a copy of the License at:
7    *
8    *   https://www.apache.org/licenses/LICENSE-2.0
9    *
10   * Unless required by applicable law or agreed to in writing, software
11   * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
12   * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
13   * License for the specific language governing permissions and limitations
14   * under the License.
15   */
16  package io.netty.handler.codec.http.websocketx.extensions.compression;
17  
18  import io.netty.buffer.ByteBuf;
19  import io.netty.buffer.CompositeByteBuf;
20  import io.netty.channel.ChannelHandler;
21  import io.netty.channel.ChannelHandlerContext;
22  import io.netty.channel.embedded.EmbeddedChannel;
23  import io.netty.handler.codec.CodecException;
24  import io.netty.handler.codec.compression.ZlibCodecFactory;
25  import io.netty.handler.codec.compression.ZlibWrapper;
26  import io.netty.handler.codec.http.websocketx.BinaryWebSocketFrame;
27  import io.netty.handler.codec.http.websocketx.ContinuationWebSocketFrame;
28  import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;
29  import io.netty.handler.codec.http.websocketx.WebSocketFrame;
30  import io.netty.handler.codec.http.websocketx.extensions.WebSocketExtensionEncoder;
31  import io.netty.handler.codec.http.websocketx.extensions.WebSocketExtensionFilter;
32  
33  import java.util.List;
34  
35  import static io.netty.handler.codec.http.websocketx.extensions.compression.PerMessageDeflateDecoder.*;
36  import static io.netty.util.internal.ObjectUtil.*;
37  
38  /**
39   * Deflate implementation of a payload compressor for
40   * <tt>io.netty.handler.codec.http.websocketx.WebSocketFrame</tt>.
41   */
42  abstract class DeflateEncoder extends WebSocketExtensionEncoder {
43  
44      static final int DEFAULT_MEM_LEVEL = 8;
45  
46      private final int compressionLevel;
47      private final int windowSize;
48      private final int memLevel;
49      private final boolean noContext;
50      private final WebSocketExtensionFilter extensionEncoderFilter;
51  
52      private EmbeddedChannel encoder;
53  
54      /**
55       * Constructor
56       * @param compressionLevel compression level of the compressor.
57       * @param windowSize maximum size of the window compressor buffer.
58       * @param noContext true to disable context takeover.
59       * @param extensionEncoderFilter extension encoder filter.
60       */
61      DeflateEncoder(int compressionLevel, int windowSize, boolean noContext,
62                     WebSocketExtensionFilter extensionEncoderFilter) {
63          this(compressionLevel, windowSize, DEFAULT_MEM_LEVEL, noContext, extensionEncoderFilter);
64      }
65  
66      /**
67       * Constructor
68       * @param compressionLevel compression level of the compressor.
69       * @param windowSize maximum size of the window compressor buffer.
70       * @param memLevel internal compression state memory level (1..9).
71       * @param noContext true to disable context takeover.
72       * @param extensionEncoderFilter extension encoder filter.
73       */
74      DeflateEncoder(int compressionLevel, int windowSize, int memLevel, boolean noContext,
75                     WebSocketExtensionFilter extensionEncoderFilter) {
76          this.compressionLevel = compressionLevel;
77          this.windowSize = windowSize;
78          this.memLevel = memLevel;
79          this.noContext = noContext;
80          this.extensionEncoderFilter = checkNotNull(extensionEncoderFilter, "extensionEncoderFilter");
81      }
82  
83      /**
84       * Returns the extension encoder filter.
85       */
86      protected WebSocketExtensionFilter extensionEncoderFilter() {
87          return extensionEncoderFilter;
88      }
89  
90      /**
91       * @param msg the current frame.
92       * @return the rsv bits to set in the compressed frame.
93       */
94      protected abstract int rsv(WebSocketFrame msg);
95  
96      /**
97       * @param msg the current frame.
98       * @return true if compressed payload tail needs to be removed.
99       */
100     protected abstract boolean removeFrameTail(WebSocketFrame msg);
101 
102     @Override
103     protected void encode(ChannelHandlerContext ctx, WebSocketFrame msg, List<Object> out) throws Exception {
104         final ByteBuf compressedContent;
105         if (msg.content().isReadable()) {
106             compressedContent = compressContent(ctx, msg);
107         } else if (msg.isFinalFragment()) {
108             // Set empty DEFLATE block manually for unknown buffer size
109             // https://tools.ietf.org/html/rfc7692#section-7.2.3.6
110             compressedContent = EMPTY_DEFLATE_BLOCK.duplicate();
111         } else {
112             throw new CodecException("cannot compress content buffer");
113         }
114 
115         final WebSocketFrame outMsg;
116         if (msg instanceof TextWebSocketFrame) {
117             outMsg = new TextWebSocketFrame(msg.isFinalFragment(), rsv(msg), compressedContent);
118         } else if (msg instanceof BinaryWebSocketFrame) {
119             outMsg = new BinaryWebSocketFrame(msg.isFinalFragment(), rsv(msg), compressedContent);
120         } else if (msg instanceof ContinuationWebSocketFrame) {
121             outMsg = new ContinuationWebSocketFrame(msg.isFinalFragment(), rsv(msg), compressedContent);
122         } else {
123             throw new CodecException("unexpected frame type: " + msg.getClass().getName());
124         }
125 
126         out.add(outMsg);
127     }
128 
129     @Override
130     public void handlerRemoved(ChannelHandlerContext ctx) throws Exception {
131         cleanup();
132         super.handlerRemoved(ctx);
133     }
134 
135     private ByteBuf compressContent(ChannelHandlerContext ctx, WebSocketFrame msg) {
136         if (encoder == null) {
137             encoder = EmbeddedChannel.builder()
138                     .handlers(ZlibCodecFactory.newZlibEncoder(
139                             ZlibWrapper.NONE, compressionLevel, windowSize, memLevel))
140                     .build();
141         }
142 
143         encoder.writeOutbound(msg.content().retain());
144 
145         CompositeByteBuf fullCompressedContent = ctx.alloc().compositeBuffer();
146         for (;;) {
147             ByteBuf partCompressedContent = encoder.readOutbound();
148             if (partCompressedContent == null) {
149                 break;
150             }
151             if (!partCompressedContent.isReadable()) {
152                 partCompressedContent.release();
153                 continue;
154             }
155             fullCompressedContent.addComponent(true, partCompressedContent);
156         }
157 
158         if (fullCompressedContent.numComponents() <= 0) {
159             fullCompressedContent.release();
160             throw new CodecException("cannot read compressed buffer");
161         }
162 
163         if (msg.isFinalFragment() && noContext) {
164             cleanup();
165         }
166 
167         ByteBuf compressedContent;
168         if (removeFrameTail(msg)) {
169             int realLength = fullCompressedContent.readableBytes() - FRAME_TAIL.readableBytes();
170             compressedContent = fullCompressedContent.slice(0, realLength);
171         } else {
172             compressedContent = fullCompressedContent;
173         }
174 
175         return compressedContent;
176     }
177 
178     private void cleanup() {
179         if (encoder != null) {
180             // Clean-up the previous encoder if not cleaned up correctly.
181             encoder.finishAndReleaseAll();
182             encoder = null;
183         }
184     }
185 }