This is an automated email from the ASF dual-hosted git repository.

jt2594838 pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/dev/1.3 by this push:
     new 22d9de463d2 [Pipe] Improve timeout diagnostics and DROP PIPE 
observability (#18297) (#18340)
22d9de463d2 is described below

commit 22d9de463d2023627d805c23be225b247477cf00
Author: Caideyipi <[email protected]>
AuthorDate: Thu Jul 30 10:42:36 2026 +0800

    [Pipe] Improve timeout diagnostics and DROP PIPE observability (#18297) 
(#18340)
    
    * Improve Pipe procedure timeout diagnostics
    
    * Add pre-delete state to drop pipe procedure
    
    * Address pipe timeout diagnostic review comments
    
    * Fix pipe timeout message i18n
---
 .../iotdb/confignode/manager/ProcedureManager.java |  24 ++-
 .../runtime/heartbeat/PipeHeartbeatParser.java     |   8 +
 .../confignode/persistence/pipe/PipeTaskInfo.java  |  26 ++-
 .../procedure/env/ConfigNodeProcedureEnv.java      |  57 ++++--
 .../procedure/impl/StateMachineProcedure.java      |  11 ++
 .../impl/pipe/AbstractOperatePipeProcedureV2.java  | 199 ++++++++++++++++++++-
 .../impl/pipe/task/DropPipeProcedureV2.java        |  81 ++++++++-
 .../confignode/procedure/scheduler/LockQueue.java  |   6 +-
 .../request/ConfigPhysicalPlanSerDeTest.java       |   2 +-
 .../consensus/response/pipe/PipeTableRespTest.java |  37 ++++
 .../runtime/heartbeat/PipeHeartbeatParserTest.java |  35 ++++
 .../pipe/PipeTaskInfoAutoRestartTest.java          |  46 +++++
 .../pipe/AbstractOperatePipeProcedureV2Test.java   |  33 ++++
 .../impl/pipe/task/DropPipeProcedureV2Test.java    | 130 ++++++++++++++
 .../commons/pipe/agent/task/PipeTaskAgent.java     |   6 +
 .../commons/pipe/agent/task/meta/PipeStatus.java   |   3 +
 .../iotdb/commons/pipe/task/PipeMetaDeSerTest.java |  13 ++
 17 files changed, 686 insertions(+), 31 deletions(-)

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 4912f2dad5c..f656444ca19 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
@@ -64,6 +64,7 @@ import 
org.apache.iotdb.confignode.procedure.impl.node.AddConfigNodeProcedure;
 import org.apache.iotdb.confignode.procedure.impl.node.RemoveAINodeProcedure;
 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.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;
@@ -1360,7 +1361,7 @@ public class ProcedureManager {
         return status;
       } else {
         return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode())
-            
.setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage()));
+            .setMessage(wrapTimeoutMessageForPipeProcedure(status, procedure));
       }
     } catch (Exception e) {
       return new 
TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage());
@@ -1376,7 +1377,7 @@ public class ProcedureManager {
         return status;
       } else {
         return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode())
-            
.setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage()));
+            .setMessage(wrapTimeoutMessageForPipeProcedure(status, procedure));
       }
     } catch (Exception e) {
       return new 
TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage());
@@ -1392,7 +1393,7 @@ public class ProcedureManager {
         return status;
       } else {
         return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode())
-            
.setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage()));
+            .setMessage(wrapTimeoutMessageForPipeProcedure(status, procedure));
       }
     } catch (Exception e) {
       return new 
TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage());
@@ -1408,7 +1409,7 @@ public class ProcedureManager {
         return status;
       } else {
         return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode())
-            
.setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage()));
+            .setMessage(wrapTimeoutMessageForPipeProcedure(status, procedure));
       }
     } catch (Exception e) {
       return new 
TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage());
@@ -1424,7 +1425,7 @@ public class ProcedureManager {
         return status;
       } else {
         return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode())
