Created
February 28, 2014 10:17
-
-
Save normanmaurer/9268632 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| /* | |
| * Copyright 2012 The Netty Project | |
| * | |
| * The Netty Project licenses this file to you under the Apache License, | |
| * version 2.0 (the "License"); you may not use this file except in compliance | |
| * with the License. You may obtain a copy of the License at: | |
| * | |
| * http://www.apache.org/licenses/LICENSE-2.0 | |
| * | |
| * Unless required by applicable law or agreed to in writing, software | |
| * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT | |
| * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the | |
| * License for the specific language governing permissions and limitations | |
| * under the License. | |
| */ | |
| package io.netty.handler.codec; | |
| import io.netty.buffer.ByteBuf; | |
| import io.netty.channel.ChannelFlushPromiseNotifier; | |
| import io.netty.channel.ChannelFuture; | |
| import io.netty.channel.ChannelFutureListener; | |
| import io.netty.channel.ChannelHandlerContext; | |
| import io.netty.channel.ChannelPromise; | |
| import io.netty.util.ReferenceCountUtil; | |
| public abstract class MessageToBufferedByteEncoder<I> extends MessageToByteEncoder<I> { | |
| private final ChannelFlushPromiseNotifier notifier = new ChannelFlushPromiseNotifier(); | |
| private final int bufferSize; | |
| private ByteBuf buffer; | |
| public MessageToBufferedByteEncoder() { | |
| this(1024); | |
| } | |
| public MessageToBufferedByteEncoder(int bufferSize) { | |
| this.bufferSize = bufferSize; | |
| } | |
| public MessageToBufferedByteEncoder(Class<? extends I> outboundMessageType) { | |
| this(outboundMessageType, 1024); | |
| } | |
| public MessageToBufferedByteEncoder(Class<? extends I> outboundMessageType, int bufferSize) { | |
| super(outboundMessageType); | |
| this.bufferSize = bufferSize; | |
| } | |
| public MessageToBufferedByteEncoder(boolean preferDirect) { | |
| this(preferDirect, 1024); | |
| } | |
| public MessageToBufferedByteEncoder(boolean preferDirect, int bufferSize) { | |
| super(preferDirect); | |
| this.bufferSize = bufferSize; | |
| } | |
| public MessageToBufferedByteEncoder(Class<? extends I> outboundMessageType, boolean preferDirect) { | |
| this(outboundMessageType, preferDirect, 1024); | |
| } | |
| public MessageToBufferedByteEncoder(Class<? extends I> outboundMessageType, boolean preferDirect, int bufferSize) { | |
| super(outboundMessageType, preferDirect); | |
| this.bufferSize = bufferSize; | |
| } | |
| @Override | |
| public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception { | |
| try { | |
| if (acceptOutboundMessage(msg)) { | |
| @SuppressWarnings("unchecked") | |
| I cast = (I) msg; | |
| if (buffer == null) { | |
| buffer = newBuffer(ctx, msg, preferDirect, bufferSize); | |
| } | |
| try { | |
| encode(ctx, cast, buffer); | |
| } finally { | |
| ReferenceCountUtil.release(cast); | |
| } | |
| } else { | |
| writeBufferedData(ctx); | |
| ctx.write(msg, promise); | |
| } | |
| } catch (EncoderException e) { | |
| throw e; | |
| } catch (Throwable e) { | |
| throw new EncoderException(e); | |
| } | |
| } | |
| protected final void writeBufferedData(ChannelHandlerContext ctx) { | |
| if (buffer == null) { | |
| return; | |
| } | |
| ByteBuf buf = buffer; | |
| final int length = buf.readableBytes(); | |
| buffer = null; | |
| ctx.write(buf).addListener(new ChannelFutureListener() { | |
| @Override | |
| public void operationComplete(ChannelFuture future) throws Exception { | |
| notifier.increaseWriteCounter(length); | |
| if (future.isSuccess()) { | |
| notifier.notifyFlushFutures(); | |
| } else { | |
| notifier.notifyFlushFutures(future.cause()); | |
| } | |
| } | |
| }); | |
| } | |
| protected ByteBuf newBuffer(ChannelHandlerContext ctx, @SuppressWarnings("unused") Object msg, | |
| boolean preferDirect, int preferSize) { | |
| if (preferDirect) { | |
| return ctx.alloc().ioBuffer(preferSize); | |
| } else { | |
| return ctx.alloc().heapBuffer(preferSize); | |
| } | |
| } | |
| @Override | |
| public void flush(ChannelHandlerContext ctx) throws Exception { | |
| writeBufferedData(ctx); | |
| super.flush(ctx); | |
| } | |
| @Override | |
| public void handlerRemoved(ChannelHandlerContext ctx) throws Exception { | |
| writeBufferedData(ctx); | |
| super.handlerRemoved(ctx); | |
| } | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment