Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion hugegraph-pd/hg-pd-cli/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@
<dependency>
<groupId>com.alipay.sofa</groupId>
<artifactId>jraft-core</artifactId>
<version>1.3.13</version>
<version>${pd-store-jraft.version}</version>
<exclusions>
<exclusion>
<groupId>org.rocksdb</groupId>
Expand Down
2 changes: 1 addition & 1 deletion hugegraph-pd/hg-pd-core/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@
<groupId>com.alipay.sofa</groupId>
<artifactId>jraft-core</artifactId>
<!-- TODO: use open source version & adopt the code later -->
<version>1.3.13</version>
<version>${pd-store-jraft.version}</version>
<exclusions>
<exclusion>
<groupId>org.rocksdb</groupId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -125,10 +125,10 @@ public void onApply(Iterator iter) {
done.run(Status.OK());
}
} catch (Throwable t) {
log.error("StateMachine meet critical error: {}.", t);
if (done != null) {
done.run(new Status(RaftError.EINTERNAL, t.getMessage()));
}
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, "%s", t.getMessage()));
return;
}
iter.next();
}
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@

@RunWith(Suite.class)
@Suite.SuiteClasses({
RaftStateMachineTest.class,
MetadataKeyHelperTest.class,
HgKVStoreImplTest.class,
PDConfigTest.class,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,175 @@
/*
* 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.error.RaftException;
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);
}

@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<Integer> applied = new ArrayList<>();
stateMachine.addTaskHandler((operation, done) -> {
int key = operation.getKey()[0];
if (fail && key == 2) {
throw new IllegalStateException("injected apply failure: %s, 100%");
}
applied.add(key);
return true;
});

LogManager logManager = Mockito.mock(LogManager.class);
ClosureQueueImpl closures = new ClosureQueueImpl("pd-apply-test");
closures.resetFirstIndex(1);
List<KVStoreClosure> 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);
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);
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<Status> 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 (!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;
Mockito.verify(callbacks.get(index), Mockito.times(1))
.run(Mockito.argThat(status -> status.getCode() == expectedCode));
Mockito.verifyNoMoreInteractions(callbacks.get(index));
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand Down
1 change: 0 additions & 1 deletion hugegraph-server/hugegraph-core/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@

<properties>
<top.level.dir>${basedir}/..</top.level.dir>
<jraft.version>1.3.11</jraft.version>
<ohc.version>0.7.4</ohc.version>
<jna.version>5.12.1</jna.version>
<lz4.version>1.8.1</lz4.version>
Expand Down
9 changes: 9 additions & 0 deletions hugegraph-store/docs/operations-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion hugegraph-store/hg-store-core/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@
<dependency>
<groupId>com.alipay.sofa</groupId>
<artifactId>jraft-core</artifactId>
<version>1.3.13</version>
<version>${pd-store-jraft.version}</version>
<exclusions>
<exclusion>
<groupId>org.rocksdb</groupId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -125,6 +126,9 @@ public class PartitionEngine implements Lifecycle<PartitionEngineOptions>, RaftS
private SnapshotHandler snapshotHandler;
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;
Expand Down Expand Up @@ -577,15 +581,28 @@ public void shutdown() {
* Restart raft engine
*/
public void restartRaftNode() {
shutdown();
log.error("Raft {} is restarting !!!", getGroupId());
this.init(this.options);
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);
}
Comment thread
imbajin marked this conversation as resolved.
}

/**
* 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 {} ",
Expand Down Expand Up @@ -821,6 +838,13 @@ 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);
return;
}
this.restartRaftNode();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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() {
Expand Down
Loading
Loading