Skip to content

Instantly share code, notes, and snippets.

@normanmaurer
Created February 28, 2014 10:17
Show Gist options
  • Select an option

  • Save normanmaurer/9268632 to your computer and use it in GitHub Desktop.

Select an option

Save normanmaurer/9268632 to your computer and use it in GitHub Desktop.
/*
* 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