From f007f1ecfcf555c13847b0c4e74fccf39f937be6 Mon Sep 17 00:00:00 2001 From: vaughn Date: Tue, 16 Sep 2025 10:14:55 +0800 Subject: [PATCH 1/9] improve: raft state machine should report error --- .../main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java | 1 + .../org/apache/hugegraph/store/raft/HgStoreStateMachine.java | 1 + 2 files changed, 2 insertions(+) diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java index c7537d30a0..2df9241609 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java @@ -94,6 +94,7 @@ public void onApply(Iterator iter) { if (done != null) { done.run(new Status(RaftError.EINTERNAL, t.getMessage())); } + iter.setErrorAndRollback(1, new Status(RaftError.ESTATEMACHINE, t.getMessage())); } iter.next(); } diff --git a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/HgStoreStateMachine.java b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/HgStoreStateMachine.java index 0f80017c53..f2560205df 100644 --- a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/HgStoreStateMachine.java +++ b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/HgStoreStateMachine.java @@ -101,6 +101,7 @@ public void onApply(Iterator inter) { done.op.getReq()); // done.run(new Status(RaftError.EINTERNAL, t.getMessage())); } + inter.setErrorAndRollback(1, new Status(RaftError.ESTATEMACHINE, t.getMessage())); } committedIndex = inter.getIndex(); From 23ad40aa6d4004c80ab302fd391935ef12b67804 Mon Sep 17 00:00:00 2001 From: imbajin Date: Thu, 24 Sep 2026 22:41:59 +0800 Subject: [PATCH 2/9] fix: stop raft apply after state machine errors - let jraft complete failed PD closures once - preserve the last successful Store apply position - cover failure boundaries with real FSMCaller tests - run Store regressions against production jraft --- .../hugegraph/pd/raft/RaftStateMachine.java | 5 +- .../hugegraph/pd/core/PDCoreSuiteTest.java | 1 + .../pd/core/RaftStateMachineTest.java | 156 ++++++++++++++++ .../store/raft/HgStoreStateMachine.java | 2 + hugegraph-store/hg-store-test/pom.xml | 3 +- .../core/raft/HgStoreStateMachineTest.java | 176 ++++++++++++++---- 6 files changed, 298 insertions(+), 45 deletions(-) create mode 100644 hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java index 2df9241609..ae18355ef4 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java @@ -91,10 +91,9 @@ public void onApply(Iterator iter) { } } catch (Throwable t) { log.error("StateMachine encountered critical error", t); - if (done != null) { - done.run(new Status(RaftError.EINTERNAL, t.getMessage())); - } + // JRaft completes the failed and remaining closures with the state machine error. iter.setErrorAndRollback(1, new Status(RaftError.ESTATEMACHINE, t.getMessage())); + return; } iter.next(); } diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java index 3d785360d0..bdee1aa7ee 100644 --- a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java @@ -26,6 +26,7 @@ @RunWith(Suite.class) @Suite.SuiteClasses({ + RaftStateMachineTest.class, MetadataKeyHelperTest.class, HgKVStoreImplTest.class, ConfigServiceTest.class, diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java new file mode 100644 index 0000000000..2f1fdbb903 --- /dev/null +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java @@ -0,0 +1,156 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF 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 org.apache.hugegraph.pd.core; + +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.hugegraph.pd.raft.KVOperation; +import org.apache.hugegraph.pd.raft.KVStoreClosure; +import org.apache.hugegraph.pd.raft.RaftStateMachine; +import org.junit.Assert; +import org.junit.Test; +import org.mockito.Mockito; + +import com.alipay.sofa.jraft.Status; +import com.alipay.sofa.jraft.closure.ClosureQueueImpl; +import com.alipay.sofa.jraft.closure.SaveSnapshotClosure; +import com.alipay.sofa.jraft.core.FSMCallerImpl; +import com.alipay.sofa.jraft.core.NodeImpl; +import com.alipay.sofa.jraft.core.NodeMetrics; +import com.alipay.sofa.jraft.entity.EnumOutter; +import com.alipay.sofa.jraft.entity.LogEntry; +import com.alipay.sofa.jraft.entity.LogId; +import com.alipay.sofa.jraft.error.RaftError; +import com.alipay.sofa.jraft.option.FSMCallerOptions; +import com.alipay.sofa.jraft.storage.LogManager; + +public class RaftStateMachineTest { + + @Test + public void testSuccessfulLeaderApply() throws Exception { + assertApply(true, false); + } + + @Test + public void testLeaderFailureStopsApply() throws Exception { + assertApply(true, true); + } + + @Test + public void testFollowerFailureStopsReplay() throws Exception { + assertApply(false, true); + } + + private void assertApply(boolean leader, boolean fail) throws Exception { + RaftStateMachine stateMachine = new RaftStateMachine(); + List applied = new ArrayList<>(); + stateMachine.addTaskHandler((operation, done) -> { + int key = operation.getKey()[0]; + if (fail && key == 2) { + throw new IllegalStateException("injected apply failure"); + } + applied.add(key); + return true; + }); + + LogManager logManager = Mockito.mock(LogManager.class); + ClosureQueueImpl closures = new ClosureQueueImpl("pd-apply-test"); + closures.resetFirstIndex(1); + List callbacks = new ArrayList<>(); + CountDownLatch completed = new CountDownLatch(leader ? 3 : 0); + for (int index = 1; index <= 4; index++) { + KVOperation operation = KVOperation.createPut(new byte[]{(byte) index}, new byte[]{1}); + LogEntry entry = new LogEntry(EnumOutter.EntryType.ENTRY_TYPE_DATA); + entry.setId(new LogId(index, 1)); + entry.setData(ByteBuffer.wrap(operation.toByteArray())); + Mockito.when(logManager.getEntry(index)).thenReturn(entry); + Mockito.when(logManager.getTerm(index)).thenReturn(1L); + if (leader && index <= 3) { + KVStoreClosure callback = Mockito.mock(KVStoreClosure.class); + Mockito.doAnswer(invocation -> { + completed.countDown(); + return null; + }).when(callback).run(Mockito.any(Status.class)); + callbacks.add(callback); + closures.appendPendingClosure(new RaftStateMachine.RaftClosureAdapter(operation, + callback)); + } + } + + NodeImpl node = Mockito.mock(NodeImpl.class); + Mockito.when(node.getNodeMetrics()).thenReturn(new NodeMetrics(false)); + Mockito.when(node.getGroupId()).thenReturn("pd-apply-test"); + FSMCallerOptions options = new FSMCallerOptions(); + options.setNode(node); + options.setFsm(stateMachine); + options.setLogManager(logManager); + options.setClosureQueue(closures); + options.setBootstrapId(new LogId(0, 0)); + options.setDisruptorBufferSize(16); + FSMCallerImpl caller = new FSMCallerImpl(); + Assert.assertTrue(caller.init(options)); + CountDownLatch batchApplied = new CountDownLatch(1); + caller.addLastAppliedLogIndexListener(index -> batchApplied.countDown()); + try { + Assert.assertTrue(caller.onCommitted(3)); + Assert.assertTrue(batchApplied.await(10, TimeUnit.SECONDS)); + Assert.assertEquals(fail ? 1 : 3, caller.getLastAppliedIndex()); + Assert.assertEquals(3, caller.getLastCommittedIndex()); + Assert.assertTrue(completed.await(10, TimeUnit.SECONDS)); + if (fail) { + // The terminal error must also reject future batches and snapshot saves. + Assert.assertTrue(caller.onCommitted(4)); + CountDownLatch snapshotCompleted = new CountDownLatch(1); + AtomicReference snapshotStatus = new AtomicReference<>(); + SaveSnapshotClosure snapshot = Mockito.mock(SaveSnapshotClosure.class); + Mockito.doAnswer(invocation -> { + snapshotStatus.set(invocation.getArgument(0)); + snapshotCompleted.countDown(); + return null; + }).when(snapshot).run(Mockito.any(Status.class)); + caller.onSnapshotSave(snapshot); + Assert.assertTrue(snapshotCompleted.await(10, TimeUnit.SECONDS)); + Assert.assertFalse(snapshotStatus.get().isOk()); + Assert.assertTrue(snapshotStatus.get().getErrorMsg() + .startsWith("FSMCaller is in bad status")); + Mockito.verify(logManager, Mockito.never()).getConfiguration(Mockito.anyLong()); + Mockito.verify(snapshot, Mockito.never()).start(Mockito.any()); + } + } finally { + caller.shutdown(); + caller.join(); + } + Assert.assertEquals(fail ? Arrays.asList(1) : Arrays.asList(1, 2, 3), applied); + Assert.assertEquals(fail ? 1 : 3, caller.getLastAppliedIndex()); + Mockito.verify(logManager).setAppliedId(new LogId(fail ? 1 : 3, 1)); + if (leader) { + for (int index = 0; index < callbacks.size(); index++) { + int expectedCode = fail && index > 0 ? RaftError.ESTATEMACHINE.getNumber() : 0; + Mockito.verify(callbacks.get(index), Mockito.times(1)) + .run(Mockito.argThat(status -> status.getCode() == expectedCode)); + Mockito.verifyNoMoreInteractions(callbacks.get(index)); + } + } + } +} diff --git a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/HgStoreStateMachine.java b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/HgStoreStateMachine.java index f2560205df..da6ba7e6ef 100644 --- a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/HgStoreStateMachine.java +++ b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/HgStoreStateMachine.java @@ -102,6 +102,8 @@ public void onApply(Iterator inter) { // done.run(new Status(RaftError.EINTERNAL, t.getMessage())); } inter.setErrorAndRollback(1, new Status(RaftError.ESTATEMACHINE, t.getMessage())); + // Do not publish the failed entry as applied. + return; } committedIndex = inter.getIndex(); diff --git a/hugegraph-store/hg-store-test/pom.xml b/hugegraph-store/hg-store-test/pom.xml index 36308f449d..6e4f50fa20 100644 --- a/hugegraph-store/hg-store-test/pom.xml +++ b/hugegraph-store/hg-store-test/pom.xml @@ -193,7 +193,7 @@ com.alipay.sofa jraft-core - 1.3.9 + 1.3.13 org.rocksdb @@ -236,6 +236,7 @@ **/CoreSuiteTest.java + **/HgStoreStateMachineTest.java diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/HgStoreStateMachineTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/HgStoreStateMachineTest.java index 752a17ea59..9b84b8126d 100644 --- a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/HgStoreStateMachineTest.java +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/HgStoreStateMachineTest.java @@ -19,10 +19,23 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; import java.nio.ByteBuffer; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicIntegerArray; +import java.util.concurrent.atomic.AtomicReference; import org.apache.hugegraph.store.raft.HgStoreStateMachine; import org.apache.hugegraph.store.raft.RaftClosure; @@ -37,15 +50,21 @@ import org.mockito.Mock; import org.mockito.junit.MockitoJUnitRunner; -import com.alipay.sofa.jraft.Closure; -import com.alipay.sofa.jraft.Iterator; import com.alipay.sofa.jraft.Status; +import com.alipay.sofa.jraft.closure.ClosureQueueImpl; +import com.alipay.sofa.jraft.closure.SaveSnapshotClosure; import com.alipay.sofa.jraft.conf.Configuration; -import com.alipay.sofa.jraft.entity.LeaderChangeContext; +import com.alipay.sofa.jraft.core.FSMCallerImpl; +import com.alipay.sofa.jraft.core.NodeImpl; +import com.alipay.sofa.jraft.core.NodeMetrics; +import com.alipay.sofa.jraft.entity.EnumOutter; +import com.alipay.sofa.jraft.entity.LogEntry; +import com.alipay.sofa.jraft.entity.LogId; import com.alipay.sofa.jraft.entity.PeerId; -import com.alipay.sofa.jraft.entity.Task; import com.alipay.sofa.jraft.error.RaftError; import com.alipay.sofa.jraft.error.RaftException; +import com.alipay.sofa.jraft.option.FSMCallerOptions; +import com.alipay.sofa.jraft.storage.LogManager; @RunWith(MockitoJUnitRunner.class) public class HgStoreStateMachineTest { @@ -115,62 +134,137 @@ public void testIsLeader() { } @Test - public void testOnApply() { - RaftOperation op = RaftOperation.create((byte) 0b0); - final Task task = new Task(); - task.setData(ByteBuffer.wrap(op.getValues())); - task.setDone(new HgStoreStateMachine.RaftClosureAdapter(op, closure -> { + public void testLeaderApplySuccess() throws Exception { + assertApply(false, true); + } - })); + @Test + public void testFollowerApplySuccess() throws Exception { + assertApply(false, false); + } - List tasks = new ArrayList<>(); - tasks.add(task); - // Setup - final Iterator inter = new Iterator() { - final java.util.Iterator iterator = tasks.iterator(); - Task task; + @Test + public void testLeaderApplyFailure() throws Exception { + assertApply(true, true); + } - @Override - public ByteBuffer getData() { - return task.getData(); - } + @Test + public void testFollowerApplyFailure() throws Exception { + assertApply(true, false); + } + private void assertApply(boolean fail, boolean leader) throws Exception { + List invoked = new ArrayList<>(); + List notified = new ArrayList<>(); + hgStoreStateMachineUnderTest.addTaskHandler(new RaftTaskHandler() { @Override - public long getIndex() { - return 0; + public boolean invoke(int groupId, byte[] request, RaftClosure response) { + return apply(request[0]); } @Override - public long getTerm() { - return 0; + public boolean invoke(int groupId, byte methodId, Object req, RaftClosure response) { + return apply(methodId); } - @Override - public Closure done() { - return null; + private boolean apply(int value) { + invoked.add(value); + if (fail && value == 2) { + throw new IllegalStateException("injected apply failure"); + } + return true; } - + }); + hgStoreStateMachineUnderTest.addStateListener(new RaftStateListener() { @Override - public void setErrorAndRollback(long ntail, Status st) { - + public void onLeaderStart(long term) { } @Override - public boolean hasNext() { - return iterator.hasNext(); + public void onError(RaftException error) { } @Override - public ByteBuffer next() { - task = iterator.next(); - return task.getData(); + public void onDataCommitted(long index) { + notified.add(index); } - }; - - // Run the test - hgStoreStateMachineUnderTest.onApply(inter); - - // Verify the results + }); + + LogManager logs = mock(LogManager.class); + when(logs.getEntry(anyLong())).thenAnswer(invocation -> { + long index = invocation.getArgument(0); + LogEntry entry = new LogEntry(EnumOutter.EntryType.ENTRY_TYPE_DATA); + entry.setId(new LogId(index, 1)); + entry.setData(ByteBuffer.wrap(new byte[]{(byte) index})); + return entry; + }); + ClosureQueueImpl closures = new ClosureQueueImpl("store-apply-test"); + closures.resetFirstIndex(1); + AtomicIntegerArray calls = new AtomicIntegerArray(3); + AtomicIntegerArray codes = new AtomicIntegerArray(3); + CountDownLatch completed = new CountDownLatch(leader ? 3 : 0); + for (int i = 0; i < 3; i++) { + final int slot = i; + closures.appendPendingClosure(leader ? new HgStoreStateMachine.RaftClosureAdapter( + RaftOperation.create((byte) (i + 1)), status -> { + codes.set(slot, status.getCode()); + calls.incrementAndGet(slot); + completed.countDown(); + }) : null); + } + when(logs.getTerm(anyLong())).thenReturn(1L); + NodeImpl node = mock(NodeImpl.class); + when(node.getNodeMetrics()).thenReturn(new NodeMetrics(false)); + when(node.getGroupId()).thenReturn("store-apply-test"); + FSMCallerOptions options = new FSMCallerOptions(); + options.setNode(node); + options.setFsm(hgStoreStateMachineUnderTest); + options.setLogManager(logs); + options.setClosureQueue(closures); + options.setBootstrapId(new LogId(0, 0)); + options.setDisruptorBufferSize(16); + FSMCallerImpl caller = new FSMCallerImpl(); + assertTrue(caller.init(options)); + CountDownLatch batchApplied = new CountDownLatch(1); + caller.addLastAppliedLogIndexListener(index -> batchApplied.countDown()); + try { + assertTrue(caller.onCommitted(3)); + assertTrue(batchApplied.await(5, TimeUnit.SECONDS)); + assertEquals(fail ? 1L : 3L, caller.getLastAppliedIndex()); + assertEquals(3L, caller.getLastCommittedIndex()); + if (fail) { + assertTrue(caller.onCommitted(4)); + CountDownLatch snapshotCompleted = new CountDownLatch(1); + AtomicReference snapshotStatus = new AtomicReference<>(); + SaveSnapshotClosure snapshot = mock(SaveSnapshotClosure.class); + doAnswer(invocation -> { + snapshotStatus.set(invocation.getArgument(0)); + snapshotCompleted.countDown(); + return null; + }).when(snapshot).run(any(Status.class)); + assertTrue(caller.onSnapshotSave(snapshot)); + assertTrue(snapshotCompleted.await(5, TimeUnit.SECONDS)); + assertFalse(snapshotStatus.get().isOk()); + assertTrue(snapshotStatus.get().getErrorMsg() + .startsWith("FSMCaller is in bad status")); + verify(logs, never()).getConfiguration(anyLong()); + verify(snapshot, never()).start(any()); + } + } finally { + caller.shutdown(); + caller.join(); + } + assertEquals(fail ? Arrays.asList(1, 2) : Arrays.asList(1, 2, 3), invoked); + assertEquals(fail ? Arrays.asList(1L) : Arrays.asList(1L, 2L, 3L), notified); + assertEquals(fail ? 1L : 3L, hgStoreStateMachineUnderTest.getCommittedIndex()); + assertEquals(fail ? 1L : 3L, caller.getLastAppliedIndex()); + verify(logs).setAppliedId(new LogId(fail ? 1 : 3, 1)); + assertTrue(completed.await(5, TimeUnit.SECONDS)); + for (int i = 0; i < 3; i++) { + assertEquals(leader ? 1 : 0, calls.get(i)); + assertEquals(leader && fail && i > 0 ? RaftError.ESTATEMACHINE.getNumber() : 0, + codes.get(i)); + } } @Test From 4fb8c0f83a32dcdd5380d8e7f0fa46a3f4ff260b Mon Sep 17 00:00:00 2001 From: imbajin Date: Thu, 24 Sep 2026 23:11:29 +0800 Subject: [PATCH 3/9] chore: centralize existing jraft versions - keep Server and PD/Store runtime versions unchanged - share the Store production version with its tests - track compatibility and future alignment in #3239 --- hugegraph-pd/hg-pd-cli/pom.xml | 2 +- hugegraph-pd/hg-pd-core/pom.xml | 2 +- hugegraph-server/hugegraph-core/pom.xml | 1 - hugegraph-store/hg-store-core/pom.xml | 2 +- hugegraph-store/hg-store-test/pom.xml | 2 +- pom.xml | 3 +++ 6 files changed, 7 insertions(+), 5 deletions(-) diff --git a/hugegraph-pd/hg-pd-cli/pom.xml b/hugegraph-pd/hg-pd-cli/pom.xml index 4920174d76..01524c7479 100644 --- a/hugegraph-pd/hg-pd-cli/pom.xml +++ b/hugegraph-pd/hg-pd-cli/pom.xml @@ -49,7 +49,7 @@ com.alipay.sofa jraft-core - 1.3.13 + ${pd-store-jraft.version} org.rocksdb diff --git a/hugegraph-pd/hg-pd-core/pom.xml b/hugegraph-pd/hg-pd-core/pom.xml index e17570d592..38837f1e7d 100644 --- a/hugegraph-pd/hg-pd-core/pom.xml +++ b/hugegraph-pd/hg-pd-core/pom.xml @@ -38,7 +38,7 @@ com.alipay.sofa jraft-core - 1.3.13 + ${pd-store-jraft.version} org.rocksdb diff --git a/hugegraph-server/hugegraph-core/pom.xml b/hugegraph-server/hugegraph-core/pom.xml index b2519633a9..4e03478a53 100644 --- a/hugegraph-server/hugegraph-core/pom.xml +++ b/hugegraph-server/hugegraph-core/pom.xml @@ -29,7 +29,6 @@ ${basedir}/.. - 1.3.11 0.7.4 5.12.1 1.8.1 diff --git a/hugegraph-store/hg-store-core/pom.xml b/hugegraph-store/hg-store-core/pom.xml index 0ecf723280..19c255ae5e 100644 --- a/hugegraph-store/hg-store-core/pom.xml +++ b/hugegraph-store/hg-store-core/pom.xml @@ -55,7 +55,7 @@ com.alipay.sofa jraft-core - 1.3.13 + ${pd-store-jraft.version} org.rocksdb diff --git a/hugegraph-store/hg-store-test/pom.xml b/hugegraph-store/hg-store-test/pom.xml index 919b3c60aa..c99541b36c 100644 --- a/hugegraph-store/hg-store-test/pom.xml +++ b/hugegraph-store/hg-store-test/pom.xml @@ -201,7 +201,7 @@ com.alipay.sofa jraft-core - 1.3.13 + ${pd-store-jraft.version} org.rocksdb diff --git a/pom.xml b/pom.xml index 3581d0346f..65dee5a154 100644 --- a/pom.xml +++ b/pom.xml @@ -89,6 +89,9 @@ 5.6.0 1.7.0 1.18.30 + + 1.3.11 + 1.3.13 hugegraph 11 11 From 67f6285277a39fe75da7055d614c1c8aa4e423ad Mon Sep 17 00:00:00 2001 From: imbajin Date: Thu, 24 Sep 2026 23:33:15 +0800 Subject: [PATCH 4/9] chore: remove obsolete jraft dependency records - remove unused jraft 1.3.9 from the dependency inventory - remove its obsolete release license entry - verify the regenerated inventory matches CI checks --- install-dist/release-docs/LICENSE | 1 - install-dist/scripts/dependency/known-dependencies.txt | 1 - 2 files changed, 2 deletions(-) diff --git a/install-dist/release-docs/LICENSE b/install-dist/release-docs/LICENSE index 06a31044eb..b721884ff0 100644 --- a/install-dist/release-docs/LICENSE +++ b/install-dist/release-docs/LICENSE @@ -268,7 +268,6 @@ non-Apache-2.0 license text is required for a bundled component. https://central.sonatype.com/artifact/com.alipay.sofa/hessian/3.3.7 -> Apache 2.0 https://central.sonatype.com/artifact/com.alipay.sofa/jraft-core/1.3.11 -> Apache 2.0 https://central.sonatype.com/artifact/com.alipay.sofa/jraft-core/1.3.13 -> Apache 2.0 - https://central.sonatype.com/artifact/com.alipay.sofa/jraft-core/1.3.9 -> Apache 2.0 https://central.sonatype.com/artifact/com.alipay.sofa/sofa-rpc-all/5.7.6 -> Apache 2.0 https://central.sonatype.com/artifact/com.alipay.sofa/tracer-core/3.0.8 -> Apache 2.0 https://central.sonatype.com/artifact/com.beust/jcommander/1.30 -> Apache 2.0 diff --git a/install-dist/scripts/dependency/known-dependencies.txt b/install-dist/scripts/dependency/known-dependencies.txt index 157da1c064..29151532a6 100644 --- a/install-dist/scripts/dependency/known-dependencies.txt +++ b/install-dist/scripts/dependency/known-dependencies.txt @@ -284,7 +284,6 @@ joda-time-2.10.8.jar joni-2.2.1.jar jraft-core-1.3.11.jar jraft-core-1.3.13.jar -jraft-core-1.3.9.jar jsonassert-1.5.0.jar json-path-2.5.0.jar json-smart-2.3.jar From c285cbc56c7286f59ac48adce8eace1f471ea035 Mon Sep 17 00:00:00 2001 From: imbajin Date: Fri, 25 Sep 2026 00:24:07 +0800 Subject: [PATCH 5/9] fix(store): assert stable compaction lock release - remove the transient compaction completion assertion - wait for path lock release after worker interruption - release the reserved range lock in test cleanup --- .../core/snapshot/HgSnapshotHandlerTest.java | 43 ++++++------------- 1 file changed, 12 insertions(+), 31 deletions(-) diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/snapshot/HgSnapshotHandlerTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/snapshot/HgSnapshotHandlerTest.java index e32f1df498..4bf691a6cd 100644 --- a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/snapshot/HgSnapshotHandlerTest.java +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/snapshot/HgSnapshotHandlerTest.java @@ -32,6 +32,7 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -332,8 +333,8 @@ public void testDbCompactionSkipsWhenRangeLockStillHeldAfterWait() throws Interr * logs), skipping unlock(path) - so every later compaction for that partition would block * until the 6-hour path-lock timeout. Interrupts a real compactionPool worker thread while * it is parked in tryLock() (identified by stack trace, since the pool is shared), then - * confirms a second dbCompaction() call is able to complete instead of hanging behind the - * still-held path lock. + * confirms the path lock returns to its available state while the snapshot range lock + * remains owned by this test. */ @Test public void testDbCompactionReleasesPathLockWhenInterruptedWaitingForRangeLock() @@ -371,13 +372,16 @@ public void testDbCompactionReleasesPathLockWhenInterruptedWaitingForRangeLock() BusinessHandler.doing, pathLockBeforeInterrupt.get()); other.interrupt(); - // Give the interrupted task time to run its InterruptedException handling and - // return. - Thread.sleep(500); + // Wait for the interrupted task to release the path lock. No second compaction + // runs here, so the available state remains stable. + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5); + while (pathLockBeforeInterrupt.get() != BusinessHandler.compactionCanStart && + System.nanoTime() < deadline) { + Thread.sleep(10); + } // The path lock must have been released by the interrupted task's - // InterruptedException handler, directly confirming the fix rather than relying - // solely on the second dbCompaction() call below to prove it indirectly. + // InterruptedException handler, directly confirming the fix. String path = businessHandler.getLockPath(partitionId); AtomicInteger pathLockState = businessHandler.getPathLockState(path); assertNotNull("path lock must have been initialized by the interrupted task", @@ -395,31 +399,8 @@ public void testDbCompactionReleasesPathLockWhenInterruptedWaitingForRangeLock() rangeLockCheck.join(); assertFalse("range lock must still belong to the snapshot save", concurrentRangeLockResult.get()); - - // The path lock, however, must have been released by the interrupted task - - // otherwise this second dbCompaction() call would block on lock(path) until the - // 6-hour timeout instead of reaching the compacting state below once the range - // lock is released. Ownership of the range lock passes to the second - // dbCompaction() call's own worker thread, which acquires and releases it itself - - // do not touch compactionRangeLock again after this point. - businessHandler.unlockCompactionRange(partitionId); - businessHandler.dbCompaction("graph0", partitionId); - - // Compaction on the test's near-empty RocksDB completes almost immediately, so the - // transient "doing" state cannot be reliably observed here - poll for the - // terminal compactionDone state instead. What this proves is that dbCompaction() - // was able to acquire lock(path) at all: before the fix, the leaked path lock - // would have made this call hang on lock(path) until the 6-hour timeout instead - // of ever reaching compactionDone. - long start = System.currentTimeMillis(); - while (businessHandler.getState(partitionId).get() != BusinessHandler.compactionDone && - System.currentTimeMillis() - start < 5000) { - Thread.sleep(50); - } - assertEquals("second dbCompaction() must complete, proving the path lock was " + - "released rather than leaked by the interrupted task", - BusinessHandler.compactionDone, businessHandler.getState(partitionId).get()); } finally { + businessHandler.unlockCompactionRange(partitionId); BusinessHandlerImpl.setCompactionRangeLockWaitMillis(originalWaitMillis); } } From 181865d5d64c1f46d2ec96fd41bec18e75fc8c84 Mon Sep 17 00:00:00 2001 From: imbajin Date: Fri, 25 Sep 2026 08:48:40 +0800 Subject: [PATCH 6/9] fix: isolate raft callback and terminal errors - keep completion exceptions out of apply rollback - latch Store state errors before recovery checks - cover callback failures and restart protection - stabilize the monitor timestamp assertion - document manual recovery after state machine errors --- .../hugegraph/pd/raft/RaftStateMachine.java | 7 +- .../pd/core/RaftStateMachineTest.java | 19 +++++ .../pd/core/StoreMonitorDataServiceTest.java | 5 +- hugegraph-store/docs/operations-guide.md | 9 +++ .../hugegraph/store/PartitionEngine.java | 14 ++++ .../store/raft/DefaultRaftClosure.java | 15 +++- .../store/raft/PartitionStateMachine.java | 14 +++- hugegraph-store/hg-store-test/pom.xml | 1 + .../store/core/PartitionEngineErrorTest.java | 79 +++++++++++++++++++ .../core/raft/PartitionStateMachineTest.java | 45 ++++++++++- 10 files changed, 199 insertions(+), 9 deletions(-) create mode 100644 hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/PartitionEngineErrorTest.java diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java index f99caa2c0f..ce5eadc692 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java @@ -351,7 +351,12 @@ public KVStoreClosure getClosure() { @Override public void run(Status status) { - closure.run(status); + try { + closure.run(status); + } catch (Throwable t) { + // Response delivery must not turn an applied entry into an apply failure. + log.error("Raft completion callback failed", t); + } } @Override diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java index 2f1fdbb903..d28e2456c8 100644 --- a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java @@ -42,6 +42,7 @@ import com.alipay.sofa.jraft.entity.LogEntry; import com.alipay.sofa.jraft.entity.LogId; import com.alipay.sofa.jraft.error.RaftError; +import com.alipay.sofa.jraft.error.RaftException; import com.alipay.sofa.jraft.option.FSMCallerOptions; import com.alipay.sofa.jraft.storage.LogManager; @@ -62,7 +63,17 @@ public void testFollowerFailureStopsReplay() throws Exception { assertApply(false, true); } + @Test + public void testCompletionFailureDoesNotStopApply() throws Exception { + assertApply(true, false, true); + } + private void assertApply(boolean leader, boolean fail) throws Exception { + assertApply(leader, fail, false); + } + + private void assertApply(boolean leader, boolean fail, boolean throwFromCallback) + throws Exception { RaftStateMachine stateMachine = new RaftStateMachine(); List applied = new ArrayList<>(); stateMachine.addTaskHandler((operation, done) -> { @@ -88,8 +99,12 @@ private void assertApply(boolean leader, boolean fail) throws Exception { Mockito.when(logManager.getTerm(index)).thenReturn(1L); if (leader && index <= 3) { KVStoreClosure callback = Mockito.mock(KVStoreClosure.class); + boolean failCallback = throwFromCallback && index == 2; Mockito.doAnswer(invocation -> { completed.countDown(); + if (failCallback && ((Status) invocation.getArgument(0)).isOk()) { + throw new IllegalStateException("injected callback failure"); + } return null; }).when(callback).run(Mockito.any(Status.class)); callbacks.add(callback); @@ -144,6 +159,10 @@ private void assertApply(boolean leader, boolean fail) throws Exception { Assert.assertEquals(fail ? Arrays.asList(1) : Arrays.asList(1, 2, 3), applied); Assert.assertEquals(fail ? 1 : 3, caller.getLastAppliedIndex()); Mockito.verify(logManager).setAppliedId(new LogId(fail ? 1 : 3, 1)); + if (!fail) { + Mockito.verify(node, Mockito.never()) + .onError(Mockito.any(RaftException.class)); + } if (leader) { for (int index = 0; index < callbacks.size(); index++) { int expectedCode = fail && index > 0 ? RaftError.ESTATEMACHINE.getNumber() : 0; diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/StoreMonitorDataServiceTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/StoreMonitorDataServiceTest.java index 9caf9eaaeb..13c9db1ea7 100644 --- a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/StoreMonitorDataServiceTest.java +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/StoreMonitorDataServiceTest.java @@ -50,8 +50,9 @@ public void test() throws InterruptedException, PDException { now = System.currentTimeMillis() / 1000; Thread.sleep(1100); } - assertTrue(this.service.getLatestStoreMonitorDataTimeStamp(1) == 0 || - this.service.getLatestStoreMonitorDataTimeStamp(1) == now); + // The one-second lookup window can expire between two reads. + long latest = this.service.getLatestStoreMonitorDataTimeStamp(1); + assertTrue(latest == 0 || latest == now); var data = this.service.getStoreMonitorData(1); assertEquals(5, data.size()); diff --git a/hugegraph-store/docs/operations-guide.md b/hugegraph-store/docs/operations-guide.md index 6eb80ce136..1df19b8be6 100644 --- a/hugegraph-store/docs/operations-guide.md +++ b/hugegraph-store/docs/operations-guide.md @@ -448,6 +448,15 @@ df -h --- +### State Machine Apply Failures + +A partition that reports a Raft state machine error stops applying logs and is +not automatically restarted, including by the activity check. Inspect the Store +error log and repair the underlying storage or apply failure before manually +restarting the Store process. Restarting without repairing the cause may replay +the same failing entry. This protection does not undo writes already performed +by the failed entry. + ## Backup and Recovery ### Backup Strategies diff --git a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java index e18e730444..ab580565dc 100644 --- a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java +++ b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java @@ -83,6 +83,7 @@ import com.alipay.sofa.jraft.core.DefaultJRaftServiceFactory; import com.alipay.sofa.jraft.core.NodeMetrics; import com.alipay.sofa.jraft.core.Replicator; +import com.alipay.sofa.jraft.entity.EnumOutter.ErrorType; import com.alipay.sofa.jraft.entity.PeerId; import com.alipay.sofa.jraft.entity.Task; import com.alipay.sofa.jraft.error.RaftException; @@ -125,6 +126,7 @@ public class PartitionEngine implements Lifecycle, RaftS private SnapshotHandler snapshotHandler; private Node raftNode; private volatile boolean started; + private volatile boolean stateMachineError; public PartitionEngine(HgStoreEngine storeEngine, ShardGroup shardGroup) { this.storeEngine = storeEngine; @@ -577,6 +579,9 @@ public void shutdown() { * Restart raft engine */ public void restartRaftNode() { + if (this.stateMachineError) { + return; + } shutdown(); log.error("Raft {} is restarting !!!", getGroupId()); this.init(this.options); @@ -586,6 +591,9 @@ public void restartRaftNode() { * Check if it is active, if not, restart it. */ public void checkActivity() { + if (this.stateMachineError) { + return; + } Utils.runInThread(() -> { if (!this.raftNode.getNodeState().isActive()) { log.error("Raft {} is not activity state is {} ", @@ -821,6 +829,12 @@ public void onDataCommitted(long index) { @Override public void onError(RaftException e) { + if (e.getType() == ErrorType.ERROR_TYPE_STATE_MACHINE) { + this.stateMachineError = true; + log.error("Raft {} stopped after a state machine error; repair the cause and " + + "restart the Store process before resuming", getGroupId(), e); + return; + } this.restartRaftNode(); } diff --git a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/DefaultRaftClosure.java b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/DefaultRaftClosure.java index e98bc16e9d..b235222c67 100644 --- a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/DefaultRaftClosure.java +++ b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/DefaultRaftClosure.java @@ -19,9 +19,12 @@ import com.alipay.sofa.jraft.Status; +import lombok.extern.slf4j.Slf4j; + /** * @date 2023/9/8 **/ +@Slf4j public class DefaultRaftClosure implements RaftClosure { private RaftOperation operation; @@ -34,7 +37,17 @@ public DefaultRaftClosure(RaftOperation op, RaftClosure closure) { @Override public void run(Status status) { - closure.run(status); + try { + closure.run(status); + } catch (Throwable t) { + // Response delivery must not turn an applied entry into an apply failure. + log.error("Raft completion callback failed", t); + } + } + + @Override + public void onLeaderChanged(Integer partId, Long storeId) { + closure.onLeaderChanged(partId, storeId); } public RaftClosure getClosure() { diff --git a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/PartitionStateMachine.java b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/PartitionStateMachine.java index bc48d0429d..3de94b7474 100644 --- a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/PartitionStateMachine.java +++ b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/PartitionStateMachine.java @@ -33,6 +33,7 @@ import com.alipay.sofa.jraft.Status; import com.alipay.sofa.jraft.conf.Configuration; import com.alipay.sofa.jraft.core.StateMachineAdapter; +import com.alipay.sofa.jraft.entity.EnumOutter.ErrorType; import com.alipay.sofa.jraft.entity.LeaderChangeContext; import com.alipay.sofa.jraft.entity.RaftOutter; import com.alipay.sofa.jraft.error.RaftError; @@ -87,7 +88,7 @@ public void onApply(Iterator iter) { // Leader branch, call locally RaftOperation operation = done.getOperation(); if (handler.invoke(groupId, operation.getOp(), operation.getReq(), - done.getClosure())) { + done)) { done.run(Status.OK()); break; } @@ -131,9 +132,14 @@ public long getLeaderTerm() { @Override public void onError(final RaftException e) { log.error(String.format("Raft %s StateMachine on error {}", groupId), e); - Utils.runInThread(() -> { - stateListeners.forEach(listener -> listener.onError(e)); - }); + Runnable notifyListeners = () -> stateListeners.forEach(listener -> listener.onError(e)); + if (e.getType() == ErrorType.ERROR_TYPE_STATE_MACHINE) { + // Latch the terminal error before JRaft marks the node inactive for activity checks. + notifyListeners.run(); + } else { + // Other errors can restart the node and must not join the FSM thread here. + Utils.runInThread(notifyListeners); + } } @Override diff --git a/hugegraph-store/hg-store-test/pom.xml b/hugegraph-store/hg-store-test/pom.xml index c99541b36c..a284c1c5c4 100644 --- a/hugegraph-store/hg-store-test/pom.xml +++ b/hugegraph-store/hg-store-test/pom.xml @@ -251,6 +251,7 @@ **/CoreSuiteTest.java **/PartitionStateMachineTest.java + **/PartitionEngineErrorTest.java **/BatchGraphIsolationTest.java diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/PartitionEngineErrorTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/PartitionEngineErrorTest.java new file mode 100644 index 0000000000..234af0226b --- /dev/null +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/PartitionEngineErrorTest.java @@ -0,0 +1,79 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF 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 org.apache.hugegraph.store.core; + +import static org.junit.Assert.assertEquals; +import static org.mockito.Mockito.mock; + +import org.apache.hugegraph.store.HgStoreEngine; +import org.apache.hugegraph.store.PartitionEngine; +import org.apache.hugegraph.store.meta.ShardGroup; +import org.apache.hugegraph.store.options.PartitionEngineOptions; +import org.junit.Test; + +import com.alipay.sofa.jraft.entity.EnumOutter.ErrorType; +import com.alipay.sofa.jraft.error.RaftException; + +public class PartitionEngineErrorTest { + + @Test + public void testStateMachineErrorPreventsAutomaticRestart() { + CountingEngine engine = new CountingEngine(); + engine.onError(new RaftException(ErrorType.ERROR_TYPE_STATE_MACHINE)); + engine.restartRaftNode(); + engine.checkActivity(); + // A later error must not clear the terminal state either. + engine.onError(new RaftException(ErrorType.ERROR_TYPE_LOG)); + assertEquals(0, engine.shutdowns); + assertEquals(0, engine.starts); + } + + @Test + public void testOtherErrorsStillRestart() { + CountingEngine engine = new CountingEngine(); + engine.onError(new RaftException(ErrorType.ERROR_TYPE_LOG)); + assertEquals(1, engine.shutdowns); + assertEquals(1, engine.starts); + } + + private static class CountingEngine extends PartitionEngine { + + private int shutdowns; + private int starts; + + CountingEngine() { + super(mock(HgStoreEngine.class), mock(ShardGroup.class)); + } + + @Override + public Integer getGroupId() { + return 1; + } + + @Override + public void shutdown() { + this.shutdowns++; + } + + @Override + public synchronized boolean init(PartitionEngineOptions options) { + this.starts++; + return true; + } + } +} diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java index 3c5916fe25..384ca30872 100644 --- a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java @@ -182,6 +182,7 @@ private static Lock getInternalLock(PartitionStateMachine stateMachine) throws E field.setAccessible(true); return (Lock) field.get(stateMachine); } + @Test public void testLeaderApplySuccess() throws Exception { assertApply(false, true); @@ -202,7 +203,38 @@ public void testFollowerApplyFailure() throws Exception { assertApply(true, false); } + @Test + public void testStateMachineErrorNotifiesListenerSynchronously() { + PartitionStateMachine stateMachine = new PartitionStateMachine(0, mock(SnapshotHandler.class)); + RaftStateListener listener = mock(RaftStateListener.class); + AtomicReference notifiedOn = new AtomicReference<>(); + doAnswer(invocation -> { + notifiedOn.set(Thread.currentThread()); + return null; + }).when(listener).onError(any(RaftException.class)); + stateMachine.addStateListener(listener); + + stateMachine.onError(new RaftException(EnumOutter.ErrorType.ERROR_TYPE_STATE_MACHINE)); + + assertEquals(Thread.currentThread(), notifiedOn.get()); + } + + @Test + public void testCompletionFailureDoesNotStopApply() throws Exception { + assertApply(false, true, true, false); + } + + @Test + public void testHandlerCallbackFailureDoesNotStopApply() throws Exception { + assertApply(false, true, true, true); + } + private void assertApply(boolean fail, boolean leader) throws Exception { + assertApply(fail, leader, false, false); + } + + private void assertApply(boolean fail, boolean leader, boolean throwFromCallback, + boolean completeInHandler) throws Exception { PartitionStateMachine stateMachine = new PartitionStateMachine(0, mock(SnapshotHandler.class)); List invoked = new ArrayList<>(); List notified = new ArrayList<>(); @@ -214,7 +246,12 @@ public boolean invoke(int groupId, byte[] request, RaftClosure response) { @Override public boolean invoke(int groupId, byte methodId, Object req, RaftClosure response) { - return apply(methodId); + boolean handled = apply(methodId); + if (completeInHandler) { + response.run(Status.OK()); + return false; + } + return handled; } private boolean apply(int value) { @@ -260,6 +297,9 @@ public void onDataCommitted(long index) { codes.set(slot, status.getCode()); calls.incrementAndGet(slot); completed.countDown(); + if (throwFromCallback && slot == 1 && status.isOk()) { + throw new IllegalStateException("injected callback failure"); + } }) : null); } when(logs.getTerm(anyLong())).thenReturn(1L); @@ -309,6 +349,9 @@ public void onDataCommitted(long index) { assertEquals(fail ? 1L : 3L, stateMachine.getCommittedIndex()); assertEquals(fail ? 1L : 3L, caller.getLastAppliedIndex()); verify(logs).setAppliedId(new LogId(fail ? 1 : 3, 1)); + if (!fail) { + verify(node, never()).onError(any(RaftException.class)); + } assertTrue(completed.await(5, TimeUnit.SECONDS)); for (int i = 0; i < 3; i++) { assertEquals(leader ? 1 : 0, calls.get(i)); From c24e5a329b0a3feb6b68ee8b7f9ce2c601268248 Mon Sep 17 00:00:00 2001 From: imbajin Date: Fri, 25 Sep 2026 13:08:11 +0800 Subject: [PATCH 7/9] fix(store): preserve response and restart contracts - retain typed gRPC results through closure decorators - serialize restarts and recheck terminal errors after join - exercise real handlers and shutdown-error interleavings - verify queued activity checks reach the restart gate --- .../hugegraph/store/PartitionEngine.java | 20 ++- .../store/node/grpc/GrpcClosure.java | 4 + .../store/core/PartitionEngineErrorTest.java | 88 +++++++++++++- .../core/raft/PartitionStateMachineTest.java | 115 ++++++++++++++++++ 4 files changed, 221 insertions(+), 6 deletions(-) diff --git a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java index ab580565dc..78ac781375 100644 --- a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java +++ b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/PartitionEngine.java @@ -127,6 +127,8 @@ public class PartitionEngine implements Lifecycle, RaftS private Node raftNode; private volatile boolean started; private volatile boolean stateMachineError; + // FSM callbacks may need the engine monitor while restart waits for them to drain. + private final Object restartLock = new Object(); public PartitionEngine(HgStoreEngine storeEngine, ShardGroup shardGroup) { this.storeEngine = storeEngine; @@ -579,12 +581,19 @@ public void shutdown() { * Restart raft engine */ public void restartRaftNode() { - if (this.stateMachineError) { - return; + synchronized (this.restartLock) { + if (this.stateMachineError) { + return; + } + shutdown(); + // shutdown joins the old FSM, including its synchronous error notifications. + // An error while it drained must prevent reopening the same logs. + if (this.stateMachineError) { + return; + } + log.error("Raft {} is restarting !!!", getGroupId()); + this.init(this.options); } - shutdown(); - log.error("Raft {} is restarting !!!", getGroupId()); - this.init(this.options); } /** @@ -830,6 +839,7 @@ public void onDataCommitted(long index) { @Override public void onError(RaftException e) { if (e.getType() == ErrorType.ERROR_TYPE_STATE_MACHINE) { + // Do not acquire restartLock here: shutdown may be joining this FSM thread. this.stateMachineError = true; log.error("Raft {} stopped after a state machine error; repair the cause and " + "restart the Store process before resuming", getGroupId(), e); diff --git a/hugegraph-store/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/GrpcClosure.java b/hugegraph-store/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/GrpcClosure.java index a16cdb3210..f91d63036b 100644 --- a/hugegraph-store/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/GrpcClosure.java +++ b/hugegraph-store/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/GrpcClosure.java @@ -22,6 +22,7 @@ import java.util.Map; import org.apache.hugegraph.store.grpc.session.FeedbackRes; +import org.apache.hugegraph.store.raft.DefaultRaftClosure; import org.apache.hugegraph.store.raft.RaftClosure; import io.grpc.stub.StreamObserver; @@ -35,6 +36,9 @@ public abstract class GrpcClosure implements RaftClosure { * Set the output result to raftClosure, for Follower, raftClosure is empty. */ public static void setResult(RaftClosure raftClosure, V result) { + while (raftClosure instanceof DefaultRaftClosure) { + raftClosure = ((DefaultRaftClosure) raftClosure).getClosure(); + } GrpcClosure closure = (GrpcClosure) raftClosure; if (closure != null) { closure.setResult(result); diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/PartitionEngineErrorTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/PartitionEngineErrorTest.java index 234af0226b..8908680451 100644 --- a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/PartitionEngineErrorTest.java +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/PartitionEngineErrorTest.java @@ -18,7 +18,16 @@ package org.apache.hugegraph.store.core; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.lang.reflect.Field; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import org.apache.hugegraph.store.HgStoreEngine; import org.apache.hugegraph.store.PartitionEngine; @@ -26,6 +35,8 @@ import org.apache.hugegraph.store.options.PartitionEngineOptions; import org.junit.Test; +import com.alipay.sofa.jraft.core.NodeImpl; +import com.alipay.sofa.jraft.core.State; import com.alipay.sofa.jraft.entity.EnumOutter.ErrorType; import com.alipay.sofa.jraft.error.RaftException; @@ -36,7 +47,6 @@ public void testStateMachineErrorPreventsAutomaticRestart() { CountingEngine engine = new CountingEngine(); engine.onError(new RaftException(ErrorType.ERROR_TYPE_STATE_MACHINE)); engine.restartRaftNode(); - engine.checkActivity(); // A later error must not clear the terminal state either. engine.onError(new RaftException(ErrorType.ERROR_TYPE_LOG)); assertEquals(0, engine.shutdowns); @@ -51,10 +61,72 @@ public void testOtherErrorsStillRestart() { assertEquals(1, engine.starts); } + @Test + public void testErrorWhileShutdownDrainsPreventsReinitialization() throws Exception { + CountingEngine engine = new CountingEngine(); + CountDownLatch shuttingDown = new CountDownLatch(1); + CountDownLatch drained = new CountDownLatch(1); + ExecutorService executor = Executors.newFixedThreadPool(2); + engine.duringShutdown = () -> { + shuttingDown.countDown(); + try { + assertTrue(drained.await(5, TimeUnit.SECONDS)); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError(e); + } + }; + try { + Future restart = executor.submit(engine::restartRaftNode); + assertTrue(shuttingDown.await(5, TimeUnit.SECONDS)); + // An FSM error must be able to latch while restart waits for that FSM to drain. + executor.submit(() -> engine.onError( + new RaftException(ErrorType.ERROR_TYPE_STATE_MACHINE))).get(5, TimeUnit.SECONDS); + drained.countDown(); + restart.get(5, TimeUnit.SECONDS); + assertEquals(1, engine.shutdowns); + assertEquals(0, engine.starts); + } finally { + drained.countDown(); + executor.shutdownNow(); + } + } + + @Test + public void testQueuedActivityCheckCannotRestartAfterTerminalError() throws Exception { + CountingEngine engine = new CountingEngine(); + NodeImpl node = mock(NodeImpl.class); + Field nodeField = PartitionEngine.class.getDeclaredField("raftNode"); + nodeField.setAccessible(true); + nodeField.set(engine, node); + CountDownLatch checkingState = new CountDownLatch(1); + CountDownLatch errorReported = new CountDownLatch(1); + engine.restartFinished = new CountDownLatch(1); + when(node.getNodeState()).thenAnswer(invocation -> { + checkingState.countDown(); + assertTrue(errorReported.await(5, TimeUnit.SECONDS)); + return State.STATE_ERROR; + }); + try { + engine.checkActivity(); + assertTrue(checkingState.await(5, TimeUnit.SECONDS)); + engine.onError(new RaftException(ErrorType.ERROR_TYPE_STATE_MACHINE)); + errorReported.countDown(); + // Wait for the real activity check to reach and return from the restart gate. + assertTrue(engine.restartFinished.await(5, TimeUnit.SECONDS)); + assertEquals(0, engine.shutdowns); + assertEquals(0, engine.starts); + } finally { + errorReported.countDown(); + } + } + private static class CountingEngine extends PartitionEngine { private int shutdowns; private int starts; + private Runnable duringShutdown; + private CountDownLatch restartFinished; CountingEngine() { super(mock(HgStoreEngine.class), mock(ShardGroup.class)); @@ -65,9 +137,23 @@ public Integer getGroupId() { return 1; } + @Override + public void restartRaftNode() { + try { + super.restartRaftNode(); + } finally { + if (this.restartFinished != null) { + this.restartFinished.countDown(); + } + } + } + @Override public void shutdown() { this.shutdowns++; + if (this.duringShutdown != null) { + this.duringShutdown.run(); + } } @Override diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java index 384ca30872..f541523f48 100644 --- a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java @@ -19,6 +19,7 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; @@ -37,11 +38,27 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicIntegerArray; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.Lock; import org.apache.hugegraph.store.HgStoreEngine; +import org.apache.hugegraph.store.grpc.common.GraphMethod; +import org.apache.hugegraph.store.grpc.common.Header; +import org.apache.hugegraph.store.grpc.common.ResCode; +import org.apache.hugegraph.store.grpc.common.TableMethod; +import org.apache.hugegraph.store.grpc.session.BatchReq; +import org.apache.hugegraph.store.grpc.session.BatchWriteReq; +import org.apache.hugegraph.store.grpc.session.CleanReq; +import org.apache.hugegraph.store.grpc.session.FeedbackRes; +import org.apache.hugegraph.store.grpc.session.GraphReq; +import org.apache.hugegraph.store.grpc.session.TableReq; +import org.apache.hugegraph.store.node.AppConfig; +import org.apache.hugegraph.store.node.grpc.GrpcClosure; +import org.apache.hugegraph.store.node.grpc.HgStoreNodeService; +import org.apache.hugegraph.store.node.grpc.HgStoreSessionImpl; +import org.apache.hugegraph.store.node.grpc.HgStoreWrapperEx; import org.apache.hugegraph.store.raft.DefaultRaftClosure; import org.apache.hugegraph.store.raft.PartitionStateMachine; import org.apache.hugegraph.store.raft.RaftClosure; @@ -55,6 +72,7 @@ import org.junit.BeforeClass; import org.junit.Test; +import com.alipay.sofa.jraft.Iterator; import com.alipay.sofa.jraft.Status; import com.alipay.sofa.jraft.closure.ClosureQueueImpl; import com.alipay.sofa.jraft.closure.SaveSnapshotClosure; @@ -69,6 +87,7 @@ import com.alipay.sofa.jraft.option.FSMCallerOptions; import com.alipay.sofa.jraft.storage.LogManager; import com.alipay.sofa.jraft.storage.snapshot.SnapshotWriter; +import com.google.protobuf.GeneratedMessageV3; /** * Covers the EC_RKDB_SNAPSHOT_SAVE_BUSY_FAIL -> RaftError.EBUSY mapping at the @@ -229,6 +248,102 @@ public void testHandlerCallbackFailureDoesNotStopApply() throws Exception { assertApply(false, true, true, true); } + @Test + public void testLeaderHandlersDeliverGrpcResults() throws Exception { + assertGrpcHandlers(false); + } + + @Test + public void testLeaderHandlersIgnoreResponseDeliveryFailure() throws Exception { + assertGrpcHandlers(true); + } + + private void assertGrpcHandlers(boolean throwFromCallback) throws Exception { + Header header = Header.newBuilder().setGraph("apply-test").build(); + byte[] methods = {HgStoreNodeService.BATCH_OP, HgStoreNodeService.TABLE_OP, + HgStoreNodeService.GRAPH_OP, HgStoreNodeService.CLEAN_OP}; + GeneratedMessageV3[] requests = { + BatchReq.newBuilder().setHeader(header).setBatchId("batch-1") + .setWriteReq(BatchWriteReq.getDefaultInstance()).build(), + TableReq.newBuilder().setHeader(header).setTableName("table") + .setMethod(TableMethod.TABLE_METHOD_CREATE).build(), + GraphReq.newBuilder().setHeader(header).setGraphName(header.getGraph()) + .setMethod(GraphMethod.GRAPH_METHOD_DELETE).build(), + CleanReq.newBuilder().setHeader(header).build() + }; + for (int i = 0; i < methods.length; i++) { + HgStoreWrapperEx wrapper = mock(HgStoreWrapperEx.class); + when(wrapper.doTable(0, TableMethod.TABLE_METHOD_CREATE, header.getGraph(), "table")) + .thenReturn(true); + when(wrapper.doGraph(0, GraphMethod.GRAPH_METHOD_DELETE, header.getGraph())) + .thenReturn(true); + when(wrapper.doClean(header.getGraph(), 0)).thenReturn(true); + HgStoreSessionImpl session = new HgStoreSessionImpl(); + setField(session, "wrapper", wrapper); + HgStoreNodeService service = new HgStoreNodeService(new AppConfig()); + setField(service, "hgStoreSession", session); + HgStoreEngine engine = mock(HgStoreEngine.class); + setField(service, "storeEngine", engine); + PartitionStateMachine stateMachine = new PartitionStateMachine( + 0, mock(SnapshotHandler.class)); + stateMachine.addTaskHandler(service); + AtomicInteger callbacks = new AtomicInteger(); + AtomicReference callbackStatus = new AtomicReference<>(); + GrpcClosure response = new GrpcClosure() { + @Override + public void run(Status status) { + callbackStatus.set(status); + callbacks.incrementAndGet(); + if (throwFromCallback) { + throw new IllegalStateException("injected response delivery failure"); + } + } + }; + RaftOperation operation = RaftOperation.create(methods[i], requests[i]); + Iterator iterator = mock(Iterator.class); + when(iterator.hasNext()).thenReturn(true, false); + when(iterator.done()).thenReturn(new DefaultRaftClosure(operation, response)); + when(iterator.getData()).thenReturn(ByteBuffer.wrap(operation.getValues())); + when(iterator.getIndex()).thenReturn(1L); + + stateMachine.onApply(iterator); + + assertNotNull("handler must retain the typed response for op " + methods[i], + response.getResult()); + assertEquals(ResCode.RES_CODE_OK, response.getResult().getStatus().getCode()); + assertEquals(1, callbacks.get()); + assertTrue(callbackStatus.get().isOk()); + assertEquals(1L, stateMachine.getCommittedIndex()); + verify(iterator, never()).setErrorAndRollback(anyLong(), any(Status.class)); + verify(iterator).next(); + switch (methods[i]) { + case HgStoreNodeService.BATCH_OP: + verify(wrapper).doBatch(header.getGraph(), 0, + ((BatchReq) requests[i]).getWriteReq().getEntryList()); + break; + case HgStoreNodeService.TABLE_OP: + verify(wrapper).doTable(0, TableMethod.TABLE_METHOD_CREATE, + header.getGraph(), "table"); + break; + case HgStoreNodeService.GRAPH_OP: + verify(engine).deletePartition(0, header.getGraph()); + verify(wrapper).doGraph(0, GraphMethod.GRAPH_METHOD_DELETE, header.getGraph()); + break; + case HgStoreNodeService.CLEAN_OP: + verify(wrapper).doClean(header.getGraph(), 0); + break; + default: + throw new AssertionError("unexpected method"); + } + } + } + + private static void setField(Object target, String name, Object value) throws Exception { + Field field = target.getClass().getDeclaredField(name); + field.setAccessible(true); + field.set(target, value); + } + private void assertApply(boolean fail, boolean leader) throws Exception { assertApply(fail, leader, false, false); } From d3c148af6a0033bcc5239a0307d33d72bcf7e912 Mon Sep 17 00:00:00 2001 From: imbajin Date: Fri, 25 Sep 2026 15:05:55 +0800 Subject: [PATCH 8/9] fix(store): move node tests into owning module - preserve real handler and callback assertions - avoid compiling against the executable node jar - retain existing CI test discovery and packaging --- .../grpc/PartitionStateMachineGrpcTest.java | 155 ++++++++++++++++++ .../core/raft/PartitionStateMachineTest.java | 115 ------------- 2 files changed, 155 insertions(+), 115 deletions(-) create mode 100644 hugegraph-store/hg-store-node/src/test/java/org/apache/hugegraph/store/node/grpc/PartitionStateMachineGrpcTest.java diff --git a/hugegraph-store/hg-store-node/src/test/java/org/apache/hugegraph/store/node/grpc/PartitionStateMachineGrpcTest.java b/hugegraph-store/hg-store-node/src/test/java/org/apache/hugegraph/store/node/grpc/PartitionStateMachineGrpcTest.java new file mode 100644 index 0000000000..575ed56369 --- /dev/null +++ b/hugegraph-store/hg-store-node/src/test/java/org/apache/hugegraph/store/node/grpc/PartitionStateMachineGrpcTest.java @@ -0,0 +1,155 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF 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 org.apache.hugegraph.store.node.grpc; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.lang.reflect.Field; +import java.nio.ByteBuffer; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.hugegraph.store.HgStoreEngine; +import org.apache.hugegraph.store.grpc.common.GraphMethod; +import org.apache.hugegraph.store.grpc.common.Header; +import org.apache.hugegraph.store.grpc.common.ResCode; +import org.apache.hugegraph.store.grpc.common.TableMethod; +import org.apache.hugegraph.store.grpc.session.BatchReq; +import org.apache.hugegraph.store.grpc.session.BatchWriteReq; +import org.apache.hugegraph.store.grpc.session.CleanReq; +import org.apache.hugegraph.store.grpc.session.FeedbackRes; +import org.apache.hugegraph.store.grpc.session.GraphReq; +import org.apache.hugegraph.store.grpc.session.TableReq; +import org.apache.hugegraph.store.node.AppConfig; +import org.apache.hugegraph.store.raft.DefaultRaftClosure; +import org.apache.hugegraph.store.raft.PartitionStateMachine; +import org.apache.hugegraph.store.raft.RaftOperation; +import org.apache.hugegraph.store.snapshot.SnapshotHandler; +import org.junit.Test; + +import com.alipay.sofa.jraft.Iterator; +import com.alipay.sofa.jraft.Status; +import com.google.protobuf.GeneratedMessageV3; + +public class PartitionStateMachineGrpcTest { + + @Test + public void testLeaderHandlersDeliverGrpcResults() throws Exception { + assertGrpcHandlers(false); + } + + @Test + public void testLeaderHandlersIgnoreResponseDeliveryFailure() throws Exception { + assertGrpcHandlers(true); + } + + private void assertGrpcHandlers(boolean throwFromCallback) throws Exception { + Header header = Header.newBuilder().setGraph("apply-test").build(); + byte[] methods = {HgStoreNodeService.BATCH_OP, HgStoreNodeService.TABLE_OP, + HgStoreNodeService.GRAPH_OP, HgStoreNodeService.CLEAN_OP}; + GeneratedMessageV3[] requests = { + BatchReq.newBuilder().setHeader(header).setBatchId("batch-1") + .setWriteReq(BatchWriteReq.getDefaultInstance()).build(), + TableReq.newBuilder().setHeader(header).setTableName("table") + .setMethod(TableMethod.TABLE_METHOD_CREATE).build(), + GraphReq.newBuilder().setHeader(header).setGraphName(header.getGraph()) + .setMethod(GraphMethod.GRAPH_METHOD_DELETE).build(), + CleanReq.newBuilder().setHeader(header).build() + }; + for (int i = 0; i < methods.length; i++) { + HgStoreWrapperEx wrapper = mock(HgStoreWrapperEx.class); + when(wrapper.doTable(0, TableMethod.TABLE_METHOD_CREATE, header.getGraph(), "table")) + .thenReturn(true); + when(wrapper.doGraph(0, GraphMethod.GRAPH_METHOD_DELETE, header.getGraph())) + .thenReturn(true); + when(wrapper.doClean(header.getGraph(), 0)).thenReturn(true); + HgStoreSessionImpl session = new HgStoreSessionImpl(); + setField(session, "wrapper", wrapper); + HgStoreNodeService service = new HgStoreNodeService(new AppConfig()); + setField(service, "hgStoreSession", session); + HgStoreEngine engine = mock(HgStoreEngine.class); + setField(service, "storeEngine", engine); + PartitionStateMachine stateMachine = new PartitionStateMachine( + 0, mock(SnapshotHandler.class)); + stateMachine.addTaskHandler(service); + AtomicInteger callbacks = new AtomicInteger(); + AtomicReference callbackStatus = new AtomicReference<>(); + GrpcClosure response = new GrpcClosure() { + @Override + public void run(Status status) { + callbackStatus.set(status); + callbacks.incrementAndGet(); + if (throwFromCallback) { + throw new IllegalStateException("injected response delivery failure"); + } + } + }; + RaftOperation operation = RaftOperation.create(methods[i], requests[i]); + Iterator iterator = mock(Iterator.class); + when(iterator.hasNext()).thenReturn(true, false); + when(iterator.done()).thenReturn(new DefaultRaftClosure(operation, response)); + when(iterator.getData()).thenReturn(ByteBuffer.wrap(operation.getValues())); + when(iterator.getIndex()).thenReturn(1L); + + stateMachine.onApply(iterator); + + assertNotNull("handler must retain the typed response for op " + methods[i], + response.getResult()); + assertEquals(ResCode.RES_CODE_OK, response.getResult().getStatus().getCode()); + assertEquals(1, callbacks.get()); + assertTrue(callbackStatus.get().isOk()); + assertEquals(1L, stateMachine.getCommittedIndex()); + verify(iterator, never()).setErrorAndRollback(anyLong(), any(Status.class)); + verify(iterator).next(); + switch (methods[i]) { + case HgStoreNodeService.BATCH_OP: + verify(wrapper).doBatch(header.getGraph(), 0, + ((BatchReq) requests[i]).getWriteReq().getEntryList()); + break; + case HgStoreNodeService.TABLE_OP: + verify(wrapper).doTable(0, TableMethod.TABLE_METHOD_CREATE, + header.getGraph(), "table"); + break; + case HgStoreNodeService.GRAPH_OP: + verify(engine).deletePartition(0, header.getGraph()); + verify(wrapper).doGraph(0, GraphMethod.GRAPH_METHOD_DELETE, header.getGraph()); + break; + case HgStoreNodeService.CLEAN_OP: + verify(wrapper).doClean(header.getGraph(), 0); + break; + default: + throw new AssertionError("unexpected method"); + } + } + } + + private static void setField(Object target, String name, Object value) throws Exception { + Field field = target.getClass().getDeclaredField(name); + field.setAccessible(true); + field.set(target, value); + } + +} diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java index f541523f48..384ca30872 100644 --- a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java @@ -19,7 +19,6 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; -import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; @@ -38,27 +37,11 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicIntegerArray; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.Lock; import org.apache.hugegraph.store.HgStoreEngine; -import org.apache.hugegraph.store.grpc.common.GraphMethod; -import org.apache.hugegraph.store.grpc.common.Header; -import org.apache.hugegraph.store.grpc.common.ResCode; -import org.apache.hugegraph.store.grpc.common.TableMethod; -import org.apache.hugegraph.store.grpc.session.BatchReq; -import org.apache.hugegraph.store.grpc.session.BatchWriteReq; -import org.apache.hugegraph.store.grpc.session.CleanReq; -import org.apache.hugegraph.store.grpc.session.FeedbackRes; -import org.apache.hugegraph.store.grpc.session.GraphReq; -import org.apache.hugegraph.store.grpc.session.TableReq; -import org.apache.hugegraph.store.node.AppConfig; -import org.apache.hugegraph.store.node.grpc.GrpcClosure; -import org.apache.hugegraph.store.node.grpc.HgStoreNodeService; -import org.apache.hugegraph.store.node.grpc.HgStoreSessionImpl; -import org.apache.hugegraph.store.node.grpc.HgStoreWrapperEx; import org.apache.hugegraph.store.raft.DefaultRaftClosure; import org.apache.hugegraph.store.raft.PartitionStateMachine; import org.apache.hugegraph.store.raft.RaftClosure; @@ -72,7 +55,6 @@ import org.junit.BeforeClass; import org.junit.Test; -import com.alipay.sofa.jraft.Iterator; import com.alipay.sofa.jraft.Status; import com.alipay.sofa.jraft.closure.ClosureQueueImpl; import com.alipay.sofa.jraft.closure.SaveSnapshotClosure; @@ -87,7 +69,6 @@ import com.alipay.sofa.jraft.option.FSMCallerOptions; import com.alipay.sofa.jraft.storage.LogManager; import com.alipay.sofa.jraft.storage.snapshot.SnapshotWriter; -import com.google.protobuf.GeneratedMessageV3; /** * Covers the EC_RKDB_SNAPSHOT_SAVE_BUSY_FAIL -> RaftError.EBUSY mapping at the @@ -248,102 +229,6 @@ public void testHandlerCallbackFailureDoesNotStopApply() throws Exception { assertApply(false, true, true, true); } - @Test - public void testLeaderHandlersDeliverGrpcResults() throws Exception { - assertGrpcHandlers(false); - } - - @Test - public void testLeaderHandlersIgnoreResponseDeliveryFailure() throws Exception { - assertGrpcHandlers(true); - } - - private void assertGrpcHandlers(boolean throwFromCallback) throws Exception { - Header header = Header.newBuilder().setGraph("apply-test").build(); - byte[] methods = {HgStoreNodeService.BATCH_OP, HgStoreNodeService.TABLE_OP, - HgStoreNodeService.GRAPH_OP, HgStoreNodeService.CLEAN_OP}; - GeneratedMessageV3[] requests = { - BatchReq.newBuilder().setHeader(header).setBatchId("batch-1") - .setWriteReq(BatchWriteReq.getDefaultInstance()).build(), - TableReq.newBuilder().setHeader(header).setTableName("table") - .setMethod(TableMethod.TABLE_METHOD_CREATE).build(), - GraphReq.newBuilder().setHeader(header).setGraphName(header.getGraph()) - .setMethod(GraphMethod.GRAPH_METHOD_DELETE).build(), - CleanReq.newBuilder().setHeader(header).build() - }; - for (int i = 0; i < methods.length; i++) { - HgStoreWrapperEx wrapper = mock(HgStoreWrapperEx.class); - when(wrapper.doTable(0, TableMethod.TABLE_METHOD_CREATE, header.getGraph(), "table")) - .thenReturn(true); - when(wrapper.doGraph(0, GraphMethod.GRAPH_METHOD_DELETE, header.getGraph())) - .thenReturn(true); - when(wrapper.doClean(header.getGraph(), 0)).thenReturn(true); - HgStoreSessionImpl session = new HgStoreSessionImpl(); - setField(session, "wrapper", wrapper); - HgStoreNodeService service = new HgStoreNodeService(new AppConfig()); - setField(service, "hgStoreSession", session); - HgStoreEngine engine = mock(HgStoreEngine.class); - setField(service, "storeEngine", engine); - PartitionStateMachine stateMachine = new PartitionStateMachine( - 0, mock(SnapshotHandler.class)); - stateMachine.addTaskHandler(service); - AtomicInteger callbacks = new AtomicInteger(); - AtomicReference callbackStatus = new AtomicReference<>(); - GrpcClosure response = new GrpcClosure() { - @Override - public void run(Status status) { - callbackStatus.set(status); - callbacks.incrementAndGet(); - if (throwFromCallback) { - throw new IllegalStateException("injected response delivery failure"); - } - } - }; - RaftOperation operation = RaftOperation.create(methods[i], requests[i]); - Iterator iterator = mock(Iterator.class); - when(iterator.hasNext()).thenReturn(true, false); - when(iterator.done()).thenReturn(new DefaultRaftClosure(operation, response)); - when(iterator.getData()).thenReturn(ByteBuffer.wrap(operation.getValues())); - when(iterator.getIndex()).thenReturn(1L); - - stateMachine.onApply(iterator); - - assertNotNull("handler must retain the typed response for op " + methods[i], - response.getResult()); - assertEquals(ResCode.RES_CODE_OK, response.getResult().getStatus().getCode()); - assertEquals(1, callbacks.get()); - assertTrue(callbackStatus.get().isOk()); - assertEquals(1L, stateMachine.getCommittedIndex()); - verify(iterator, never()).setErrorAndRollback(anyLong(), any(Status.class)); - verify(iterator).next(); - switch (methods[i]) { - case HgStoreNodeService.BATCH_OP: - verify(wrapper).doBatch(header.getGraph(), 0, - ((BatchReq) requests[i]).getWriteReq().getEntryList()); - break; - case HgStoreNodeService.TABLE_OP: - verify(wrapper).doTable(0, TableMethod.TABLE_METHOD_CREATE, - header.getGraph(), "table"); - break; - case HgStoreNodeService.GRAPH_OP: - verify(engine).deletePartition(0, header.getGraph()); - verify(wrapper).doGraph(0, GraphMethod.GRAPH_METHOD_DELETE, header.getGraph()); - break; - case HgStoreNodeService.CLEAN_OP: - verify(wrapper).doClean(header.getGraph(), 0); - break; - default: - throw new AssertionError("unexpected method"); - } - } - } - - private static void setField(Object target, String name, Object value) throws Exception { - Field field = target.getClass().getDeclaredField(name); - field.setAccessible(true); - field.set(target, value); - } - private void assertApply(boolean fail, boolean leader) throws Exception { assertApply(fail, leader, false, false); } From 3e028d30078de78ba5d79cf71b067e7df1f3d10d Mon Sep 17 00:00:00 2001 From: imbajin Date: Fri, 25 Sep 2026 23:56:17 +0800 Subject: [PATCH 9/9] fix: treat raft error messages as data - use a fixed Status format in PD and Store - exercise percent directives in existing failure tests - preserve rollback reporting for arbitrary messages --- .../java/org/apache/hugegraph/pd/raft/RaftStateMachine.java | 2 +- .../java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java | 2 +- .../org/apache/hugegraph/store/raft/PartitionStateMachine.java | 2 +- .../hugegraph/store/core/raft/PartitionStateMachineTest.java | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java index ce5eadc692..d07ab75f8c 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java @@ -127,7 +127,7 @@ public void onApply(Iterator iter) { } catch (Throwable t) { log.error("StateMachine encountered critical error", t); // JRaft completes the failed and remaining closures with the state machine error. - iter.setErrorAndRollback(1, new Status(RaftError.ESTATEMACHINE, t.getMessage())); + iter.setErrorAndRollback(1, new Status(RaftError.ESTATEMACHINE, "%s", t.getMessage())); return; } iter.next(); diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java index d28e2456c8..90efbb1dc2 100644 --- a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/RaftStateMachineTest.java @@ -79,7 +79,7 @@ private void assertApply(boolean leader, boolean fail, boolean throwFromCallback stateMachine.addTaskHandler((operation, done) -> { int key = operation.getKey()[0]; if (fail && key == 2) { - throw new IllegalStateException("injected apply failure"); + throw new IllegalStateException("injected apply failure: %s, 100%"); } applied.add(key); return true; diff --git a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/PartitionStateMachine.java b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/PartitionStateMachine.java index 3de94b7474..cf74810aec 100644 --- a/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/PartitionStateMachine.java +++ b/hugegraph-store/hg-store-core/src/main/java/org/apache/hugegraph/store/raft/PartitionStateMachine.java @@ -106,7 +106,7 @@ public void onApply(Iterator iter) { done.getOperation().getOp(), done.getOperation().getReq()); } - iter.setErrorAndRollback(1, new Status(RaftError.ESTATEMACHINE, t.getMessage())); + iter.setErrorAndRollback(1, new Status(RaftError.ESTATEMACHINE, "%s", t.getMessage())); // Do not publish the failed entry as applied. return; } diff --git a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java index 384ca30872..469f5e08c7 100644 --- a/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java +++ b/hugegraph-store/hg-store-test/src/main/java/org/apache/hugegraph/store/core/raft/PartitionStateMachineTest.java @@ -257,7 +257,7 @@ public boolean invoke(int groupId, byte methodId, Object req, RaftClosure respon private boolean apply(int value) { invoked.add(value); if (fail && value == 2) { - throw new IllegalStateException("injected apply failure"); + throw new IllegalStateException("injected apply failure: %s, 100%"); } return true; }