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

924060929 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new fde71764978 [fix](coordinator) Preserve cancellation before 
coordinator setup (#68283)
fde71764978 is described below

commit fde71764978b3995ba96237f1daab518fbf9e30d
Author: 924060929 <[email protected]>
AuthorDate: Mon Sep 28 17:42:46 2026 +0800

    [fix](coordinator) Preserve cancellation before coordinator setup (#68283)
    
    ### What problem does this PR solve?
    
    INSERT and row-level DML statements expose their query ID while planning
    so they remain visible to `SHOW PROCESSLIST` and `KILL QUERY`. Their
    coordinator is created later, after planning and—in some
    paths—transaction or sink preparation.
    
    Previously, a timeout or cancellation in that interval had no
    coordinator to notify. The statement could later publish a new
    coordinator and continue setup or fragment dispatch even though it had
    already been terminated. A second cancellation could also overwrite the
    reason delivered to that coordinator.
    
    This PR keeps query IDs visible and closes the coordinator-publication
    race at one transactional lifecycle boundary:
    
    - `StmtExecutor` retains the first terminal status and replays that same
    status when a coordinator is published.
    - All `AbstractInsertExecutor` paths publish their coordinator at the
    common execution boundary, covering regular INSERT, row-level DML,
    rewrite-table, and connector-rewrite execution. The retained terminal
    status is fenced immediately after publication, before executor-specific
    setup runs, and again at the common pre-commit completion boundary.
    - Legacy and Nereids coordinators publish terminal status and internal
    cancellation before fallible scan cleanup, then reject execution and
    fragment dispatch with the original timeout or cancellation message.
    Scan cleanup is exhaustive and non-throwing, so one failing scan cannot
    skip the remaining scans, the owner's close, or mask the retained
    reason. Nereids cleanup also runs when the status publication races a
    partially initialized load processor.
    - A queue wait unblocked by TIMEOUT/KILL reports the retained terminal
    reason instead of the token's generic "query is cancelled".
    - A cancelled connector-rewrite group is reported as a failed task
    (queued or running), executor publication and task cancellation share
    one atomic handoff, the owner drains live groups (cancel plus a single
    shared terminal budget) before rolling the shared transaction back,
    queued cancellation is conclusively terminal, and a task whose scheduler
    publication fails is unregistered instead of retained.
    - The distributed rewrite owner receives the outer statement's sticky
    cancellation, cancels/drains live groups, and serializes cancellation
    with the source-registration/final-commit decision so a terminated
    statement cannot commit the shared rewrite.
    - Dictionary INSERT failure handling accepts the resulting general
    `UserException` instead of assuming every failure is a `DdlException`.
    
    INSERT INTO TVF is intentionally unchanged. It performs
    non-transactional FE-side recursive replacement of external files;
    making cancellation atomic with that destructive phase requires a
    separate staged or atomic replacement design, not coordinator
    publication alone.
    
    Cross-FE forwarding, point queries, FE-only results, result-file
    cleanup, and TVF replacement keep their existing cancellation contracts.
    
    ### Release note
    
    Fix transactional INSERT and DML statements continuing execution after
    cancellation during coordinator setup.
---
 .../trees/plans/commands/RowLevelDmlCommand.java   |   1 -
 .../commands/execute/ConnectorExecuteAction.java   |   2 +-
 .../commands/execute/ConnectorRewriteDriver.java   | 236 +++++++++++++---
 .../execute/ConnectorRewriteGroupTask.java         |  56 +++-
 .../commands/insert/AbstractInsertExecutor.java    |  18 ++
 .../commands/insert/DictionaryInsertExecutor.java  |  12 +-
 .../commands/insert/InsertIntoTableCommand.java    |   3 +-
 .../main/java/org/apache/doris/qe/Coordinator.java |  58 +++-
 .../org/apache/doris/qe/NereidsCoordinator.java    |  71 +++--
 .../java/org/apache/doris/qe/StmtExecutor.java     |  42 ++-
 .../org/apache/doris/qe/runtime/LoadProcessor.java |   5 +-
 .../doris/qe/runtime/PipelineExecutionTask.java    |   4 +
 .../doris/scheduler/disruptor/TaskDisruptor.java   |   4 +-
 .../scheduler/manager/TransientTaskManager.java    |  11 +-
 .../plans/commands/RowLevelDmlCommandTest.java     |   3 +-
 .../execute/ConnectorRewriteDriverTest.java        | 308 ++++++++++++++++++++-
 .../execute/ConnectorRewriteGroupTaskTest.java     | 117 ++++++++
 .../insert/DictionaryInsertTargetDropRaceTest.java |  23 ++
 .../commands/insert/OlapInsertExecutorTest.java    |  93 +++++++
 .../apache/doris/qe/NereidsCoordinatorTest.java    | 102 +++++++
 .../org/apache/doris/qe/OldCoordinatorTest.java    |  74 +++++
 .../java/org/apache/doris/qe/StmtExecutorTest.java |  72 +++++
 .../qe/runtime/PipelineExecutionTaskTest.java      |  38 +++
 .../manager/TransientTaskManagerTest.java          |  60 ++++
 24 files changed, 1315 insertions(+), 98 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/RowLevelDmlCommand.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/RowLevelDmlCommand.java
index c992f3a5a52..498a580e714 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/RowLevelDmlCommand.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/RowLevelDmlCommand.java
@@ -140,7 +140,6 @@ public class RowLevelDmlCommand {
             applyWriteConstraintIfPresent(transform, insertExecutor, 
analyzedPlan, table);
             transform.finalizeSink(insertExecutor, op, fragment, dataSink, 
physicalSink);
             
insertExecutor.getCoordinator().setTxnId(insertExecutor.getTxnId());
-            stmtExecutor.setCoord(insertExecutor.getCoordinator());
         } catch (Throwable e) {
             // the abortTxn in onFail need to acquire table write lock
             insertExecutor.onFail(e);
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorExecuteAction.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorExecuteAction.java
index f69170d638f..89e0ea8654e 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorExecuteAction.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorExecuteAction.java
@@ -161,7 +161,7 @@ public class ConnectorExecuteAction implements 
ExecuteAction {
                     : null;
             ConnectorRewriteDriver driver = new 
ConnectorRewriteDriver(ConnectContext.get(), table, catalog,
                     metadata, procedureOps, session, tableHandle, actionType, 
properties, partitionNames,
-                    loweredWhere);
+                    loweredWhere, ConnectContext.get() == null ? null : 
ConnectContext.get().getExecutor());
             try {
                 ConnectorProcedureResult result = driver.run();
                 return wrapResult(result);
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteDriver.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteDriver.java
index ee61e2b4cfd..996a1183681 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteDriver.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteDriver.java
@@ -18,6 +18,7 @@
 package org.apache.doris.nereids.trees.plans.commands.execute;
 
 import org.apache.doris.catalog.Env;
+import org.apache.doris.common.Status;
 import org.apache.doris.common.UserException;
 import org.apache.doris.connector.spi.ConnectorMetadata;
 import org.apache.doris.connector.spi.ConnectorSession;
@@ -33,20 +34,22 @@ import 
org.apache.doris.connector.spi.pushdown.ConnectorPredicate;
 import org.apache.doris.datasource.ExternalTable;
 import org.apache.doris.datasource.plugin.PluginDrivenExternalCatalog;
 import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.StmtExecutor;
 import org.apache.doris.scheduler.exception.JobException;
-import org.apache.doris.scheduler.executor.TransientTaskExecutor;
 import org.apache.doris.transaction.PluginDrivenTransactionManager;
 
 import com.google.common.collect.Lists;
 import org.apache.logging.log4j.LogManager;
 import org.apache.logging.log4j.Logger;
 
+import java.util.Collections;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
 
 /**
@@ -83,6 +86,14 @@ public class ConnectorRewriteDriver {
     // The engine-lowered WHERE restricting which files to rewrite, or null 
when there is no WHERE. Passed
     // straight through to the connector's planRewrite (the connector scopes 
the rewrite to the matching files).
     private final ConnectorPredicate whereCondition;
+    // The outer statement executor, or null in tests/standalone use. 
Cancellation of the outer statement is
+    // handed to this driver so the distributed rewrite does not commit after 
the statement is terminal.
+    private final StmtExecutor owner;
+    // Serializes the outer cancellation handoff with the 
source-registration/final-commit decision.
+    private final Object cancelLock = new Object();
+    private volatile Status outerCancelReason;
+    private volatile List<ConnectorRewriteGroupTask> submittedGroups = 
Collections.emptyList();
+    private final AtomicBoolean cancelHandoffInstalled = new 
AtomicBoolean(false);
 
     /**
      * Builds a driver bound to one {@code ALTER TABLE ... EXECUTE 
rewrite_data_files} invocation; all of
@@ -91,7 +102,7 @@ public class ConnectorRewriteDriver {
     public ConnectorRewriteDriver(ConnectContext ctx, ExternalTable table, 
PluginDrivenExternalCatalog catalog,
             ConnectorMetadata metadata, ConnectorProcedureOps procedureOps, 
ConnectorSession session,
             ConnectorTableHandle tableHandle, String procedureName, 
Map<String, String> properties,
-            List<String> partitionNames, ConnectorPredicate whereCondition) {
+            List<String> partitionNames, ConnectorPredicate whereCondition, 
StmtExecutor owner) {
         this.ctx = ctx;
         this.table = table;
         this.catalog = catalog;
@@ -103,6 +114,79 @@ public class ConnectorRewriteDriver {
         this.properties = properties;
         this.partitionNames = partitionNames;
         this.whereCondition = whereCondition;
+        this.owner = owner;
+    }
+
+    /**
+     * Installs the sticky outer-executor cancellation handoff: pick up a 
cancellation that landed before the
+     * driver existed, then register for the ones that arrive while the 
rewrite runs.
+     */
+    private void installCancelHandoff() {
+        StmtExecutor outer = this.owner;
+        if (outer == null) {
+            return;
+        }
+        Status prior = outer.getPendingCancelReason();
+        if (prior != null && !prior.ok()) {
+            synchronized (cancelLock) {
+                if (outerCancelReason == null) {
+                    outerCancelReason = prior;
+                }
+            }
+        }
+        outer.setCancelDelegate(this::onOuterCancel);
+        cancelHandoffInstalled.set(true);
+    }
+
+    private void clearCancelHandoff() {
+        if (cancelHandoffInstalled.compareAndSet(true, false) && owner != 
null) {
+            owner.clearCancelDelegate();
+        }
+    }
+
+    /**
+     * Runs on the cancelling thread. Records the sticky reason and stops the 
live groups; the owner thread
+     * drains them (with a shared budget) before it rolls the shared 
transaction back.
+     */
+    private void onOuterCancel(Status reason) {
+        synchronized (cancelLock) {
+            if (outerCancelReason == null) {
+                outerCancelReason = reason;
+            }
+        }
+        for (ConnectorRewriteGroupTask task : submittedGroups) {
+            try {
+                task.cancel();
+            } catch (Exception e) {
+                LOG.warn("Failed to cancel rewrite task {}: {}", task.getId(), 
e.getMessage());
+            }
+        }
+    }
+
+    /**
+     * Sticky: picks up the owner's first terminal reason the first time it is 
observed. Polled by the wait
+     * loop and re-checked at the register/commit decision, so it also covers 
a cancellation whose delegate
+     * notification was missed.
+     */
+    private boolean isOuterCancelled() {
+        StmtExecutor outer = this.owner;
+        if (outer != null && outerCancelReason == null) {
+            Status current = outer.getPendingCancelReason();
+            if (current != null && !current.ok()) {
+                synchronized (cancelLock) {
+                    if (outerCancelReason == null) {
+                        outerCancelReason = current;
+                    }
+                }
+            }
+        }
+        return outerCancelReason != null;
+    }
+
+    private UserException cancelledException() {
+        Status reason = outerCancelReason;
+        return new UserException("Rewrite is cancelled: "
+                + (reason == null ? "statement terminated" : 
reason.getErrorMsg()));
     }
 
     /**
@@ -146,31 +230,49 @@ public class ConnectorRewriteDriver {
         }
         RewriteCapableTransaction rewriteTx = (RewriteCapableTransaction) 
connectorTx;
 
+        installCancelHandoff();
         try {
-            // STEP 2: run one INSERT-SELECT per group concurrently, all 
sharing the transaction.
-            runGroups(groups, txnId, connectorTx);
-
-            // STEP 3: register the UNION of every group's source data files 
in a SINGLE call. The connector
-            // re-derives them from the table at the pinned OCC snapshot with 
ONE planFiles() scan; the former
-            // per-group loop repeated that full-table scan once per group (G 
groups = G+1 scans). Ordering is
-            // unchanged — still AFTER the groups ran (so the first group's 
write loaded the table + pinned the
-            // OCC snapshot that the connector re-derives against) and BEFORE 
commit (which consumes the
-            // registered files in the RewriteFiles op). Every per-group call 
scanned the SAME pinned snapshot, so
-            // one union scan is equivalent; the connector's registration 
accumulates and dedups by path, and the
-            // planner emits path-DISJOINT groups (iceberg: planFiles() yields 
one task per data file, bin-packed
-            // into disjoint groups), so the union reconstructs exactly the 
per-group calls' accumulated file set.
-            rewriteTx.registerRewriteSourceFiles(unionSourceFilePaths(groups));
-        } catch (Exception e) {
-            txnManager.rollback(txnId);
-            if (e instanceof UserException) {
-                throw (UserException) e;
+            try {
+                // STEP 2: run one INSERT-SELECT per group concurrently, all 
sharing the transaction.
+                runGroups(groups, txnId, connectorTx);
+
+                // STEP 3: register the UNION of every group's source data 
files in a SINGLE call. The connector
+                // re-derives them from the table at the pinned OCC snapshot 
with ONE planFiles() scan; the former
+                // per-group loop repeated that full-table scan once per group 
(G groups = G+1 scans). Ordering is
+                // unchanged — still AFTER the groups ran (so the first 
group's write loaded the table + pinned the
+                // OCC snapshot that the connector re-derives against) and 
BEFORE commit (which consumes the
+                // registered files in the RewriteFiles op). Every per-group 
call scanned the SAME pinned snapshot, so
+                // one union scan is equivalent; the connector's registration 
accumulates and dedups by path, and the
+                // planner emits path-DISJOINT groups (iceberg: planFiles() 
yields one task per data file, bin-packed
+                // into disjoint groups), so the union reconstructs exactly 
the per-group calls' accumulated file set.
+                synchronized (cancelLock) {
+                    if (isOuterCancelled()) {
+                        throw cancelledException();
+                    }
+                    
rewriteTx.registerRewriteSourceFiles(unionSourceFilePaths(groups));
+                }
+            } catch (Exception e) {
+                txnManager.rollback(txnId);
+                if (e instanceof UserException) {
+                    throw (UserException) e;
+                }
+                throw new UserException("Failed to rewrite data files: " + 
e.getMessage(), e);
             }
-            throw new UserException("Failed to rewrite data files: " + 
e.getMessage(), e);
-        }
 
-        // STEP 4: commit once. The manager deregisters the transaction on 
both success and failure, so a
-        // failed commit needs no rollback (it would find nothing) — surface 
it directly.
-        txnManager.commit(txnId);
+            // STEP 4: commit once, serialized on cancelLock with the outer 
cancellation handoff. Cancellation
+            // that wins the lock first rolls back instead of committing; 
cancellation that arrives while the
+            // critical section runs is linearized after the commit. The 
manager deregisters the transaction on
+            // both success and failure, so a failed commit needs no rollback 
— surface it directly.
+            synchronized (cancelLock) {
+                if (isOuterCancelled()) {
+                    txnManager.rollback(txnId);
+                    throw cancelledException();
+                }
+                txnManager.commit(txnId);
+            }
+        } finally {
+            clearCancelHandoff();
+        }
 
         // The rewrite is committed. Persist follower replay identity and 
refresh leader caches before the
         // post-commit statistics and result construction below, which can 
fail independently of the mutation.
@@ -224,27 +326,85 @@ public class ConnectorRewriteDriver {
             tasks.add(task);
         }
 
+        List<ConnectorRewriteGroupTask> submitted = Lists.newArrayList();
         try {
-            for (TransientTaskExecutor task : tasks) {
+            for (ConnectorRewriteGroupTask task : tasks) {
                 
Env.getCurrentEnv().getTransientTaskManager().addMemoryTask(task);
+                submitted.add(task);
             }
         } catch (JobException e) {
+            // Groups submitted before the failing call already bind the 
shared transaction; drain them so the
+            // caller never rolls that transaction back while a live group 
still reports into it.
+            drain(submitted, drainBudgetNanos());
             throw new UserException("Failed to submit rewrite tasks: " + 
e.getMessage(), e);
         }
+        submittedGroups = submitted;
 
-        int maxWaitTime = ctx.getSessionVariable().getInsertTimeoutS();
-        try {
-            boolean completed = collector.await(maxWaitTime, TimeUnit.SECONDS);
-            if (!completed) {
-                throw new UserException("Rewrite tasks did not complete within 
timeout");
+        long maxWaitTime = ctx.getSessionVariable().getInsertTimeoutS();
+        long deadlineNanos = System.nanoTime() + 
TimeUnit.SECONDS.toNanos(maxWaitTime);
+        boolean completed = false;
+        while (!completed) {
+            // A TIMEOUT/KILL on the outer statement only reaches this driver 
through the handoff; poll it so
+            // live groups are stopped instead of finishing and committing 
after the statement is terminal.
+            if (isOuterCancelled()) {
+                drain(submitted, drainBudgetNanos());
+                throw cancelledException();
             }
-            if (collector.getFirstError() != null) {
-                throw new UserException("Some rewrite tasks failed: " + 
collector.getFirstError().getMessage(),
-                        collector.getFirstError());
+            long remainingMs = TimeUnit.NANOSECONDS.toMillis(deadlineNanos - 
System.nanoTime());
+            if (remainingMs <= 0) {
+                completed = collector.isDone();
+                break;
+            }
+            try {
+                completed = collector.await(Math.min(200L, remainingMs), 
TimeUnit.MILLISECONDS);
+            } catch (InterruptedException e) {
+                // The interrupt flag was cleared by the exception, so the 
drain below can still wait.
+                drain(submitted, drainBudgetNanos());
+                Thread.currentThread().interrupt();
+                throw new UserException("Wait for rewrite tasks completion was 
interrupted", e);
+            }
+        }
+        if (!completed) {
+            // The owner gave up waiting: stop every live group and wait for 
its terminal callback so no group
+            // is still reporting into the shared transaction when the caller 
rolls it back.
+            drain(submitted, drainBudgetNanos());
+            throw new UserException("Rewrite tasks did not complete within 
timeout");
+        }
+        if (collector.getFirstError() != null) {
+            throw new UserException("Some rewrite tasks failed: " + 
collector.getFirstError().getMessage(),
+                    collector.getFirstError());
+        }
+    }
+
+    private long drainBudgetNanos() {
+        return TimeUnit.SECONDS.toNanos(Math.max(1, 
ctx.getSessionVariable().getInsertTimeoutS()));
+    }
+
+    /**
+     * Cancels every submitted group and waits (bounded) for its terminal 
callback within ONE shared budget,
+     * so a shared transaction is never rolled back while a live group still 
has commit data flowing into it,
+     * and G never-terminal groups cannot multiply the drain deadline into G * 
insert_timeout.
+     */
+    private void drain(List<ConnectorRewriteGroupTask> submitted, long 
budgetNanos) {
+        for (ConnectorRewriteGroupTask task : submitted) {
+            try {
+                task.cancel();
+            } catch (Exception e) {
+                LOG.warn("Failed to cancel rewrite task {}: {}", task.getId(), 
e.getMessage());
+            }
+        }
+        long deadline = System.nanoTime() + budgetNanos;
+        for (ConnectorRewriteGroupTask task : submitted) {
+            long remaining = deadline - System.nanoTime();
+            if (remaining <= 0) {
+                return;
+            }
+            try {
+                task.awaitTerminal(remaining, TimeUnit.NANOSECONDS);
+            } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+                return;
             }
-        } catch (InterruptedException e) {
-            Thread.currentThread().interrupt();
-            throw new UserException("Wait for rewrite tasks completion was 
interrupted", e);
         }
     }
 
@@ -299,6 +459,10 @@ public class ConnectorRewriteDriver {
             return completionLatch.await(timeout, unit);
         }
 
+        public boolean isDone() {
+            return completionLatch.getCount() == 0;
+        }
+
         public Exception getFirstError() {
             return firstError;
         }
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteGroupTask.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteGroupTask.java
index 56f61fe2d4c..ead91a50b3e 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteGroupTask.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteGroupTask.java
@@ -34,6 +34,7 @@ import 
org.apache.doris.nereids.trees.plans.commands.insert.ConnectorRewriteExec
 import 
org.apache.doris.nereids.trees.plans.commands.insert.RewriteTableCommand;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.qe.OriginStatement;
+import org.apache.doris.qe.QueryState.MysqlStateType;
 import org.apache.doris.qe.StmtExecutor;
 import org.apache.doris.qe.VariableMgr;
 import org.apache.doris.scheduler.exception.JobException;
@@ -50,6 +51,8 @@ import java.util.ArrayList;
 import java.util.List;
 import java.util.Optional;
 import java.util.UUID;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 
 /**
@@ -74,9 +77,14 @@ public class ConnectorRewriteGroupTask implements 
TransientTaskExecutor {
     private final Long taskId;
     private final AtomicBoolean isCanceled;
     private final AtomicBoolean isFinished;
+    // Set once the scheduler has actually invoked execute(); a cancel that 
wins before this is conclusively
+    // terminal (execute() will throw at its entry check without touching the 
shared transaction).
+    private final AtomicBoolean started = new AtomicBoolean(false);
+    // Counted down when execute() reaches a terminal state, so a failure 
owner can drain live groups.
+    private final CountDownLatch terminalLatch = new CountDownLatch(1);
 
-    // for canceling the task
-    private StmtExecutor stmtExecutor;
+    // for canceling the task; volatile so the cancel/publication handoff in 
executeGroup() is safe
+    private volatile StmtExecutor stmtExecutor;
 
     /**
      * Builds a task for one bin-packed rewrite group, sharing {@code 
transactionId} /
@@ -106,14 +114,18 @@ public class ConnectorRewriteGroupTask implements 
TransientTaskExecutor {
 
     @Override
     public void execute() throws JobException {
-        if (isCanceled.get()) {
-            throw new JobException("Rewrite task has been canceled, task id: " 
+ taskId);
-        }
         if (isFinished.get()) {
             return;
         }
+        started.set(true);
 
         try {
+            // A queued task may have been canceled before the scheduler ran 
it. Report that through the
+            // same failure path as a running task, otherwise the collector 
would wait for the full insert
+            // timeout because terminating early used to skip the callback 
entirely.
+            if (isCanceled.get()) {
+                throw new JobException("Rewrite task has been canceled, task 
id: " + taskId);
+            }
             // Step 1: Build a fresh ConnectContext for this group and stash 
the per-group scan scope + the
             // shared connector transaction (read back during planning by 
pinRewriteFileScope / finalizeSink).
             ConnectContext taskConnectContext = buildConnectContext();
@@ -140,9 +152,18 @@ public class ConnectorRewriteGroupTask implements 
TransientTaskExecutor {
             throw new JobException("Rewrite group execution failed: " + 
e.getMessage(), e);
         } finally {
             isFinished.set(true);
+            terminalLatch.countDown();
         }
     }
 
+    /**
+     * Waits (bounded) until this task has reached a terminal state, so a 
caller can drain live groups
+     * before it rolls their shared transaction back.
+     */
+    public boolean awaitTerminal(long timeout, TimeUnit unit) throws 
InterruptedException {
+        return terminalLatch.await(timeout, unit);
+    }
+
     @Override
     public void cancel() throws JobException {
         if (isFinished.get()) {
@@ -152,17 +173,28 @@ public class ConnectorRewriteGroupTask implements 
TransientTaskExecutor {
         if (stmtExecutor != null) {
             stmtExecutor.cancel(new Status(TStatusCode.CANCELLED, "rewrite 
task cancelled"));
         }
+        if (!started.get()) {
+            // A queued task fails at execute() entry without writing 
anything, so it is already conclusively
+            // terminal: make it visible to a drain now instead of burning a 
full timeout on it.
+            terminalLatch.countDown();
+        }
         LOG.info("[Connector Rewrite Task] taskId: {} cancelled", taskId);
     }
 
     private void executeGroup(ConnectContext taskConnectContext,
             RewriteTableCommand taskLogicalPlan,
             StatementBase taskParsedStmt) throws Exception {
-        stmtExecutor = new StmtExecutor(taskConnectContext, taskParsedStmt);
+        StmtExecutor taskStmtExecutor = new StmtExecutor(taskConnectContext, 
taskParsedStmt);
+        // Publish under the cancel handoff: cancel() sets isCanceled before 
reading this field, so assigning
+        // first and re-checking here guarantees at least one side observes 
the other.
+        stmtExecutor = taskStmtExecutor;
+        if (isCanceled.get()) {
+            throw new JobException("Rewrite task has been canceled, task id: " 
+ taskId);
+        }
 
         // initPlan finalizes the sink (ConnectorRewriteExecutor.finalizeSink 
binds the shared transaction
         // onto the sink session BEFORE planWrite reads it).
-        AbstractInsertExecutor insertExecutor = 
taskLogicalPlan.initPlan(taskConnectContext, stmtExecutor);
+        AbstractInsertExecutor insertExecutor = 
taskLogicalPlan.initPlan(taskConnectContext, taskStmtExecutor);
         Preconditions.checkState(insertExecutor instanceof 
ConnectorRewriteExecutor,
                 "Expected ConnectorRewriteExecutor, got: " + 
insertExecutor.getClass());
 
@@ -170,7 +202,15 @@ public class ConnectorRewriteGroupTask implements 
TransientTaskExecutor {
         // accumulate on the one rewrite transaction (mirrors legacy 
RewriteGroupTask).
         insertExecutor.getCoordinator().setTxnId(transactionId);
 
-        insertExecutor.executeSingleInsert(stmtExecutor);
+        insertExecutor.executeSingleInsert(taskStmtExecutor);
+
+        // executeSingleInsert turns a retained coordinator cancellation into 
QueryState.ERR and returns
+        // normally. Surface it as a failure so the collector rolls back the 
shared transaction instead of
+        // reporting the group complete and committing a partial rewrite.
+        if (taskConnectContext.getState().getStateType() == 
MysqlStateType.ERR) {
+            throw new JobException("Rewrite group failed: "
+                    + taskConnectContext.getState().getErrorMessage());
+        }
     }
 
     private RewriteTableCommand buildRewriteLogicalPlan() {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/AbstractInsertExecutor.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/AbstractInsertExecutor.java
index a8946f2c93d..bc501b5816b 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/AbstractInsertExecutor.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/AbstractInsertExecutor.java
@@ -272,6 +272,16 @@ public abstract class AbstractInsertExecutor {
      */
     public void executeSingleInsert(StmtExecutor executor) throws Exception {
         try {
+            // Every statement-owned insert coordinator is published at the 
common execution boundary.
+            // Cancellation retained during planning is replayed before any 
executor-specific setup or dispatch.
+            executor.setCoord(coordinator);
+            // Publication synchronously replays any cancellation retained 
during planning. Fence on that
+            // terminal status before executor-specific setup runs, so a later 
setup failure cannot mask the
+            // original timeout/cancel reason.
+            Status execStatus = coordinator.getExecStatus();
+            if (!execStatus.ok()) {
+                throw new UserException(execStatus.getErrorMsg());
+            }
             // Pre-execution work may register external resources, so it must 
share the transaction cleanup scope.
             beforeExec();
             executor.updateProfile(false);
@@ -285,6 +295,14 @@ public abstract class AbstractInsertExecutor {
             for (InsertExecutorListener listener : listeners) {
                 listener.beforeComplete(this, executor, jobId);
             }
+            // Every transaction-owning executor commits inside onComplete(). 
Re-fence the first terminal
+            // status here so a cancellation/TIMEOUT that landed after 
execImpl()'s last status read cannot
+            // be committed by a path such as row-level UPDATE/DELETE/MERGE, 
which has no command-level
+            // cancellation listener of its own.
+            Status preCommitStatus = coordinator.getExecStatus();
+            if (!preCommitStatus.ok()) {
+                throw new UserException(preCommitStatus.getErrorMsg());
+            }
             onComplete();
             for (InsertExecutorListener listener : listeners) {
                 try {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/DictionaryInsertExecutor.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/DictionaryInsertExecutor.java
index f045a4c19b4..243031a30fe 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/DictionaryInsertExecutor.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/DictionaryInsertExecutor.java
@@ -18,7 +18,7 @@
 package org.apache.doris.nereids.trees.plans.commands.insert;
 
 import org.apache.doris.catalog.DatabaseIf;
-import org.apache.doris.common.DdlException;
+import org.apache.doris.common.ErrorCode;
 import org.apache.doris.common.UserException;
 import org.apache.doris.common.util.DebugUtil;
 import org.apache.doris.dictionary.Dictionary;
@@ -75,11 +75,9 @@ public class DictionaryInsertExecutor extends 
AbstractInsertExecutor {
 
     @Override
     protected void onFail(Throwable t) {
-        // must from AbstractInsertExecutor so got DdlException
-        DdlException ddlException = (DdlException) t;
-        errMsg = t.getMessage() == null ? "unknown reason" : 
ddlException.getMessage();
+        errMsg = t.getMessage() == null ? "unknown reason" : t.getMessage();
         String queryId = DebugUtil.printId(ctx.queryId());
-        LOG.warn("dictionary insert [{}] with query id {} failed", labelName, 
queryId, ddlException);
+        LOG.warn("dictionary insert [{}] with query id {} failed", labelName, 
queryId, t);
 
         String finalErrorMsg = InsertUtils.getFinalErrorMsg(
                 errMsg,
@@ -87,7 +85,9 @@ public class DictionaryInsertExecutor extends 
AbstractInsertExecutor {
                 coordinator.getTrackingUrl()
         );
         // we should set the context to make the caller know the command failed
-        ctx.getState().setError(ddlException.getMysqlErrorCode(), 
finalErrorMsg);
+        ErrorCode errorCode = t instanceof UserException
+                ? ((UserException) t).getMysqlErrorCode() : 
ErrorCode.ERR_UNKNOWN_ERROR;
+        ctx.getState().setError(errorCode, finalErrorMsg);
     }
 
     @Override
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertIntoTableCommand.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertIntoTableCommand.java
index 0882a19a7a0..04fea4794ce 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertIntoTableCommand.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertIntoTableCommand.java
@@ -419,9 +419,8 @@ public class InsertIntoTableCommand extends Command
             }
             stmtExecutor.setProfileType(ProfileType.LOAD);
             // We exposed @StmtExecutor#cancel as a unified entry point for 
statement interruption,
-            // so we need to set this here
+            // so executeSingleInsert publishes the coordinator at its common 
execution boundary.
             
insertExecutor.getCoordinator().setTxnId(insertExecutor.getTxnId());
-            stmtExecutor.setCoord(insertExecutor.getCoordinator());
             if (needsExternalDmlAuditBarrier(insertExecutor)) {
                 // The resolved executor is the invariant that distinguishes 
an external write;
                 // logical sink roots are rewritten and are not a stable audit 
classification.
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
index 0acfe3d7635..b78e6aeab4a 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
@@ -767,6 +767,10 @@ public class Coordinator implements CoordInterface {
     // A call to Exec() must precede all other member function calls.
     @Override
     public void exec() throws Exception {
+        Status status = getQueryStatus();
+        if (!status.ok()) {
+            throw new UserException(status.getErrorMsg());
+        }
         // LoadTask does not have context, not controlled by queue now
         if (context != null) {
             if (Config.enable_workload_group) {
@@ -785,8 +789,12 @@ public class Coordinator implements CoordInterface {
                     // AllBackendComputeGroup may assocatiate with multiple 
workload groups
                     queryQueue = wgs.get(0).getQueryQueue();
                     queueToken = 
queryQueue.getToken(context.getSessionVariable().wgQuerySlotCount);
-                    queueToken.get(DebugUtil.printId(queryId),
-                            this.queryOptions.getExecutionTimeout() * 1000);
+                    try {
+                        queueToken.get(DebugUtil.printId(queryId),
+                                this.queryOptions.getExecutionTimeout() * 
1000);
+                    } catch (UserException e) {
+                        throw preferTerminalReason(e);
+                    }
                 }
                 context.setWorkloadGroupName(wgs.get(0).getName());
             } else {
@@ -922,6 +930,9 @@ public class Coordinator implements CoordInterface {
     protected void sendPipelineCtx() throws Exception {
         lock();
         try {
+            if (!queryStatus.ok()) {
+                throw new UserException(queryStatus.getErrorMsg());
+            }
             Multiset<TNetworkAddress> hostCounter = HashMultiset.create();
             for (FragmentExecParams params : fragmentExecParamsMap.values()) {
                 for (FInstanceExecParam fi : params.instanceExecParams) {
@@ -1409,12 +1420,6 @@ public class Coordinator implements CoordInterface {
 
     @Override
     public void cancel(Status cancelReason) {
-        if (queueToken != null) {
-            queueToken.cancel();
-        }
-        for (ScanNode scanNode : scanNodes) {
-            scanNode.stop();
-        }
         if (cancelReason.ok()) {
             throw new RuntimeException("Should use correct cancel reason, but 
it is "
                     + cancelReason.toString());
@@ -1439,6 +1444,21 @@ public class Coordinator implements CoordInterface {
         } finally {
             unlock();
         }
+        if (queueToken != null) {
+            queueToken.cancel();
+        }
+        // Scan cleanup is best-effort and must never escape: the terminal 
status and interval cancellation
+        // above are already published, and a throwing scan would otherwise 
skip the remaining scans (and the
+        // caller's coordinator close), masking the retained reason. A scan 
whose first stop() threw still
+        // removes its own sources on the close-time retry because 
SplitAssignment.stop() is idempotent.
+        for (ScanNode scanNode : scanNodes) {
+            try {
+                scanNode.stop();
+            } catch (Throwable t) {
+                LOG.error("error happens when scannode stop during cancel, 
query id: {}",
+                        DebugUtil.printId(queryId), t);
+            }
+        }
     }
 
     public boolean isQueryCancelled() {
@@ -1450,6 +1470,28 @@ public class Coordinator implements CoordInterface {
         }
     }
 
+    protected Status getQueryStatus() {
+        lock();
+        try {
+            return new Status(queryStatus);
+        } finally {
+            unlock();
+        }
+    }
+
+    /**
+     * A queue wait can be unblocked by {@code cancel()}, which records the 
real TIMEOUT/KILL reason on the
+     * coordinator before cancelling the token, while {@link QueueToken#get} 
can only report a generic
+     * "query is cancelled". Prefer the retained terminal reason when one is 
present.
+     */
+    protected UserException preferTerminalReason(UserException queueFailure) {
+        Status current = getQueryStatus();
+        if (!current.ok()) {
+            return new UserException(current.getErrorMsg());
+        }
+        return queueFailure;
+    }
+
     private void cancelLatch() {
         if (instancesDoneLatch != null) {
             instancesDoneLatch.countDownToZero(new Status());
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java
index 8c53bca65f7..b7715396048 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java
@@ -154,9 +154,17 @@ public class NereidsCoordinator extends Coordinator {
 
     @Override
     public void exec() throws Exception {
+        Status status = getQueryStatus();
+        if (!status.ok()) {
+            throw new UserException(status.getErrorMsg());
+        }
         enqueue(coordinatorContext.connectContext);
 
         processTopSink(coordinatorContext, 
coordinatorContext.topDistributedPlan);
+        status = getQueryStatus();
+        if (!status.ok()) {
+            throw new UserException(status.getErrorMsg());
+        }
 
         QeProcessorImpl.INSTANCE.registerInstances(coordinatorContext.queryId, 
coordinatorContext.instanceNum.get());
 
@@ -173,31 +181,47 @@ public class NereidsCoordinator extends Coordinator {
 
     @Override
     public void cancel(Status cancelReason) {
-        coordinatorContext.getQueueToken().ifPresent(QueueToken::cancel);
-
-        for (ScanNode scanNode : coordinatorContext.scanNodes) {
-            scanNode.stop();
-        }
-
         if (cancelReason.ok()) {
             throw new RuntimeException("Should use correct cancel reason, but 
it is " + cancelReason);
         }
 
         TUniqueId queryId = coordinatorContext.queryId;
-        Status originQueryStatus = 
coordinatorContext.updateStatusIfOk(cancelReason);
-        if (!originQueryStatus.ok()) {
-            if (LOG.isDebugEnabled()) {
-                // Print an error stack here to know why send cancel again.
-                LOG.warn("Query {} already in abnormal status {}, but received 
cancel again,"
-                                + "so that send cancel to BE again",
-                        DebugUtil.printId(queryId), 
originQueryStatus.toString(),
-                        new Exception("cancel failed"));
+        try {
+            Status originQueryStatus = 
coordinatorContext.updateStatusIfOk(cancelReason);
+            if (!originQueryStatus.ok()) {
+                if (LOG.isDebugEnabled()) {
+                    // Print an error stack here to know why send cancel again.
+                    LOG.warn("Query {} already in abnormal status {}, but 
received cancel again,"
+                                    + "so that send cancel to BE again",
+                            DebugUtil.printId(queryId), 
originQueryStatus.toString(),
+                            new Exception("cancel failed"));
+                }
+            } else {
+                LOG.warn("Cancel execution of query {}, this is a outside 
invoke, cancelReason {}",
+                        DebugUtil.printId(queryId), cancelReason);
+            }
+        } finally {
+            // Publishing the status above can itself cancel a partially 
initialized processor. Start the
+            // non-throwing cleanup scope before that publication so the queue 
token, the scan nodes, and
+            // the final internal cancel are never skipped, even if the 
publication was only half wired.
+            try {
+                
coordinatorContext.getQueueToken().ifPresent(QueueToken::cancel);
+                // Scan cleanup is best-effort and must never escape: a 
throwing scan would otherwise skip
+                // the remaining scans (and the caller's coordinator close), 
masking the retained reason. A
+                // scan whose first stop() threw still removes its own sources 
on the close-time retry
+                // because SplitAssignment.stop() is idempotent.
+                for (ScanNode scanNode : coordinatorContext.scanNodes) {
+                    try {
+                        scanNode.stop();
+                    } catch (Throwable t) {
+                        LOG.error("error happens when scannode stop during 
cancel, query id: {}",
+                                DebugUtil.printId(queryId), t);
+                    }
+                }
+            } finally {
+                cancelInternal(cancelReason);
             }
-        } else {
-            LOG.warn("Cancel execution of query {}, this is a outside invoke, 
cancelReason {}",
-                    DebugUtil.printId(queryId), cancelReason);
         }
-        cancelInternal(cancelReason);
     }
 
     public QueryProcessor asQueryProcessor() {
@@ -227,6 +251,11 @@ public class NereidsCoordinator extends Coordinator {
         return coordinatorContext.readCloneStatus().isCancelled();
     }
 
+    @Override
+    protected Status getQueryStatus() {
+        return coordinatorContext.readCloneStatus();
+    }
+
     @Override
     public RowBatch getNext() throws Exception {
         return coordinatorContext.asQueryProcessor().getNext();
@@ -576,7 +605,11 @@ public class NereidsCoordinator extends Coordinator {
                     QueueToken queueToken = 
queryQueue.getToken(context.getSessionVariable().wgQuerySlotCount);
                     int queryTimeout = 
coordinatorContext.queryOptions.getExecutionTimeout() * 1000;
                     coordinatorContext.setQueueInfo(queryQueue, queueToken);
-                    
queueToken.get(DebugUtil.printId(coordinatorContext.queryId), queryTimeout);
+                    try {
+                        
queueToken.get(DebugUtil.printId(coordinatorContext.queryId), queryTimeout);
+                    } catch (UserException e) {
+                        throw preferTerminalReason(e);
+                    }
                 }
                 context.setWorkloadGroupName(wgs.get(0).getName());
             } else {
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
index 8227950bcd2..9a8ae03257f 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
@@ -148,7 +148,6 @@ import com.google.common.collect.Lists;
 import com.google.common.collect.Maps;
 import com.google.common.collect.Sets;
 import com.google.protobuf.ByteString;
-import lombok.Setter;
 import org.apache.commons.codec.digest.DigestUtils;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.logging.log4j.LogManager;
@@ -166,6 +165,7 @@ import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.Future;
 import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.Consumer;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
@@ -200,8 +200,11 @@ public class StmtExecutor {
     // never be finished. Null until decided.
     private volatile Boolean profileEnabled;
 
-    @Setter
     private volatile Coordinator coord = null;
+    // A statement can be cancelled while it is still planning and has no 
coordinator yet.
+    // Keep this state scoped to the coordinator publication handoff: other 
execution targets
+    // retain their existing cancellation contracts.
+    private final AtomicReference<Status> pendingCoordinatorCancelReason = new 
AtomicReference<>();
     private volatile Coordinator externalDmlAuditCoordinator = null;
     // Arrow Flight SQL: when true, this query's coordinator is kept alive 
past GetFlightInfo and
     // is finalized later by ConnectContext (see #62259), so the eager close 
in executeAndSendResult
@@ -1444,6 +1447,8 @@ public class StmtExecutor {
     }
 
     public void cancel(Status cancelReason, boolean needWaitCancelComplete) {
+        pendingCoordinatorCancelReason.compareAndSet(null, cancelReason);
+        Status coordinatorCancelReason = pendingCoordinatorCancelReason.get();
         Consumer<Status> delegate = cancelDelegate;
         if (delegate != null) {
             delegate.accept(cancelReason);
@@ -1463,7 +1468,7 @@ public class StmtExecutor {
         }
         Coordinator coordRef = coord;
         if (coordRef != null) {
-            coordRef.cancel(cancelReason);
+            coordRef.cancel(coordinatorCancelReason);
         }
         if (mysqlLoadId != null) {
             
Env.getCurrentEnv().getLoadManager().getMysqlLoadManager().cancelMySqlLoad(mysqlLoadId);
@@ -1474,6 +1479,23 @@ public class StmtExecutor {
         }
     }
 
+    public void setCoord(Coordinator coordinator) {
+        coord = coordinator;
+        Status cancelReason = pendingCoordinatorCancelReason.get();
+        if (coordinator != null && cancelReason != null) {
+            coordinator.cancel(cancelReason);
+        }
+    }
+
+    /**
+     * The first terminal status delivered to this executor, or null when it 
has not been cancelled. Sticky:
+     * a later cancellation never replaces the first one. Used by owners (such 
as the distributed rewrite
+     * driver) that execute outside the coordinator publication handoff.
+     */
+    public Status getPendingCancelReason() {
+        return pendingCoordinatorCancelReason.get();
+    }
+
     public void cancel(Status cancelReason) {
         cancel(cancelReason, true);
     }
@@ -1660,15 +1682,15 @@ public class StmtExecutor {
                     
context.getSessionVariable().getMaxMsgSizeOfResultReceiver());
             context.getState().setIsQuery(true);
         } else if (planner instanceof NereidsPlanner && ((NereidsPlanner) 
planner).getDistributedPlans() != null) {
-            coord = new NereidsCoordinator(context,
-                    (NereidsPlanner) planner, 
context.getStatsErrorEstimator());
+            setCoord(new NereidsCoordinator(context,
+                    (NereidsPlanner) planner, 
context.getStatsErrorEstimator()));
             profile.addExecutionProfile(coord.getExecutionProfile());
             QeProcessorImpl.INSTANCE.registerQuery(context.queryId(),
                     new QueryInfo(context, originStmt.originStmt, coord));
             coordBase = coord;
         } else {
-            coord = EnvFactory.getInstance().createCoordinator(
-                    context, planner, context.getStatsErrorEstimator());
+            setCoord(EnvFactory.getInstance().createCoordinator(
+                    context, planner, context.getStatsErrorEstimator()));
             profile.addExecutionProfile(coord.getExecutionProfile());
             QeProcessorImpl.INSTANCE.registerQuery(context.queryId(),
                     new QueryInfo(context, originStmt.originStmt, coord));
@@ -1820,7 +1842,7 @@ public class StmtExecutor {
             LOG.warn(internalErrorSt.getErrorMsg());
             coordBase.cancel(internalErrorSt);
             // set to null so that the retry logic will generate a new 
coordinator
-            this.coord = null;
+            setCoord(null);
             throw e;
         } finally {
             // For deferred Arrow Flight queries the coordinator is closed 
later by ConnectContext
@@ -2137,8 +2159,8 @@ public class StmtExecutor {
             if (Config.enable_collect_internal_query_profile) {
                 context.getSessionVariable().enableProfile = true;
             }
-            coord = EnvFactory.getInstance().createCoordinator(context,
-                    planner, context.getStatsErrorEstimator());
+            setCoord(EnvFactory.getInstance().createCoordinator(context,
+                    planner, context.getStatsErrorEstimator()));
             profile.addExecutionProfile(coord.getExecutionProfile());
             try {
                 QeProcessorImpl.INSTANCE.registerQuery(context.queryId(),
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/LoadProcessor.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/LoadProcessor.java
index d4878ff99a3..8a314b21ccd 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/LoadProcessor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/LoadProcessor.java
@@ -114,7 +114,10 @@ public class LoadProcessor extends AbstractJobProcessor {
             for (MultiFragmentsPipelineTask fragmentsTask : 
executionTask.get().getChildrenTasks().values()) {
                 fragmentsTask.cancelExecute(cancelReason);
             }
-            latch.get().countDownToZero(new Status());
+            // setPipelineExecutionTask publishes this task before 
afterSetPipelineExecutionTask builds the
+            // latch. A cancel that crosses that publication must not throw 
from an empty latch, which would
+            // escape the coordinator's cleanup scope and skip the queue/scan 
teardown.
+            latch.ifPresent(l -> l.countDownToZero(new Status()));
         }
     }
 
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/PipelineExecutionTask.java
 
b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/PipelineExecutionTask.java
index 6ea7842562c..ef5a2df7079 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/PipelineExecutionTask.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/PipelineExecutionTask.java
@@ -97,6 +97,10 @@ public class PipelineExecutionTask extends 
AbstractRuntimeTask<BackendWorker, Mu
     @Override
     public void execute() throws Exception {
         coordinatorContext.withLock(() -> {
+            Status status = coordinatorContext.readCloneStatus();
+            if (!status.ok()) {
+                throw new UserException(status.getErrorMsg());
+            }
             dispatchedBackendIdsForAudit.clear();
             sendAndWaitPhaseOneRpc();
             if (coordinatorContext.twoPhaseExecution()) {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/scheduler/disruptor/TaskDisruptor.java
 
b/fe/fe-core/src/main/java/org/apache/doris/scheduler/disruptor/TaskDisruptor.java
index 9be30124c9e..434c6ecf8b3 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/scheduler/disruptor/TaskDisruptor.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/scheduler/disruptor/TaskDisruptor.java
@@ -122,8 +122,10 @@ public class TaskDisruptor implements Closeable {
      */
     public void tryPublishTask(Long taskId) throws JobException {
         if (isClosed) {
+            // Fail loudly: silently returning would leave the caller's 
already-registered task with no event
+            // to run or remove it.
             log.info("tryPublish failed, disruptor is closed, taskId: {}", 
taskId);
-            return;
+            throw new JobException("Disruptor is closed, cannot publish 
transient task: " + taskId);
         }
         // We reserve two slots in the ring buffer
         // to prevent it from becoming stuck due to competition between 
producers and consumers.
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/scheduler/manager/TransientTaskManager.java
 
b/fe/fe-core/src/main/java/org/apache/doris/scheduler/manager/TransientTaskManager.java
index feaf3365a35..3f86e0f9efd 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/scheduler/manager/TransientTaskManager.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/scheduler/manager/TransientTaskManager.java
@@ -56,8 +56,17 @@ public class TransientTaskManager {
 
     public Long addMemoryTask(TransientTaskExecutor executor) throws 
JobException {
         Long taskId = executor.getId();
+        // Registration must precede publication so the handler can resolve 
the task, but publication is
+        // fallible (ring buffer saturated or disruptor closed). TaskHandler 
only removes the task it actually
+        // runs, so a task whose event was never published must be 
unregistered here or it stays reachable for
+        // the FE lifetime.
         taskExecutorMap.put(taskId, executor);
-        disruptor.tryPublishTask(taskId);
+        try {
+            disruptor.tryPublishTask(taskId);
+        } catch (JobException e) {
+            taskExecutorMap.remove(taskId);
+            throw e;
+        }
         LOG.info("add memory task, taskId: {}", taskId);
         return taskId;
     }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/RowLevelDmlCommandTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/RowLevelDmlCommandTest.java
index f98da4b7b77..81ef3a894d0 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/RowLevelDmlCommandTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/RowLevelDmlCommandTest.java
@@ -90,7 +90,7 @@ public class RowLevelDmlCommandTest {
     }
 
     @Test
-    public void successPathWiresCoordinatorWithoutRollback() {
+    public void successPathPreparesCoordinatorWithoutRollback() {
         Coordinator coordinator = Mockito.mock(Coordinator.class);
         Mockito.when(executor.getCoordinator()).thenReturn(coordinator);
         Mockito.when(executor.getTxnId()).thenReturn(42L);
@@ -98,7 +98,6 @@ public class RowLevelDmlCommandTest {
         invoke();
 
         Mockito.verify(coordinator).setTxnId(42L);
-        Mockito.verify(stmtExecutor).setCoord(coordinator);
         Mockito.verify(executor, Mockito.never()).onFail(Mockito.any());
     }
 }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteDriverTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteDriverTest.java
index 75c5b6366f1..ca0f0ebe7e9 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteDriverTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteDriverTest.java
@@ -19,6 +19,7 @@ package org.apache.doris.nereids.trees.plans.commands.execute;
 
 import org.apache.doris.catalog.Env;
 import org.apache.doris.catalog.RefreshManager;
+import org.apache.doris.common.Status;
 import org.apache.doris.common.UserException;
 import org.apache.doris.connector.spi.ConnectorColumn;
 import org.apache.doris.connector.spi.ConnectorMetadata;
@@ -38,7 +39,10 @@ import org.apache.doris.datasource.ExternalTable;
 import org.apache.doris.datasource.plugin.PluginDrivenExternalCatalog;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.qe.SessionVariable;
+import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.scheduler.exception.JobException;
 import org.apache.doris.scheduler.manager.TransientTaskManager;
+import org.apache.doris.thrift.TStatusCode;
 import org.apache.doris.transaction.PluginDrivenTransactionManager;
 
 import com.google.common.collect.ImmutableSet;
@@ -52,7 +56,10 @@ import org.mockito.Mockito;
 
 import java.util.Arrays;
 import java.util.Collections;
+import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.atomic.AtomicReference;
 
 /**
@@ -84,7 +91,8 @@ public class ConnectorRewriteDriverTest {
                 "rewrite_data_files",
                 Collections.emptyMap(),
                 Collections.emptyList(),
-                where);
+                where,
+                null);
     }
 
     @Test
@@ -189,7 +197,7 @@ public class ConnectorRewriteDriverTest {
 
         ConnectorRewriteDriver driver = new ConnectorRewriteDriver(
                 context, table, catalog, metadata, procedureOps, session, 
tableHandle,
-                "rewrite_data_files", Collections.emptyMap(), 
Collections.emptyList(), null);
+                "rewrite_data_files", Collections.emptyMap(), 
Collections.emptyList(), null, null);
 
         try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
                 MockedConstruction<ConnectorRewriteGroupTask> taskConstruction 
= Mockito.mockConstruction(
@@ -215,6 +223,302 @@ public class ConnectorRewriteDriverTest {
         }
     }
 
+    @Test
+    public void groupFailureRollsBackTheSharedTransactionWithoutCommit() 
throws Exception {
+        // A failed group must reach the collector as a failure so the driver 
rolls the shared transaction
+        // back. If a group's cancellation is swallowed and reported as 
completed, the driver would register
+        // every group's sources and commit a partial rewrite. MUTATION: 
reporting onTaskCompleted for a
+        // cancelled/failed group is killed here.
+        ConnectorProcedureOps procedureOps = 
Mockito.mock(ConnectorProcedureOps.class);
+        ConnectorMetadata metadata = Mockito.mock(ConnectorMetadata.class);
+        ConnectorSession session = Mockito.mock(ConnectorSession.class);
+        ConnectorTableHandle tableHandle = 
Mockito.mock(ConnectorTableHandle.class);
+        ConnectorTransaction connectorTx = 
Mockito.mock(ConnectorTransaction.class,
+                
Mockito.withSettings().extraInterfaces(RewriteCapableTransaction.class));
+        PluginDrivenTransactionManager txnManager = 
Mockito.mock(PluginDrivenTransactionManager.class);
+        PluginDrivenExternalCatalog catalog = 
Mockito.mock(PluginDrivenExternalCatalog.class);
+        ExternalTable table = Mockito.mock(ExternalTable.class);
+        ConnectContext context = Mockito.mock(ConnectContext.class);
+        SessionVariable sessionVariable = Mockito.mock(SessionVariable.class);
+        ConnectorRewriteGroup first = new ConnectorRewriteGroup(
+                ImmutableSet.of("s3://bucket/table/a.parquet"), 1, 1024L, 0);
+        ConnectorRewriteGroup second = new ConnectorRewriteGroup(
+                ImmutableSet.of("s3://bucket/table/b.parquet"), 1, 1024L, 0);
+
+        Mockito.when(procedureOps.planRewrite(Mockito.any(), Mockito.any(), 
Mockito.any(), Mockito.any(),
+                Mockito.any(), Mockito.any())).thenReturn(Arrays.asList(first, 
second));
+        Mockito.when(metadata.beginTransaction(session, 
tableHandle)).thenReturn(connectorTx);
+        Mockito.when(catalog.getTransactionManager()).thenReturn(txnManager);
+        Mockito.when(txnManager.begin(connectorTx)).thenReturn(7L);
+        Mockito.when(context.getSessionVariable()).thenReturn(sessionVariable);
+        Mockito.when(sessionVariable.getInsertTimeoutS()).thenReturn(1);
+
+        Env env = Mockito.mock(Env.class);
+        TransientTaskManager transientTaskManager = 
Mockito.mock(TransientTaskManager.class);
+        
Mockito.when(env.getTransientTaskManager()).thenReturn(transientTaskManager);
+        AtomicInteger ids = new AtomicInteger(10);
+        Map<ConnectorRewriteGroupTask, 
ConnectorRewriteGroupTask.RewriteResultCallback> callbacks =
+                new ConcurrentHashMap<>();
+
+        ConnectorRewriteDriver driver = new ConnectorRewriteDriver(
+                context, table, catalog, metadata, procedureOps, session, 
tableHandle,
+                "rewrite_data_files", Collections.emptyMap(), 
Collections.emptyList(), null, null);
+
+        try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
+                MockedConstruction<ConnectorRewriteGroupTask> taskConstruction 
= Mockito.mockConstruction(
+                        ConnectorRewriteGroupTask.class, (task, 
constructionContext) -> {
+                            Mockito.when(task.getId()).thenReturn((long) 
ids.incrementAndGet());
+                            callbacks.put(task, 
(ConnectorRewriteGroupTask.RewriteResultCallback)
+                                    constructionContext.arguments().get(5));
+                        })) {
+            envStatic.when(Env::getCurrentEnv).thenReturn(env);
+            
Mockito.when(transientTaskManager.addMemoryTask(Mockito.any())).thenAnswer(invocation
 -> {
+                ConnectorRewriteGroupTask task = invocation.getArgument(0);
+                callbacks.get(task).onTaskFailed(task.getId(), new 
JobException("group failed"));
+                return task.getId();
+            });
+
+            UserException ex = Assertions.assertThrows(UserException.class, 
driver::run);
+            Assertions.assertTrue(ex.getMessage().contains("Some rewrite tasks 
failed"),
+                    "the first group failure must surface, got: " + 
ex.getMessage());
+
+            Assertions.assertEquals(2, taskConstruction.constructed().size());
+            Mockito.verify(txnManager).rollback(7L);
+            Mockito.verify(txnManager, Mockito.never()).commit(7L);
+        }
+    }
+
+    @Test
+    public void timeoutDrainsSubmittedGroupsBeforeSharedTransactionRollback() 
throws Exception {
+        RewriteDrainHarness harness = newDrainHarness();
+        AtomicInteger ids = new AtomicInteger(10);
+        
Mockito.when(harness.transientTaskManager.addMemoryTask(Mockito.any())).thenReturn(1L);
+        try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
+                MockedConstruction<ConnectorRewriteGroupTask> taskConstruction 
= Mockito.mockConstruction(
+                        ConnectorRewriteGroupTask.class, (task, 
constructionContext) -> {
+                            Mockito.when(task.getId()).thenReturn((long) 
ids.incrementAndGet());
+                            try {
+                                
Mockito.when(task.awaitTerminal(Mockito.anyLong(), 
Mockito.any())).thenReturn(true);
+                            } catch (InterruptedException e) {
+                                throw new RuntimeException(e);
+                            }
+                        })) {
+            envStatic.when(Env::getCurrentEnv).thenReturn(harness.env);
+
+            UserException ex = Assertions.assertThrows(UserException.class, 
harness.driver::run);
+            Assertions.assertTrue(ex.getMessage().contains("did not complete 
within timeout"),
+                    "the timeout must surface, got: " + ex.getMessage());
+
+            Assertions.assertEquals(2, taskConstruction.constructed().size());
+            for (ConnectorRewriteGroupTask task : 
taskConstruction.constructed()) {
+                Mockito.verify(task, Mockito.atLeastOnce()).cancel();
+            }
+            Mockito.verify(harness.txnManager).rollback(7L);
+            Mockito.verify(harness.txnManager, Mockito.never()).commit(7L);
+        }
+    }
+
+    @Test
+    public void submissionFailureDrainsAlreadySubmittedGroupsBeforeRollback() 
throws Exception {
+        RewriteDrainHarness harness = newDrainHarness();
+        AtomicInteger ids = new AtomicInteger(20);
+        AtomicInteger submissions = new AtomicInteger();
+        
Mockito.when(harness.transientTaskManager.addMemoryTask(Mockito.any())).thenAnswer(invocation
 -> {
+            if (submissions.incrementAndGet() == 2) {
+                throw new JobException("submit boom");
+            }
+            return 1L;
+        });
+        try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
+                MockedConstruction<ConnectorRewriteGroupTask> taskConstruction 
= Mockito.mockConstruction(
+                        ConnectorRewriteGroupTask.class, (task, 
constructionContext) -> {
+                            Mockito.when(task.getId()).thenReturn((long) 
ids.incrementAndGet());
+                            try {
+                                
Mockito.when(task.awaitTerminal(Mockito.anyLong(), 
Mockito.any())).thenReturn(true);
+                            } catch (InterruptedException e) {
+                                throw new RuntimeException(e);
+                            }
+                        })) {
+            envStatic.when(Env::getCurrentEnv).thenReturn(harness.env);
+
+            UserException ex = Assertions.assertThrows(UserException.class, 
harness.driver::run);
+            Assertions.assertTrue(ex.getMessage().contains("Failed to submit 
rewrite tasks"),
+                    "the submission failure must surface, got: " + 
ex.getMessage());
+
+            Assertions.assertEquals(2, taskConstruction.constructed().size());
+            // Only the first group made it into the manager; it still has to 
be drained before rollback.
+            Mockito.verify(taskConstruction.constructed().get(0), 
Mockito.atLeastOnce()).cancel();
+            Mockito.verify(harness.txnManager).rollback(7L);
+            Mockito.verify(harness.txnManager, Mockito.never()).commit(7L);
+        }
+    }
+
+    @Test
+    public void interruptionDrainsSubmittedGroupsBeforeRollback() throws 
Exception {
+        RewriteDrainHarness harness = newDrainHarness();
+        AtomicInteger ids = new AtomicInteger(30);
+        
Mockito.when(harness.transientTaskManager.addMemoryTask(Mockito.any())).thenAnswer(invocation
 -> {
+            Thread.currentThread().interrupt();
+            return 1L;
+        });
+        try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
+                MockedConstruction<ConnectorRewriteGroupTask> taskConstruction 
= Mockito.mockConstruction(
+                        ConnectorRewriteGroupTask.class, (task, 
constructionContext) -> {
+                            Mockito.when(task.getId()).thenReturn((long) 
ids.incrementAndGet());
+                            try {
+                                
Mockito.when(task.awaitTerminal(Mockito.anyLong(), 
Mockito.any())).thenReturn(true);
+                            } catch (InterruptedException e) {
+                                throw new RuntimeException(e);
+                            }
+                        })) {
+            envStatic.when(Env::getCurrentEnv).thenReturn(harness.env);
+
+            UserException ex = Assertions.assertThrows(UserException.class, 
harness.driver::run);
+            Assertions.assertTrue(ex.getMessage().contains("interrupted"),
+                    "the interruption must surface, got: " + ex.getMessage());
+
+            Assertions.assertEquals(2, taskConstruction.constructed().size());
+            for (ConnectorRewriteGroupTask task : 
taskConstruction.constructed()) {
+                Mockito.verify(task, Mockito.atLeastOnce()).cancel();
+            }
+            Mockito.verify(harness.txnManager).rollback(7L);
+            Mockito.verify(harness.txnManager, Mockito.never()).commit(7L);
+        } finally {
+            Thread.interrupted();
+        }
+    }
+
+    private static final class RewriteDrainHarness {
+        final ConnectorRewriteDriver driver;
+        final PluginDrivenTransactionManager txnManager;
+        final TransientTaskManager transientTaskManager;
+        final Env env;
+        final RewriteCapableTransaction rewriteTx;
+
+        RewriteDrainHarness(ConnectorRewriteDriver driver, 
PluginDrivenTransactionManager txnManager,
+                TransientTaskManager transientTaskManager, Env env, 
RewriteCapableTransaction rewriteTx) {
+            this.driver = driver;
+            this.txnManager = txnManager;
+            this.transientTaskManager = transientTaskManager;
+            this.env = env;
+            this.rewriteTx = rewriteTx;
+        }
+    }
+
+    private RewriteDrainHarness newDrainHarness() {
+        return newDrainHarness(null);
+    }
+
+    private RewriteDrainHarness newDrainHarness(StmtExecutor owner) {
+        ConnectorProcedureOps procedureOps = 
Mockito.mock(ConnectorProcedureOps.class);
+        ConnectorMetadata metadata = Mockito.mock(ConnectorMetadata.class);
+        ConnectorSession session = Mockito.mock(ConnectorSession.class);
+        ConnectorTableHandle tableHandle = 
Mockito.mock(ConnectorTableHandle.class);
+        ConnectorTransaction connectorTx = 
Mockito.mock(ConnectorTransaction.class,
+                
Mockito.withSettings().extraInterfaces(RewriteCapableTransaction.class));
+        PluginDrivenTransactionManager txnManager = 
Mockito.mock(PluginDrivenTransactionManager.class);
+        PluginDrivenExternalCatalog catalog = 
Mockito.mock(PluginDrivenExternalCatalog.class);
+        ExternalTable table = Mockito.mock(ExternalTable.class);
+        ConnectContext context = Mockito.mock(ConnectContext.class);
+        SessionVariable sessionVariable = Mockito.mock(SessionVariable.class);
+        ConnectorRewriteGroup first = new ConnectorRewriteGroup(
+                ImmutableSet.of("s3://bucket/table/a.parquet"), 1, 1024L, 0);
+        ConnectorRewriteGroup second = new ConnectorRewriteGroup(
+                ImmutableSet.of("s3://bucket/table/b.parquet"), 1, 1024L, 0);
+
+        Mockito.when(procedureOps.planRewrite(Mockito.any(), Mockito.any(), 
Mockito.any(), Mockito.any(),
+                Mockito.any(), Mockito.any())).thenReturn(Arrays.asList(first, 
second));
+        Mockito.when(metadata.beginTransaction(session, 
tableHandle)).thenReturn(connectorTx);
+        Mockito.when(catalog.getTransactionManager()).thenReturn(txnManager);
+        Mockito.when(txnManager.begin(connectorTx)).thenReturn(7L);
+        Mockito.when(context.getSessionVariable()).thenReturn(sessionVariable);
+        Mockito.when(sessionVariable.getInsertTimeoutS()).thenReturn(1);
+
+        Env env = Mockito.mock(Env.class);
+        TransientTaskManager transientTaskManager = 
Mockito.mock(TransientTaskManager.class);
+        
Mockito.when(env.getTransientTaskManager()).thenReturn(transientTaskManager);
+
+        ConnectorRewriteDriver driver = new ConnectorRewriteDriver(
+                context, table, catalog, metadata, procedureOps, session, 
tableHandle,
+                "rewrite_data_files", Collections.emptyMap(), 
Collections.emptyList(), null, owner);
+        return new RewriteDrainHarness(driver, txnManager, 
transientTaskManager, env,
+                (RewriteCapableTransaction) connectorTx);
+    }
+
+    @Test
+    public void 
outerCancellationBeforeCommitRollsBackWithoutRegisteringOrCommitting() throws 
Exception {
+        StmtExecutor owner = Mockito.mock(StmtExecutor.class);
+        // Cancellation is observed only once the groups have finished and the 
commit decision is taken:
+        // installCancelHandoff polls once (call 1), the wait loop polls once 
(call 2), and the
+        // registration/commit decision polls once more (call 3) and must see 
the terminal status.
+        Mockito.when(owner.getPendingCancelReason())
+                .thenReturn(null, null, new Status(TStatusCode.CANCELLED, 
"outer kill before commit"));
+        RewriteDrainHarness harness = newDrainHarness(owner);
+        AtomicInteger ids = new AtomicInteger(40);
+        Map<ConnectorRewriteGroupTask, 
ConnectorRewriteGroupTask.RewriteResultCallback> callbacks =
+                new ConcurrentHashMap<>();
+        try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
+                MockedConstruction<ConnectorRewriteGroupTask> taskConstruction 
= Mockito.mockConstruction(
+                        ConnectorRewriteGroupTask.class, (task, 
constructionContext) -> {
+                            Mockito.when(task.getId()).thenReturn((long) 
ids.incrementAndGet());
+                            callbacks.put(task, 
(ConnectorRewriteGroupTask.RewriteResultCallback)
+                                    constructionContext.arguments().get(5));
+                        })) {
+            envStatic.when(Env::getCurrentEnv).thenReturn(harness.env);
+            
Mockito.when(harness.transientTaskManager.addMemoryTask(Mockito.any())).thenAnswer(invocation
 -> {
+                ConnectorRewriteGroupTask task = invocation.getArgument(0);
+                // The group completes successfully; the cancellation only 
lands at the commit decision.
+                callbacks.get(task).onTaskCompleted(task.getId());
+                return task.getId();
+            });
+
+            UserException ex = Assertions.assertThrows(UserException.class, 
harness.driver::run);
+            Assertions.assertTrue(ex.getMessage().contains("Rewrite is 
cancelled"),
+                    "the outer cancellation must surface, got: " + 
ex.getMessage());
+
+            Mockito.verify(harness.txnManager).rollback(7L);
+            Mockito.verify(harness.txnManager, Mockito.never()).commit(7L);
+            Mockito.verify(harness.rewriteTx, Mockito.never())
+                    .registerRewriteSourceFiles(Mockito.any());
+        }
+    }
+
+    @Test
+    public void outerCancellationWhileWaitingDrainsLiveGroupsBeforeRollback() 
throws Exception {
+        StmtExecutor owner = Mockito.mock(StmtExecutor.class);
+        AtomicReference<Status> ownerReason = new AtomicReference<>();
+        Mockito.when(owner.getPendingCancelReason()).thenAnswer(invocation -> 
ownerReason.get());
+        RewriteDrainHarness harness = newDrainHarness(owner);
+        AtomicInteger ids = new AtomicInteger(50);
+        try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
+                MockedConstruction<ConnectorRewriteGroupTask> taskConstruction 
= Mockito.mockConstruction(
+                        ConnectorRewriteGroupTask.class, (task, 
constructionContext) -> {
+                            Mockito.when(task.getId()).thenReturn((long) 
ids.incrementAndGet());
+                            try {
+                                
Mockito.when(task.awaitTerminal(Mockito.anyLong(), 
Mockito.any())).thenReturn(true);
+                            } catch (InterruptedException e) {
+                                throw new RuntimeException(e);
+                            }
+                        })) {
+            envStatic.when(Env::getCurrentEnv).thenReturn(harness.env);
+            
Mockito.when(harness.transientTaskManager.addMemoryTask(Mockito.any())).thenAnswer(invocation
 -> {
+                // The groups are live; the outer TIMEOUT/KILL lands while the 
owner waits for them.
+                ownerReason.set(new Status(TStatusCode.CANCELLED, "outer kill 
while waiting"));
+                return 1L;
+            });
+
+            UserException ex = Assertions.assertThrows(UserException.class, 
harness.driver::run);
+            Assertions.assertTrue(ex.getMessage().contains("Rewrite is 
cancelled"));
+
+            Assertions.assertEquals(2, taskConstruction.constructed().size());
+            for (ConnectorRewriteGroupTask task : 
taskConstruction.constructed()) {
+                Mockito.verify(task, Mockito.atLeastOnce()).cancel();
+            }
+            Mockito.verify(harness.txnManager).rollback(7L);
+            Mockito.verify(harness.txnManager, Mockito.never()).commit(7L);
+        }
+    }
+
     @Test
     public void unionSourceFilePathsMergesAllGroupsAndDedupsByPath() {
         // STEP 3 registers the UNION of every group's source files in ONE 
connector call (one planFiles() scan)
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteGroupTaskTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteGroupTaskTest.java
new file mode 100644
index 00000000000..643cec801e7
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/execute/ConnectorRewriteGroupTaskTest.java
@@ -0,0 +1,117 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.nereids.trees.plans.commands.execute;
+
+import org.apache.doris.common.Status;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.connector.spi.handle.ConnectorTransaction;
+import org.apache.doris.connector.spi.procedure.ConnectorRewriteGroup;
+import org.apache.doris.datasource.ExternalTable;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.scheduler.exception.JobException;
+
+import com.google.common.collect.ImmutableSet;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
+
+/**
+ * Guards the cancellation/publication handoff of {@link 
ConnectorRewriteGroupTask}. The distributed write
+ * path needs a live cluster, but the handoff itself is unit-testable: a 
queued task must report its
+ * cancellation through the collector, and a cancel that lands after the 
executor is published must reach it.
+ */
+public class ConnectorRewriteGroupTaskTest {
+
+    private ConnectorRewriteGroupTask newTask(
+            AtomicBoolean completed,
+            AtomicReference<Long> failedId,
+            AtomicReference<Exception> failure) {
+        ConnectorRewriteGroup group = new ConnectorRewriteGroup(
+                ImmutableSet.of("s3://bucket/table/a.parquet"), 1, 1024L, 0);
+        return new ConnectorRewriteGroupTask(
+                group,
+                7L,
+                Mockito.mock(ConnectorTransaction.class),
+                Mockito.mock(ExternalTable.class),
+                Mockito.mock(ConnectContext.class),
+                new ConnectorRewriteGroupTask.RewriteResultCallback() {
+                    @Override
+                    public void onTaskCompleted(Long taskId) {
+                        completed.set(true);
+                    }
+
+                    @Override
+                    public void onTaskFailed(Long taskId, Exception error) {
+                        failedId.set(taskId);
+                        failure.set(error);
+                    }
+                });
+    }
+
+    @Test
+    public void queuedCancellationIsReportedToTheCollector() throws Exception {
+        AtomicBoolean completed = new AtomicBoolean(false);
+        AtomicReference<Long> failedId = new AtomicReference<>();
+        AtomicReference<Exception> failure = new AtomicReference<>();
+        ConnectorRewriteGroupTask task = newTask(completed, failedId, failure);
+
+        // Cancelled while still queued: the scheduler may still invoke 
execute().
+        task.cancel();
+        Assertions.assertThrows(JobException.class, task::execute);
+
+        Assertions.assertFalse(completed.get(), "a cancelled task must never 
report completion");
+        Assertions.assertEquals(task.getId(), failedId.get(), "the collector 
must be told this task failed");
+        Assertions.assertNotNull(failure.get());
+    }
+
+    @Test
+    public void cancelAfterExecutorPublicationCancelsTheExecutor() throws 
Exception {
+        AtomicBoolean completed = new AtomicBoolean(false);
+        AtomicReference<Long> failedId = new AtomicReference<>();
+        AtomicReference<Exception> failure = new AtomicReference<>();
+        ConnectorRewriteGroupTask task = newTask(completed, failedId, failure);
+
+        // Simulate a running task whose executor has been published.
+        StmtExecutor stmtExecutor = Mockito.mock(StmtExecutor.class);
+        Deencapsulation.setField(task, "stmtExecutor", stmtExecutor);
+
+        task.cancel();
+
+        Mockito.verify(stmtExecutor).cancel(Mockito.any(Status.class));
+        Assertions.assertFalse(completed.get());
+    }
+
+    @Test
+    public void cancelAfterFinishIsIgnored() throws Exception {
+        AtomicBoolean completed = new AtomicBoolean(false);
+        AtomicReference<Long> failedId = new AtomicReference<>();
+        AtomicReference<Exception> failure = new AtomicReference<>();
+        ConnectorRewriteGroupTask task = newTask(completed, failedId, failure);
+        ((AtomicBoolean) Deencapsulation.getField(task, 
"isFinished")).set(true);
+        StmtExecutor stmtExecutor = Mockito.mock(StmtExecutor.class);
+        Deencapsulation.setField(task, "stmtExecutor", stmtExecutor);
+
+        task.cancel();
+
+        Mockito.verify(stmtExecutor, 
Mockito.never()).cancel(Mockito.any(Status.class));
+    }
+}
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/DictionaryInsertTargetDropRaceTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/DictionaryInsertTargetDropRaceTest.java
index df31a7b4f4e..4fedc988f5b 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/DictionaryInsertTargetDropRaceTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/DictionaryInsertTargetDropRaceTest.java
@@ -22,13 +22,17 @@ import org.apache.doris.catalog.Database;
 import org.apache.doris.catalog.Env;
 import org.apache.doris.catalog.TableIf;
 import org.apache.doris.common.Config;
+import org.apache.doris.common.UserException;
+import org.apache.doris.common.jmockit.Deencapsulation;
 import org.apache.doris.common.util.DebugPointUtil;
 import org.apache.doris.common.util.DebugPointUtil.DebugPoint;
 import org.apache.doris.dictionary.Dictionary;
 import org.apache.doris.nereids.StatementContext;
 import org.apache.doris.nereids.parser.NereidsParser;
 import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.Coordinator;
 import org.apache.doris.qe.OriginStatement;
+import org.apache.doris.qe.QueryState.MysqlStateType;
 import org.apache.doris.qe.StmtExecutor;
 import org.apache.doris.thrift.TUniqueId;
 import org.apache.doris.utframe.TestWithFeService;
@@ -37,6 +41,8 @@ import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
+import org.mockito.Answers;
+import org.mockito.Mockito;
 
 import java.util.List;
 import java.util.concurrent.CountDownLatch;
@@ -140,6 +146,23 @@ class DictionaryInsertTargetDropRaceTest extends 
TestWithFeService {
         Env.getCurrentInternalCatalog().dropDb(sourceDbName, false, true);
     }
 
+    @Test
+    void reportsCoordinatorCancellationWithoutAssumingDdlException() {
+        ConnectContext ctx = new ConnectContext();
+        ctx.setQueryId(new TUniqueId(3, 4));
+        Coordinator coordinator = Mockito.mock(Coordinator.class);
+        Mockito.when(coordinator.getTrackingUrl()).thenReturn("");
+        DictionaryInsertExecutor executor = 
Mockito.mock(DictionaryInsertExecutor.class, Answers.CALLS_REAL_METHODS);
+        Deencapsulation.setField(executor, "ctx", ctx);
+        Deencapsulation.setField(executor, "coordinator", coordinator);
+        Deencapsulation.setField(executor, "labelName", 
"dictionary_cancel_test");
+
+        Assertions.assertDoesNotThrow(() -> executor.onFail(new 
UserException("dictionary insert timeout")));
+
+        Assertions.assertEquals(MysqlStateType.ERR, 
ctx.getState().getStateType());
+        
Assertions.assertTrue(ctx.getState().getErrorMessage().contains("dictionary 
insert timeout"));
+    }
+
     private void createSourceTable() throws Exception {
         createTable("CREATE TABLE source_table (id INT NOT NULL, city 
VARCHAR(32) NOT NULL) "
                 + "DISTRIBUTED BY HASH(id) BUCKETS 1 PROPERTIES 
('replication_num' = '1')");
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapInsertExecutorTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapInsertExecutorTest.java
index d2d792541dc..793b9d60fa6 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapInsertExecutorTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/OlapInsertExecutorTest.java
@@ -57,6 +57,7 @@ import org.mockito.Mockito;
 
 import java.util.List;
 import java.util.Optional;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 /**
  * Tests for publish-timeout behaviors in {@link OlapInsertExecutor}.
@@ -94,6 +95,8 @@ class OlapInsertExecutorTest {
             executor.txnId = 10001L;
             executor.executeSingleInsert(stmtExecutor);
 
+            Mockito.verify(stmtExecutor).setCoord(coordinator);
+
             Assertions.assertEquals(TransactionStatus.COMMITTED, 
executor.txnStatus);
             Assertions.assertEquals(MysqlStateType.ERR, 
ctx.getState().getStateType());
             Assertions.assertTrue(ctx.getState().getErrorMessage().contains(
@@ -238,6 +241,39 @@ class OlapInsertExecutorTest {
         }
     }
 
+    @Test
+    void testPendingCoordinatorTimeoutFencesBeforeExecSetup() throws Exception 
{
+        ConnectContext ctx = createExecutorContext();
+        Coordinator coordinator = createCoordinator();
+        Mockito.when(coordinator.getExecStatus())
+                .thenReturn(new Status(TStatusCode.TIMEOUT, "timeout before 
coordinator setup"));
+        GlobalTransactionMgrIface txnMgr = 
Mockito.mock(GlobalTransactionMgrIface.class);
+        TransactionState txnState = Mockito.mock(TransactionState.class);
+        LoadManager loadManager = Mockito.mock(LoadManager.class);
+        Env currentEnv = createCurrentEnv(loadManager);
+        StmtExecutor stmtExecutor = createStmtExecutor();
+
+        try (MockedStatic<EnvFactory> envFactoryMock = 
Mockito.mockStatic(EnvFactory.class);
+                MockedStatic<Env> envMock = Mockito.mockStatic(Env.class)) {
+            prepareFactoryMocks(envFactoryMock, envMock, coordinator, txnMgr, 
txnState, currentEnv);
+            ctx.setEnv(currentEnv);
+
+            AtomicBoolean beforeExecRan = new AtomicBoolean(false);
+            OlapInsertExecutor executor = 
createExecutorWithBeforeExecProbe(ctx, beforeExecRan);
+            executor.txnId = 10006L;
+
+            Assertions.assertDoesNotThrow(() -> 
executor.executeSingleInsert(stmtExecutor));
+
+            // The retained terminal reason is fenced immediately after 
coordinator publication, before any
+            // executor-specific setup can run or fail with a different error.
+            Assertions.assertFalse(beforeExecRan.get(), "beforeExec must not 
run after a retained timeout");
+            Assertions.assertEquals(MysqlStateType.ERR, 
ctx.getState().getStateType());
+            
Assertions.assertTrue(ctx.getState().getErrorMessage().contains("timeout before 
coordinator setup"));
+            Mockito.verify(coordinator, Mockito.never()).exec();
+            Mockito.verify(coordinator).close();
+        }
+    }
+
     @Test
     void testEmptyStreamInsertCommitsWithoutCoordinatorExecution() throws 
Exception {
         ConnectContext ctx = createExecutorContext();
@@ -272,6 +308,44 @@ class OlapInsertExecutorTest {
         }
     }
 
+    @Test
+    void testCancellationAfterLastStatusReadFencesBeforeCommit() throws 
Exception {
+        ConnectContext ctx = createExecutorContext();
+        Coordinator coordinator = createCoordinator();
+        // execImpl() reads an OK status; the completion handoff then flips it 
before onComplete() commits,
+        // exactly the window row-level UPDATE/DELETE/MERGE has no 
command-level listener for.
+        AtomicBoolean cancelled = new AtomicBoolean(false);
+        Mockito.when(coordinator.getExecStatus()).thenAnswer(invocation -> 
cancelled.get()
+                ? new Status(TStatusCode.TIMEOUT, "timeout after last status 
read")
+                : new Status(TStatusCode.OK, ""));
+        GlobalTransactionMgrIface txnMgr = 
Mockito.mock(GlobalTransactionMgrIface.class);
+        TransactionState txnState = Mockito.mock(TransactionState.class);
+        LoadManager loadManager = Mockito.mock(LoadManager.class);
+        Env currentEnv = createCurrentEnv(loadManager);
+        StmtExecutor stmtExecutor = createStmtExecutor();
+
+        try (MockedStatic<EnvFactory> envFactoryMock = 
Mockito.mockStatic(EnvFactory.class);
+                MockedStatic<Env> envMock = Mockito.mockStatic(Env.class)) {
+            prepareFactoryMocks(envFactoryMock, envMock, coordinator, txnMgr, 
txnState, currentEnv);
+            ctx.setEnv(currentEnv);
+
+            OlapInsertExecutor executor = createExecutor(ctx);
+            executor.txnId = 10007L;
+            executor.registerListener(new 
AbstractInsertExecutor.InsertExecutorListener() {
+                @Override
+                public void beforeComplete(AbstractInsertExecutor 
insertExecutor, StmtExecutor executor, long jobId) {
+                    cancelled.set(true);
+                }
+            });
+
+            Assertions.assertDoesNotThrow(() -> 
executor.executeSingleInsert(stmtExecutor));
+
+            Assertions.assertEquals(MysqlStateType.ERR, 
ctx.getState().getStateType());
+            
Assertions.assertTrue(ctx.getState().getErrorMessage().contains("timeout after 
last status read"));
+            Mockito.verify(txnMgr).abortTransaction(Mockito.eq(1L), 
Mockito.eq(10007L), Mockito.anyString());
+        }
+    }
+
     // Build a fresh context per case so insertResult and QueryState do not 
leak between tests.
     private ConnectContext createExecutorContext() {
         ConnectContext ctx = new ConnectContext();
@@ -371,6 +445,25 @@ class OlapInsertExecutorTest {
         };
     }
 
+    private OlapInsertExecutor 
createExecutorWithBeforeExecProbe(ConnectContext ctx, AtomicBoolean 
beforeExecRan) {
+        Database database = Mockito.mock(Database.class);
+        Mockito.when(database.getFullName()).thenReturn("test_db");
+        Mockito.when(database.getId()).thenReturn(1L);
+
+        OlapTable table = Mockito.mock(OlapTable.class);
+        Mockito.when(table.getDatabase()).thenReturn(database);
+        Mockito.when(table.getName()).thenReturn("test_tbl");
+        Mockito.when(table.getId()).thenReturn(2L);
+
+        return new OlapInsertExecutor(ctx, table, "label_test", 
Mockito.mock(NereidsPlanner.class),
+                Optional.empty(), false, 0L) {
+            @Override
+            protected void beforeExec() {
+                beforeExecRan.set(true);
+            }
+        };
+    }
+
     // Redirect coordinator creation and transaction access to mocks so the 
test stays deterministic.
     private void prepareFactoryMocks(MockedStatic<EnvFactory> envFactoryMock, 
MockedStatic<Env> envMock,
             Coordinator coordinator, GlobalTransactionMgrIface txnMgr, 
TransactionState txnState, Env currentEnv) {
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/qe/NereidsCoordinatorTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/qe/NereidsCoordinatorTest.java
index f36b1db9dd5..27fdac75ef8 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/qe/NereidsCoordinatorTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/qe/NereidsCoordinatorTest.java
@@ -18,19 +18,29 @@
 package org.apache.doris.qe;
 
 import org.apache.doris.catalog.EnvFactory;
+import org.apache.doris.common.AnalysisException;
 import org.apache.doris.common.FeConstants;
+import org.apache.doris.common.Status;
+import org.apache.doris.common.UserException;
 import org.apache.doris.nereids.NereidsPlanner;
+import org.apache.doris.nereids.trees.plans.distribute.PipelineDistributedPlan;
 import org.apache.doris.nereids.util.PlanChecker;
 import org.apache.doris.planner.PlanFragment;
+import org.apache.doris.planner.ScanNode;
+import org.apache.doris.thrift.TStatusCode;
 import org.apache.doris.thrift.TUniqueId;
 import org.apache.doris.utframe.TestWithFeService;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
+import org.mockito.Mockito;
 
 import java.io.IOException;
 import java.util.UUID;
+import java.util.concurrent.atomic.AtomicInteger;
 
 public class NereidsCoordinatorTest extends TestWithFeService {
     @BeforeAll
@@ -81,6 +91,98 @@ public class NereidsCoordinatorTest extends 
TestWithFeService {
         }
     }
 
+    @ParameterizedTest
+    @EnumSource(value = TStatusCode.class, names = {"CANCELLED", "TIMEOUT"})
+    public void testTerminateBeforeFragmentDispatch(TStatusCode statusCode) 
throws Exception {
+        ConnectContext context = createDefaultCtx();
+        NereidsPlanner planner = plan("select * from test.tbl", context);
+        Status cancelReason = new Status(statusCode, "terminate before 
fragment dispatch");
+        NereidsCoordinator coordinator = new NereidsCoordinator(context, 
planner, null) {
+            @Override
+            protected void processTopSink(CoordinatorContext 
coordinatorContext,
+                    PipelineDistributedPlan topPlan) throws AnalysisException {
+                cancel(cancelReason);
+            }
+        };
+
+        UserException exception = Assertions.assertThrows(UserException.class, 
coordinator::exec);
+        Assertions.assertTrue(exception.getMessage().contains("terminate 
before fragment dispatch"));
+    }
+
+    @Test
+    public void testCancelPublishesStatusBeforeScanCleanupFailure() throws 
Exception {
+        ConnectContext context = createDefaultCtx();
+        NereidsPlanner planner = plan("select * from test.tbl", context);
+        Status cancelReason = new Status(TStatusCode.TIMEOUT, "timeout before 
scan cleanup");
+        ScanNode failingScan = Mockito.mock(ScanNode.class);
+        Mockito.doThrow(new RuntimeException("scan cleanup 
failed")).when(failingScan).stop();
+        ScanNode remainingScan = Mockito.mock(ScanNode.class);
+        AtomicInteger cancelInternalCalls = new AtomicInteger();
+        NereidsCoordinator coordinator = new NereidsCoordinator(context, 
planner, null) {
+            @Override
+            protected void cancelInternal(Status status) {
+                cancelInternalCalls.incrementAndGet();
+            }
+        };
+        coordinator.coordinatorContext.scanNodes.clear();
+        coordinator.coordinatorContext.scanNodes.add(failingScan);
+        coordinator.coordinatorContext.scanNodes.add(remainingScan);
+
+        // A fallible scan cleanup must neither escape to the caller nor skip 
the remaining scans; the
+        // terminal status is published before cleanup and cancelInternal() 
still runs in finally.
+        Assertions.assertDoesNotThrow(() -> coordinator.cancel(cancelReason));
+
+        Assertions.assertEquals(TStatusCode.TIMEOUT, 
coordinator.getExecStatus().getErrorCode());
+        Assertions.assertEquals("timeout before scan cleanup", 
coordinator.getExecStatus().getErrorMsg());
+        Assertions.assertTrue(cancelInternalCalls.get() >= 1);
+        Mockito.verify(remainingScan).stop();
+    }
+
+    @Test
+    public void testQueueCancellationPrefersRetainedTerminalReason() throws 
Exception {
+        ConnectContext context = createDefaultCtx();
+        NereidsPlanner planner = plan("select * from test.tbl", context);
+        NereidsCoordinator coordinator = new NereidsCoordinator(context, 
planner, null) {
+            @Override
+            protected void cancelInternal(Status status) {
+            }
+        };
+
+        Assertions.assertTrue(coordinator.preferTerminalReason(new 
UserException("query is cancelled"))
+                .getMessage().contains("query is cancelled"));
+
+        coordinator.cancel(new Status(TStatusCode.TIMEOUT, "retained queue 
timeout"));
+        Assertions.assertTrue(coordinator.preferTerminalReason(new 
UserException("query is cancelled"))
+                .getMessage().contains("retained queue timeout"));
+    }
+
+    @Test
+    public void testCancelStillCleansUpWhenStatusPublicationFails() throws 
Exception {
+        ConnectContext context = createDefaultCtx();
+        NereidsPlanner planner = plan("select * from test.tbl", context);
+        Status cancelReason = new Status(TStatusCode.TIMEOUT, "timeout before 
cleanup");
+        ScanNode scanNode = Mockito.mock(ScanNode.class);
+        AtomicInteger cancelInternalCalls = new AtomicInteger();
+        NereidsCoordinator coordinator = new NereidsCoordinator(context, 
planner, null) {
+            @Override
+            protected void cancelInternal(Status status) {
+                // updateStatusIfOk calls this while publishing; simulate a 
partially initialized processor
+                // whose cancel throws before the cleanup scope used to start.
+                if (cancelInternalCalls.incrementAndGet() == 1) {
+                    throw new RuntimeException("partially initialized 
processor cancel");
+                }
+            }
+        };
+        coordinator.coordinatorContext.scanNodes.clear();
+        coordinator.coordinatorContext.scanNodes.add(scanNode);
+
+        Assertions.assertThrows(RuntimeException.class, () -> 
coordinator.cancel(cancelReason));
+        // The publication threw, but the scan cleanup and the final internal 
cancel still ran.
+        Mockito.verify(scanNode).stop();
+        Assertions.assertTrue(cancelInternalCalls.get() >= 2, "the final 
internal cancel must still be resent");
+        Assertions.assertEquals(TStatusCode.TIMEOUT, 
coordinator.getExecStatus().getErrorCode());
+    }
+
     private NereidsPlanner plan(String sql) throws IOException {
         return plan(sql, connectContext);
     }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/qe/OldCoordinatorTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/qe/OldCoordinatorTest.java
index b8cb5b74165..975f2c35894 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/qe/OldCoordinatorTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/qe/OldCoordinatorTest.java
@@ -17,16 +17,27 @@
 
 package org.apache.doris.qe;
 
+import org.apache.doris.analysis.DescriptorTable;
+import org.apache.doris.common.Status;
+import org.apache.doris.common.UserException;
 import org.apache.doris.nereids.rules.RuleType;
 import org.apache.doris.planner.OlapScanNode;
 import org.apache.doris.planner.PlanFragment;
 import org.apache.doris.planner.PlanFragmentId;
 import org.apache.doris.planner.PlanNode;
+import org.apache.doris.planner.ScanNode;
+import org.apache.doris.thrift.TStatusCode;
+import org.apache.doris.thrift.TUniqueId;
 import org.apache.doris.utframe.TestWithFeService;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
+import org.mockito.Mockito;
 
+import java.util.Arrays;
+import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.atomic.AtomicBoolean;
@@ -99,4 +110,67 @@ public class OldCoordinatorTest extends TestWithFeService {
         }.test();
         Assertions.assertTrue(shuffleFragmentHasMultiInstances.get());
     }
+
+    @ParameterizedTest
+    @EnumSource(value = TStatusCode.class, names = {"CANCELLED", "TIMEOUT"})
+    public void 
testTerminateBetweenExecutionAdmissionAndFragmentDispatch(TStatusCode 
statusCode) {
+        Status cancelReason = new Status(statusCode, "terminate before 
fragment dispatch");
+        Coordinator coordinator = new Coordinator(0L, new TUniqueId(1L, 1L),
+                new DescriptorTable(), Collections.emptyList(), 
Collections.emptyList(), "UTC", false, false) {
+            @Override
+            protected void execInternal() throws Exception {
+                cancel(cancelReason);
+                sendPipelineCtx();
+            }
+        };
+
+        UserException exception = Assertions.assertThrows(UserException.class, 
coordinator::exec);
+        Assertions.assertTrue(exception.getMessage().contains("terminate 
before fragment dispatch"));
+    }
+
+    @Test
+    public void testCancelPublishesStatusBeforeScanCleanupFailure() {
+        Status cancelReason = new Status(TStatusCode.TIMEOUT, "timeout before 
scan cleanup");
+        ScanNode failingScan = Mockito.mock(ScanNode.class);
+        Mockito.doThrow(new RuntimeException("scan cleanup 
failed")).when(failingScan).stop();
+        ScanNode remainingScan = Mockito.mock(ScanNode.class);
+        AtomicBoolean cancelInternalCalled = new AtomicBoolean(false);
+        Coordinator coordinator = new Coordinator(0L, new TUniqueId(1L, 1L),
+                new DescriptorTable(), Collections.emptyList(), 
Arrays.asList(failingScan, remainingScan),
+                "UTC", false, false) {
+            @Override
+            protected void cancelInternal(Status status) {
+                cancelInternalCalled.set(true);
+            }
+        };
+
+        // A fallible scan cleanup must neither escape to the caller (which 
would skip the owner's close and
+        // mask the retained reason) nor skip the remaining scans. The 
terminal status is published first.
+        Assertions.assertDoesNotThrow(() -> coordinator.cancel(cancelReason));
+
+        Assertions.assertEquals(TStatusCode.TIMEOUT, 
coordinator.getExecStatus().getErrorCode());
+        Assertions.assertEquals("timeout before scan cleanup", 
coordinator.getExecStatus().getErrorMsg());
+        Assertions.assertTrue(cancelInternalCalled.get());
+        Mockito.verify(remainingScan).stop();
+    }
+
+    @Test
+    public void testQueueCancellationPrefersRetainedTerminalReason() {
+        Coordinator coordinator = new Coordinator(0L, new TUniqueId(1L, 1L),
+                new DescriptorTable(), Collections.emptyList(), 
Collections.emptyList(),
+                "UTC", false, false) {
+            @Override
+            protected void cancelInternal(Status status) {
+            }
+        };
+
+        // Without a retained status the queue token's own message stays.
+        Assertions.assertTrue(coordinator.preferTerminalReason(new 
UserException("query is cancelled"))
+                .getMessage().contains("query is cancelled"));
+
+        coordinator.cancel(new Status(TStatusCode.TIMEOUT, "retained queue 
timeout"));
+        // A TIMEOUT/KILL that unblocked the queue wait must win over the 
token's generic message.
+        Assertions.assertTrue(coordinator.preferTerminalReason(new 
UserException("query is cancelled"))
+                .getMessage().contains("retained queue timeout"));
+    }
 }
diff --git a/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
index 95b7d17734f..86f36224ad0 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
@@ -37,6 +37,7 @@ import org.apache.doris.planner.PlanFragment;
 import org.apache.doris.planner.Planner;
 import org.apache.doris.planner.ResultFileSink;
 import org.apache.doris.thrift.TQueryOptions;
+import org.apache.doris.thrift.TStatusCode;
 import org.apache.doris.thrift.TUniqueId;
 import org.apache.doris.utframe.TestWithFeService;
 
@@ -44,11 +45,20 @@ import com.google.common.collect.Lists;
 import com.google.common.collect.Maps;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
 import org.mockito.MockedConstruction;
 import org.mockito.Mockito;
 
 import java.lang.reflect.Field;
 import java.lang.reflect.Method;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 
 public class StmtExecutorTest extends TestWithFeService {
@@ -100,6 +110,68 @@ public class StmtExecutorTest extends TestWithFeService {
         Assertions.assertEquals(QueryState.MysqlStateType.OK, 
connectContext.getState().getStateType());
     }
 
+    @ParameterizedTest
+    @EnumSource(value = TStatusCode.class, names = {"CANCELLED", "TIMEOUT"})
+    public void testTerminateBeforeCoordinatorIsPublished(TStatusCode 
statusCode) {
+        StmtExecutor stmtExecutor = new StmtExecutor(connectContext, "");
+        Status cancelReason = new Status(statusCode, "terminate before 
coordinator publication");
+        Coordinator coordinator = Mockito.mock(Coordinator.class);
+
+        stmtExecutor.cancel(cancelReason, false);
+        stmtExecutor.setCoord(coordinator);
+
+        Mockito.verify(coordinator).cancel(cancelReason);
+    }
+
+    @Test
+    public void testCancelAfterCoordinatorIsPublished() {
+        StmtExecutor stmtExecutor = new StmtExecutor(connectContext, "");
+        Status cancelReason = new Status(TStatusCode.CANCELLED, "cancel after 
coordinator publication");
+        Coordinator coordinator = Mockito.mock(Coordinator.class);
+
+        stmtExecutor.setCoord(coordinator);
+        stmtExecutor.cancel(cancelReason, false);
+
+        Mockito.verify(coordinator).cancel(cancelReason);
+    }
+
+    @Test
+    public void testFirstTerminalReasonWinsAcrossCoordinatorPublication() 
throws Exception {
+        StmtExecutor stmtExecutor = new StmtExecutor(connectContext, "");
+        Status timeout = new Status(TStatusCode.TIMEOUT, "first timeout");
+        Status cancelled = new Status(TStatusCode.CANCELLED, "later 
cancellation");
+        Coordinator coordinator = Mockito.mock(Coordinator.class);
+        CountDownLatch publicationReachedCoordinator = new CountDownLatch(1);
+        CountDownLatch finishPublication = new CountDownLatch(1);
+        AtomicInteger calls = new AtomicInteger();
+        List<Status> delivered = new CopyOnWriteArrayList<>();
+        Mockito.doAnswer(invocation -> {
+            delivered.add(invocation.getArgument(0));
+            if (calls.getAndIncrement() == 0) {
+                publicationReachedCoordinator.countDown();
+                Assertions.assertTrue(finishPublication.await(10, 
TimeUnit.SECONDS));
+            }
+            return null;
+        }).when(coordinator).cancel(Mockito.any(Status.class));
+
+        stmtExecutor.cancel(timeout, false);
+        ExecutorService executorService = Executors.newSingleThreadExecutor();
+        try {
+            Future<?> publication = executorService.submit(() -> 
stmtExecutor.setCoord(coordinator));
+            Assertions.assertTrue(publicationReachedCoordinator.await(10, 
TimeUnit.SECONDS));
+            stmtExecutor.cancel(cancelled, false);
+            finishPublication.countDown();
+            publication.get(10, TimeUnit.SECONDS);
+        } finally {
+            finishPublication.countDown();
+            executorService.shutdownNow();
+        }
+
+        Assertions.assertEquals(2, delivered.size());
+        Assertions.assertTrue(delivered.stream().allMatch(status -> 
status.getErrorCode() == TStatusCode.TIMEOUT));
+        Assertions.assertTrue(delivered.stream().allMatch(status -> "first 
timeout".equals(status.getErrorMsg())));
+    }
+
     // The deferral gate (#67503): a coordinator is kept alive past 
GetFlightInfo only when the BE
     // still depends on it (Coordinator.mustOutliveDispatch), and the 
execution timeout it
     // ran with is frozen at that moment. SET_VAR hint values are reverted 
when execute() ends, so
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/qe/runtime/PipelineExecutionTaskTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/qe/runtime/PipelineExecutionTaskTest.java
index 83c61861e2c..25bc54ad0d6 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/qe/runtime/PipelineExecutionTaskTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/qe/runtime/PipelineExecutionTaskTest.java
@@ -26,6 +26,7 @@ import org.apache.doris.qe.CoordinatorContext;
 import org.apache.doris.rpc.BackendServiceProxy;
 import org.apache.doris.system.Backend;
 import org.apache.doris.thrift.TQueryOptions;
+import org.apache.doris.thrift.TStatusCode;
 import org.apache.doris.thrift.TUniqueId;
 
 import org.junit.jupiter.api.Assertions;
@@ -53,6 +54,7 @@ class PipelineExecutionTaskTest {
         Deencapsulation.setField(coordinatorContext, "timeoutDeadline", 
(Supplier<Long>) () -> 0L);
         
Mockito.when(coordinatorContext.withLock(ArgumentMatchers.<Callable<Object>>any()))
                 .thenAnswer(invocation -> 
invocation.<Callable<Object>>getArgument(0).call());
+        
Mockito.when(coordinatorContext.readCloneStatus()).thenReturn(Status.OK);
         Mockito.when(coordinatorContext.twoPhaseExecution()).thenReturn(false);
 
         MultiFragmentsPipelineTask fragmentsTask = 
Mockito.mock(MultiFragmentsPipelineTask.class);
@@ -83,6 +85,7 @@ class PipelineExecutionTaskTest {
         Deencapsulation.setField(coordinatorContext, "timeoutDeadline", 
(Supplier<Long>) () -> Long.MAX_VALUE);
         
Mockito.when(coordinatorContext.withLock(ArgumentMatchers.<Callable<Object>>any()))
                 .thenAnswer(invocation -> 
invocation.<Callable<Object>>getArgument(0).call());
+        
Mockito.when(coordinatorContext.readCloneStatus()).thenReturn(Status.OK);
         Mockito.when(coordinatorContext.twoPhaseExecution()).thenReturn(false);
 
         MultiFragmentsPipelineTask first = mockFragmentTask(10001L);
@@ -105,6 +108,41 @@ class PipelineExecutionTaskTest {
         Mockito.verify(neverAttempted, Mockito.never()).sendPhaseOneRpc(false);
     }
 
+    @Test
+    void timeoutStatusPreventsFragmentDispatch() throws Exception {
+        CoordinatorContext coordinatorContext = 
Mockito.mock(CoordinatorContext.class);
+        Deencapsulation.setField(coordinatorContext, "timeoutDeadline", 
(Supplier<Long>) () -> Long.MAX_VALUE);
+        
Mockito.when(coordinatorContext.withLock(ArgumentMatchers.<Callable<Object>>any()))
+                .thenAnswer(invocation -> 
invocation.<Callable<Object>>getArgument(0).call());
+        Mockito.when(coordinatorContext.readCloneStatus())
+                .thenReturn(new Status(TStatusCode.TIMEOUT, "timeout before 
fragment dispatch"));
+
+        MultiFragmentsPipelineTask fragmentsTask = mockFragmentTask(10001L);
+        PipelineExecutionTask executionTask = new PipelineExecutionTask(
+                coordinatorContext,
+                Mockito.mock(BackendServiceProxy.class),
+                Collections.singletonMap(Mockito.mock(BackendWorker.class), 
fragmentsTask));
+
+        UserException exception = Assertions.assertThrows(UserException.class, 
executionTask::execute);
+
+        Assertions.assertTrue(exception.getMessage().contains("timeout before 
fragment dispatch"));
+        Mockito.verify(fragmentsTask, 
Mockito.never()).sendPhaseOneRpc(ArgumentMatchers.anyBoolean());
+    }
+
+    @Test
+    void loadProcessorCancelToleratesTaskPublishedBeforeLatch() {
+        LoadProcessor processor = Mockito.mock(LoadProcessor.class, 
Mockito.CALLS_REAL_METHODS);
+        PipelineExecutionTask task = Mockito.mock(PipelineExecutionTask.class);
+        
Mockito.when(task.getChildrenTasks()).thenReturn(Collections.emptyMap());
+        Deencapsulation.setField(processor, "executionTask", 
java.util.Optional.of(task));
+        // AfterSetPipelineExecutionTask has not created the latch yet 
(Objenesis skipped the constructor).
+        Deencapsulation.setField(processor, "latch", 
java.util.Optional.empty());
+        // latch deliberately stays Optional.empty(): the window between 
executionTask publication and
+        // afterSetPipelineExecutionTask creating the latch must not throw out 
of the coordinator cleanup.
+        Assertions.assertDoesNotThrow(() ->
+                processor.cancel(new Status(TStatusCode.TIMEOUT, "cancel 
during publication")));
+    }
+
     private static MultiFragmentsPipelineTask mockFragmentTask(long backendId) 
{
         MultiFragmentsPipelineTask task = 
Mockito.mock(MultiFragmentsPipelineTask.class);
         Backend backend = Mockito.mock(Backend.class);
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/scheduler/manager/TransientTaskManagerTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/scheduler/manager/TransientTaskManagerTest.java
new file mode 100644
index 00000000000..8e111232b56
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/scheduler/manager/TransientTaskManagerTest.java
@@ -0,0 +1,60 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.scheduler.manager;
+
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.scheduler.disruptor.TaskDisruptor;
+import org.apache.doris.scheduler.exception.JobException;
+import org.apache.doris.scheduler.executor.TransientTaskExecutor;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+/**
+ * Guards the register-then-publish contract of {@link TransientTaskManager}: 
{@code addMemoryTask}
+ * inserts the task before it publishes the event, but {@code TaskHandler} 
only removes tasks it actually
+ * runs. A publication failure must therefore unregister the task instead of 
leaking it for the FE lifetime.
+ */
+public class TransientTaskManagerTest {
+
+    @Test
+    public void publicationFailureUnregistersTheRegisteredTask() throws 
Exception {
+        TransientTaskManager manager = new TransientTaskManager();
+        TaskDisruptor disruptor = Mockito.mock(TaskDisruptor.class);
+        Mockito.doThrow(new JobException("There is not enough available 
capacity in the RingBuffer."))
+                .when(disruptor).tryPublishTask(Mockito.anyLong());
+        Deencapsulation.setField(manager, "disruptor", disruptor);
+        TransientTaskExecutor executor = 
Mockito.mock(TransientTaskExecutor.class);
+        Mockito.when(executor.getId()).thenReturn(42L);
+
+        // addMemoryTask performs the real map insertion first; only the event 
publication is stubbed to fail.
+        Assertions.assertThrows(JobException.class, () -> 
manager.addMemoryTask(executor));
+        Assertions.assertNull(manager.getMemoryTaskExecutor(42L),
+                "a task whose scheduler publication failed must not stay 
registered");
+    }
+
+    @Test
+    public void closedDisruptorFailsInsteadOfSilentlyDroppingTheTask() {
+        TaskDisruptor disruptor = Mockito.mock(TaskDisruptor.class, 
Mockito.CALLS_REAL_METHODS);
+        Deencapsulation.setField(disruptor, "isClosed", true);
+
+        // A closed disruptor must fail the publish rather than silently leave 
the caller's task registered.
+        Assertions.assertThrows(JobException.class, () -> 
disruptor.tryPublishTask(7L));
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to