diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/processor/RequestProcessor.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/processor/RequestProcessor.java
index 50b97e906ba..5a4238e64d6 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/processor/RequestProcessor.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/processor/RequestProcessor.java
@@ -20,7 +20,7 @@
*/
package org.apache.bookkeeper.processor;
-import io.netty.channel.Channel;
+import org.apache.bookkeeper.proto.BookieRequestHandler;
/**
* A request processor that is used for processing requests at bookie side.
@@ -41,5 +41,5 @@ public interface RequestProcessor extends AutoCloseable {
* @param channel
* channel received the given request r
*/
- void processRequest(Object r, Channel channel);
+ void processRequest(Object r, BookieRequestHandler channel);
}
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java
index 93be69cd377..c9d65a73174 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java
@@ -32,20 +32,27 @@
/**
* Serverside handler for bookkeeper requests.
*/
-class BookieRequestHandler extends ChannelInboundHandlerAdapter {
+public class BookieRequestHandler extends ChannelInboundHandlerAdapter {
private static final Logger LOG = LoggerFactory.getLogger(BookieRequestHandler.class);
private final RequestProcessor requestProcessor;
private final ChannelGroup allChannels;
+ private ChannelHandlerContext ctx;
+
BookieRequestHandler(ServerConfiguration conf, RequestProcessor processor, ChannelGroup allChannels) {
this.requestProcessor = processor;
this.allChannels = allChannels;
}
+ public ChannelHandlerContext ctx() {
+ return ctx;
+ }
+
@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
LOG.info("Channel connected {}", ctx.channel());
+ this.ctx = ctx;
super.channelActive(ctx);
}
@@ -75,6 +82,6 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception
ctx.fireChannelRead(msg);
return;
}
- requestProcessor.processRequest(msg, ctx.channel());
+ requestProcessor.processRequest(msg, this);
}
}
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java
index f7f4eceda30..9237c451ed6 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java
@@ -299,7 +299,8 @@ private void shutdownExecutor(OrderedExecutor service) {
}
@Override
- public void processRequest(Object msg, Channel c) {
+ public void processRequest(Object msg, BookieRequestHandler requestHandler) {
+ Channel channel = requestHandler.ctx().channel();
// If we can decode this packet as a Request protobuf packet, process
// it as a version 3 packet. Else, just use the old protocol.
if (msg instanceof BookkeeperProtocol.Request) {
@@ -309,16 +310,16 @@ public void processRequest(Object msg, Channel c) {
BookkeeperProtocol.BKPacketHeader header = r.getHeader();
switch (header.getOperation()) {
case ADD_ENTRY:
- processAddRequestV3(r, c);
+ processAddRequestV3(r, requestHandler);
break;
case READ_ENTRY:
- processReadRequestV3(r, c);
+ processReadRequestV3(r, requestHandler);
break;
case FORCE_LEDGER:
- processForceLedgerRequestV3(r, c);
+ processForceLedgerRequestV3(r, requestHandler);
break;
case AUTH:
- LOG.info("Ignoring auth operation from client {}", c.remoteAddress());
+ LOG.info("Ignoring auth operation from client {}", channel.remoteAddress());
BookkeeperProtocol.AuthMessage message = BookkeeperProtocol.AuthMessage
.newBuilder()
.setAuthPluginName(AuthProviderFactoryFactory.AUTHENTICATION_DISABLED_PLUGIN_NAME)
@@ -328,29 +329,29 @@ public void processRequest(Object msg, Channel c) {
.newBuilder().setHeader(r.getHeader())
.setStatus(BookkeeperProtocol.StatusCode.EOK)
.setAuthResponse(message);
- c.writeAndFlush(authResponse.build());
+ channel.writeAndFlush(authResponse.build());
break;
case WRITE_LAC:
- processWriteLacRequestV3(r, c);
+ processWriteLacRequestV3(r, requestHandler);
break;
case READ_LAC:
- processReadLacRequestV3(r, c);
+ processReadLacRequestV3(r, requestHandler);
break;
case GET_BOOKIE_INFO:
- processGetBookieInfoRequestV3(r, c);
+ processGetBookieInfoRequestV3(r, requestHandler);
break;
case START_TLS:
- processStartTLSRequestV3(r, c);
+ processStartTLSRequestV3(r, requestHandler);
break;
case GET_LIST_OF_ENTRIES_OF_LEDGER:
- processGetListOfEntriesOfLedgerProcessorV3(r, c);
+ processGetListOfEntriesOfLedgerProcessorV3(r, requestHandler);
break;
default:
LOG.info("Unknown operation type {}", header.getOperation());
BookkeeperProtocol.Response.Builder response =
BookkeeperProtocol.Response.newBuilder().setHeader(r.getHeader())
.setStatus(BookkeeperProtocol.StatusCode.EBADREQ);
- c.writeAndFlush(response.build());
+ channel.writeAndFlush(response.build());
if (statsEnabled) {
bkStats.getOpStats(BKStats.STATS_UNKNOWN).incrementFailedOps();
}
@@ -365,26 +366,27 @@ public void processRequest(Object msg, Channel c) {
switch (r.getOpCode()) {
case BookieProtocol.ADDENTRY:
checkArgument(r instanceof BookieProtocol.ParsedAddRequest);
- processAddRequest((BookieProtocol.ParsedAddRequest) r, c);
+ processAddRequest((BookieProtocol.ParsedAddRequest) r, requestHandler);
break;
case BookieProtocol.READENTRY:
checkArgument(r instanceof BookieProtocol.ReadRequest);
- processReadRequest((BookieProtocol.ReadRequest) r, c);
+ processReadRequest((BookieProtocol.ReadRequest) r, requestHandler);
break;
case BookieProtocol.AUTH:
- LOG.info("Ignoring auth operation from client {}", c.remoteAddress());
+ LOG.info("Ignoring auth operation from client {}",
+ requestHandler.ctx().channel().remoteAddress());
BookkeeperProtocol.AuthMessage message = BookkeeperProtocol.AuthMessage
.newBuilder()
.setAuthPluginName(AuthProviderFactoryFactory.AUTHENTICATION_DISABLED_PLUGIN_NAME)
.setPayload(ByteString.copyFrom(AuthToken.NULL.getData()))
.build();
- c.writeAndFlush(new BookieProtocol.AuthResponse(
+ channel.writeAndFlush(new BookieProtocol.AuthResponse(
BookieProtocol.CURRENT_PROTOCOL_VERSION, message));
break;
default:
LOG.error("Unknown op type {}, sending error", r.getOpCode());
- c.writeAndFlush(ResponseBuilder.buildErrorResponse(BookieProtocol.EBADREQ, r));
+ channel.writeAndFlush(ResponseBuilder.buildErrorResponse(BookieProtocol.EBADREQ, r));
if (statsEnabled) {
bkStats.getOpStats(BKStats.STATS_UNKNOWN).incrementFailedOps();
}
@@ -402,8 +404,9 @@ private void restoreMdcContextFromRequest(BookkeeperProtocol.Request req) {
}
}
- private void processWriteLacRequestV3(final BookkeeperProtocol.Request r, final Channel c) {
- WriteLacProcessorV3 writeLac = new WriteLacProcessorV3(r, c, this);
+ private void processWriteLacRequestV3(final BookkeeperProtocol.Request r,
+ final BookieRequestHandler requestHandler) {
+ WriteLacProcessorV3 writeLac = new WriteLacProcessorV3(r, requestHandler, this);
if (null == writeThreadPool) {
writeLac.run();
} else {
@@ -411,8 +414,9 @@ private void processWriteLacRequestV3(final BookkeeperProtocol.Request r, final
}
}
- private void processReadLacRequestV3(final BookkeeperProtocol.Request r, final Channel c) {
- ReadLacProcessorV3 readLac = new ReadLacProcessorV3(r, c, this);
+ private void processReadLacRequestV3(final BookkeeperProtocol.Request r,
+ final BookieRequestHandler requestHandler) {
+ ReadLacProcessorV3 readLac = new ReadLacProcessorV3(r, requestHandler, this);
if (null == readThreadPool) {
readLac.run();
} else {
@@ -420,8 +424,8 @@ private void processReadLacRequestV3(final BookkeeperProtocol.Request r, final C
}
}
- private void processAddRequestV3(final BookkeeperProtocol.Request r, final Channel c) {
- WriteEntryProcessorV3 write = new WriteEntryProcessorV3(r, c, this);
+ private void processAddRequestV3(final BookkeeperProtocol.Request r, final BookieRequestHandler requestHandler) {
+ WriteEntryProcessorV3 write = new WriteEntryProcessorV3(r, requestHandler, this);
final OrderedExecutor threadPool;
if (RequestUtils.isHighPriority(r)) {
@@ -455,8 +459,9 @@ private void processAddRequestV3(final BookkeeperProtocol.Request r, final Chann
}
}
- private void processForceLedgerRequestV3(final BookkeeperProtocol.Request r, final Channel c) {
- ForceLedgerProcessorV3 forceLedger = new ForceLedgerProcessorV3(r, c, this);
+ private void processForceLedgerRequestV3(final BookkeeperProtocol.Request r,
+ final BookieRequestHandler requestHandler) {
+ ForceLedgerProcessorV3 forceLedger = new ForceLedgerProcessorV3(r, requestHandler, this);
final OrderedExecutor threadPool;
if (RequestUtils.isHighPriority(r)) {
@@ -492,19 +497,20 @@ private void processForceLedgerRequestV3(final BookkeeperProtocol.Request r, fin
}
}
- private void processReadRequestV3(final BookkeeperProtocol.Request r, final Channel c) {
- ExecutorService fenceThread = null == highPriorityThreadPool ? null : highPriorityThreadPool.chooseThread(c);
+ private void processReadRequestV3(final BookkeeperProtocol.Request r, final BookieRequestHandler requestHandler) {
+ ExecutorService fenceThread = null == highPriorityThreadPool ? null :
+ highPriorityThreadPool.chooseThread(requestHandler.ctx());
final ReadEntryProcessorV3 read;
final OrderedExecutor threadPool;
if (RequestUtils.isLongPollReadRequest(r.getReadRequest())) {
- ExecutorService lpThread = longPollThreadPool.chooseThread(c);
+ ExecutorService lpThread = longPollThreadPool.chooseThread(requestHandler.ctx());
- read = new LongPollReadEntryProcessorV3(r, c, this, fenceThread,
+ read = new LongPollReadEntryProcessorV3(r, requestHandler, this, fenceThread,
lpThread, requestTimer);
threadPool = longPollThreadPool;
} else {
- read = new ReadEntryProcessorV3(r, c, this, fenceThread);
+ read = new ReadEntryProcessorV3(r, requestHandler, this, fenceThread);
// If it's a high priority read (fencing or as part of recovery process), we want to make sure it
// gets executed as fast as possible, so bypass the normal readThreadPool
@@ -544,13 +550,16 @@ private void processReadRequestV3(final BookkeeperProtocol.Request r, final Chan
}
}
- private void processStartTLSRequestV3(final BookkeeperProtocol.Request r, final Channel c) {
+ private void processStartTLSRequestV3(final BookkeeperProtocol.Request r,
+ final BookieRequestHandler requestHandler) {
BookkeeperProtocol.Response.Builder response = BookkeeperProtocol.Response.newBuilder();
BookkeeperProtocol.BKPacketHeader.Builder header = BookkeeperProtocol.BKPacketHeader.newBuilder();
header.setVersion(BookkeeperProtocol.ProtocolVersion.VERSION_THREE);
header.setOperation(r.getHeader().getOperation());
header.setTxnId(r.getHeader().getTxnId());
response.setHeader(header.build());
+ final Channel c = requestHandler.ctx().channel();
+
if (shFactory == null) {
LOG.error("Got StartTLS request but TLS not configured");
response.setStatus(BookkeeperProtocol.StatusCode.EBADREQ);
@@ -596,8 +605,9 @@ public void operationComplete(Future future) throws Exception {
}
}
- private void processGetBookieInfoRequestV3(final BookkeeperProtocol.Request r, final Channel c) {
- GetBookieInfoProcessorV3 getBookieInfo = new GetBookieInfoProcessorV3(r, c, this);
+ private void processGetBookieInfoRequestV3(final BookkeeperProtocol.Request r,
+ final BookieRequestHandler requestHandler) {
+ GetBookieInfoProcessorV3 getBookieInfo = new GetBookieInfoProcessorV3(r, requestHandler, this);
if (null == readThreadPool) {
getBookieInfo.run();
} else {
@@ -605,9 +615,10 @@ private void processGetBookieInfoRequestV3(final BookkeeperProtocol.Request r, f
}
}
- private void processGetListOfEntriesOfLedgerProcessorV3(final BookkeeperProtocol.Request r, final Channel c) {
- GetListOfEntriesOfLedgerProcessorV3 getListOfEntriesOfLedger = new GetListOfEntriesOfLedgerProcessorV3(r, c,
- this);
+ private void processGetListOfEntriesOfLedgerProcessorV3(final BookkeeperProtocol.Request r,
+ final BookieRequestHandler requestHandler) {
+ GetListOfEntriesOfLedgerProcessorV3 getListOfEntriesOfLedger =
+ new GetListOfEntriesOfLedgerProcessorV3(r, requestHandler, this);
if (null == readThreadPool) {
getListOfEntriesOfLedger.run();
} else {
@@ -615,8 +626,8 @@ private void processGetListOfEntriesOfLedgerProcessorV3(final BookkeeperProtocol
}
}
- private void processAddRequest(final BookieProtocol.ParsedAddRequest r, final Channel c) {
- WriteEntryProcessor write = WriteEntryProcessor.create(r, c, this);
+ private void processAddRequest(final BookieProtocol.ParsedAddRequest r, final BookieRequestHandler requestHandler) {
+ WriteEntryProcessor write = WriteEntryProcessor.create(r, requestHandler, this);
// If it's a high priority add (usually as part of recovery process), we want to make sure it gets
// executed as fast as possible, so bypass the normal writeThreadPool and execute in highPriorityThreadPool
@@ -647,10 +658,11 @@ private void processAddRequest(final BookieProtocol.ParsedAddRequest r, final Ch
}
}
- private void processReadRequest(final BookieProtocol.ReadRequest r, final Channel c) {
+ private void processReadRequest(final BookieProtocol.ReadRequest r, final BookieRequestHandler requestHandler) {
ExecutorService fenceThreadPool =
- null == highPriorityThreadPool ? null : highPriorityThreadPool.chooseThread(c);
- ReadEntryProcessor read = ReadEntryProcessor.create(r, c, this, fenceThreadPool, throttleReadResponses);
+ null == highPriorityThreadPool ? null : highPriorityThreadPool.chooseThread(requestHandler.ctx());
+ ReadEntryProcessor read = ReadEntryProcessor.create(r, requestHandler,
+ this, fenceThreadPool, throttleReadResponses);
// If it's a high priority read (fencing or as part of recovery process), we want to make sure it
// gets executed as fast as possible, so bypass the normal readThreadPool
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3.java
index de73f950118..c1627579c6e 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3.java
@@ -22,7 +22,6 @@
import static com.google.common.base.Preconditions.checkArgument;
-import io.netty.channel.Channel;
import java.util.concurrent.TimeUnit;
import org.apache.bookkeeper.bookie.BookieImpl;
import org.apache.bookkeeper.net.BookieId;
@@ -39,9 +38,9 @@
class ForceLedgerProcessorV3 extends PacketProcessorBaseV3 implements Runnable {
private static final Logger logger = LoggerFactory.getLogger(ForceLedgerProcessorV3.class);
- public ForceLedgerProcessorV3(Request request, Channel channel,
+ public ForceLedgerProcessorV3(Request request, BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor) {
- super(request, channel, requestProcessor);
+ super(request, requestHandler, requestProcessor);
}
// Returns null if there is no exception thrown
@@ -98,7 +97,7 @@ private ForceLedgerResponse getForceLedgerResponse() {
};
StatusCode status = null;
try {
- requestProcessor.getBookie().forceLedger(ledgerId, wcb, channel);
+ requestProcessor.getBookie().forceLedger(ledgerId, wcb, requestHandler);
status = StatusCode.EOK;
} catch (Throwable t) {
logger.error("Unexpected exception while forcing ledger {} : ", ledgerId, t);
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/GetBookieInfoProcessorV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/GetBookieInfoProcessorV3.java
index 6f242555864..8795263a5b5 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/GetBookieInfoProcessorV3.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/GetBookieInfoProcessorV3.java
@@ -20,7 +20,6 @@
*/
package org.apache.bookkeeper.proto;
-import io.netty.channel.Channel;
import java.io.IOException;
import java.util.concurrent.TimeUnit;
import org.apache.bookkeeper.proto.BookkeeperProtocol.GetBookieInfoRequest;
@@ -38,9 +37,9 @@
public class GetBookieInfoProcessorV3 extends PacketProcessorBaseV3 implements Runnable {
private static final Logger LOG = LoggerFactory.getLogger(GetBookieInfoProcessorV3.class);
- public GetBookieInfoProcessorV3(Request request, Channel channel,
+ public GetBookieInfoProcessorV3(Request request, BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor) {
- super(request, channel, requestProcessor);
+ super(request, requestHandler, requestProcessor);
}
private GetBookieInfoResponse getGetBookieInfoResponse() {
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/GetListOfEntriesOfLedgerProcessorV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/GetListOfEntriesOfLedgerProcessorV3.java
index 57f72208d81..90c850841d8 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/GetListOfEntriesOfLedgerProcessorV3.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/GetListOfEntriesOfLedgerProcessorV3.java
@@ -21,7 +21,6 @@
package org.apache.bookkeeper.proto;
import com.google.protobuf.ByteString;
-import io.netty.channel.Channel;
import java.io.IOException;
import java.util.concurrent.TimeUnit;
import org.apache.bookkeeper.bookie.Bookie;
@@ -44,9 +43,9 @@ public class GetListOfEntriesOfLedgerProcessorV3 extends PacketProcessorBaseV3 i
protected final GetListOfEntriesOfLedgerRequest getListOfEntriesOfLedgerRequest;
protected final long ledgerId;
- public GetListOfEntriesOfLedgerProcessorV3(Request request, Channel channel,
+ public GetListOfEntriesOfLedgerProcessorV3(Request request, BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor) {
- super(request, channel, requestProcessor);
+ super(request, requestHandler, requestProcessor);
this.getListOfEntriesOfLedgerRequest = request.getGetListOfEntriesOfLedgerRequest();
this.ledgerId = getListOfEntriesOfLedgerRequest.getLedgerId();
}
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/LongPollReadEntryProcessorV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/LongPollReadEntryProcessorV3.java
index f61a2688ba6..658c37c5949 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/LongPollReadEntryProcessorV3.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/LongPollReadEntryProcessorV3.java
@@ -18,7 +18,6 @@
package org.apache.bookkeeper.proto;
import com.google.common.base.Stopwatch;
-import io.netty.channel.Channel;
import io.netty.util.HashedWheelTimer;
import io.netty.util.Timeout;
import java.io.IOException;
@@ -55,12 +54,12 @@ class LongPollReadEntryProcessorV3 extends ReadEntryProcessorV3 implements Watch
private boolean shouldReadEntry = false;
LongPollReadEntryProcessorV3(Request request,
- Channel channel,
+ BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor,
ExecutorService fenceThreadPool,
ExecutorService longPollThreadPool,
HashedWheelTimer requestTimer) {
- super(request, channel, requestProcessor, fenceThreadPool);
+ super(request, requestHandler, requestProcessor, fenceThreadPool);
this.previousLAC = readRequest.getPreviousLAC();
this.longPollThreadPool = longPollThreadPool;
this.requestTimer = requestTimer;
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java
index 8d079504b1c..c9798156c25 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java
@@ -35,20 +35,20 @@
abstract class PacketProcessorBase implements Runnable {
private static final Logger logger = LoggerFactory.getLogger(PacketProcessorBase.class);
T request;
- Channel channel;
+ BookieRequestHandler requestHandler;
BookieRequestProcessor requestProcessor;
long enqueueNanos;
- protected void init(T request, Channel channel, BookieRequestProcessor requestProcessor) {
+ protected void init(T request, BookieRequestHandler requestHandler, BookieRequestProcessor requestProcessor) {
this.request = request;
- this.channel = channel;
+ this.requestHandler = requestHandler;
this.requestProcessor = requestProcessor;
this.enqueueNanos = MathUtils.nowInNano();
}
protected void reset() {
request = null;
- channel = null;
+ requestHandler = null;
requestProcessor = null;
enqueueNanos = -1;
}
@@ -82,8 +82,10 @@ protected void sendReadReqResponse(int rc, Object response, OpStatsLogger statsL
protected void sendResponse(int rc, Object response, OpStatsLogger statsLogger) {
final long writeNanos = MathUtils.nowInNano();
-
final long timeOut = requestProcessor.getWaitTimeoutOnBackpressureMillis();
+
+ Channel channel = requestHandler.ctx().channel();
+
if (timeOut >= 0 && !channel.isWritable()) {
if (!requestProcessor.isBlacklisted(channel)) {
synchronized (channel) {
@@ -120,18 +122,23 @@ protected void sendResponse(int rc, Object response, OpStatsLogger statsLogger)
}
if (channel.isActive()) {
- ChannelPromise promise = channel.newPromise().addListener(future -> {
- if (!future.isSuccess()) {
- logger.debug("Netty channel write exception. ", future.cause());
- }
- });
+ ChannelPromise promise = channel.voidPromise();
+ if (logger.isDebugEnabled()) {
+ promise = channel.newPromise().addListener(future -> {
+ if (!future.isSuccess()) {
+ logger.debug("Netty channel write exception. ", future.cause());
+ }
+ });
+ }
channel.writeAndFlush(response, promise);
} else {
if (response instanceof BookieProtocol.Response) {
((BookieProtocol.Response) response).release();
}
+ if (logger.isDebugEnabled()) {
logger.debug("Netty channel {} is inactive, "
+ "hence bypassing netty channel writeAndFlush during sendResponse", channel);
+ }
}
if (BookieProtocol.EOK == rc) {
statsLogger.registerSuccessfulEvent(MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS);
@@ -149,6 +156,7 @@ protected void sendResponse(int rc, Object response, OpStatsLogger statsLogger)
*/
protected void sendResponseAndWait(int rc, Object response, OpStatsLogger statsLogger) {
try {
+ Channel channel = requestHandler.ctx().channel();
ChannelFuture future = channel.writeAndFlush(response);
if (!channel.eventLoop().inEventLoop()) {
future.get();
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java
index ccc452ae60a..dac454933c1 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java
@@ -40,14 +40,14 @@
public abstract class PacketProcessorBaseV3 implements Runnable {
final Request request;
- final Channel channel;
+ final BookieRequestHandler requestHandler;
final BookieRequestProcessor requestProcessor;
final long enqueueNanos;
- public PacketProcessorBaseV3(Request request, Channel channel,
+ public PacketProcessorBaseV3(Request request, BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor) {
this.request = request;
- this.channel = channel;
+ this.requestHandler = requestHandler;
this.requestProcessor = requestProcessor;
this.enqueueNanos = MathUtils.nowInNano();
}
@@ -55,6 +55,7 @@ public PacketProcessorBaseV3(Request request, Channel channel,
protected void sendResponse(StatusCode code, Object response, OpStatsLogger statsLogger) {
final long writeNanos = MathUtils.nowInNano();
+ Channel channel = requestHandler.ctx().channel();
final long timeOut = requestProcessor.getWaitTimeoutOnBackpressureMillis();
if (timeOut >= 0 && !channel.isWritable()) {
if (!requestProcessor.isBlacklisted(channel)) {
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessor.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessor.java
index 71ee51f3fa6..6935ca8be60 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessor.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessor.java
@@ -18,7 +18,6 @@
package org.apache.bookkeeper.proto;
import io.netty.buffer.ByteBuf;
-import io.netty.channel.Channel;
import io.netty.util.Recycler;
import io.netty.util.ReferenceCountUtil;
import java.io.IOException;
@@ -44,15 +43,15 @@ class ReadEntryProcessor extends PacketProcessorBase {
private boolean throttleReadResponses;
public static ReadEntryProcessor create(ReadRequest request,
- Channel channel,
+ BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor,
ExecutorService fenceThreadPool,
boolean throttleReadResponses) {
ReadEntryProcessor rep = RECYCLER.get();
- rep.init(request, channel, requestProcessor);
+ rep.init(request, requestHandler, requestProcessor);
rep.fenceThreadPool = fenceThreadPool;
rep.throttleReadResponses = throttleReadResponses;
- requestProcessor.onReadRequestStart(channel);
+ requestProcessor.onReadRequestStart(requestHandler.ctx().channel());
return rep;
}
@@ -61,9 +60,9 @@ protected void processPacket() {
if (LOG.isDebugEnabled()) {
LOG.debug("Received new read request: {}", request);
}
- if (!channel.isOpen()) {
+ if (!requestHandler.ctx().channel().isOpen()) {
if (LOG.isDebugEnabled()) {
- LOG.debug("Dropping read request for closed channel: {}", channel);
+ LOG.debug("Dropping read request for closed channel: {}", requestHandler.ctx().channel());
}
requestProcessor.onReadRequestFinish();
return;
@@ -74,7 +73,8 @@ protected void processPacket() {
try {
CompletableFuture fenceResult = null;
if (request.isFencing()) {
- LOG.warn("Ledger: {} fenced by: {}", request.getLedgerId(), channel.remoteAddress());
+ LOG.warn("Ledger: {} fenced by: {}", request.getLedgerId(),
+ requestHandler.ctx().channel().remoteAddress());
if (request.hasMasterKey()) {
fenceResult = requestProcessor.getBookie().fenceLedger(request.getLedgerId(),
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessorV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessorV3.java
index 4672d592d88..999b8095db6 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessorV3.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessorV3.java
@@ -57,11 +57,11 @@ class ReadEntryProcessorV3 extends PacketProcessorBaseV3 {
protected final OpStatsLogger reqStats;
public ReadEntryProcessorV3(Request request,
- Channel channel,
+ BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor,
ExecutorService fenceThreadPool) {
- super(request, channel, requestProcessor);
- requestProcessor.onReadRequestStart(channel);
+ super(request, requestHandler, requestProcessor);
+ requestProcessor.onReadRequestStart(requestHandler.ctx().channel());
this.readRequest = request.getReadRequest();
this.ledgerId = readRequest.getLedgerId();
@@ -194,6 +194,7 @@ protected ReadResponse readEntry(ReadResponse.Builder readResponseBuilder,
protected ReadResponse getReadResponse() {
final Stopwatch startTimeSw = Stopwatch.createStarted();
+ final Channel channel = requestHandler.ctx().channel();
final ReadResponse.Builder readResponse = ReadResponse.newBuilder()
.setLedgerId(ledgerId)
@@ -249,9 +250,9 @@ protected ReadResponse getReadResponse() {
public void run() {
requestProcessor.getRequestStats().getReadEntrySchedulingDelayStats().registerSuccessfulEvent(
MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS);
- if (!channel.isOpen()) {
+ if (!requestHandler.ctx().channel().isOpen()) {
if (LOG.isDebugEnabled()) {
- LOG.debug("Dropping read request for closed channel: {}", channel);
+ LOG.debug("Dropping read request for closed channel: {}", requestHandler.ctx().channel());
}
requestProcessor.onReadRequestFinish();
return;
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadLacProcessorV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadLacProcessorV3.java
index abb7d616468..25fe9530ad7 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadLacProcessorV3.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadLacProcessorV3.java
@@ -22,7 +22,6 @@
import com.google.protobuf.ByteString;
import io.netty.buffer.ByteBuf;
-import io.netty.channel.Channel;
import io.netty.util.ReferenceCountUtil;
import java.io.IOException;
import java.util.concurrent.TimeUnit;
@@ -43,9 +42,9 @@
class ReadLacProcessorV3 extends PacketProcessorBaseV3 implements Runnable {
private static final Logger logger = LoggerFactory.getLogger(ReadLacProcessorV3.class);
- public ReadLacProcessorV3(Request request, Channel channel,
+ public ReadLacProcessorV3(Request request, BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor) {
- super(request, channel, requestProcessor);
+ super(request, requestHandler, requestProcessor);
}
// Returns null if there is no exception thrown
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteEntryProcessor.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteEntryProcessor.java
index 531dfecccc2..7e8f9fa768d 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteEntryProcessor.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteEntryProcessor.java
@@ -19,7 +19,6 @@
import com.google.common.annotations.VisibleForTesting;
import io.netty.buffer.ByteBuf;
-import io.netty.channel.Channel;
import io.netty.util.Recycler;
import java.io.IOException;
import java.util.concurrent.TimeUnit;
@@ -47,11 +46,11 @@ protected void reset() {
startTimeNanos = -1L;
}
- public static WriteEntryProcessor create(ParsedAddRequest request, Channel channel,
+ public static WriteEntryProcessor create(ParsedAddRequest request, BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor) {
WriteEntryProcessor wep = RECYCLER.get();
- wep.init(request, channel, requestProcessor);
- requestProcessor.onAddRequestStart(channel);
+ wep.init(request, requestHandler, requestProcessor);
+ requestProcessor.onAddRequestStart(requestHandler.ctx().channel());
return wep;
}
@@ -74,9 +73,10 @@ protected void processPacket() {
ByteBuf addData = request.getData();
try {
if (request.isRecoveryAdd()) {
- requestProcessor.getBookie().recoveryAddEntry(addData, this, channel, request.getMasterKey());
+ requestProcessor.getBookie().recoveryAddEntry(addData, this, requestHandler, request.getMasterKey());
} else {
- requestProcessor.getBookie().addEntry(addData, false, this, channel, request.getMasterKey());
+ requestProcessor.getBookie().addEntry(addData, false, this,
+ requestHandler, request.getMasterKey());
}
} catch (OperationRejectedException e) {
requestProcessor.getRequestStats().getAddEntryRejectedCounter().inc();
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteEntryProcessorV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteEntryProcessorV3.java
index 7d598587324..36aff7ad924 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteEntryProcessorV3.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteEntryProcessorV3.java
@@ -22,7 +22,6 @@
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
-import io.netty.channel.Channel;
import java.io.IOException;
import java.util.EnumSet;
import java.util.concurrent.TimeUnit;
@@ -43,10 +42,10 @@
class WriteEntryProcessorV3 extends PacketProcessorBaseV3 {
private static final Logger logger = LoggerFactory.getLogger(WriteEntryProcessorV3.class);
- public WriteEntryProcessorV3(Request request, Channel channel,
+ public WriteEntryProcessorV3(Request request, BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor) {
- super(request, channel, requestProcessor);
- requestProcessor.onAddRequestStart(channel);
+ super(request, requestHandler, requestProcessor);
+ requestProcessor.onAddRequestStart(requestHandler.ctx().channel());
}
// Returns null if there is no exception thrown
@@ -118,9 +117,11 @@ public void writeComplete(int rc, long ledgerId, long entryId,
ByteBuf entryToAdd = Unpooled.wrappedBuffer(addRequest.getBody().asReadOnlyByteBuffer());
try {
if (RequestUtils.hasFlag(addRequest, AddRequest.Flag.RECOVERY_ADD)) {
- requestProcessor.getBookie().recoveryAddEntry(entryToAdd, wcb, channel, masterKey);
+ requestProcessor.getBookie().recoveryAddEntry(entryToAdd, wcb,
+ requestHandler.ctx().channel(), masterKey);
} else {
- requestProcessor.getBookie().addEntry(entryToAdd, ackBeforeSync, wcb, channel, masterKey);
+ requestProcessor.getBookie().addEntry(entryToAdd, ackBeforeSync, wcb,
+ requestHandler.ctx().channel(), masterKey);
}
status = StatusCode.EOK;
} catch (OperationRejectedException e) {
diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteLacProcessorV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteLacProcessorV3.java
index d8a427c9059..293cea3bb0c 100644
--- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteLacProcessorV3.java
+++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteLacProcessorV3.java
@@ -21,7 +21,6 @@
package org.apache.bookkeeper.proto;
import io.netty.buffer.Unpooled;
-import io.netty.channel.Channel;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.concurrent.TimeUnit;
@@ -40,9 +39,9 @@
class WriteLacProcessorV3 extends PacketProcessorBaseV3 implements Runnable {
private static final Logger logger = LoggerFactory.getLogger(WriteLacProcessorV3.class);
- public WriteLacProcessorV3(Request request, Channel channel,
+ public WriteLacProcessorV3(Request request, BookieRequestHandler requestHandler,
BookieRequestProcessor requestProcessor) {
- super(request, channel, requestProcessor);
+ super(request, requestHandler, requestProcessor);
}
// Returns null if there is no exception thrown
@@ -103,7 +102,8 @@ public void writeComplete(int rc, long ledgerId, long entryId, BookieId addr, Ob
byte[] masterKey = writeLacRequest.getMasterKey().toByteArray();
try {
- requestProcessor.bookie.setExplicitLac(Unpooled.wrappedBuffer(lacToAdd), writeCallback, channel, masterKey);
+ requestProcessor.bookie.setExplicitLac(Unpooled.wrappedBuffer(lacToAdd),
+ writeCallback, requestHandler, masterKey);
status = StatusCode.EOK;
} catch (IOException e) {
logger.error("Error saving lac {} for ledger:{}",
diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3Test.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3Test.java
index 54460770f16..3bc9cbee427 100644
--- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3Test.java
+++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3Test.java
@@ -30,6 +30,7 @@
import static org.mockito.Mockito.when;
import io.netty.channel.Channel;
+import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelPromise;
import io.netty.channel.DefaultChannelPromise;
import java.util.concurrent.CountDownLatch;
@@ -55,6 +56,8 @@ public class ForceLedgerProcessorV3Test {
private Request request;
private ForceLedgerProcessorV3 processor;
+
+ private BookieRequestHandler requestHandler;
private Channel channel;
private BookieRequestProcessor requestProcessor;
private Bookie bookie;
@@ -71,17 +74,25 @@ public void setup() {
.setLedgerId(System.currentTimeMillis())
.build())
.build();
+
+
channel = mock(Channel.class);
when(channel.isOpen()).thenReturn(true);
+ when(channel.isActive()).thenReturn(true);
+
+ requestHandler = mock(BookieRequestHandler.class);
+ ChannelHandlerContext ctx = mock(ChannelHandlerContext.class);
+ when(ctx.channel()).thenReturn(channel);
+ when(requestHandler.ctx()).thenReturn(ctx);
+
bookie = mock(Bookie.class);
requestProcessor = mock(BookieRequestProcessor.class);
when(requestProcessor.getBookie()).thenReturn(bookie);
when(requestProcessor.getWaitTimeoutOnBackpressureMillis()).thenReturn(-1L);
when(requestProcessor.getRequestStats()).thenReturn(new RequestStats(NullStatsLogger.INSTANCE));
- when(channel.isActive()).thenReturn(true);
processor = new ForceLedgerProcessorV3(
request,
- channel,
+ requestHandler,
requestProcessor);
}
@@ -102,7 +113,7 @@ public void testForceLedger() throws Exception {
}).when(bookie).forceLedger(
eq(request.getForceLedgerRequest().getLedgerId()),
any(WriteCallback.class),
- same(channel));
+ same(requestHandler));
ChannelPromise promise = new DefaultChannelPromise(channel);
AtomicReference