diff --git a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ProcedureMessages.java b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ProcedureMessages.java index 8c0be448f7f8b..12f1bed3ea54f 100644 --- a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ProcedureMessages.java +++ b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ProcedureMessages.java @@ -1408,4 +1408,28 @@ private ProcedureMessages() {} public static final String MESSAGE_UNEXPECTED_DATAPARTITIONTABLEINTEGRITYCHECKPROCEDURESTATE_ARG_WHEN_SHOWING_PROGRESS_D3C07BA1 = "Unexpected DataPartitionTableIntegrityCheckProcedureState {} when showing progress"; + public static final String MESSAGE_NO_PROCEDURE_WORKER_IS_CURRENTLY_AVAILABLE_WORKERS_MAY_BE_BUSY_OR_BLOCKED_BY_OTHER_PROCEDURES_AB0B1595 = + "no Procedure worker is currently available; workers may be busy or blocked by other procedures."; + public static final String MESSAGE_PIPE_OPERATION_ARG_TIMED_OUT_PROCEDUREID_ARG_STUCK_AT_ARG_REASON_ARG_THE_PROCEDURE_IS_STILL_RUNNING_7EEAC50E = + "Pipe operation %s timed out (procedureId=%d). Stuck at %s. Reason: %s. The procedure is still running."; + public static final String MESSAGE_WAITING_TO_ACQUIRE_THE_PIPETASKCOORDINATOR_LOCK_BECAUSE_ANOTHER_PIPE_OPERATION_IS_HOLDING_IT_25A3B6B8 = + "waiting to acquire the PipeTaskCoordinator lock because another Pipe operation is holding it."; + public static final String MESSAGE_WAITING_TO_ACQUIRE_THE_CONFIGNODE_NODE_LOCK_BECAUSE_ANOTHER_NODE_PROCEDURE_IS_HOLDING_IT_56494E86 = + "waiting to acquire the ConfigNode node lock because another node procedure is holding it."; + public static final String MESSAGE_PIPE_REQUEST_OR_PLUGIN_VALIDATION_HAS_NOT_COMPLETED_A_PLUGIN_CHECK_OR_METADATA_ACCESS_MAY_BE_SLOW_57C36CEF = + "Pipe request or plugin validation has not completed; a plugin check or metadata access may be slow."; + public static final String MESSAGE_PIPE_METADATA_CALCULATION_HAS_NOT_COMPLETED_METADATA_ACCESS_OR_LOCAL_CALCULATION_MAY_BE_SLOW_DEBF2504 = + "Pipe metadata calculation has not completed; metadata access or local calculation may be slow."; + public static final String MESSAGE_THE_CONFIGNODE_CONSENSUS_WRITE_HAS_NOT_RETURNED_THE_CONSENSUS_GROUP_MAY_BE_UNAVAILABLE_OR_SLOW_F8911CE7 = + "the ConfigNode consensus write has not returned; the consensus group may be unavailable or slow."; + public static final String MESSAGE_ONE_OR_MORE_DATANODES_HAVE_NOT_RESPONDED_TO_THE_PIPE_METADATA_PUSH_THEY_MAY_BE_UNAVAILABLE_OR_SLOW_11BBB333 = + "one or more DataNodes have not responded to the Pipe metadata push; they may be unavailable or slow."; + public static final String MESSAGE_THE_PREVIOUS_ATTEMPT_FAILED_WITH_ARG_AND_THIS_STATE_IS_BEING_RETRIED_7A541F27 = + "the previous attempt failed with '%s' and this state is being retried."; + public static final String MESSAGE_THE_STATE_FAILED_WITH_ARG_AND_ROLLBACK_IS_PENDING_E7B43829 = + "the state failed with '%s' and rollback is pending."; + public static final String MESSAGE_ROLLING_BACK_AFTER_FAILURE_ARG_474DF456 = + "rolling back after failure: %s."; + public static final String MESSAGE_ROLLING_BACK_AFTER_AN_EARLIER_FAILURE_850D0AF5 = + "rolling back after an earlier failure."; } diff --git a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ProcedureMessages.java b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ProcedureMessages.java index 91bf304b3266f..1b7ca6a0fde58 100644 --- a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ProcedureMessages.java +++ b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ProcedureMessages.java @@ -1484,4 +1484,38 @@ private ProcedureMessages() {} public static final String MESSAGE_UNEXPECTED_DATAPARTITIONTABLEINTEGRITYCHECKPROCEDURESTATE_ARG_WHEN_SHOWING_PROGRESS_D3C07BA1 = "非预期的 DataPartitionTableIntegrityCheckProcedureState {}(显示进度时)"; + public static final String + MESSAGE_PIPE_OPERATION_ARG_TIMED_OUT_PROCEDUREID_ARG_STUCK_AT_ARG_REASON_ARG_THE_PROCEDURE_IS_STILL_RUNNING_7EEAC50E = + "Pipe 操作 %s 超时(procedureId=%d)。卡在 %s。原因:%s。该 Procedure 仍在运行。"; + public static final String + MESSAGE_NO_PROCEDURE_WORKER_IS_CURRENTLY_AVAILABLE_WORKERS_MAY_BE_BUSY_OR_BLOCKED_BY_OTHER_PROCEDURES_AB0B1595 = + "当前没有可用的 Procedure worker;worker 可能正忙或被其他 Procedure 阻塞。"; + public static final String + MESSAGE_WAITING_TO_ACQUIRE_THE_PIPETASKCOORDINATOR_LOCK_BECAUSE_ANOTHER_PIPE_OPERATION_IS_HOLDING_IT_25A3B6B8 = + "正在等待获取 PipeTaskCoordinator 锁,因为另一个 Pipe 操作正在持有该锁。"; + public static final String + MESSAGE_WAITING_TO_ACQUIRE_THE_CONFIGNODE_NODE_LOCK_BECAUSE_ANOTHER_NODE_PROCEDURE_IS_HOLDING_IT_56494E86 = + "正在等待获取 ConfigNode 节点锁,因为另一个节点 Procedure 正在持有该锁。"; + public static final String + MESSAGE_PIPE_REQUEST_OR_PLUGIN_VALIDATION_HAS_NOT_COMPLETED_A_PLUGIN_CHECK_OR_METADATA_ACCESS_MAY_BE_SLOW_57C36CEF = + "Pipe 请求或插件校验尚未完成;插件检查或元数据访问可能过慢。"; + public static final String + MESSAGE_PIPE_METADATA_CALCULATION_HAS_NOT_COMPLETED_METADATA_ACCESS_OR_LOCAL_CALCULATION_MAY_BE_SLOW_DEBF2504 = + "Pipe 元数据计算尚未完成;元数据访问或本地计算可能过慢。"; + public static final String + MESSAGE_THE_CONFIGNODE_CONSENSUS_WRITE_HAS_NOT_RETURNED_THE_CONSENSUS_GROUP_MAY_BE_UNAVAILABLE_OR_SLOW_F8911CE7 = + "ConfigNode 共识写尚未返回;共识组可能不可用或响应过慢。"; + public static final String + MESSAGE_ONE_OR_MORE_DATANODES_HAVE_NOT_RESPONDED_TO_THE_PIPE_METADATA_PUSH_THEY_MAY_BE_UNAVAILABLE_OR_SLOW_11BBB333 = + "一个或多个 DataNode 尚未响应 Pipe 元数据推送;这些 DataNode 可能不可用或响应过慢。"; + public static final String + MESSAGE_THE_PREVIOUS_ATTEMPT_FAILED_WITH_ARG_AND_THIS_STATE_IS_BEING_RETRIED_7A541F27 = + "上一次尝试因“%s”失败,正在重试此状态。"; + public static final String + MESSAGE_THE_STATE_FAILED_WITH_ARG_AND_ROLLBACK_IS_PENDING_E7B43829 = + "此状态因“%s”失败,正在等待回滚。"; + public static final String MESSAGE_ROLLING_BACK_AFTER_FAILURE_ARG_474DF456 = + "正在回滚,失败原因:%s。"; + public static final String MESSAGE_ROLLING_BACK_AFTER_AN_EARLIER_FAILURE_850D0AF5 = + "正在回滚此前发生的失败。"; } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java index 8f5b27b7e2fa4..894c40df2ecb4 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java @@ -74,6 +74,7 @@ import org.apache.iotdb.confignode.procedure.impl.node.RemoveConfigNodeProcedure; import org.apache.iotdb.confignode.procedure.impl.node.RemoveDataNodesProcedure; import org.apache.iotdb.confignode.procedure.impl.partition.DataPartitionTableIntegrityCheckProcedure; +import org.apache.iotdb.confignode.procedure.impl.pipe.AbstractOperatePipeProcedureV2; import org.apache.iotdb.confignode.procedure.impl.pipe.plugin.CreatePipePluginProcedure; import org.apache.iotdb.confignode.procedure.impl.pipe.plugin.DropPipePluginProcedure; import org.apache.iotdb.confignode.procedure.impl.pipe.runtime.PipeHandleLeaderChangeProcedure; @@ -1661,7 +1662,7 @@ public TSStatus createPipe(TCreatePipeReq req) { return status; } else { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()) - .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage())); + .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage(), procedure)); } } catch (final Exception e) { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage()); @@ -1677,7 +1678,7 @@ public TSStatus alterPipe(final TAlterPipeReq req) { return status; } else { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()) - .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage())); + .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage(), procedure)); } } catch (final Exception e) { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage()); @@ -1708,7 +1709,7 @@ public TSStatus startPipe(String pipeName, boolean isTableModel) { return status; } else { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()) - .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage())); + .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage(), procedure)); } } catch (Exception e) { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage()); @@ -1739,7 +1740,7 @@ public TSStatus stopPipe(String pipeName, boolean isTableModel) { return status; } else { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()) - .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage())); + .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage(), procedure)); } } catch (Exception e) { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage()); @@ -1786,7 +1787,7 @@ public TSStatus dropPipe(String pipeName, boolean isTableModel) { return status; } else { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()) - .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage())); + .setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage(), procedure)); } } catch (Exception e) { return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage()); @@ -2266,6 +2267,13 @@ private static String wrapTimeoutMessageForPipeProcedure(String message) { return message; } + private static String wrapTimeoutMessageForPipeProcedure( + final String message, final AbstractOperatePipeProcedureV2 procedure) { + return PROCEDURE_TIMEOUT_MESSAGE.equals(message) + ? procedure.getTimeoutDiagnosticMessage() + : message; + } + public static void sleepWithoutInterrupt(final long timeToSleep) { long currentTime = System.currentTimeMillis(); final long endTime = timeToSleep + currentTime; diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java index ca68dee4f7bb8..c33beed69c062 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java @@ -142,6 +142,11 @@ private void parseHeartbeatAndSaveMetaChangeLocally( final int nodeId, final PipeHeartbeat pipeHeartbeat) { for (final PipeMeta pipeMetaFromCoordinator : pipeTaskInfo.get().getPipeMetaList()) { + if (PipeStatus.PRE_DELETE.equals( + pipeMetaFromCoordinator.getRuntimeMeta().getStatus().get())) { + continue; + } + final PipeStaticMeta staticMeta = pipeMetaFromCoordinator.getStaticMeta(); final PipeMeta pipeMetaFromAgent = pipeHeartbeat.getPipeMeta(staticMeta); if (pipeMetaFromAgent == null) { @@ -292,6 +297,9 @@ private void parseHeartbeatAndSaveMetaChangeLocally( } final PipeRuntimeMeta runtimeMeta = pipeMeta.getRuntimeMeta(); + if (PipeStatus.PRE_DELETE.equals(runtimeMeta.getStatus().get())) { + return; + } if (!runtimeMeta.getStatus().get().equals(PipeStatus.STOPPED)) { // Record the connector exception for each pipe affected Map exceptionMap = diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java index b301944417225..63de8cf16dae1 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java @@ -224,7 +224,9 @@ public void checkAndUpdateRequestBeforeAlterPipe(final TAlterPipeReq alterPipeRe private void checkAndUpdateRequestBeforeAlterPipeInternal(final TAlterPipeReq alterPipeRequest) throws PipeException { - if (!isPipeExisted(alterPipeRequest.getPipeName(), alterPipeRequest.isTableModel)) { + if (!isPipeExisted(alterPipeRequest.getPipeName(), alterPipeRequest.isTableModel) + || PipeStatus.PRE_DELETE.equals( + getPipeStatus(alterPipeRequest.getPipeName(), alterPipeRequest.isTableModel))) { final String exceptionMessage = String.format( "Failed to alter pipe %s, %s", alterPipeRequest.getPipeName(), PIPE_NOT_EXIST_MSG); @@ -365,7 +367,7 @@ private void checkBeforeStartPipeInternal(final String pipeName) throws PipeExce } final PipeStatus pipeStatus = getPipeStatus(pipeName); - if (pipeStatus == PipeStatus.DROPPED) { + if (pipeStatus == PipeStatus.DROPPED || pipeStatus == PipeStatus.PRE_DELETE) { final String exceptionMessage = String.format( ConfigNodeMessages.FAILED_TO_START_PIPE_BECAUSE_PIPE_IS_ALREADY_DROPPED, pipeName); @@ -385,7 +387,7 @@ private void checkBeforeStartPipeInternal(final String pipeName, final boolean i } final PipeStatus pipeStatus = getPipeStatus(pipeName, isTableModel); - if (pipeStatus == PipeStatus.DROPPED) { + if (pipeStatus == PipeStatus.DROPPED || pipeStatus == PipeStatus.PRE_DELETE) { final String exceptionMessage = String.format( ConfigNodeMessages.FAILED_TO_START_PIPE_BECAUSE_PIPE_IS_ALREADY_DROPPED, pipeName); @@ -423,7 +425,7 @@ private void checkBeforeStopPipeInternal(final String pipeName) throws PipeExcep } final PipeStatus pipeStatus = getPipeStatus(pipeName); - if (pipeStatus == PipeStatus.DROPPED) { + if (pipeStatus == PipeStatus.DROPPED || pipeStatus == PipeStatus.PRE_DELETE) { final String exceptionMessage = String.format( ConfigNodeMessages.FAILED_TO_STOP_PIPE_BECAUSE_PIPE_IS_ALREADY_DROPPED, pipeName); @@ -443,7 +445,7 @@ private void checkBeforeStopPipeInternal(final String pipeName, final boolean is } final PipeStatus pipeStatus = getPipeStatus(pipeName, isTableModel); - if (pipeStatus == PipeStatus.DROPPED) { + if (pipeStatus == PipeStatus.DROPPED || pipeStatus == PipeStatus.PRE_DELETE) { final String exceptionMessage = String.format( ConfigNodeMessages.FAILED_TO_STOP_PIPE_BECAUSE_PIPE_IS_ALREADY_DROPPED, pipeName); @@ -1125,6 +1127,10 @@ private boolean recordDataNodePushPipeMetaExceptionsInternal( final PipeRuntimeMeta runtimeMeta = pipeMeta.getRuntimeMeta(); + if (PipeStatus.PRE_DELETE.equals(runtimeMeta.getStatus().get())) { + continue; + } + // Keep user-stopped pipes out of the auto-restart flow. Otherwise, a failed STOPPED meta // sync can turn a manually stopped pipe into a runtime-stopped one and the next // PipeMetaSyncer round will restart it automatically. @@ -1177,7 +1183,8 @@ private boolean autoRestartInternal() { .forEach( pipeMeta -> { final PipeRuntimeMeta runtimeMeta = pipeMeta.getRuntimeMeta(); - if (runtimeMeta.getIsStoppedByRuntimeException()) { + if (!PipeStatus.PRE_DELETE.equals(runtimeMeta.getStatus().get()) + && runtimeMeta.getIsStoppedByRuntimeException()) { runtimeMeta.setExceptionsClearTime(exceptionsClearTime); runtimeMeta.getStatus().set(PipeStatus.RUNNING); diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/StateMachineProcedure.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/StateMachineProcedure.java index 8f3c1c93633c3..d0cd62e5ef00b 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/StateMachineProcedure.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/StateMachineProcedure.java @@ -132,6 +132,17 @@ protected void setNextState(final TState state) { setNextState(getStateId(state)); } + /** + * Returns whether the specified state is already present in the persisted state history. + * + *

The current state is included once it has been scheduled. This is useful when an append-only + * state is added to a procedure and the new execution path needs to coexist with procedures + * persisted by an older version. + */ + protected final boolean hasReachedState(final TState state) { + return states.contains(getStateId(state)); + } + /** * Add a child procedure to execute. * diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java index 26e84e3ebe47d..a339ca9272cc6 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java @@ -103,6 +103,12 @@ public abstract class AbstractOperatePipeProcedureV2 // recovered procedure is already re-scheduled by the procedure framework. private transient boolean shouldYieldAfterExecution; + // These fields are only used to report where a running Pipe procedure is blocked when the caller + // times out. They do not affect procedure execution and do not need to be persisted. + private volatile PipeProcedureExecutionStage executionStage = + PipeProcedureExecutionStage.WAITING_FOR_PROCEDURE_WORKER; + private volatile String lastExecutionExceptionMessage; + private static final String SKIP_PIPE_PROCEDURE_MESSAGE = "Try to start a RUNNING pipe or stop a STOPPED pipe, do nothing."; @@ -118,6 +124,7 @@ protected AtomicReference acquireLockInternal( @Override protected ProcedureLockState acquireLock(ConfigNodeProcedureEnv configNodeProcedureEnv) { LOGGER.debug(ProcedureMessages.PROCEDUREID_TRY_TO_ACQUIRE_PIPE_LOCK, getProcId()); + executionStage = PipeProcedureExecutionStage.WAITING_FOR_PIPE_TASK_COORDINATOR_LOCK; pipeTaskInfo = acquireLockInternal(configNodeProcedureEnv); if (pipeTaskInfo == null) { LOGGER.warn(ProcedureMessages.PROCEDUREID_FAILED_TO_ACQUIRE_PIPE_LOCK, getProcId()); @@ -125,9 +132,11 @@ protected ProcedureLockState acquireLock(ConfigNodeProcedureEnv configNodeProced LOGGER.debug(ProcedureMessages.PROCEDUREID_ACQUIRED_PIPE_LOCK, getProcId()); } + executionStage = PipeProcedureExecutionStage.WAITING_FOR_NODE_LOCK; final ProcedureLockState procedureLockState = super.acquireLock(configNodeProcedureEnv); switch (procedureLockState) { case LOCK_ACQUIRED: + updateExecutionStage(getCurrentState(), false); if (pipeTaskInfo == null) { LOGGER.warn( ProcedureMessages @@ -192,6 +201,9 @@ protected void releaseLock(ConfigNodeProcedureEnv configNodeProcedureEnv) { .updateTimer(this.getOperation().getName(), this.elapsedTime()); } releasePipeTaskCoordinatorLock(configNodeProcedureEnv); + if (!isFinished()) { + executionStage = PipeProcedureExecutionStage.WAITING_FOR_PROCEDURE_WORKER; + } } } @@ -215,6 +227,16 @@ private void releasePipeTaskCoordinatorLock(ConfigNodeProcedureEnv configNodePro /** Execute at state {@link OperatePipeTaskState#CALCULATE_INFO_FOR_TASK}. */ public abstract void executeFromCalculateInfoForTask(ConfigNodeProcedureEnv env); + /** Whether this procedure uses the append-only {@link OperatePipeTaskState#PRE_DELETE} state. */ + protected boolean shouldExecutePreDeleteState() { + return false; + } + + /** Execute at state {@link OperatePipeTaskState#PRE_DELETE}. */ + public void executeFromPreDelete(ConfigNodeProcedureEnv env) { + // Do nothing by default + } + /** * Execute at state {@link OperatePipeTaskState#WRITE_CONFIG_NODE_CONSENSUS}. * @@ -236,6 +258,7 @@ public abstract void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) protected Flow executeFromState(ConfigNodeProcedureEnv env, OperatePipeTaskState state) throws InterruptedException { shouldYieldAfterExecution = false; + updateExecutionStage(state, false); if (pipeTaskInfo == null) { LOGGER.warn( ProcedureMessages.PROCEDUREID_PIPE_LOCK_IS_NOT_ACQUIRED_EXECUTEFROMSTATE_S_EXECUTION_WILL, @@ -247,6 +270,7 @@ protected Flow executeFromState(ConfigNodeProcedureEnv env, OperatePipeTaskState switch (state) { case VALIDATE_TASK: if (!executeFromValidateTask(env)) { + lastExecutionExceptionMessage = null; LOGGER.info(ProcedureMessages.PROCEDUREID, getProcId(), SKIP_PIPE_PROCEDURE_MESSAGE); // On client side, the message returned after the successful execution of the pipe // command corresponding to this procedure is "Msg: The statement is executed @@ -258,21 +282,39 @@ protected Flow executeFromState(ConfigNodeProcedureEnv env, OperatePipeTaskState break; case CALCULATE_INFO_FOR_TASK: executeFromCalculateInfoForTask(env); - setNextState(OperatePipeTaskState.WRITE_CONFIG_NODE_CONSENSUS); + setNextState( + shouldExecutePreDeleteState() + ? OperatePipeTaskState.PRE_DELETE + : OperatePipeTaskState.WRITE_CONFIG_NODE_CONSENSUS); + break; + case PRE_DELETE: + executeFromPreDelete(env); + setNextState(OperatePipeTaskState.OPERATE_ON_DATA_NODES); break; case WRITE_CONFIG_NODE_CONSENSUS: executeFromWriteConfigNodeConsensus(env); + if (shouldExecutePreDeleteState() && hasReachedState(OperatePipeTaskState.PRE_DELETE)) { + lastExecutionExceptionMessage = null; + return Flow.NO_MORE_STATE; + } setNextState(OperatePipeTaskState.OPERATE_ON_DATA_NODES); break; case OPERATE_ON_DATA_NODES: executeFromOperateOnDataNodes(env); + if (shouldExecutePreDeleteState() && hasReachedState(OperatePipeTaskState.PRE_DELETE)) { + setNextState(OperatePipeTaskState.WRITE_CONFIG_NODE_CONSENSUS); + break; + } + lastExecutionExceptionMessage = null; return Flow.NO_MORE_STATE; default: throw new UnsupportedOperationException( String.format( ProcedureMessages.UNKNOWN_STATE_DURING_EXECUTING_OPERATEPIPEPROCEDURE, state)); } + lastExecutionExceptionMessage = null; } catch (Exception e) { + lastExecutionExceptionMessage = getExceptionMessage(e); // Retry before rollback if (getCycles() < RETRY_THRESHOLD) { LOGGER.warn( @@ -319,6 +361,7 @@ protected boolean isRollbackSupported(OperatePipeTaskState state) { @Override protected void rollbackState(ConfigNodeProcedureEnv env, OperatePipeTaskState state) throws IOException, InterruptedException, ProcedureException { + updateExecutionStage(state, true); if (pipeTaskInfo == null) { LOGGER.warn( ProcedureMessages.PROCEDUREID_PIPE_LOCK_IS_NOT_ACQUIRED_ROLLBACKSTATE_S_EXECUTION_WILL, @@ -351,6 +394,9 @@ protected void rollbackState(ConfigNodeProcedureEnv env, OperatePipeTaskState st e); } break; + case PRE_DELETE: + rollbackFromPreDelete(env); + break; case WRITE_CONFIG_NODE_CONSENSUS: try { // rollbackFromWriteConfigNodeConsensus can be called before @@ -392,6 +438,10 @@ protected void rollbackState(ConfigNodeProcedureEnv env, OperatePipeTaskState st public abstract void rollbackFromCalculateInfoForTask(ConfigNodeProcedureEnv env); + public void rollbackFromPreDelete(ConfigNodeProcedureEnv env) { + // Do nothing by default + } + public abstract void rollbackFromWriteConfigNodeConsensus(ConfigNodeProcedureEnv env); public abstract void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env) @@ -412,6 +462,151 @@ protected OperatePipeTaskState getInitialState() { return OperatePipeTaskState.VALIDATE_TASK; } + public final String getTimeoutDiagnosticMessage() { + final PipeProcedureExecutionStage currentExecutionStage = executionStage; + return String.format( + ProcedureMessages + .MESSAGE_PIPE_OPERATION_ARG_TIMED_OUT_PROCEDUREID_ARG_STUCK_AT_ARG_REASON_ARG_THE_PROCEDURE_IS_STILL_RUNNING_7EEAC50E, + getOperation().name(), + getProcId(), + currentExecutionStage.name(), + getTimeoutReason(currentExecutionStage)); + } + + private String getTimeoutReason(final PipeProcedureExecutionStage currentExecutionStage) { + final String failureMessage = getFailureMessage(); + if (currentExecutionStage.isRollback()) { + return failureMessage == null + ? ProcedureMessages.MESSAGE_ROLLING_BACK_AFTER_AN_EARLIER_FAILURE_850D0AF5 + : String.format( + ProcedureMessages.MESSAGE_ROLLING_BACK_AFTER_FAILURE_ARG_474DF456, failureMessage); + } + if (isFailed() && failureMessage != null) { + return String.format( + ProcedureMessages.MESSAGE_THE_STATE_FAILED_WITH_ARG_AND_ROLLBACK_IS_PENDING_E7B43829, + failureMessage); + } + if (lastExecutionExceptionMessage != null) { + return String.format( + ProcedureMessages + .MESSAGE_THE_PREVIOUS_ATTEMPT_FAILED_WITH_ARG_AND_THIS_STATE_IS_BEING_RETRIED_7A541F27, + lastExecutionExceptionMessage); + } + + switch (currentExecutionStage) { + case WAITING_FOR_PROCEDURE_WORKER: + return ProcedureMessages + .MESSAGE_NO_PROCEDURE_WORKER_IS_CURRENTLY_AVAILABLE_WORKERS_MAY_BE_BUSY_OR_BLOCKED_BY_OTHER_PROCEDURES_AB0B1595; + case WAITING_FOR_PIPE_TASK_COORDINATOR_LOCK: + return ProcedureMessages + .MESSAGE_WAITING_TO_ACQUIRE_THE_PIPETASKCOORDINATOR_LOCK_BECAUSE_ANOTHER_PIPE_OPERATION_IS_HOLDING_IT_25A3B6B8; + case WAITING_FOR_NODE_LOCK: + return ProcedureMessages + .MESSAGE_WAITING_TO_ACQUIRE_THE_CONFIGNODE_NODE_LOCK_BECAUSE_ANOTHER_NODE_PROCEDURE_IS_HOLDING_IT_56494E86; + case VALIDATE_TASK: + return ProcedureMessages + .MESSAGE_PIPE_REQUEST_OR_PLUGIN_VALIDATION_HAS_NOT_COMPLETED_A_PLUGIN_CHECK_OR_METADATA_ACCESS_MAY_BE_SLOW_57C36CEF; + case CALCULATE_INFO_FOR_TASK: + return ProcedureMessages + .MESSAGE_PIPE_METADATA_CALCULATION_HAS_NOT_COMPLETED_METADATA_ACCESS_OR_LOCAL_CALCULATION_MAY_BE_SLOW_DEBF2504; + case PRE_DELETE: + case WRITE_CONFIG_NODE_CONSENSUS: + return ProcedureMessages + .MESSAGE_THE_CONFIGNODE_CONSENSUS_WRITE_HAS_NOT_RETURNED_THE_CONSENSUS_GROUP_MAY_BE_UNAVAILABLE_OR_SLOW_F8911CE7; + case OPERATE_ON_DATA_NODES: + return ProcedureMessages + .MESSAGE_ONE_OR_MORE_DATANODES_HAVE_NOT_RESPONDED_TO_THE_PIPE_METADATA_PUSH_THEY_MAY_BE_UNAVAILABLE_OR_SLOW_11BBB333; + default: + return ProcedureMessages + .MESSAGE_NO_PROCEDURE_WORKER_IS_CURRENTLY_AVAILABLE_WORKERS_MAY_BE_BUSY_OR_BLOCKED_BY_OTHER_PROCEDURES_AB0B1595; + } + } + + private String getFailureMessage() { + if (getException() != null) { + return getException().getMessage(); + } + return lastExecutionExceptionMessage; + } + + private static String getExceptionMessage(final Exception exception) { + return exception.getMessage() == null || exception.getMessage().isEmpty() + ? exception.getClass().getSimpleName() + : exception.getMessage(); + } + + private void updateExecutionStage(final OperatePipeTaskState state, final boolean isRollback) { + if (state == null) { + executionStage = PipeProcedureExecutionStage.WAITING_FOR_PROCEDURE_WORKER; + return; + } + switch (state) { + case VALIDATE_TASK: + executionStage = + isRollback + ? PipeProcedureExecutionStage.ROLLBACK_VALIDATE_TASK + : PipeProcedureExecutionStage.VALIDATE_TASK; + break; + case CALCULATE_INFO_FOR_TASK: + executionStage = + isRollback + ? PipeProcedureExecutionStage.ROLLBACK_CALCULATE_INFO_FOR_TASK + : PipeProcedureExecutionStage.CALCULATE_INFO_FOR_TASK; + break; + case PRE_DELETE: + executionStage = + isRollback + ? PipeProcedureExecutionStage.ROLLBACK_PRE_DELETE + : PipeProcedureExecutionStage.PRE_DELETE; + break; + case WRITE_CONFIG_NODE_CONSENSUS: + executionStage = + isRollback + ? PipeProcedureExecutionStage.ROLLBACK_WRITE_CONFIG_NODE_CONSENSUS + : PipeProcedureExecutionStage.WRITE_CONFIG_NODE_CONSENSUS; + break; + case OPERATE_ON_DATA_NODES: + executionStage = + isRollback + ? PipeProcedureExecutionStage.ROLLBACK_OPERATE_ON_DATA_NODES + : PipeProcedureExecutionStage.OPERATE_ON_DATA_NODES; + break; + default: + executionStage = PipeProcedureExecutionStage.WAITING_FOR_PROCEDURE_WORKER; + break; + } + } + + private enum PipeProcedureExecutionStage { + WAITING_FOR_PROCEDURE_WORKER, + WAITING_FOR_PIPE_TASK_COORDINATOR_LOCK, + WAITING_FOR_NODE_LOCK, + VALIDATE_TASK, + CALCULATE_INFO_FOR_TASK, + PRE_DELETE, + WRITE_CONFIG_NODE_CONSENSUS, + OPERATE_ON_DATA_NODES, + ROLLBACK_VALIDATE_TASK(true), + ROLLBACK_CALCULATE_INFO_FOR_TASK(true), + ROLLBACK_PRE_DELETE(true), + ROLLBACK_WRITE_CONFIG_NODE_CONSENSUS(true), + ROLLBACK_OPERATE_ON_DATA_NODES(true); + + private final boolean rollback; + + PipeProcedureExecutionStage() { + this(false); + } + + PipeProcedureExecutionStage(final boolean rollback) { + this.rollback = rollback; + } + + private boolean isRollback() { + return rollback; + } + } + /** * Pushing all the pipeMeta's to all the dataNodes, forcing an update to the pipe's runtime state. * diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeMetaSyncProcedure.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeMetaSyncProcedure.java index d371b7e9b203b..f18cfb6d5453f 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeMetaSyncProcedure.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeMetaSyncProcedure.java @@ -22,6 +22,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMeta; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta; import org.apache.iotdb.commons.pipe.config.PipeConfig; import org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant; @@ -125,6 +126,9 @@ public void executeFromCalculateInfoForTask(ConfigNodeProcedureEnv env) { .getPipeMetaList() .forEach( pipeMeta -> { + if (PipeStatus.PRE_DELETE.equals(pipeMeta.getRuntimeMeta().getStatus().get())) { + return; + } if (!pipeMeta.getStaticMeta().isSourceExternal()) { return; } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2.java index 80044cf04c873..f07586613f2e7 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2.java @@ -23,6 +23,7 @@ import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus; import org.apache.iotdb.confignode.consensus.request.write.pipe.task.DropPipePlanV2; +import org.apache.iotdb.confignode.consensus.request.write.pipe.task.SetPipeStatusPlanV2; import org.apache.iotdb.confignode.i18n.ConfigNodeMessages; import org.apache.iotdb.confignode.i18n.ProcedureMessages; import org.apache.iotdb.confignode.persistence.pipe.PipeTaskInfo; @@ -130,6 +131,49 @@ public void executeFromCalculateInfoForTask(ConfigNodeProcedureEnv env) throws P : pipeTaskInfo.get().getPipeMetaByPipeName(pipeName); } + @Override + protected boolean shouldExecutePreDeleteState() { + return true; + } + + @Override + public void executeFromPreDelete(ConfigNodeProcedureEnv env) throws PipeException { + // Legacy procedures created without an explicit model do not persist pipeMetaToDrop. Restore + // it from PipeTaskInfo so a procedure recovered between CALCULATE_INFO_FOR_TASK and PRE_DELETE + // still exposes PRE_DELETE through SHOW PIPES. + if (!restorePipeMetaToDropIfNecessary()) { + return; + } + + TSStatus response; + try { + response = + env.getConfigManager() + .getConsensusManager() + .write( + isTableModelSet + ? new SetPipeStatusPlanV2(pipeName, PipeStatus.PRE_DELETE, isTableModel) + : new SetPipeStatusPlanV2(pipeName, PipeStatus.PRE_DELETE)); + } catch (ConsensusException e) { + LOGGER.warn(ConfigNodeMessages.FAILED_IN_THE_WRITE_API_EXECUTING_THE_CONSENSUS_LAYER_DUE, e); + response = new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode()); + response.setMessage(e.getMessage()); + } + if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) { + throw new PipeException(response.getMessage()); + } + } + + boolean restorePipeMetaToDropIfNecessary() { + if (pipeMetaToDrop == null) { + pipeMetaToDrop = + isTableModelSet + ? pipeTaskInfo.get().getPipeMetaByPipeName(pipeName, isTableModel) + : pipeTaskInfo.get().getPipeMetaByPipeName(pipeName); + } + return pipeMetaToDrop != null; + } + @Override public void executeFromWriteConfigNodeConsensus(ConfigNodeProcedureEnv env) throws PipeException { LOGGER.info( @@ -194,6 +238,11 @@ public void rollbackFromCalculateInfoForTask(ConfigNodeProcedureEnv env) { // Do nothing } + @Override + public void rollbackFromPreDelete(ConfigNodeProcedureEnv env) { + // Keep PRE_DELETE so SHOW PIPES continues to expose that the drop did not complete. + } + @Override public void rollbackFromWriteConfigNodeConsensus(ConfigNodeProcedureEnv env) { LOGGER.info( diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/pipe/task/OperatePipeTaskState.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/pipe/task/OperatePipeTaskState.java index fdbc65fe4dd8a..10f00ea7682fe 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/pipe/task/OperatePipeTaskState.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/pipe/task/OperatePipeTaskState.java @@ -23,5 +23,7 @@ public enum OperatePipeTaskState { VALIDATE_TASK, CALCULATE_INFO_FOR_TASK, OPERATE_ON_DATA_NODES, - WRITE_CONFIG_NODE_CONSENSUS + WRITE_CONFIG_NODE_CONSENSUS, + // Keep new states appended to preserve persisted state ordinals. + PRE_DELETE } diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanSerDeTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanSerDeTest.java index e844dcf6910e9..35537a2cea2b8 100644 --- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanSerDeTest.java +++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanSerDeTest.java @@ -1014,7 +1014,7 @@ public void AlterPipePlanV2CurrentStaticMetaTest() throws IOException { public void SetPipeStatusPlanV2Test() throws IOException { final SetPipeStatusPlanV2 setPipeStatusPlanV2 = new SetPipeStatusPlanV2( - "pipe", org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus.RUNNING, true); + "pipe", org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus.PRE_DELETE, true); final SetPipeStatusPlanV2 setPipeStatusPlanV21 = (SetPipeStatusPlanV2) ConfigPhysicalPlan.Factory.create(setPipeStatusPlanV2.serializeToByteBuffer()); diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/response/pipe/PipeTableRespTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/response/pipe/PipeTableRespTest.java index 093e5989bb0b9..f4437fd2c34e1 100644 --- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/response/pipe/PipeTableRespTest.java +++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/response/pipe/PipeTableRespTest.java @@ -23,6 +23,7 @@ import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeRuntimeMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMetaInCoordinator; import org.apache.iotdb.commons.pipe.config.constant.SystemConstant; @@ -168,6 +169,17 @@ public void testConvertToTShowPipeRespIncludesDegradedStatus() { Assert.assertFalse(showPipeResult.get(2).isSetIsDegraded()); } + @Test + public void testConvertToTShowPipeRespIncludesPreDeleteStatus() { + final PipeTableResp pipeTableResp = constructPipeTableResp(); + pipeTableResp.getAllPipeMeta().get(0).getRuntimeMeta().getStatus().set(PipeStatus.PRE_DELETE); + + final List showPipeResult = + pipeTableResp.convertToTShowPipeResp().getPipeInfoList(); + + Assert.assertEquals(PipeStatus.PRE_DELETE.name(), showPipeResult.get(0).getState()); + } + @Test public void testFilterByModelBeforeWhereClause() { TSStatus status = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()); diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java index d77992c18be9b..04a898629fd81 100644 --- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java +++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java @@ -214,6 +214,42 @@ public void testParseHeartbeatTracksExceptionsAfterClearTime() throws Exception verify(context.procedureManager, times(1)).pipeHandleMetaChange(true, false); } + @Test + public void testParseHeartbeatDoesNotOverwritePreDeleteStatus() throws Exception { + CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false); + + final String pipeName = "preDeletePipe"; + final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo(); + createPipe(pipeTaskInfo, pipeName, PipeStatus.RUNNING); + + final PipeMeta pipeMeta = pipeTaskInfo.getPipeMetaByPipeName(pipeName); + final PipeRuntimeMeta runtimeMeta = pipeMeta.getRuntimeMeta(); + runtimeMeta.getStatus().set(PipeStatus.PRE_DELETE); + + final PipeTaskMeta agentTaskMeta = + new PipeTaskMeta(MinimumProgressIndex.INSTANCE, DATA_NODE_ID); + agentTaskMeta.trackExceptionMessage(new PipeRuntimeCriticalException("fresh failure", 300L)); + final ConcurrentMap agentPipeTasks = new ConcurrentHashMap<>(); + agentPipeTasks.put(DATA_NODE_ID, agentTaskMeta); + final PipeHeartbeat heartbeat = + new PipeHeartbeat( + Collections.singletonList( + new PipeMeta(pipeMeta.getStaticMeta(), new PipeRuntimeMeta(agentPipeTasks)) + .serialize()), + Collections.singletonList(false), + Collections.singletonList(0L), + Collections.singletonList(0D), + null); + + final ParserTestContext context = createParserTestContext(1, pipeTaskInfo); + context.parser.parseHeartbeat(DATA_NODE_ID, heartbeat); + + Assert.assertEquals(PipeStatus.PRE_DELETE, runtimeMeta.getStatus().get()); + Assert.assertFalse( + runtimeMeta.getConsensusGroupId2TaskMetaMap().get(DATA_NODE_ID).hasExceptionMessages()); + verify(context.procedureManager, never()).pipeHandleMetaChange(anyBoolean(), anyBoolean()); + } + @Test public void testParseHeartbeatRecordsPipeDegradedStatus() throws Exception { CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false); diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2Test.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2Test.java index 5993e73ec42a9..accab5f59024a 100644 --- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2Test.java +++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2Test.java @@ -19,6 +19,7 @@ package org.apache.iotdb.confignode.procedure.impl.pipe; +import org.apache.iotdb.confignode.i18n.ProcedureMessages; import org.apache.iotdb.confignode.persistence.pipe.PipeTaskInfo; import org.apache.iotdb.confignode.procedure.Procedure; import org.apache.iotdb.confignode.procedure.impl.StateMachineProcedure; @@ -29,6 +30,8 @@ import org.junit.Test; import java.io.IOException; +import java.util.ArrayList; +import java.util.List; import java.util.concurrent.atomic.AtomicReference; public class AbstractOperatePipeProcedureV2Test { @@ -89,12 +92,86 @@ public void testRetryStateYieldsOnlyBeforeRetryThreshold() throws Exception { Assert.assertEquals(2, procedure.calculateExecutionCount); } + @Test + public void testTimeoutDiagnosticReportsCurrentStateAndRetryReason() throws Exception { + final TestOperatePipeProcedure procedure = new TestOperatePipeProcedure(); + procedure.failValidation = true; + + procedure.executeFromState(null, OperatePipeTaskState.VALIDATE_TASK); + + final String diagnosticMessage = procedure.getTimeoutDiagnosticMessage(); + Assert.assertTrue(diagnosticMessage.contains("START_PIPE")); + Assert.assertTrue(diagnosticMessage.contains("VALIDATE_TASK")); + Assert.assertTrue(diagnosticMessage.contains("retry")); + } + + @Test + public void testTimeoutDiagnosticReportsDataNodeOperation() throws Exception { + final TestOperatePipeProcedure procedure = new TestOperatePipeProcedure(); + + procedure.executeFromState(null, OperatePipeTaskState.OPERATE_ON_DATA_NODES); + + final String diagnosticMessage = procedure.getTimeoutDiagnosticMessage(); + Assert.assertTrue(diagnosticMessage.contains("OPERATE_ON_DATA_NODES")); + Assert.assertTrue( + diagnosticMessage.contains( + ProcedureMessages + .MESSAGE_ONE_OR_MORE_DATANODES_HAVE_NOT_RESPONDED_TO_THE_PIPE_METADATA_PUSH_THEY_MAY_BE_UNAVAILABLE_OR_SLOW_11BBB333)); + } + + @Test + public void testPreDeleteFlowAndStateOrdinalCompatibility() throws Exception { + Assert.assertEquals(0, OperatePipeTaskState.VALIDATE_TASK.ordinal()); + Assert.assertEquals(1, OperatePipeTaskState.CALCULATE_INFO_FOR_TASK.ordinal()); + Assert.assertEquals(2, OperatePipeTaskState.OPERATE_ON_DATA_NODES.ordinal()); + Assert.assertEquals(3, OperatePipeTaskState.WRITE_CONFIG_NODE_CONSENSUS.ordinal()); + Assert.assertEquals(4, OperatePipeTaskState.PRE_DELETE.ordinal()); + + final TestOperatePipeProcedure procedure = new TestOperatePipeProcedure(); + procedure.preDeleteEnabled = true; + + for (int i = 0; i < 5; i++) { + procedure.runOnce(); + } + + Assert.assertEquals( + List.of( + OperatePipeTaskState.VALIDATE_TASK, + OperatePipeTaskState.CALCULATE_INFO_FOR_TASK, + OperatePipeTaskState.PRE_DELETE, + OperatePipeTaskState.OPERATE_ON_DATA_NODES, + OperatePipeTaskState.WRITE_CONFIG_NODE_CONSENSUS), + procedure.executionOrder); + } + + @Test + public void testLegacyDropFlowWithoutPreDeleteHistoryStillOperatesOnDataNodes() throws Exception { + final TestOperatePipeProcedure procedure = new TestOperatePipeProcedure(); + + // Schedule WRITE_CONFIG_NODE_CONSENSUS with the legacy flow, then continue with the new logic. + procedure.runOnce(); + procedure.runOnce(); + procedure.preDeleteEnabled = true; + procedure.runOnce(); + procedure.runOnce(); + + Assert.assertEquals( + List.of( + OperatePipeTaskState.VALIDATE_TASK, + OperatePipeTaskState.CALCULATE_INFO_FOR_TASK, + OperatePipeTaskState.WRITE_CONFIG_NODE_CONSENSUS, + OperatePipeTaskState.OPERATE_ON_DATA_NODES), + procedure.executionOrder); + } + private static class TestOperatePipeProcedure extends AbstractOperatePipeProcedureV2 { private int validateExecutionCount; private int calculateExecutionCount; private boolean failValidation; private boolean failCalculation; + private boolean preDeleteEnabled; + private final List executionOrder = new ArrayList<>(); private TestOperatePipeProcedure() { pipeTaskInfo = new AtomicReference<>(new PipeTaskInfo()); @@ -113,6 +190,7 @@ protected PipeTaskOperation getOperation() { public boolean executeFromValidateTask( final org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv env) throws PipeException { + executionOrder.add(OperatePipeTaskState.VALIDATE_TASK); validateExecutionCount++; if (failValidation) { throw new PipeException("retry"); @@ -123,22 +201,34 @@ public boolean executeFromValidateTask( @Override public void executeFromCalculateInfoForTask( final org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv env) { + executionOrder.add(OperatePipeTaskState.CALCULATE_INFO_FOR_TASK); calculateExecutionCount++; if (failCalculation) { throw new RuntimeException("retry"); } } + @Override + protected boolean shouldExecutePreDeleteState() { + return preDeleteEnabled; + } + + @Override + public void executeFromPreDelete( + final org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv env) { + executionOrder.add(OperatePipeTaskState.PRE_DELETE); + } + @Override public void executeFromWriteConfigNodeConsensus( final org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv env) { - // Do nothing + executionOrder.add(OperatePipeTaskState.WRITE_CONFIG_NODE_CONSENSUS); } @Override public void executeFromOperateOnDataNodes( final org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv env) { - // Do nothing + executionOrder.add(OperatePipeTaskState.OPERATE_ON_DATA_NODES); } @Override diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2Test.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2Test.java index 1317b83883074..61b868137554d 100644 --- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2Test.java +++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2Test.java @@ -19,19 +19,59 @@ package org.apache.iotdb.confignode.procedure.impl.pipe.task; +import org.apache.iotdb.common.rpc.thrift.TSStatus; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeRuntimeMeta; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus; +import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlan; +import org.apache.iotdb.confignode.consensus.request.write.pipe.task.CreatePipePlanV2; +import org.apache.iotdb.confignode.consensus.request.write.pipe.task.SetPipeStatusPlanV2; +import org.apache.iotdb.confignode.manager.ConfigManager; +import org.apache.iotdb.confignode.manager.consensus.ConsensusManager; +import org.apache.iotdb.confignode.persistence.pipe.PipeTaskInfo; +import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv; import org.apache.iotdb.confignode.procedure.store.ProcedureFactory; +import org.apache.iotdb.pipe.api.exception.PipeException; +import org.apache.iotdb.rpc.TSStatusCode; import org.apache.tsfile.utils.PublicBAOS; import org.junit.Test; +import org.mockito.ArgumentCaptor; +import org.mockito.Mockito; import java.io.DataOutputStream; import java.nio.ByteBuffer; +import java.util.Collections; +import java.util.concurrent.atomic.AtomicReference; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; public class DropPipeProcedureV2Test { + + private static class TestDropPipeProcedureV2 extends DropPipeProcedureV2 { + + private TestDropPipeProcedureV2() { + super(); + } + + private TestDropPipeProcedureV2(final String pipeName) throws PipeException { + super(pipeName); + } + + private TestDropPipeProcedureV2(final String pipeName, final boolean isTableModel) + throws PipeException { + super(pipeName, isTableModel); + } + + private void setPipeTaskInfo(final PipeTaskInfo pipeTaskInfo) { + this.pipeTaskInfo = new AtomicReference<>(pipeTaskInfo); + } + } + @Test public void serializeDeserializeTest() { PublicBAOS byteArrayOutputStream = new PublicBAOS(); @@ -72,4 +112,91 @@ public void serializeDeserializeLegacyFormatTest() { fail(); } } + + @Test + public void testPreDeleteWritesStatusBeforeFinalDrop() throws Exception { + final String pipeName = "testPipe"; + final PipeTaskInfo pipeTaskInfo = createPipeTaskInfo(pipeName); + final TestDropPipeProcedureV2 proc = new TestDropPipeProcedureV2(pipeName, false); + proc.setPipeTaskInfo(pipeTaskInfo); + proc.executeFromCalculateInfoForTask(Mockito.mock(ConfigNodeProcedureEnv.class)); + + final ConfigNodeProcedureEnv env = Mockito.mock(ConfigNodeProcedureEnv.class); + final ConfigManager configManager = Mockito.mock(ConfigManager.class); + final ConsensusManager consensusManager = Mockito.mock(ConsensusManager.class); + Mockito.when(env.getConfigManager()).thenReturn(configManager); + Mockito.when(configManager.getConsensusManager()).thenReturn(consensusManager); + Mockito.when(consensusManager.write(Mockito.any(ConfigPhysicalPlan.class))) + .thenReturn(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())); + + proc.executeFromPreDelete(env); + + final ArgumentCaptor planCaptor = + ArgumentCaptor.forClass(ConfigPhysicalPlan.class); + Mockito.verify(consensusManager).write(planCaptor.capture()); + assertEquals( + new SetPipeStatusPlanV2(pipeName, PipeStatus.PRE_DELETE, false), planCaptor.getValue()); + } + + @Test + public void testRecoveredLegacyProcedureRestoresPipeMetaBeforePreDelete() throws Exception { + final String pipeName = "testPipe"; + final PipeTaskInfo pipeTaskInfo = createPipeTaskInfo(pipeName); + final TestDropPipeProcedureV2 proc = new TestDropPipeProcedureV2(pipeName); + proc.setPipeTaskInfo(pipeTaskInfo); + proc.executeFromCalculateInfoForTask(Mockito.mock(ConfigNodeProcedureEnv.class)); + + final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + proc.serialize(new DataOutputStream(byteArrayOutputStream)); + final ByteBuffer byteBuffer = + ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size()); + byteBuffer.getShort(); + final TestDropPipeProcedureV2 recoveredProc = new TestDropPipeProcedureV2(); + recoveredProc.deserialize(byteBuffer); + recoveredProc.setPipeTaskInfo(pipeTaskInfo); + + assertFalse(recoveredProc.isTableModelSet()); + assertNull(recoveredProc.getPipeMetaToDrop()); + assertTrue(recoveredProc.restorePipeMetaToDropIfNecessary()); + assertEquals(pipeTaskInfo.getPipeMetaByPipeName(pipeName), recoveredProc.getPipeMetaToDrop()); + } + + @Test + public void testPreDeletePipeIsNotAutoRestarted() { + final String pipeName = "testPipe"; + final PipeTaskInfo pipeTaskInfo = createPipeTaskInfo(pipeName); + pipeTaskInfo + .getPipeMetaByPipeName(pipeName) + .getRuntimeMeta() + .getStatus() + .set(PipeStatus.PRE_DELETE); + pipeTaskInfo + .getPipeMetaByPipeName(pipeName) + .getRuntimeMeta() + .setIsStoppedByRuntimeException(true); + + assertFalse(pipeTaskInfo.autoRestart()); + assertEquals( + PipeStatus.PRE_DELETE, + pipeTaskInfo.getPipeMetaByPipeName(pipeName).getRuntimeMeta().getStatus().get()); + assertTrue( + pipeTaskInfo + .getPipeMetaByPipeName(pipeName) + .getRuntimeMeta() + .getIsStoppedByRuntimeException()); + } + + private PipeTaskInfo createPipeTaskInfo(final String pipeName) { + final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo(); + pipeTaskInfo.createPipe( + new CreatePipePlanV2( + new PipeStaticMeta( + pipeName, + System.currentTimeMillis(), + Collections.emptyMap(), + Collections.emptyMap(), + Collections.emptyMap()), + new PipeRuntimeMeta())); + return pipeTaskInfo; + } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/PipeTaskAgent.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/PipeTaskAgent.java index ee465ada69b8e..f2f912cbc61a4 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/PipeTaskAgent.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/PipeTaskAgent.java @@ -220,6 +220,12 @@ private void executeSinglePipeMetaChanges(final PipeMeta metaFromCoordinator) return; } + // PRE_DELETE is a coordinator-only marker. The drop procedure will push DROPPED explicitly + // after the marker is persisted, so task agents should retain their current runtime state here. + if (metaFromCoordinator.getRuntimeMeta().getStatus().get() == PipeStatus.PRE_DELETE) { + return; + } + if (metaFromCoordinator.getRuntimeMeta().getStatus().get() == PipeStatus.DROPPED) { dropPipe(pipeName, creationTime); return; diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeStatus.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeStatus.java index 0feca59bc8593..cdd3a491fa4d5 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeStatus.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeStatus.java @@ -25,6 +25,7 @@ public enum PipeStatus { RUNNING((byte) 0), STOPPED((byte) 1), DROPPED((byte) 2), + PRE_DELETE((byte) 3), ; private final byte type; @@ -45,6 +46,8 @@ public static PipeStatus getPipeStatus(byte type) { return PipeStatus.STOPPED; case 2: return PipeStatus.DROPPED; + case 3: + return PipeStatus.PRE_DELETE; default: throw new IllegalArgumentException(SchemaMessages.SCHEMA_INVALID_INPUT + type); } diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/task/PipeMetaDeSerTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/task/PipeMetaDeSerTest.java index 3b30ec62a7f3c..eee842621bb84 100644 --- a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/task/PipeMetaDeSerTest.java +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/task/PipeMetaDeSerTest.java @@ -50,6 +50,14 @@ public class PipeMetaDeSerTest { + @Test + public void testPipeStatusTypeCompatibility() { + Assert.assertEquals((byte) 0, PipeStatus.RUNNING.getType()); + Assert.assertEquals((byte) 1, PipeStatus.STOPPED.getType()); + Assert.assertEquals((byte) 2, PipeStatus.DROPPED.getType()); + Assert.assertEquals((byte) 3, PipeStatus.PRE_DELETE.getType()); + } + @Test public void test() throws IOException { final PipeStaticMeta pipeStaticMeta = @@ -157,6 +165,11 @@ public void test() throws IOException { pipeRuntimeMeta1 = PipeRuntimeMeta.deserialize(runtimeByteBuffer); Assert.assertEquals(pipeRuntimeMeta, pipeRuntimeMeta1); + pipeRuntimeMeta.getStatus().set(PipeStatus.PRE_DELETE); + runtimeByteBuffer = pipeRuntimeMeta.serialize(); + pipeRuntimeMeta1 = PipeRuntimeMeta.deserialize(runtimeByteBuffer); + Assert.assertEquals(pipeRuntimeMeta, pipeRuntimeMeta1); + final PipeMeta pipeMeta = new PipeMeta(pipeStaticMeta, pipeRuntimeMeta); final ByteBuffer byteBuffer = pipeMeta.serialize(); final PipeMeta pipeMeta1 = PipeMeta.deserialize4Coordinator(byteBuffer);