From 926ca7da94d38fe6acfab8d660c3740750bdec0c Mon Sep 17 00:00:00 2001 From: ZhenyuLi <893652269@qq.com> Date: Wed, 16 Sep 2026 23:21:22 -0400 Subject: [PATCH 1/2] Fix state reuse after interrupted TRUNC sync --- .../zookeeper/server/quorum/Learner.java | 8 ++- .../zookeeper/server/quorum/LearnerTest.java | 58 +++++++++++++++++++ 2 files changed, 63 insertions(+), 3 deletions(-) diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java index adf0ef6e510..cead3d37381 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java @@ -919,9 +919,11 @@ public void shutdown() { closeSocket(); // shutdown previous zookeeper if (zk != null) { - // If we haven't finished SNAP sync, force fully shutdown - // to avoid potential inconsistency - zk.shutdown(self.getSyncMode().equals(QuorumPeer.SyncMode.SNAP)); + QuorumPeer.SyncMode syncMode = self.getSyncMode(); + // SNAP and TRUNC sync can apply transactions directly to the in-memory + // database before the state is persisted in a snapshot. If sync fails, + // discard that database so the next attempt reloads the state on disk. + zk.shutdown(syncMode == QuorumPeer.SyncMode.SNAP || syncMode == QuorumPeer.SyncMode.TRUNC); } } diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/LearnerTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/LearnerTest.java index d64d051b093..026e8d949ce 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/LearnerTest.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/LearnerTest.java @@ -24,6 +24,9 @@ import static org.hamcrest.CoreMatchers.is; import static org.hamcrest.MatcherAssert.assertThat; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; @@ -48,6 +51,7 @@ import org.apache.zookeeper.common.X509Exception; import org.apache.zookeeper.data.ACL; import org.apache.zookeeper.server.ExitCode; +import org.apache.zookeeper.server.Request; import org.apache.zookeeper.server.ZKDatabase; import org.apache.zookeeper.server.persistence.FileTxnSnapLog; import org.apache.zookeeper.txn.CreateTxn; @@ -356,4 +360,58 @@ public void accept(Integer exitCode) { assertThat("System.exit() should have been called", exitProcCalled[0], is(true)); } + + @Test + public void incompleteTruncSyncClearsInMemoryDatabase(@TempDir File tmpDir) throws Exception { + FileTxnSnapLog txnSnapLog = new FileTxnSnapLog(tmpDir, tmpDir); + SimpleLearner learner = new SimpleLearner(txnSnapLog); + ZKDatabase zkDb = learner.zk.getZKDatabase(); + + txnSnapLog.save(zkDb.getDataTree(), zkDb.getSessionWithTimeOuts(), false); + appendCreate(txnSnapLog, 1, "/retained"); + appendCreate(txnSnapLog, 2, "/discarded"); + txnSnapLog.commit(); + assertEquals(2, zkDb.loadDataBase()); + + ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + BinaryOutputArchive leaderOutput = BinaryOutputArchive.getArchive(bytes); + leaderOutput.writeRecord(new QuorumPacket(Leader.TRUNC, 1, null, null), null); + + TxnHeader replacementHeader = new TxnHeader(1, 3, 2, 3, ZooDefs.OpCode.create); + CreateTxn replacementTxn = new CreateTxn( + "/replacement", + new byte[0], + ZooDefs.Ids.OPEN_ACL_UNSAFE, + false, + 2); + ByteArrayOutputStream proposalBytes = new ByteArrayOutputStream(); + BinaryOutputArchive proposalOutput = BinaryOutputArchive.getArchive(proposalBytes); + replacementHeader.serialize(proposalOutput, "hdr"); + replacementTxn.serialize(proposalOutput, "txn"); + leaderOutput.writeRecord(new QuorumPacket(Leader.PROPOSAL, 2, proposalBytes.toByteArray(), null), null); + leaderOutput.writeRecord(new QuorumPacket(Leader.COMMIT, 2, null, null), null); + + learner.leaderIs = BinaryInputArchive.getArchive(new ByteArrayInputStream(bytes.toByteArray())); + learner.leaderOs = BinaryOutputArchive.getArchive(new ByteArrayOutputStream()); + learner.bufferedOutput = new BufferedOutputStream(new ByteArrayOutputStream()); + learner.sock = new Socket(); + + assertThrows(EOFException.class, () -> learner.syncWithLeader(3)); + assertEquals(QuorumPeer.SyncMode.TRUNC, learner.self.getSyncMode()); + assertNotNull(zkDb.getNode("/replacement")); + assertNull(zkDb.getNode("/discarded")); + + learner.shutdown(); + + assertFalse(zkDb.isInitialized()); + assertEquals(1, zkDb.loadDataBase()); + assertNotNull(zkDb.getNode("/retained")); + assertNull(zkDb.getNode("/replacement")); + } + + private static void appendCreate(FileTxnSnapLog txnSnapLog, long zxid, String path) throws IOException { + TxnHeader header = new TxnHeader(1, (int) zxid, zxid, zxid, ZooDefs.OpCode.create); + CreateTxn txn = new CreateTxn(path, new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, false, (int) zxid); + txnSnapLog.append(new Request(header, txn, null)); + } } From 46c8b43a17b65de4a8c935cb3bc0d0dac78e9337 Mon Sep 17 00:00:00 2001 From: ZhenyuLi <893652269@qq.com> Date: Thu, 17 Sep 2026 20:14:42 -0400 Subject: [PATCH 2/2] Improve interrupted TRUNC sync regression test --- .../zookeeper/server/quorum/LearnerTest.java | 31 +++++++++++++++---- 1 file changed, 25 insertions(+), 6 deletions(-) diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/LearnerTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/LearnerTest.java index 026e8d949ce..64fe57ff995 100644 --- a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/LearnerTest.java +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/LearnerTest.java @@ -23,8 +23,8 @@ import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.CoreMatchers.is; import static org.hamcrest.MatcherAssert.assertThat; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -362,7 +362,7 @@ public void accept(Integer exitCode) { } @Test - public void incompleteTruncSyncClearsInMemoryDatabase(@TempDir File tmpDir) throws Exception { + public void incompleteTruncSyncDoesNotCreateTxnLogGapOnReconnect(@TempDir File tmpDir) throws Exception { FileTxnSnapLog txnSnapLog = new FileTxnSnapLog(tmpDir, tmpDir); SimpleLearner learner = new SimpleLearner(txnSnapLog); ZKDatabase zkDb = learner.zk.getZKDatabase(); @@ -403,10 +403,29 @@ public void incompleteTruncSyncClearsInMemoryDatabase(@TempDir File tmpDir) thro learner.shutdown(); - assertFalse(zkDb.isInitialized()); - assertEquals(1, zkDb.loadDataBase()); - assertNotNull(zkDb.getNode("/retained")); - assertNull(zkDb.getNode("/replacement")); + // Mirror the lazy reload performed by QuorumPeer.getLastLoggedZxid() on + // the next connection without depending on how shutdown invalidates the database. + if (!zkDb.isInitialized()) { + zkDb.loadDataBase(); + } + long nextZxid = zkDb.getDataTreeLastProcessedZxid() + 1; + + // Persist the transaction the leader would send next. If the incomplete + // in-memory state was reused, nextZxid is 3 and the on-disk log has a gap. + appendCreate(txnSnapLog, nextZxid, "/continued"); + txnSnapLog.commit(); + zkDb.close(); + + // Simulate a process restart. Replaying the log must not detect a zxid gap. + FileTxnSnapLog restartedTxnSnapLog = new FileTxnSnapLog(tmpDir, tmpDir); + ZKDatabase restartedDb = new ZKDatabase(restartedTxnSnapLog); + long restoredZxid = assertDoesNotThrow(restartedDb::loadDataBase); + assertEquals(2, restoredZxid); + assertNotNull(restartedDb.getNode("/retained")); + assertNotNull(restartedDb.getNode("/continued")); + assertNull(restartedDb.getNode("/discarded")); + assertNull(restartedDb.getNode("/replacement")); + restartedDb.close(); } private static void appendCreate(FileTxnSnapLog txnSnapLog, long zxid, String path) throws IOException {