diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java index 78dc8fadccc..0d17ce2bc6a 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerCnxn.java @@ -101,6 +101,7 @@ public enum DisconnectReason { CLOSE_CONNECTION_COMMAND("close_connection_command"), CLEAN_UP("clean_up"), CONNECTION_MODE_CHANGED("connection_mode_changed"), + RENEW_GLOBAL_SESSION_IN_RO_MODE("renew a global session in readonly mode"), // Below reasons are NettyServerCnxnFactory only CHANNEL_DISCONNECTED("channel disconnected"), CHANNEL_CLOSED_EXCEPTION("channel_closed_exception"), @@ -298,7 +299,7 @@ void disableRecv() { protected ZooKeeperSaslServer zooKeeperSaslServer = null; - protected static class CloseRequestException extends IOException { + public static class CloseRequestException extends IOException { private static final long serialVersionUID = -7854505709816442681L; private DisconnectReason reason; diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ZooKeeperServer.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ZooKeeperServer.java index c391425a309..8a0321c445a 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ZooKeeperServer.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ZooKeeperServer.java @@ -1413,13 +1413,13 @@ public void processConnectRequest(ServerCnxn cnxn, ByteBuffer incomingBuffer) connReq.getTimeOut(), cnxn.getRemoteSocketAddress()); } else { - long clientSessionId = connReq.getSessionId(); - LOG.debug( - "Client attempting to renew session: session = 0x{}, zxid = 0x{}, timeout = {}, address = {}", - Long.toHexString(clientSessionId), - Long.toHexString(connReq.getLastZxidSeen()), - connReq.getTimeOut(), - cnxn.getRemoteSocketAddress()); + validateSession(cnxn, sessionId); + LOG.debug( + "Client attempting to renew session: session = 0x{}, zxid = 0x{}, timeout = {}, address = {}", + Long.toHexString(sessionId), + Long.toHexString(connReq.getLastZxidSeen()), + connReq.getTimeOut(), + cnxn.getRemoteSocketAddress()); if (serverCnxnFactory != null) { serverCnxnFactory.closeSession(sessionId, ServerCnxn.DisconnectReason.CLIENT_RECONNECT); } @@ -1433,6 +1433,17 @@ public void processConnectRequest(ServerCnxn cnxn, ByteBuffer incomingBuffer) } } + /** + * Validate if a particular session can be reestablished. + * + * @param cnxn + * @param sessionId + */ + protected void validateSession(ServerCnxn cnxn, long sessionId) + throws IOException { + // do nothing + } + public boolean shouldThrottle(long outStandingCount) { int globalOutstandingLimit = getGlobalOutstandingLimit(); if (globalOutstandingLimit < getInflight() || globalOutstandingLimit < getInProcess()) { diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyRequestProcessor.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyRequestProcessor.java index c50dd539f4e..3aec97f72e2 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyRequestProcessor.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyRequestProcessor.java @@ -86,16 +86,14 @@ public void run() { case OpCode.setACL: case OpCode.multi: case OpCode.check: - ReplyHeader hdr = new ReplyHeader( - request.cxid, - zks.getZKDatabase().getDataTreeLastProcessedZxid(), - Code.NOTREADONLY.intValue()); - try { - request.cnxn.sendResponse(hdr, null, null); - } catch (IOException e) { - LOG.error("IO exception while sending response", e); - } + sendErrorResponse(request); continue; + case OpCode.closeSession: + case OpCode.createSession: + if (!request.isLocalSession()) { + sendErrorResponse(request); + continue; + } } // proceed to the next processor @@ -109,6 +107,18 @@ public void run() { LOG.info("ReadOnlyRequestProcessor exited loop!"); } + private void sendErrorResponse(Request request) { + ReplyHeader hdr = new ReplyHeader( + request.cxid, + zks.getZKDatabase().getDataTreeLastProcessedZxid(), + Code.NOTREADONLY.intValue()); + try { + request.cnxn.sendResponse(hdr, null, null); + } catch (IOException e) { + LOG.error("IO exception while sending response", e); + } + } + @Override public void processRequest(Request request) { if (!finished) { diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyZooKeeperServer.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyZooKeeperServer.java index f8517eb9b02..7f3084c5376 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyZooKeeperServer.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/ReadOnlyZooKeeperServer.java @@ -18,6 +18,7 @@ package org.apache.zookeeper.server.quorum; +import java.io.IOException; import java.io.PrintWriter; import java.util.Objects; import java.util.stream.Collectors; @@ -26,6 +27,7 @@ import org.apache.zookeeper.server.FinalRequestProcessor; import org.apache.zookeeper.server.PrepRequestProcessor; import org.apache.zookeeper.server.RequestProcessor; +import org.apache.zookeeper.server.ServerCnxn; import org.apache.zookeeper.server.ZKDatabase; import org.apache.zookeeper.server.ZooKeeperServer; import org.apache.zookeeper.server.ZooKeeperServerBean; @@ -80,6 +82,15 @@ public synchronized void startup() { LOG.info("Read-only server started"); } + @Override + protected void validateSession(ServerCnxn cnxn, long sessionId) throws IOException { + if (((LearnerSessionTracker) sessionTracker).isGlobalSession(sessionId)) { + String msg = "Refusing global session reconnection in RO mode " + cnxn.getRemoteSocketAddress(); + LOG.info(msg); + throw new ServerCnxn.CloseRequestException(msg, ServerCnxn.DisconnectReason.RENEW_GLOBAL_SESSION_IN_RO_MODE); + } + } + @Override protected void registerJMX() { // register with JMX