-            
.setMessage(wrapTimeoutMessageForPipeProcedure(status.getMessage()));
+            .setMessage(wrapTimeoutMessageForPipeProcedure(status, procedure));
       }
     } catch (Exception e) {
       return new 
TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage());
@@ -1676,7 +1677,7 @@ public class ProcedureManager {
     if (!procedure.isFinished()) {
       // The procedure is still executing
       status =
-          RpcUtils.getStatus(TSStatusCode.OVERLAP_WITH_EXISTING_TASK, 
PROCEDURE_TIMEOUT_MESSAGE);
+          RpcUtils.getStatus(TSStatusCode.INTERNAL_REQUEST_TIME_OUT, 
PROCEDURE_TIMEOUT_MESSAGE);
     } else {
       if (procedure.isSuccess()) {
         if (procedure.getResult() != null) {
@@ -1714,6 +1715,17 @@ public class ProcedureManager {
     return message;
   }
 
+  private static boolean isProcedureTimeout(final TSStatus status) {
+    return status.getCode() == 
TSStatusCode.INTERNAL_REQUEST_TIME_OUT.getStatusCode();
+  }
+
+  private static String wrapTimeoutMessageForPipeProcedure(
+      final TSStatus status, final AbstractOperatePipeProcedureV2 procedure) {
+    return isProcedureTimeout(status)
+        ? procedure.getTimeoutDiagnosticMessage()
+        : status.getMessage();
+  }
+
   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 60ca3b8580a..e67a832ea3e 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
@@ -136,6 +136,11 @@ public class PipeHeartbeatParser {
       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) {
@@ -271,6 +276,9 @@ public class PipeHeartbeatParser {
                   .filter(true, pipeName).getAllPipeMeta().stream()
                       .filter(pipeMeta -> 
!pipeMeta.getStaticMeta().getPipeName().equals(pipeName))
                       .map(PipeMeta::getRuntimeMeta)
+                      .filter(
+                          runtimeMeta ->
+                              
!PipeStatus.PRE_DELETE.equals(runtimeMeta.getStatus().get()))
                       .filter(
                           runtimeMeta -> 
!runtimeMeta.getStatus().get().equals(PipeStatus.STOPPED))
                       .forEach(
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 01b043d57a2..e1d1eaacba2 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
@@ -213,6 +213,13 @@ public class PipeTaskInfo implements SnapshotProcessor {
       LOGGER.info(exceptionMessage);
       throw new PipeException(exceptionMessage);
     }
+    if 
(PipeStatus.PRE_DELETE.equals(getPipeStatus(alterPipeRequest.getPipeName()))) {
+      final String exceptionMessage =
+          String.format(
+              "Failed to alter pipe %s, the pipe is being dropped", 
alterPipeRequest.getPipeName());
+      LOGGER.info(exceptionMessage);
+      throw new PipeException(exceptionMessage);
+    }
 
     final PipeStaticMeta pipeStaticMetaFromCoordinator =
         getPipeMetaByPipeName(alterPipeRequest.getPipeName()).getStaticMeta();
@@ -298,6 +305,12 @@ public class PipeTaskInfo implements SnapshotProcessor {
       LOGGER.warn(exceptionMessage);
       throw new PipeException(exceptionMessage);
     }
+    if (pipeStatus == PipeStatus.PRE_DELETE) {
+      final String exceptionMessage =
+          String.format("Failed to start pipe %s, the pipe is being dropped", 
pipeName);
+      LOGGER.warn(exceptionMessage);
+      throw new PipeException(exceptionMessage);
+    }
   }
 
   public void checkBeforeStopPipe(final String pipeName) throws PipeException {
@@ -324,6 +337,12 @@ public class PipeTaskInfo implements SnapshotProcessor {
       LOGGER.warn(exceptionMessage);
       throw new PipeException(exceptionMessage);
     }
+    if (pipeStatus == PipeStatus.PRE_DELETE) {
+      final String exceptionMessage =
+          String.format("Failed to stop pipe %s, the pipe is being dropped", 
pipeName);
+      LOGGER.warn(exceptionMessage);
+      throw new PipeException(exceptionMessage);
+    }
   }
 
   public void checkBeforeDropPipe(final String pipeName) {
@@ -836,6 +855,10 @@ public class PipeTaskInfo implements SnapshotProcessor {
                   final PipeRuntimeMeta runtimeMeta =
                       
pipeMetaKeeper.getPipeMeta(message.getPipeName()).getRuntimeMeta();
 
+                  if 
(PipeStatus.PRE_DELETE.equals(runtimeMeta.getStatus().get())) {
+                    return;
+                  }
+
                   // 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.
@@ -889,7 +912,8 @@ public class PipeTaskInfo implements SnapshotProcessor {
         .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/env/ConfigNodeProcedureEnv.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
index a96d12af11c..174bc2ca6dc 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
@@ -99,9 +99,11 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
 import java.util.TreeMap;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.locks.ReentrantLock;
+import java.util.function.Consumer;
 import java.util.stream.Collectors;
 
 public class ConfigNodeProcedureEnv {
@@ -109,6 +111,7 @@ public class ConfigNodeProcedureEnv {
   private static final Logger LOG = 
LoggerFactory.getLogger(ConfigNodeProcedureEnv.class);
 
   private static final int RUNTIME_META_PUSH_RETRY_NUM = 1;
+  private static final Consumer<Set<Integer>> NO_OP_PENDING_DATA_NODE_TRACKER 
= ignored -> {};
 
   /** Add or remove node lock. */
   private final LockQueue nodeLock = new LockQueue();
@@ -683,6 +686,11 @@ public class ConfigNodeProcedureEnv {
 
   public Map<Integer, TPushPipeMetaResp> pushAllPipeMetaToDataNodes(
       List<ByteBuffer> pipeMetaBinaryList) {
+    return pushAllPipeMetaToDataNodes(pipeMetaBinaryList, 
NO_OP_PENDING_DATA_NODE_TRACKER);
+  }
+
+  public Map<Integer, TPushPipeMetaResp> pushAllPipeMetaToDataNodes(
+      List<ByteBuffer> pipeMetaBinaryList, final Consumer<Set<Integer>> 
pendingDataNodeTracker) {
     final Map<Integer, TDataNodeLocation> dataNodeLocationMap =
         configManager.getNodeManager().getRegisteredDataNodeLocations();
     final TPushPipeMetaReq request = new 
TPushPipeMetaReq().setPipeMetas(pipeMetaBinaryList);
@@ -691,13 +699,17 @@ public class ConfigNodeProcedureEnv {
         new DataNodeAsyncRequestContext<>(
             CnToDnAsyncRequestType.PIPE_PUSH_ALL_META, request, 
dataNodeLocationMap);
     final long timeoutInMs = getRequiredPipeMetadataRequestTimeoutInMs();
-    sendRequiredMetadataRequest(clientHandler, timeoutInMs);
-    fillMissingPipePushMetaResponses(clientHandler, timeoutInMs);
-    return clientHandler.getResponseMap();
+    return sendPipeMetaRequest(clientHandler, timeoutInMs, false, 
pendingDataNodeTracker);
   }
 
   public Map<Integer, TPushPipeMetaResp> pushAllPipeMetaToDataNodesBestEffort(
       List<ByteBuffer> pipeMetaBinaryList) {
+    return pushAllPipeMetaToDataNodesBestEffort(
+        pipeMetaBinaryList, NO_OP_PENDING_DATA_NODE_TRACKER);
+  }
+
+  public Map<Integer, TPushPipeMetaResp> pushAllPipeMetaToDataNodesBestEffort(
+      List<ByteBuffer> pipeMetaBinaryList, final Consumer<Set<Integer>> 
pendingDataNodeTracker) {
     final Map<Integer, TDataNodeLocation> dataNodeLocationMap =
         configManager.getNodeManager().getRegisteredDataNodeLocations();
     final TPushPipeMetaReq request = new 
TPushPipeMetaReq().setPipeMetas(pipeMetaBinaryList);
@@ -705,12 +717,16 @@ public class ConfigNodeProcedureEnv {
     final DataNodeAsyncRequestContext<TPushPipeMetaReq, TPushPipeMetaResp> 
clientHandler =
         new DataNodeAsyncRequestContext<>(
             CnToDnAsyncRequestType.PIPE_PUSH_ALL_META, request, 
dataNodeLocationMap);
-    final long timeoutInMs = sendBestEffortRuntimeMetaRequest(clientHandler);
-    fillMissingPipePushMetaResponses(clientHandler, timeoutInMs);
-    return clientHandler.getResponseMap();
+    return sendPipeMetaRequest(
+        clientHandler, getRuntimeMetaPushTimeoutInMs(), true, 
pendingDataNodeTracker);
   }
 
   public Map<Integer, TPushPipeMetaResp> 
pushSinglePipeMetaToDataNodes(ByteBuffer pipeMetaBinary) {
+    return pushSinglePipeMetaToDataNodes(pipeMetaBinary, 
NO_OP_PENDING_DATA_NODE_TRACKER);
+  }
+
+  public Map<Integer, TPushPipeMetaResp> pushSinglePipeMetaToDataNodes(
+      ByteBuffer pipeMetaBinary, final Consumer<Set<Integer>> 
pendingDataNodeTracker) {
     final Map<Integer, TDataNodeLocation> dataNodeLocationMap =
         configManager.getNodeManager().getRegisteredDataNodeLocations();
     final TPushSinglePipeMetaReq request = new 
TPushSinglePipeMetaReq().setPipeMeta(pipeMetaBinary);
@@ -719,12 +735,15 @@ public class ConfigNodeProcedureEnv {
         new DataNodeAsyncRequestContext<>(
             CnToDnAsyncRequestType.PIPE_PUSH_SINGLE_META, request, 
dataNodeLocationMap);
     final long timeoutInMs = getRequiredPipeMetadataRequestTimeoutInMs();
-    sendRequiredMetadataRequest(clientHandler, timeoutInMs);
-    fillMissingPipePushMetaResponses(clientHandler, timeoutInMs);
-    return clientHandler.getResponseMap();
+    return sendPipeMetaRequest(clientHandler, timeoutInMs, false, 
pendingDataNodeTracker);
   }
 
   public Map<Integer, TPushPipeMetaResp> dropSinglePipeOnDataNodes(String 
pipeNameToDrop) {
+    return dropSinglePipeOnDataNodes(pipeNameToDrop, 
NO_OP_PENDING_DATA_NODE_TRACKER);
+  }
+
+  public Map<Integer, TPushPipeMetaResp> dropSinglePipeOnDataNodes(
+      String pipeNameToDrop, final Consumer<Set<Integer>> 
pendingDataNodeTracker) {
     final Map<Integer, TDataNodeLocation> dataNodeLocationMap =
         configManager.getNodeManager().getRegisteredDataNodeLocations();
     final TPushSinglePipeMetaReq request =
@@ -734,9 +753,7 @@ public class ConfigNodeProcedureEnv {
         new DataNodeAsyncRequestContext<>(
             CnToDnAsyncRequestType.PIPE_PUSH_SINGLE_META, request, 
dataNodeLocationMap);
     final long timeoutInMs = getRequiredPipeMetadataRequestTimeoutInMs();
-    sendRequiredMetadataRequest(clientHandler, timeoutInMs);
-    fillMissingPipePushMetaResponses(clientHandler, timeoutInMs);
-    return clientHandler.getResponseMap();
+    return sendPipeMetaRequest(clientHandler, timeoutInMs, false, 
pendingDataNodeTracker);
   }
 
   public Map<Integer, TPushPipeMetaResp> pushMultiPipeMetaToDataNodes(
@@ -936,6 +953,22 @@ public class ConfigNodeProcedureEnv {
         .collect(Collectors.toList());
   }
 
+  private static Map<Integer, TPushPipeMetaResp> sendPipeMetaRequest(
+      final DataNodeAsyncRequestContext<?, TPushPipeMetaResp> clientHandler,
+      final long timeoutInMs,
+      final boolean keepSilent,
+      final Consumer<Set<Integer>> pendingDataNodeTracker) {
+    pendingDataNodeTracker.accept(
+        
Collections.unmodifiableSet(clientHandler.getNodeLocationMap().keySet()));
+    try {
+      sendRuntimeMetaRequest(clientHandler, keepSilent, timeoutInMs);
+      fillMissingPipePushMetaResponses(clientHandler, timeoutInMs);
+      return clientHandler.getResponseMap();
+    } finally {
+      pendingDataNodeTracker.accept(Collections.emptySet());
+    }
+  }
+
   private static long sendBestEffortRuntimeMetaRequest(
       final DataNodeAsyncRequestContext<?, ?> clientHandler) {
     final long timeoutInMs = getRuntimeMetaPushTimeoutInMs();
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 588a04dacb1..2f74bddfed8 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
@@ -131,6 +131,17 @@ public abstract class StateMachineProcedure<Env, TState> 
extends Procedure<Env>
     setNextState(getStateId(state));
   }
 
+  /**
+   * Returns whether the specified state is already present in the persisted 
state history.
+   *
+   * <p>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 0bfe1612f07..5587b51af8d 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
@@ -23,6 +23,7 @@ import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMeta;
 import org.apache.iotdb.commons.pipe.config.constant.SystemConstant;
 import 
org.apache.iotdb.confignode.manager.pipe.metric.overview.PipeProcedureMetrics;
 import org.apache.iotdb.confignode.persistence.pipe.PipeTaskInfo;
+import org.apache.iotdb.confignode.procedure.Procedure;
 import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
 import org.apache.iotdb.confignode.procedure.exception.ProcedureException;
 import org.apache.iotdb.confignode.procedure.impl.node.AbstractNodeProcedure;
@@ -46,6 +47,8 @@ import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
+import java.util.TreeSet;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicReference;
 
@@ -96,6 +99,14 @@ public abstract class AbstractOperatePipeProcedureV2
   // recovered procedure is already re-scheduled by the procedure framework.
   private transient boolean shouldYieldAfterExecution;
 
+  // These fields only describe where a running Pipe procedure is blocked when 
the caller times out.
+  // They do not affect execution and do not need to be persisted.
+  private volatile PipeProcedureExecutionStage executionStage =
+      PipeProcedureExecutionStage.WAITING_FOR_PROCEDURE_WORKER;
+  private volatile String lastExecutionExceptionMessage;
+  private volatile Procedure<?> nodeLockOwnerProcedure;
+  private volatile Set<Integer> pendingDataNodeIds = Collections.emptySet();
+
   private static final String SKIP_PIPE_PROCEDURE_MESSAGE =
       "Try to start a RUNNING pipe or stop a STOPPED pipe, do nothing.";
 
@@ -110,6 +121,7 @@ public abstract class AbstractOperatePipeProcedureV2
 
   @Override
   protected ProcedureLockState acquireLock(ConfigNodeProcedureEnv 
configNodeProcedureEnv) {
+    executionStage = 
PipeProcedureExecutionStage.WAITING_FOR_PIPE_TASK_COORDINATOR_LOCK;
     LOGGER.debug("ProcedureId {} try to acquire pipe lock.", getProcId());
     pipeTaskInfo = acquireLockInternal(configNodeProcedureEnv);
     if (pipeTaskInfo == null) {
@@ -118,9 +130,12 @@ public abstract class AbstractOperatePipeProcedureV2
       LOGGER.debug("ProcedureId {} acquired pipe lock.", getProcId());
     }
 
+    executionStage = PipeProcedureExecutionStage.WAITING_FOR_NODE_LOCK;
     final ProcedureLockState procedureLockState = 
super.acquireLock(configNodeProcedureEnv);
     switch (procedureLockState) {
       case LOCK_ACQUIRED:
+        nodeLockOwnerProcedure = null;
+        updateExecutionStage(getCurrentState(), false);
         if (pipeTaskInfo == null) {
           LOGGER.warn(
               "ProcedureId {}: LOCK_ACQUIRED. The following procedure should 
not be executed without pipe lock.",
@@ -132,6 +147,7 @@ public abstract class AbstractOperatePipeProcedureV2
         }
         break;
       case LOCK_EVENT_WAIT:
+        nodeLockOwnerProcedure = 
configNodeProcedureEnv.getNodeLock().getLockOwnerProcedure();
         if (pipeTaskInfo == null) {
           LOGGER.warn("ProcedureId {}: LOCK_EVENT_WAIT. Without acquiring pipe 
lock.", getProcId());
         } else {
@@ -177,6 +193,9 @@ public abstract class AbstractOperatePipeProcedureV2
             .updateTimer(this.getOperation().getName(), this.elapsedTime());
       }
       releasePipeTaskCoordinatorLock(configNodeProcedureEnv);
+      if (!isFinished()) {
+        executionStage = 
PipeProcedureExecutionStage.WAITING_FOR_PROCEDURE_WORKER;
+      }
     }
   }
 
@@ -221,6 +240,7 @@ public abstract class AbstractOperatePipeProcedureV2
   protected Flow executeFromState(ConfigNodeProcedureEnv env, 
OperatePipeTaskState state)
       throws InterruptedException {
     shouldYieldAfterExecution = false;
+    updateExecutionStage(state, false);
     if (pipeTaskInfo == null) {
       LOGGER.warn(
           "ProcedureId {}: Pipe lock is not acquired, executeFromState's 
execution will be skipped.",
@@ -232,6 +252,7 @@ public abstract class AbstractOperatePipeProcedureV2
       switch (state) {
         case VALIDATE_TASK:
           if (!executeFromValidateTask(env)) {
+            lastExecutionExceptionMessage = null;
             LOGGER.info("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
@@ -251,12 +272,15 @@ public abstract class AbstractOperatePipeProcedureV2
           break;
         case OPERATE_ON_DATA_NODES:
           executeFromOperateOnDataNodes(env);
+          lastExecutionExceptionMessage = null;
           return Flow.NO_MORE_STATE;
         default:
           throw new UnsupportedOperationException(
               String.format("Unknown state during executing 
operatePipeProcedure, %s", state));
       }
+      lastExecutionExceptionMessage = null;
     } catch (Exception e) {
+      lastExecutionExceptionMessage = getExceptionMessage(e);
       // Retry before rollback
       if (getCycles() < RETRY_THRESHOLD) {
         LOGGER.warn(
@@ -301,6 +325,7 @@ public abstract class AbstractOperatePipeProcedureV2
   @Override
   protected void rollbackState(ConfigNodeProcedureEnv env, 
OperatePipeTaskState state)
       throws IOException, InterruptedException, ProcedureException {
+    updateExecutionStage(state, true);
     if (pipeTaskInfo == null) {
       LOGGER.warn(
           "ProcedureId {}: Pipe lock is not acquired, rollbackState({})'s 
execution will be skipped.",
@@ -387,6 +412,167 @@ public abstract class AbstractOperatePipeProcedureV2
     return OperatePipeTaskState.VALIDATE_TASK;
   }
 
+  public final String getTimeoutDiagnosticMessage() {
+    final PipeProcedureExecutionStage currentExecutionStage = executionStage;
+    return String.format(
+        "Pipe operation %s timed out (procedureId=%d). Stuck at %s. Reason: 
%s. "
+            + "The procedure is still running.",
+        getOperation().name(),
+        getProcId(),
+        currentExecutionStage.name(),
+        getTimeoutReason(currentExecutionStage));
+  }
+
+  private String getTimeoutReason(final PipeProcedureExecutionStage 
currentExecutionStage) {
+    final String failureMessage = getFailureMessage();
+    if (currentExecutionStage.isRollback()) {
+      return getRollbackTimeoutReason(failureMessage);
+    }
+    if (isFailed() && failureMessage != null) {
+      return String.format("the state failed with '%s' and rollback is 
pending.", failureMessage);
+    }
+    if (lastExecutionExceptionMessage != null) {
+      return String.format(
+          "the previous attempt failed with '%s' and this state is being 
retried.",
+          lastExecutionExceptionMessage);
+    }
+
+    switch (currentExecutionStage) {
+      case WAITING_FOR_PROCEDURE_WORKER:
+        return "no Procedure worker is currently available; "
+            + "workers may be busy or blocked by other procedures.";
+      case WAITING_FOR_PIPE_TASK_COORDINATOR_LOCK:
+        return "waiting to acquire the PipeTaskCoordinator lock because "
+            + "another Pipe operation is holding it.";
+      case WAITING_FOR_NODE_LOCK:
+        return getNodeLockTimeoutReason();
+      case VALIDATE_TASK:
+        return "Pipe request or plugin validation has not completed; "
+            + "a plugin check or metadata access may be slow.";
+      case CALCULATE_INFO_FOR_TASK:
+        return "Pipe metadata calculation has not completed; "
+            + "metadata access or local calculation may be slow.";
+      case WRITE_CONFIG_NODE_CONSENSUS:
+        return "the ConfigNode consensus write has not returned; "
+            + "run SHOW CLUSTER to check node status.";
+      case OPERATE_ON_DATA_NODES:
+        return getDataNodeTimeoutReason();
+      default:
+        return getRollbackTimeoutReason(failureMessage);
+    }
+  }
+
+  private static String getRollbackTimeoutReason(final String failureMessage) {
+    return failureMessage == null
+        ? "rolling back after an earlier failure."
+        : String.format("rolling back after failure: %s.", failureMessage);
+  }
+
+  private String getNodeLockTimeoutReason() {
+    final Procedure<?> lockOwnerProcedure = nodeLockOwnerProcedure;
+    if (lockOwnerProcedure == null) {
+      return "waiting to acquire the ConfigNode node lock because another node 
procedure is holding it.";
+    }
+    final String operation =
+        lockOwnerProcedure instanceof AbstractOperatePipeProcedureV2
+            ? ((AbstractOperatePipeProcedureV2) 
lockOwnerProcedure).getOperation().name()
+            : lockOwnerProcedure.getClass().getSimpleName();
+    return String.format(
+        "waiting to acquire the ConfigNode node lock held by %s 
(procedureId=%d).",
+        operation, lockOwnerProcedure.getProcId());
+  }
+
+  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 String getDataNodeTimeoutReason() {
+    final Set<Integer> currentPendingDataNodeIds = new 
TreeSet<>(pendingDataNodeIds);
+    return currentPendingDataNodeIds.isEmpty()
+        ? "the Pipe metadata push has not completed; run SHOW CLUSTER to check 
DataNode status."
+        : String.format(
+            "DataNodes %s have not responded to the Pipe metadata push; "
+                + "run SHOW CLUSTER to check their status.",
+            currentPendingDataNodeIds);
+  }
+
+  protected final 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 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;
+    }
+  }
+
+  protected final void setPendingDataNodeIds(final Set<Integer> 
pendingDataNodeIds) {
+    this.pendingDataNodeIds = pendingDataNodeIds;
+  }
+
+  private enum PipeProcedureExecutionStage {
+    WAITING_FOR_PROCEDURE_WORKER,
+    WAITING_FOR_PIPE_TASK_COORDINATOR_LOCK,
+    WAITING_FOR_NODE_LOCK,
+    VALIDATE_TASK,
+    CALCULATE_INFO_FOR_TASK,
+    WRITE_CONFIG_NODE_CONSENSUS,
+    OPERATE_ON_DATA_NODES,
+    ROLLBACK_VALIDATE_TASK(true),
+    ROLLBACK_CALCULATE_INFO_FOR_TASK(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.
    *
@@ -401,7 +587,7 @@ public abstract class AbstractOperatePipeProcedureV2
       pipeMetaBinaryList.add(pipeMeta.serialize());
     }
 
-    return env.pushAllPipeMetaToDataNodes(pipeMetaBinaryList);
+    return env.pushAllPipeMetaToDataNodes(pipeMetaBinaryList, 
this::setPendingDataNodeIds);
   }
 
   /**
@@ -527,7 +713,8 @@ public abstract class AbstractOperatePipeProcedureV2
     for (final PipeMeta pipeMeta : pipeTaskInfo.get().getPipeMetaList()) {
       pipeMetaBinaryList.add(pipeMeta.serialize());
     }
-    return env.pushAllPipeMetaToDataNodesBestEffort(pipeMetaBinaryList);
+    return env.pushAllPipeMetaToDataNodesBestEffort(
+        pipeMetaBinaryList, this::setPendingDataNodeIds);
   }
 
   protected void pushPipeMetaToDataNodesBestEffort(ConfigNodeProcedureEnv env) 
{
@@ -549,7 +736,8 @@ public abstract class AbstractOperatePipeProcedureV2
   protected Map<Integer, TPushPipeMetaResp> pushSinglePipeMetaToDataNodes(
       String pipeName, ConfigNodeProcedureEnv env) throws IOException {
     return env.pushSinglePipeMetaToDataNodes(
-        pipeTaskInfo.get().getPipeMetaByPipeName(pipeName).serialize());
+        pipeTaskInfo.get().getPipeMetaByPipeName(pipeName).serialize(),
+        this::setPendingDataNodeIds);
   }
 
   protected Map<Integer, TPushPipeMetaResp> 
pushSinglePipeMetaToDataNodes4Realtime(
@@ -565,7 +753,8 @@ public abstract class AbstractOperatePipeProcedureV2
                 Collections.singletonMap(
                     SystemConstant.RESTART_OR_NEWLY_ADDED_KEY, 
Boolean.FALSE.toString())));
     return env.pushSinglePipeMetaToDataNodes(
-        pipeTaskInfo.get().getPipeMetaByPipeName(pipeName).serialize());
+        pipeTaskInfo.get().getPipeMetaByPipeName(pipeName).serialize(),
+        this::setPendingDataNodeIds);
   }
 
   /**
@@ -577,7 +766,7 @@ public abstract class AbstractOperatePipeProcedureV2
    */
   protected Map<Integer, TPushPipeMetaResp> dropSinglePipeOnDataNodes(
       String pipeName, ConfigNodeProcedureEnv env) {
-    return env.dropSinglePipeOnDataNodes(pipeName);
+    return env.dropSinglePipeOnDataNodes(pipeName, 
this::setPendingDataNodeIds);
   }
 
   @Override
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 0f9a8bc0d8c..62db3c53b6d 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
@@ -20,11 +20,15 @@
 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.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.persistence.pipe.PipeTaskInfo;
 import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
 import 
org.apache.iotdb.confignode.procedure.impl.pipe.AbstractOperatePipeProcedureV2;
 import org.apache.iotdb.confignode.procedure.impl.pipe.PipeTaskOperation;
+import 
org.apache.iotdb.confignode.procedure.state.pipe.task.OperatePipeTaskState;
 import org.apache.iotdb.confignode.procedure.store.ProcedureType;
 import org.apache.iotdb.consensus.exception.ConsensusException;
 import org.apache.iotdb.pipe.api.exception.PipeException;
@@ -43,8 +47,10 @@ import java.util.concurrent.atomic.AtomicReference;
 public class DropPipeProcedureV2 extends AbstractOperatePipeProcedureV2 {
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(DropPipeProcedureV2.class);
+  private static final int SERIALIZATION_VERSION_MAGIC = 0x44505632;
 
   private String pipeName;
+  private PipeMeta pipeMetaToDrop;
 
   public DropPipeProcedureV2() {
     super();
@@ -67,6 +73,10 @@ public class DropPipeProcedureV2 extends 
AbstractOperatePipeProcedureV2 {
     return pipeName;
   }
 
+  PipeMeta getPipeMetaToDrop() {
+    return pipeMetaToDrop;
+  }
+
   @Override
   protected PipeTaskOperation getOperation() {
     return PipeTaskOperation.DROP_PIPE;
@@ -84,13 +94,41 @@ public class DropPipeProcedureV2 extends 
AbstractOperatePipeProcedureV2 {
   @Override
   public void executeFromCalculateInfoForTask(ConfigNodeProcedureEnv env) 
throws PipeException {
     LOGGER.info("DropPipeProcedureV2: executeFromCalculateInfoForTask({})", 
pipeName);
-    // Do nothing
+    pipeMetaToDrop = pipeTaskInfo.get().getPipeMetaByPipeName(pipeName);
   }
 
   @Override
   public void executeFromWriteConfigNodeConsensus(ConfigNodeProcedureEnv env) 
throws PipeException {
     LOGGER.info("DropPipeProcedureV2: 
executeFromWriteConfigNodeConsensus({})", pipeName);
 
+    if (!restorePipeMetaToDropIfNecessary()) {
+      return;
+    }
+
+    TSStatus response;
+    try {
+      response =
+          env.getConfigManager()
+              .getConsensusManager()
+              .write(new SetPipeStatusPlanV2(pipeName, PipeStatus.PRE_DELETE));
+    } catch (ConsensusException e) {
+      LOGGER.warn("Failed in the write API executing the consensus layer due 
to: ", 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 = pipeTaskInfo.get().getPipeMetaByPipeName(pipeName);
+    }
+    return pipeMetaToDrop != null;
+  }
+
+  private void dropPipeOnConfigNode(final ConfigNodeProcedureEnv env) throws 
PipeException {
     TSStatus response;
     try {
       response = env.getConfigManager().getConsensusManager().write(new 
DropPipePlanV2(pipeName));
@@ -105,17 +143,34 @@ public class DropPipeProcedureV2 extends 
AbstractOperatePipeProcedureV2 {
   }
 
   @Override
-  public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) {
+  public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) throws 
PipeException {
     LOGGER.info("DropPipeProcedureV2: executeFromOperateOnDataNodes({})", 
pipeName);
 
-    final String exceptionMessage =
-        parsePushPipeMetaExceptionForPipe(pipeName, 
dropSinglePipeOnDataNodes(pipeName, env));
-    if (!exceptionMessage.isEmpty()) {
+    String exceptionMessage;
+    try {
+      if (pipeMetaToDrop == null) {
+        exceptionMessage =
+            parsePushPipeMetaExceptionForPipe(pipeName, 
dropSinglePipeOnDataNodes(pipeName, env));
+      } else {
+        final PipeMeta droppedPipeMeta = pipeMetaToDrop.deepCopy4TaskAgent();
+        droppedPipeMeta.getRuntimeMeta().getStatus().set(PipeStatus.DROPPED);
+        exceptionMessage =
+            parsePushPipeMetaExceptionForPipe(
+                pipeName,
+                env.pushSinglePipeMetaToDataNodes(
+                    droppedPipeMeta.serialize(), this::setPendingDataNodeIds));
+      }
+    } catch (final IOException e) {
+      exceptionMessage = e.getMessage();
+    }
+    if (exceptionMessage != null && !exceptionMessage.isEmpty()) {
       LOGGER.warn(
           "Failed to drop pipe {}, details: {}, metadata will be synchronized 
later.",
           pipeName,
           exceptionMessage);
     }
+    updateExecutionStage(OperatePipeTaskState.WRITE_CONFIG_NODE_CONSENSUS, 
false);
+    dropPipeOnConfigNode(env);
   }
 
   @Override
@@ -147,12 +202,28 @@ public class DropPipeProcedureV2 extends 
AbstractOperatePipeProcedureV2 {
     stream.writeShort(ProcedureType.DROP_PIPE_PROCEDURE_V2.getTypeCode());
     super.serialize(stream);
     ReadWriteIOUtils.write(pipeName, stream);
+    ReadWriteIOUtils.write(SERIALIZATION_VERSION_MAGIC, stream);
+    ReadWriteIOUtils.write(pipeMetaToDrop != null, stream);
+    if (pipeMetaToDrop != null) {
+      pipeMetaToDrop.serialize(stream);
+    }
   }
 
   @Override
   public void deserialize(ByteBuffer byteBuffer) {
     super.deserialize(byteBuffer);
     pipeName = ReadWriteIOUtils.readString(byteBuffer);
+    if (byteBuffer.remaining() < Integer.BYTES) {
+      return;
+    }
+    final int position = byteBuffer.position();
+    if (ReadWriteIOUtils.readInt(byteBuffer) != SERIALIZATION_VERSION_MAGIC) {
+      byteBuffer.position(position);
+      return;
+    }
+    if (byteBuffer.hasRemaining() && ReadWriteIOUtils.readBool(byteBuffer)) {
+      pipeMetaToDrop = PipeMeta.deserialize4Coordinator(byteBuffer);
+    }
   }
 
   @Override
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/scheduler/LockQueue.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/scheduler/LockQueue.java
index e2f5935a909..57b0817940e 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/scheduler/LockQueue.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/scheduler/LockQueue.java
@@ -27,7 +27,7 @@ import java.util.ArrayDeque;
 public class LockQueue {
   private final ArrayDeque<Procedure<?>> deque = new ArrayDeque<>();
 
-  private Procedure<?> lockOwnerProcedure = null;
+  private volatile Procedure<?> lockOwnerProcedure = null;
 
   public boolean tryLock(Procedure<?> procedure) {
     if (lockOwnerProcedure == null) {
@@ -45,6 +45,10 @@ public class LockQueue {
     return true;
   }
 
+  public Procedure<?> getLockOwnerProcedure() {
+    return lockOwnerProcedure;
+  }
+
   public void waitProcedure(Procedure<?> procedure, ProcedureScheduler 
procedureScheduler) {
     if (lockOwnerProcedure == null) {
       procedureScheduler.addFront(procedure);
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 1c73fe323ce..c231ff18ead 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
@@ -895,7 +895,7 @@ public class ConfigPhysicalPlanSerDeTest {
   public void SetPipeStatusPlanV2Test() throws IOException {
     final SetPipeStatusPlanV2 setPipeStatusPlanV2 =
         new SetPipeStatusPlanV2(
-            "pipe", 
org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus.RUNNING);
+            "pipe", 
org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus.PRE_DELETE);
     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 94189a19d99..072f98a7500 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,12 +23,19 @@ 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.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.confignode.consensus.response.pipe.task.PipeTableResp;
+import org.apache.iotdb.confignode.manager.ConfigManager;
+import org.apache.iotdb.confignode.manager.node.NodeManager;
+import org.apache.iotdb.confignode.service.ConfigNode;
 import org.apache.iotdb.rpc.TSStatusCode;
 
+import org.junit.After;
 import org.junit.Assert;
+import org.junit.Before;
 import org.junit.Test;
+import org.mockito.Mockito;
 
 import java.util.ArrayList;
 import java.util.HashMap;
@@ -39,6 +46,26 @@ import java.util.concurrent.ConcurrentMap;
 
 public class PipeTableRespTest {
 
+  private ConfigNode oldInstance;
+
+  @Before
+  public void setUp() {
+    oldInstance = ConfigNode.getInstance();
+    final ConfigNode configNode = Mockito.mock(ConfigNode.class);
+    final ConfigManager configManager = Mockito.mock(ConfigManager.class);
+    final NodeManager nodeManager = Mockito.mock(NodeManager.class);
+
+    Mockito.when(configNode.getConfigManager()).thenReturn(configManager);
+    Mockito.when(configManager.getNodeManager()).thenReturn(nodeManager);
+    Mockito.when(nodeManager.getRegisteredDataNodeCount()).thenReturn(2);
+    ConfigNode.setInstance(configNode);
+  }
+
+  @After
+  public void tearDown() {
+    ConfigNode.setInstance(oldInstance);
+  }
+
   public PipeTableResp constructPipeTableResp() {
     TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
     List<PipeMeta> pipeMetaList = new ArrayList<>();
@@ -118,4 +145,14 @@ public class PipeTableRespTest {
     PipeTableResp allPipeTableResp = pipeTableResp.filter(true, null);
     Assert.assertEquals(3, allPipeTableResp.getAllPipeMeta().size());
   }
+
+  @Test
+  public void testConvertToTShowPipeRespIncludesPreDeleteStatus() {
+    final PipeTableResp pipeTableResp = constructPipeTableResp();
+    
pipeTableResp.getAllPipeMeta().get(0).getRuntimeMeta().getStatus().set(PipeStatus.PRE_DELETE);
+
+    Assert.assertEquals(
+        PipeStatus.PRE_DELETE.name(),
+        
pipeTableResp.convertToTShowPipeResp().getPipeInfoList().get(0).getState());
+  }
 }
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 0c68881584e..b080f8f51b1 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
@@ -208,6 +208,41 @@ public class PipeHeartbeatParserTest {
     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<Integer, PipeTaskMeta> 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));
+
+    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());
+  }
+
   private ParserTestContext createParserTestContext(final int 
registeredDataNodeCount) {
     return createParserTestContext(registeredDataNodeCount, new 
PipeTaskInfo());
   }
diff --git 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfoAutoRestartTest.java
 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfoAutoRestartTest.java
index 0947637812d..6d933dbc413 100644
--- 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfoAutoRestartTest.java
+++ 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfoAutoRestartTest.java
@@ -27,8 +27,10 @@ import 
org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
 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.rpc.thrift.TAlterPipeReq;
 import org.apache.iotdb.mpp.rpc.thrift.TPushPipeMetaResp;
 import org.apache.iotdb.mpp.rpc.thrift.TPushPipeMetaRespExceptionMessage;
+import org.apache.iotdb.pipe.api.exception.PipeException;
 import org.apache.iotdb.rpc.TSStatusCode;
 
 import org.junit.Assert;
@@ -116,6 +118,50 @@ public class PipeTaskInfoAutoRestartTest {
     Assert.assertEquals(exceptionsClearTime, 
runtimeMeta.getExceptionsClearTime());
   }
 
+  @Test
+  public void testPreDeletePipeIsNotOverwrittenOrAutoRestarted() {
+    final String pipeName = "droppingPipe";
+    createPipe(pipeName, PipeStatus.STOPPED);
+    pipeTaskInfo.setPipeStatus(new SetPipeStatusPlanV2(pipeName, 
PipeStatus.PRE_DELETE));
+
+    Assert.assertTrue(
+        
pipeTaskInfo.recordDataNodePushPipeMetaExceptions(createErrorRespMap(pipeName)));
+
+    final PipeRuntimeMeta runtimeMeta =
+        pipeTaskInfo.getPipeMetaByPipeName(pipeName).getRuntimeMeta();
+    Assert.assertEquals(PipeStatus.PRE_DELETE, runtimeMeta.getStatus().get());
+    Assert.assertFalse(runtimeMeta.getIsStoppedByRuntimeException());
+    Assert.assertFalse(pipeTaskInfo.autoRestart());
+    Assert.assertEquals(PipeStatus.PRE_DELETE, runtimeMeta.getStatus().get());
+  }
+
+  @Test
+  public void testPreDeletePipeRejectsAlterStartAndStop() {
+    final String pipeName = "droppingPipe";
+    createPipe(pipeName, PipeStatus.STOPPED);
+    pipeTaskInfo.setPipeStatus(new SetPipeStatusPlanV2(pipeName, 
PipeStatus.PRE_DELETE));
+
+    final TAlterPipeReq alterPipeRequest =
+        new TAlterPipeReq(pipeName, new HashMap<>(), new HashMap<>(), false, 
false);
+    alterPipeRequest.setExtractorAttributes(new HashMap<>());
+    alterPipeRequest.setIsReplaceAllExtractorAttributes(false);
+
+    assertPipeBeingDropped(
+        pipeName, () -> 
pipeTaskInfo.checkAndUpdateRequestBeforeAlterPipe(alterPipeRequest));
+    assertPipeBeingDropped(pipeName, () -> 
pipeTaskInfo.checkBeforeStartPipe(pipeName));
+    assertPipeBeingDropped(pipeName, () -> 
pipeTaskInfo.checkBeforeStopPipe(pipeName));
+  }
+
+  private static void assertPipeBeingDropped(final String pipeName, final 
Runnable pipeOperation) {
+    try {
+      pipeOperation.run();
+      Assert.fail();
+    } catch (final PipeException e) {
+      Assert.assertTrue(e.getMessage().contains(pipeName));
+      Assert.assertTrue(e.getMessage().contains("being dropped"));
+    }
+  }
+
   private Map<Integer, TPushPipeMetaResp> createErrorRespMap(final String 
pipeName) {
     final TPushPipeMetaRespExceptionMessage exceptionMessage =
         new TPushPipeMetaRespExceptionMessage(
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 5993e73ec42..39da89ffe76 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
@@ -29,6 +29,9 @@ import org.junit.Assert;
 import org.junit.Test;
 
 import java.io.IOException;
+import java.util.Arrays;
+import java.util.HashSet;
+import java.util.Set;
 import java.util.concurrent.atomic.AtomicReference;
 
 public class AbstractOperatePipeProcedureV2Test {
@@ -89,6 +92,32 @@ public class AbstractOperatePipeProcedureV2Test {
     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.setPendingDataNodeIdsForTest(new HashSet<>(Arrays.asList(3, 1)));
+
+    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("DataNodes [1, 3]"));
+    Assert.assertTrue(diagnosticMessage.contains("SHOW CLUSTER"));
+  }
+
   private static class TestOperatePipeProcedure extends 
AbstractOperatePipeProcedureV2 {
 
     private int validateExecutionCount;
@@ -104,6 +133,10 @@ public class AbstractOperatePipeProcedureV2Test {
       return execute(null);
     }
 
+    private void setPendingDataNodeIdsForTest(final Set<Integer> 
pendingDataNodeIds) {
+      setPendingDataNodeIds(pendingDataNodeIds);
+    }
+
     @Override
     protected PipeTaskOperation getOperation() {
       return PipeTaskOperation.START_PIPE;
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 032288df93b..ddbac8da9aa 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,18 +19,54 @@
 
 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.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.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.DropPipePlanV2;
+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.assertNotNull;
 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 void setPipeTaskInfo(final PipeTaskInfo pipeTaskInfo) {
+      this.pipeTaskInfo = new AtomicReference<>(pipeTaskInfo);
+    }
+  }
+
   @Test
   public void serializeDeserializeTest() {
     PublicBAOS byteArrayOutputStream = new PublicBAOS();
@@ -46,8 +82,102 @@ public class DropPipeProcedureV2Test {
           (DropPipeProcedureV2) ProcedureFactory.getInstance().create(buffer);
 
       assertEquals(proc, proc2);
+
+      final ByteBuffer legacyBuffer =
+          ByteBuffer.wrap(
+              byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size() 
- Integer.BYTES - 1);
+      final DropPipeProcedureV2 legacyProc =
+          (DropPipeProcedureV2) 
ProcedureFactory.getInstance().create(legacyBuffer);
+      assertEquals(proc, legacyProc);
     } catch (Exception e) {
       fail();
     }
   }
+
+  @Test
+  public void testWriteConsensusMarksPreDelete() throws Exception {
+    final String pipeName = "testPipe";
+    final PipeTaskInfo pipeTaskInfo = createPipeTaskInfo(pipeName);
+    final TestDropPipeProcedureV2 procedure = new 
TestDropPipeProcedureV2(pipeName);
+    procedure.setPipeTaskInfo(pipeTaskInfo);
+    
procedure.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()));
+
+    procedure.executeFromWriteConfigNodeConsensus(env);
+
+    final ArgumentCaptor<ConfigPhysicalPlan> planCaptor =
+        ArgumentCaptor.forClass(ConfigPhysicalPlan.class);
+    Mockito.verify(consensusManager).write(planCaptor.capture());
+    assertEquals(new SetPipeStatusPlanV2(pipeName, PipeStatus.PRE_DELETE), 
planCaptor.getValue());
+  }
+
+  @Test
+  public void testDataNodeStagePushesDroppedMetaBeforeFinalDrop() throws 
Exception {
+    final String pipeName = "testPipe";
+    final PipeTaskInfo pipeTaskInfo = createPipeTaskInfo(pipeName);
+    final TestDropPipeProcedureV2 procedure = new 
TestDropPipeProcedureV2(pipeName);
+    procedure.setPipeTaskInfo(pipeTaskInfo);
+    
procedure.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(env.pushSinglePipeMetaToDataNodes(Mockito.any(ByteBuffer.class), 
Mockito.any()))
+        .thenReturn(Collections.emptyMap());
+    Mockito.when(consensusManager.write(Mockito.any(ConfigPhysicalPlan.class)))
+        .thenReturn(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
+
+    procedure.executeFromOperateOnDataNodes(env);
+
+    final ArgumentCaptor<ByteBuffer> metaCaptor = 
ArgumentCaptor.forClass(ByteBuffer.class);
+    Mockito.verify(env).pushSinglePipeMetaToDataNodes(metaCaptor.capture(), 
Mockito.any());
+    final PipeMeta droppedPipeMeta = 
PipeMeta.deserialize4TaskAgent(metaCaptor.getValue());
+    assertEquals(PipeStatus.DROPPED, 
droppedPipeMeta.getRuntimeMeta().getStatus().get());
+
+    final ArgumentCaptor<ConfigPhysicalPlan> planCaptor =
+        ArgumentCaptor.forClass(ConfigPhysicalPlan.class);
+    Mockito.verify(consensusManager).write(planCaptor.capture());
+    assertEquals(new DropPipePlanV2(pipeName), planCaptor.getValue());
+  }
+
+  @Test
+  public void testPipeMetaToDropSurvivesSerialization() throws Exception {
+    final String pipeName = "testPipe";
+    final TestDropPipeProcedureV2 procedure = new 
TestDropPipeProcedureV2(pipeName);
+    procedure.setPipeTaskInfo(createPipeTaskInfo(pipeName));
+    
procedure.executeFromCalculateInfoForTask(Mockito.mock(ConfigNodeProcedureEnv.class));
+
+    final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+    procedure.serialize(new DataOutputStream(byteArrayOutputStream));
+    final ByteBuffer buffer =
+        ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, 
byteArrayOutputStream.size());
+    final DropPipeProcedureV2 recoveredProcedure =
+        (DropPipeProcedureV2) ProcedureFactory.getInstance().create(buffer);
+
+    assertNotNull(recoveredProcedure.getPipeMetaToDrop());
+    assertEquals(pipeName, 
recoveredProcedure.getPipeMetaToDrop().getStaticMeta().getPipeName());
+  }
+
+  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 03906e7e977..1eef328bc7e 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
@@ -214,6 +214,12 @@ public abstract class PipeTaskAgent {
       return;
     }
 
+    // PRE_DELETE is a coordinator-only marker. The drop procedure pushes 
DROPPED explicitly after
+    // persisting the marker, so task agents should retain their current 
runtime state here.
+    if (metaFromCoordinator.getRuntimeMeta().getStatus().get() == 
PipeStatus.PRE_DELETE) {
+      return;
+    }
+
     final PipeMeta metaInAgent = pipeMetaKeeper.getPipeMeta(pipeName);
 
     // If pipe meta does not exist on local agent, create a new pipe
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 7898266c712..93330151f2f 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
@@ -23,6 +23,7 @@ public enum PipeStatus {
   RUNNING((byte) 0),
   STOPPED((byte) 1),
   DROPPED((byte) 2),
+  PRE_DELETE((byte) 3),
   ;
 
   private final byte type;
@@ -43,6 +44,8 @@ public enum PipeStatus {
         return PipeStatus.STOPPED;
       case 2:
         return PipeStatus.DROPPED;
+      case 3:
+        return PipeStatus.PRE_DELETE;
       default:
         throw new IllegalArgumentException("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 fa0ac2c1a37..1e2ecc1bfd3 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
@@ -49,6 +49,14 @@ import java.util.concurrent.ConcurrentHashMap;
 
 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 =
@@ -146,6 +154,11 @@ public class PipeMetaDeSerTest {
     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);

Reply via email to