Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -2146,7 +2146,10 @@ protected TSStatus waitingProcedureFinished(
&& System.currentTimeMillis() - startTimeForCurrentProcedure < procedureWaitRetryTimeout) {
sleepWithoutInterrupt(PROCEDURE_WAIT_RETRY_TIMEOUT);
}
if (!procedure.isFinished()) {
// Conflict checks can fail before a procedure is assigned a persistent id. In that case the
// procedure remains at NO_PROC_ID but already carries the user-facing failure.
if (!procedure.isFinished()
&& !(procedure.getProcId() == Procedure.NO_PROC_ID && procedure.isFailed())) {
// The procedure is still executing
status =
RpcUtils.getStatus(TSStatusCode.INTERNAL_REQUEST_TIME_OUT, PROCEDURE_TIMEOUT_MESSAGE);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -267,6 +267,15 @@ protected boolean isYieldAfterExecution(Env env) {
return false;
}

/**
* Delay in milliseconds before retrying a failed rollback. A negative value keeps the default
* behavior of completing rollback after an exception. Procedures opting in must make rollback
* idempotent and retain the state needed by the next attempt.
*/
protected long getRollbackRetryTimeout() {
return -1;
}

// -------------------------Internal methods - called by the procedureExecutor------------------
final boolean tryAcquireExecution() {
return executing.compareAndSet(false, true);
Expand Down Expand Up @@ -696,6 +705,11 @@ protected synchronized void setFailure(final ProcedureException exception) {
*/
protected synchronized boolean setTimeoutFailure(Env env) {
if (state == ProcedureState.WAITING_TIMEOUT) {
if (exception != null && getRollbackRetryTimeout() >= 0) {
// Resume compensation without replacing the original execution failure.
setState(ProcedureState.FAILED);
return true;
}
long timeDiff = System.currentTimeMillis() - lastUpdate;
setFailure(
"ProcedureExecutor",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -738,15 +738,30 @@ private ProcedureLockState acquireLock(Procedure<Env> proc) {
* @return procedure lock state
*/
private ProcedureLockState executeRollback(Procedure<Env> procedure) {
boolean rollbackFailed = false;
try {
procedure.doRollback(this.environment);
} catch (IOException e) {
rollbackFailed = true;
LOG.error(ProcedureMessages.ROLL_BACK_FAILED_FOR, procedure, e);
} catch (InterruptedException e) {
rollbackFailed = true;
LOG.warn(ProcedureMessages.INTERRUPTED_EXCEPTION_OCCURRED_FOR, procedure, e);
} catch (Throwable t) {
rollbackFailed = true;
LOG.error(ProcedureMessages.CODE_BUG_RUNTIME_EXCEPTION_FOR, procedure, t);
}
if (rollbackFailed && procedure.getRollbackRetryTimeout() >= 0) {
procedure.setTimeout(procedure.getRollbackRetryTimeout());
procedure.setState(ProcedureState.WAITING_TIMEOUT);
try {
store.update(procedure);
} catch (Exception e) {
LOG.warn(ProcedureMessages.FAILED_TO_UPDATE_PROCEDURE, procedure, e);
}
timeoutExecutor.add(procedure);
return ProcedureLockState.LOCK_EVENT_WAIT;
}
cleanupAfterRollback(procedure);
return ProcedureLockState.LOCK_ACQUIRED;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -235,11 +235,15 @@ protected void rollback(final Env env)
states.removeLast();
}

boolean completed = false;
try {
updateTimestamp();
rollbackState(env, getCurrentState());
completed = true;
} finally {
states.removeLast();
if ((completed || getRollbackRetryTimeout() < 0) && !states.isEmpty()) {
states.removeLast();
}
updateTimestamp();
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
import org.apache.iotdb.confignode.procedure.ProcedureExecutor;
import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
import org.apache.iotdb.confignode.procedure.env.RemoveDataNodeHandler;
import org.apache.iotdb.confignode.procedure.exception.ProcedureException;
import org.apache.iotdb.confignode.procedure.impl.node.RemoveDataNodesProcedure;
import org.apache.iotdb.confignode.procedure.impl.region.RegionMigrateProcedure;
import org.apache.iotdb.confignode.procedure.impl.region.RegionMigrationPlan;
Expand All @@ -52,6 +53,7 @@
import java.util.concurrent.ConcurrentHashMap;

import static org.apache.iotdb.db.service.RegionMigrateService.isFailed;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.when;

Expand Down Expand Up @@ -149,6 +151,19 @@
Assert.assertTrue(isFailed(status));
}

@Test
public void testPreSubmissionFailureWithNoProcedureIdIsReportedImmediately() {
final Procedure<ConfigNodeProcedureEnv> procedure = mock(Procedure.class);
when(procedure.getProcId()).thenReturn(Procedure.NO_PROC_ID);
when(procedure.isFinished()).thenReturn(false);
when(procedure.isFailed()).thenReturn(true);
when(procedure.getException()).thenReturn(new ProcedureException("conflict"));

Assert.assertEquals(
TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode(),

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / zh-locale-compile

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / zh-locale-compile

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-basic (17, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-basic (17, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-enhanced (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-enhanced (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-basic (17, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-basic (17, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / Ubuntu

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / Ubuntu

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / unit-test (17, ubuntu-latest, others)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / unit-test (17, ubuntu-latest, others)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-basic (17, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-basic (17, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-basic (17, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-basic (17, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-basic (17, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-basic (17, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-manual (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-manual (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-enhanced (17, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-enhanced (17, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-enhanced (17, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-enhanced (17, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / Ubuntu

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / Ubuntu

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-basic (17, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-basic (17, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / single (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / single (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-manual (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-manual (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-enhanced (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-enhanced (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-manual (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-manual (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 0)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-enhanced (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-tree-auto-enhanced (17, HighPerformanceMode, HighPerformanceMode, ubuntu-latest, 1)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / triple (17, ScalableSingleNodeMode, ScalableSingleNodeMode, ScalableSingleNodeMode, ubuntu-latest)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / triple (17, ScalableSingleNodeMode, ScalableSingleNodeMode, ScalableSingleNodeMode, ubuntu-latest)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-enhanced (17, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / dual-table-manual-enhanced (17, HighPerformanceMode, ubuntu-latest, 2)

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / Ubuntu

package TSStatusCode does not exist

Check failure on line 163 in iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java

View workflow job for this annotation

GitHub Actions / Ubuntu

package TSStatusCode does not exist
PROCEDURE_MANAGER.waitingProcedureFinished(procedure, 0).getCode());
}

@Test
public void testCheckRemoveDataNodeWithConflictRegionMigrateProcedure() {
RegionMigrateProcedure regionMigrateProcedure =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,15 +21,25 @@

import org.apache.iotdb.confignode.procedure.entity.SimpleSTMProcedure;
import org.apache.iotdb.confignode.procedure.env.TestProcEnv;
import org.apache.iotdb.confignode.procedure.exception.ProcedureException;
import org.apache.iotdb.confignode.procedure.impl.StateMachineProcedure;
import org.apache.iotdb.confignode.procedure.state.ProcedureState;
import org.apache.iotdb.confignode.procedure.util.ProcedureTestUtil;

import org.junit.Assert;
import org.junit.Test;

import java.io.ByteArrayOutputStream;
import java.io.DataOutputStream;
import java.io.IOException;
import java.lang.reflect.Field;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.concurrent.ConcurrentLinkedDeque;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

public class STMProcedureTest extends TestProcedureBase {
Expand Down Expand Up @@ -60,6 +70,104 @@ public void testRolledBackProcedure() {
Assert.assertEquals(1 + success - rolledback, acc.get());
}

@Test
public void testFailedBeforeExecutionCanRollbackWithoutState() throws Exception {
final RetryingRollbackProcedure procedure = new RetryingRollbackProcedure();
procedure.setFailure(procedure.originalFailure);
procedure.doRollback(env);
Assert.assertEquals(Arrays.asList(0), procedure.attemptedStates);
Assert.assertSame(procedure.originalFailure, procedure.getException());
}

@Test
public void testFailedRollbackRetainsStateAfterSerialization() throws Exception {
final RetryingRollbackProcedure procedure = new RetryingRollbackProcedure();
procedure.setState(ProcedureState.RUNNABLE);
procedure.doExecute(env);
procedure.doExecute(env);
try {
procedure.doRollback(env);
Assert.fail("Compensation should fail while the Region is unavailable");
} catch (IOException expected) {
Assert.assertEquals(Arrays.asList(1), procedure.attemptedStates);
}
procedure.setTimeout(1000);
procedure.setState(ProcedureState.WAITING_TIMEOUT);
final ByteArrayOutputStream bytes = new ByteArrayOutputStream();
procedure.serialize(new DataOutputStream(bytes));
final RetryingRollbackProcedure restored = new RetryingRollbackProcedure();
restored.deserialize(ByteBuffer.wrap(bytes.toByteArray()));
Assert.assertTrue(restored.isFailed());
Assert.assertEquals(
procedure.originalFailure.getMessage(), restored.getException().getMessage());
restored.attemptedStates.addAll(Arrays.asList(1, 1));
restored.doRollback(env);
restored.doRollback(env);
Assert.assertEquals(Arrays.asList(1, 1, 1, 0), restored.attemptedStates);
}

@Test
public void testFailedRollbackRetriesSameStateWithoutOverwritingFailure() throws Exception {
final RetryingRollbackProcedure procedure = new RetryingRollbackProcedure();
final long procId = procExecutor.submitProcedure(procedure);
Assert.assertTrue(procedure.firstFailure.await(5, TimeUnit.SECONDS));
Assert.assertTrue(procedure.completed.await(5, TimeUnit.SECONDS));
ProcedureTestUtil.waitForProcedure(procExecutor, procId);
Assert.assertTrue(procedure.isFinished());
Assert.assertEquals(Arrays.asList(1, 1, 1, 0), procedure.attemptedStates);
Assert.assertSame(procedure.originalFailure, procedure.getException());
}

private static class RetryingRollbackProcedure
extends StateMachineProcedure<TestProcEnv, Integer> {
private final CountDownLatch firstFailure = new CountDownLatch(1);
private final CountDownLatch completed = new CountDownLatch(1);
private final List<Integer> attemptedStates = new ArrayList<>();
private final ProcedureException originalFailure = new ProcedureException("Execution failed");

@Override
protected Flow executeFromState(TestProcEnv env, Integer state) {
if (state == 0) {
setNextState(1);
return Flow.HAS_MORE_STATE;
}
setFailure(originalFailure);
return Flow.NO_MORE_STATE;
}

@Override
protected void rollbackState(TestProcEnv env, Integer state) throws IOException {
attemptedStates.add(state);
if (state == 1 && attemptedStates.size() < 3) {
firstFailure.countDown();
throw new IOException("Region temporarily unavailable");
}
if (state == 0) {
completed.countDown();
}
}

@Override
protected long getRollbackRetryTimeout() {
return 50;
}

@Override
protected Integer getState(int stateId) {
return stateId;
}

@Override
protected int getStateId(Integer state) {
return state;
}

@Override
protected Integer getInitialState() {
return 0;
}
}

@Test
public void testEofStateReexecutionDoesNotCallExecuteFromState() throws Exception {
EofReexecutionProcedure procedure = new EofReexecutionProcedure();
Expand Down
Loading