github-actions[bot] commented on code in PR #67754:
URL: https://github.com/apache/doris/pull/67754#discussion_r4140477863


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobManager.java:
##########
@@ -498,6 +543,122 @@ private void downgradeRunningRefresh(LanceIndexJob 
sweepCandidate) {
         }
     }
 
+    /**
+     * Durable FORCE_RELEASE of an UNKNOWN job, the operator escape hatch for a
+     * mutation whose outcome cannot be established. Validation order is fixed:
+     * unknown job -> idempotent short-circuit on an already released record
+     * (deliberately ahead of the revision CAS, so an operator retry carrying 
the
+     * pre-release revision still observes the existing release record, unlike
+     * the stale-callback convention) -> revision CAS -> UNKNOWN state 
gate
+     * (a null mutation state is treated as UNKNOWN, the safe direction, 
mirroring
+     * {@link LanceIndexJob#isUnresolved()}). The short-circuit only honors a
+     * released record in a null/UNKNOWN state: a record that is force-released
+     * on a non-UNKNOWN state can only come from a corrupt journal or image, so
+     * it falls through to the CAS and the state gate instead (fail-closed).
+     *
+     * <p>The staged copy sets only the five FORCE audit fields and bumps
+     * revision and update time; {@code possibleLiveOwned} and
+     * {@code terminationProof} keep their values for audit. Swapping the 
record
+     * in flips {@link LanceIndexJob#isUnresolved()} and
+     * {@link LanceIndexJob#holdsPossibleLiveSlot()} to false, so
+     * {@link #applyToMemory(LanceIndexJob)} releases the fence, the quota
+     * charge, and the corrupt-admission blocker with no explicit teardown.
+     *
+     * @return false (with a warning) on an unknown job id, a revision 
mismatch,
+     * or a non-UNKNOWN state; true when the job is released or already was
+     * @throws DdlException when a bounded FORCE text field exceeds its limit
+     */
+    public boolean forceRelease(long jobId, long expectedRevision, String 
actor, String note, String warning)
+            throws DdlException {
+        writeLock();
+        try {
+            LanceIndexJob current = jobs.get(jobId);
+            if (current == null) {
+                LOG.warn("reject force release of unknown lance index job {}", 
jobId);
+                return false;
+            }
+            if (current.isForceReleased() && (current.getMutationState() == 
null
+                    || current.getMutationState() == 
LanceIndexJobMutationState.UNKNOWN)) {
+                return true;
+            }
+            if (current.getRevision() != expectedRevision) {
+                LOG.warn("reject force release of lance index job {}: expected 
revision {}, current {}",
+                        jobId, expectedRevision, current);
+                return false;
+            }
+            if (current.getMutationState() != null
+                    && current.getMutationState() != 
LanceIndexJobMutationState.UNKNOWN) {
+                LOG.warn("reject force release of lance index job {} in 
mutation state {}: only UNKNOWN may be"
+                        + " force-released", jobId, 
current.getMutationState());
+                return false;
+            }
+            long now = System.currentTimeMillis();
+            LanceIndexJob updated = new LanceIndexJob(current);
+            try {
+                updated.setForceReleased(true);
+                updated.setForceActor(actor);
+                updated.setForceTimeMs(now);
+                updated.setForceNote(note);
+                updated.setForceWarning(warning);
+            } catch (IllegalArgumentException e) {
+                throw new DdlException("invalid lance index job force release: 
" + e.getMessage(), e);
+            }
+            updated.setRevision(current.getRevision() + 1);
+            updated.setUpdateTimeMs(now);
+            writeEditLog(updated);
+            applyToMemory(updated);
+            return true;
+        } finally {
+            writeUnlock();
+        }
+    }
+
+    /**
+     * Retention GC of resolved job records, master-only: every record that
+     * stayed resolved for longer than {@code keepMs} (measured on the durable
+     * update time, bumped by every durable transition including the force
+     * release) is removed on all FEs through one batch edit-log record per
+     * round, so every FE serves the same SHOW LANCE INDEX JOBS view. An
+     * unresolved record is never removed, no matter its age (fail-closed). At
+     * most {@code maxPerRound} jobs are removed per round, oldest first.
+     *
+     * @return the ids actually removed, in oldest-first order
+     */
+    public List<Long> removeResolvedJobsOlderThan(long keepMs, int 
maxPerRound) {
+        writeLock();
+        try {
+            long now = System.currentTimeMillis();
+            List<LanceIndexJob> expired = new ArrayList<>();
+            for (LanceIndexJob job : jobs.values()) {
+                if (job != null && !job.isUnresolved() && now - 
job.getUpdateTimeMs() > keepMs) {

Review Comment:
   [P2] Keep slot-owning jobs out of retention GC. A normal NATIVE_OK report 
can finish refresh without a CHILD_REAPED proof, leaving 
holdsPossibleLiveSlot() true even though isUnresolved() is false. Once the keep 
window expires, this predicate deletes the only record counted by the per-BE 
capacity limit; a later proof is also dropped because the job is gone. With 
keep_max_second=1, this can occur as soon as that window elapses. Retain the 
record or a separate slot tombstone until the worker is proven gone, and age 
retention from that final transition.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobCleaner.java:
##########
@@ -0,0 +1,65 @@
+// 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.datasource.lance.job;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.util.MasterDaemon;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+
+/**
+ * Master-only retention sweeper of resolved Lance index job records. Resolved
+ * records serve audit only; once one has been resolved for longer than
+ * {@code Config.lance_index_job_keep_max_second} it is removed on all FEs
+ * through one batch edit-log record per round, so every FE serves the same
+ * SHOW LANCE INDEX JOBS view. Unresolved records are never removed, no matter
+ * their age (fail-closed).
+ *
+ * <p>MasterDaemon itself performs no master check: the master-only semantics
+ * come entirely from being started in {@code 
Env.startMasterOnlyDaemonThreads}.
+ * Both retention configs are mutable and re-read every round.

Review Comment:
   [P2] Drain expired jobs at a rate the dispatcher can sustain. This cap 
removes only 1,024 records per default one-hour cleaner interval. Even with two 
BEs on separate hosts and the default two-slot cap, the current stub can 
resolve four clean rejections every 10-second round, or 1,440 jobs/hour. After 
retention, at least 416 stale records/hour accumulate indefinitely; the audit 
map, image, and every dispatcher scan keep growing. Use bounded repeated 
batches or a cleanup rate matched to supported throughput.



##########
fe/fe-common/src/main/java/org/apache/doris/common/Config.java:
##########
@@ -4276,4 +4276,72 @@ public void handle(Field field, String value) throws 
Exception {
                     "Static upper bound for num_sub_vectors of Lance IVF_PQ 
indexes."})
     public static int lance_index_max_num_sub_vectors = 256;
 
+    @ConfField(mutable = true, masterOnly = true,
+            callback = 
LanceIndexConfigValidator.PositiveIntConfigHandler.class,
+            description = {"Lance 索引 job 派发器(含 deadline/possible-live 扫掠与 
refresh 驱动)的轮询周期(秒)。",
+                    "Polling interval in seconds of the Lance index job 
dispatcher "
+                    + "(dispatch sweep, deadline/possible-live sweeps, and 
refresh driver)."})
+    public static int lance_index_job_dispatch_interval_second = 10;
+
+    @ConfField(mutable = true, masterOnly = true,
+            callback = 
LanceIndexConfigValidator.PositiveLongConfigHandler.class,
+            description = {"单个 Lance 索引 job 派发后的结果等待上限(秒)。到期仍无完整可信结果即收敛为 
UNKNOWN;"
+                    + "该期限只限定等待,不证明终止,也不释放 possible-live 槽位。",
+                    "Wait bound in seconds for the result of one dispatched 
Lance index job. Expiry without "
+                    + "a complete trusted result converges the job to UNKNOWN; 
the deadline bounds the wait "
+                    + "only, never proves termination, and never releases a 
possible-live slot."})
+    public static long lance_index_job_execute_deadline_second = 3600;
+
+    @ConfField(mutable = true, masterOnly = true,
+            callback = 
LanceIndexConfigValidator.PositiveIntConfigHandler.class,
+            description = {"派发器单轮最多新派发的 Lance 索引 job 数(背压上限)。",
+                    "Maximum number of Lance index jobs newly dispatched per 
dispatcher round (backpressure)."})
+    public static int lance_index_job_max_dispatch_per_round = 16;
+
+    @ConfField(mutable = true, masterOnly = true,
+            callback = 
LanceIndexConfigValidator.PositiveIntConfigHandler.class,
+            description = {"单个 BE 上允许同时在途(RUNNING)的 Lance 索引 job 数上限。",
+                    "Maximum number of in-flight (RUNNING) Lance index jobs 
per backend."})
+    public static int lance_index_job_max_inflight_per_backend = 2;
+
+    @ConfField(mutable = true, masterOnly = true,
+            callback = 
LanceIndexConfigValidator.PositiveIntConfigHandler.class,
+            description = {"refresh 失败的 Lance 索引 job 的最小重试间隔(秒);首次刷新不受此间隔限制。",
+                    "Minimum retry interval in seconds for a terminal Lance 
index job whose metadata "
+                    + "refresh FAILED; the first refresh attempt is never 
delayed by this interval."})
+    public static int lance_index_job_refresh_retry_second = 300;
+
+    @ConfField(mutable = true, masterOnly = true, description = {
+            "暂停 Lance 索引 job 的派发阶段(运维与测试屏障,默认关闭)。暂停只影响派发:deadline 与 "
+                    + "possible-live 扫掠、refresh 驱动照常运行;派发器在派发阶段入口和每个 job 尝试前"
+                    + "检查本开关,被跳过的 job 保持 PENDING 且不消耗单轮派发额度;恢复后继续派发。",
+            "Pause switch for the dispatch phase of Lance index jobs (operator 
and test barrier, "
+                    + "disabled by default). Pausing affects dispatch only: 
the deadline and "
+                    + "possible-live sweeps and the refresh driver keep 
running. The dispatcher checks "
+                    + "this switch at the dispatch-phase entry and before 
every job attempt; skipped "
+                    + "jobs stay PENDING and never consume the per-round 
dispatch budget. Dispatch "
+                    + "resumes once unpaused."})
+    public static boolean lance_index_job_dispatcher_paused = false;

Review Comment:
   [P2] Give the pause switch cross-thread visibility. ADMIN SET writes this 
plain static boolean under ConfigBase's monitor, but the dispatcher reads it 
without that monitor or a volatile/atomic field. A completed pause SET 
therefore does not establish the documented hard barrier: a daemon may keep 
reading stale false and dispatch a job admitted afterward. The local-file 
assertion has the same issue. Use a volatile/atomic field or a shared 
synchronization path for these runtime switches.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ResolveLanceIndexJobCommand.java:
##########
@@ -0,0 +1,367 @@
+// 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;
+
+import org.apache.doris.catalog.DatabaseIf;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.TableIf;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.DdlException;
+import org.apache.doris.common.ErrorCode;
+import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.CatalogMgr;
+import org.apache.doris.datasource.ExternalDatabase;
+import org.apache.doris.datasource.ExternalTable;
+import org.apache.doris.datasource.lance.LanceExternalCatalog;
+import org.apache.doris.datasource.lance.LanceIndexMutationValidator;
+import org.apache.doris.datasource.lance.job.LanceIndexJob;
+import org.apache.doris.datasource.lance.job.LanceIndexJobManager;
+import org.apache.doris.datasource.lance.job.LanceIndexJobMutationState;
+import org.apache.doris.mysql.privilege.PrivPredicate;
+import org.apache.doris.nereids.trees.plans.PlanType;
+import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.StmtExecutor;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.nio.charset.StandardCharsets;
+
+/**
+ * RESOLVE LANCE INDEX JOB &lt;jobId&gt; AS FORCE_RELEASE COMMENT 
'&lt;note&gt;' — the operator
+ * escape hatch that durably releases a job whose mutation outcome is UNKNOWN 
(design section
+ * 7.1). RESOLVE is deliberately not gated by {@code 
enable_lance_index_mutation}: the gate
+ * controls mutation admission, while FORCE must stay available exactly when 
the gate is off.
+ *
+ * <p>The release protocol keeps the fence, the quota charge and the 
possible-live slot while
+ * it performs one authoritative latest-metadata read and one external-table 
refresh with the
+ * current credentials of the surviving catalog, both outside every 
catalog/manager lock; only
+ * then does the durable release transfer inside the admission critical section
+ * ({@code captureLanceIndexTarget} → lock-free read/refresh → {@code 
withLanceIndexAdmission}
+ * recheck → manager write lock), serialized against DROP CATALOG and identity 
ALTER exactly
+ * like admission. Any failure before the transfer is the typed
+ * {@code ERR_LANCE_INDEX_JOB_RESOLUTION_INCOMPLETE}: nothing is written, 
nothing is released,
+ * and the operator fixes the cause and retries the same statement.
+ *
+ * <p>Target resolution (design section 7.1 step 1) is three-valued. RESOLVED 
means the
+ * persisted names resolve and the catalog's current durable dataset locator 
still matches
+ * the job's — the same revalidation SHOW LANCE INDEX JOBS applies, so a 
repointed dataset
+ * reusing the same names never turns a stale name into table-level 
authorization. MISSING
+ * means the catalog, database or table is verifiably absent, or the locator 
positively
+ * points at a different dataset: that is the orphan family — a full orphan 
(catalog gone)
+ * has no credentials to read with and nothing to invalidate, so it is 
released directly
+ * after global ADMIN authorization, while a half-orphan skips the 
authoritative read and
+ * refreshes with {@code ignoreIfNotExists=true} as a best-effort 
invalidation. FAILED means
+ * a resolution that errors out, or a locator that cannot be resolved right 
now: never an
+ * orphan verdict — after ADMIN authorization the statement fails with the 
typed 5105 so the
+ * fence is kept when "table gone" cannot be told apart from "network down". 
SHOW fails the
+ * same uncertainty closed by hiding the row; RESOLVE fails it closed by not 
releasing.
+ *
+ * <p>Non-disclosure (design section 8): the job is loaded first and 
authorized against its
+ * persisted target — table-level ALTER when the target resolves, global ADMIN 
otherwise — and
+ * a missing job and an unauthorized job share the same fixed 
ERR_LANCE_INDEX_JOB_NOT_FOUND
+ * response naming only the job id. The 5104 state rejection and the 5105 
resolution failure
+ * are only visible to an already authorized caller.
+ *
+ * <p>Success returns an OK packet carrying one warning row with {@link 
#LATE_COMMIT_WARNING},
+ * the same text persisted as the job's durable {@code forceWarning}: the old 
worker may still
+ * overwrite, remove, or reintroduce the index name; the mutation outcome 
remains UNKNOWN.
+ * Retrying FORCE on an already released job is an idempotent success 
returning the existing
+ * release record, never an error.
+ */
+public class ResolveLanceIndexJobCommand extends Command implements 
ForwardWithSync {
+    /**
+     * The late-commit warning (design section 7.1), returned in the OK packet 
and persisted
+     * verbatim as the durable {@code forceWarning}; bounded well under
+     * {@link LanceIndexJob#MAX_FORCE_TEXT_BYTES}.
+     */
+    static final String LATE_COMMIT_WARNING =
+            "the old worker may still overwrite, remove, or reintroduce the 
index name; "
+                    + "the mutation outcome remains UNKNOWN";
+
+    private static final Logger LOG = 
LogManager.getLogger(ResolveLanceIndexJobCommand.class);
+
+    private final long jobId;
+    private final String comment;
+
+    public ResolveLanceIndexJobCommand(long jobId, String comment) {
+        super(PlanType.RESOLVE_LANCE_INDEX_JOB_COMMAND);
+        this.jobId = jobId;
+        this.comment = comment;
+    }
+
+    public long getJobId() {
+        return jobId;
+    }
+
+    public String getComment() {
+        return comment;
+    }
+
+    @Override
+    public void run(ConnectContext ctx, StmtExecutor executor) throws 
Exception {
+        Env env = Env.getCurrentEnv();
+        LanceIndexJobManager manager = env.getLanceIndexJobManager();
+        // 1. Load the job without disclosing any field (design section 7.1 
step 1).
+        LanceIndexJob job = manager.getJob(jobId);
+        if (job == null) {
+            throw notFound();
+        }
+        // 2. Resolve and authorize against the persisted target before any 
state is revealed:
+        //    table-level ALTER when the target resolves, global ADMIN for the 
orphan family
+        //    and for a target whose resolution failed outright.
+        CatalogMgr catalogMgr = env.getCatalogMgr();
+        CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog = 
catalogMgr.getCatalog(job.getCatalogId());
+        TargetResolution resolution = resolveTarget(catalog, job);
+        boolean authorized = resolution == TargetResolution.RESOLVED
+                ? env.getAccessManager().checkTblPriv(ctx, catalog.getName(), 
job.getDbName(), job.getTableName(),
+                        PrivPredicate.ALTER)
+                : env.getAccessManager().checkGlobalPriv(ctx, 
PrivPredicate.ADMIN);
+        if (!authorized) {
+            throw notFound();
+        }
+        // 3. Idempotent replay: a retry returns the existing release record 
(section 7.1).
+        //    This deliberately precedes the resolution-failure rejection: 
once the release
+        //    has landed, a retry during a provider outage is a success, not a 
5105.
+        if (job.isForceReleased()) {
+            ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING);
+            return;
+        }
+        // 4. Only UNKNOWN may be force-released; a null mutation state reads 
as UNKNOWN,
+        //    same as the manager's own gate. The state rejection also 
precedes the
+        //    resolution-failure rejection: for a terminal job the accurate 
answer is 5104,
+        //    not a 5105 claiming the job still holds its fence.
+        if (job.getMutationState() != null && job.getMutationState() != 
LanceIndexJobMutationState.UNKNOWN) {
+            throw new 
AnalysisException(ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN.formatErrorMsg(jobId),
+                    ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN);
+        }
+        if (resolution == TargetResolution.FAILED) {
+            // Never an orphan verdict: "table gone" cannot be told apart from 
"network down",
+            // so nothing is released and nothing beyond the typed error is 
disclosed; the
+            // operator fixes the cause and retries the same statement (design 
7.1 step 4).
+            throw incompleteResolution("the persisted target could not be 
resolved with current catalog"
+                    + " metadata; see fe.log for the cause");
+        }
+        boolean targetResolves = resolution == TargetResolution.RESOLVED;
+        // 5. The grammar makes COMMENT mandatory; here the note must also be 
non-empty after
+        //    trimming and fit the durable force text bound.
+        String note = comment == null ? "" : comment.trim();
+        if (note.isEmpty()) {
+            throw new AnalysisException("force release note must not be empty",
+                    ErrorCode.ERR_LANCE_INDEX_INVALID);
+        }
+        if (note.getBytes(StandardCharsets.UTF_8).length > 
LanceIndexJob.MAX_FORCE_TEXT_BYTES) {
+            throw new AnalysisException("force release note exceeds " + 
LanceIndexJob.MAX_FORCE_TEXT_BYTES
+                    + " UTF-8 bytes", ErrorCode.ERR_LANCE_INDEX_INVALID);
+        }
+        // 6-9. Branch on the orphan state, then the durable release transfer.
+        String actor = ctx.getQualifiedUser();
+        boolean released;
+        if (catalog == null) {
+            // Full orphan: no credentials survive to read with and nothing 
can be
+            // invalidated, so the release goes straight to the manager write 
lock.
+            released = manager.forceRelease(jobId, job.getRevision(), actor, 
note, LATE_COMMIT_WARNING);
+        } else {
+            released = releaseWithLiveCatalog(env, catalogMgr, manager, 
catalog, targetResolves, job, actor, note);
+        }
+        if (!released) {
+            // 10. The expected-revision transfer lost a race. A concurrent 
FORCE_RELEASE
+            //     that already landed makes this an idempotent success; 
anything else means
+            //     the job left UNKNOWN concurrently (UNKNOWN has no other 
outgoing
+            //     transition), so the pinned not-UNKNOWN wording stays 
accurate.
+            LanceIndexJob reread = manager.getJob(jobId);
+            if (reread != null && reread.isForceReleased()) {
+                ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING);
+                return;
+            }
+            throw new 
AnalysisException(ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN.formatErrorMsg(jobId),
+                    ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN);
+        }
+        // 11. The OK packet carries the late-commit warning; it survives the 
forward chain
+        //     byte-identically (proxyExecute serializes the master state).
+        ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING);

Review Comment:
   [P2] Make the late-commit warning retrievable when advertising one warning. 
This success path sets warningRows=1 and puts the text only in the OK packet 
info field; QueryState records no warning entry, and Nereids SHOW WARNINGS 
always returns an empty result set. A client that follows the reported warning 
count with SHOW WARNINGS therefore gets no row containing the risk that the old 
worker may still commit. Return a real warning row or use a response contract 
that consistently exposes the warning text.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,657 @@
+// 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.datasource.lance.job;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.ClientPool;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.util.MasterDaemon;
+import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.lance.LanceExternalCatalog;
+import org.apache.doris.datasource.lance.storage.LanceStorageOptions;
+import org.apache.doris.persist.gson.GsonUtils;
+import org.apache.doris.system.Backend;
+import org.apache.doris.system.BeSelectionPolicy;
+import org.apache.doris.system.SystemInfoService;
+import org.apache.doris.thrift.BackendService;
+import org.apache.doris.thrift.TLanceIndexJobDispatch;
+import org.apache.doris.thrift.TLanceIndexMutationType;
+import org.apache.doris.thrift.TNetworkAddress;
+import org.apache.doris.thrift.TStatus;
+import org.apache.doris.thrift.TStatusCode;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.apache.thrift.TApplicationException;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+import java.util.function.Supplier;
+
+/**
+ * Master-only daemon that drives the durable Lance index job records through
+ * the lifecycle after admission. Each round runs in a fixed order: converge
+ * expired RUNNING jobs to UNKNOWN, release possible-live slots whose backend
+ * process was replaced, drive the refresh a terminal job still owes, then
+ * dispatch PENDING jobs. Every durable transition goes through
+ * {@link LanceIndexJobManager} under its own lock; the daemon holds no catalog
+ * or manager lock across any call.
+ *
+ * <p>The daemon does not read the admission gate: a job that is already 
durable
+ * must be driven to its terminal state, whatever the gate says now, so the
+ * thread runs unconditionally on the master and simply finds nothing to do
+ * while no jobs exist. An idle round writes no journal record.
+ *
+ * <p>Dispatch follows the durable-before-send boundary: the whole request is
+ * prepared first (so a preparation failure just leaves the job PENDING), then
+ * the markRunning edit log is written and re-read before the first byte of
+ * network I/O, and the invocation id of an attempt that lost the 
compare-and-set
+ * is never reused. After a successful markRunning there is exactly one send;
+ * from that point a job converges only through a matching result callback, the
+ * deadline sweep, or the epoch sweep, never through a resend. A failure that
+ * still proves the dispatch was never enqueued (a clean pre-enqueue error
+ * status, a client-pool borrow failure, or an UNKNOWN_METHOD answer from an
+ * old backend) converges it NOT_COMMITTED through the no-enqueue channel,
+ * which releases the possible-live slot in the same durable transition;
+ * anything ambiguous after the invocation may have started converges UNKNOWN
+ * with the slot retained.
+ *
+ * <p>The manager is resolved from the supplier once per round rather than
+ * captured at construction: {@code Env.loadLanceIndexJobManager} replaces the
+ * Env-owned manager with a brand-new object on every image load, so a cached
+ * reference would keep scanning the abandoned pre-image manager after an FE
+ * restart while replay, admission and SHOW all move on to the restored one.
+ * Every phase of one round shares the single resolved instance.
+ *
+ * <p>The sleep between rounds is sliced at {@link #MAX_SLEEP_SLICE_MS} so a
+ * shortened polling interval takes effect within one slice (see the field
+ * javadoc), and {@link Config#lance_index_job_dispatcher_paused} suspends only
+ * the dispatch phase (see {@link #dispatchPendingJobs}).
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    /**
+     * Upper bound of one sleep slice, equal to the shipped default interval. 
The
+     * daemon never sleeps longer than this, so a shortened
+     * {@link Config#lance_index_job_dispatch_interval_second} takes effect 
within
+     * one slice instead of waiting out a previously adopted long sleep: the
+     * elapsed check in {@link #runAfterCatalogReady} is re-evaluated against 
the
+     * current config at every wake. Slices bound only the sleep; rounds still
+     * honor the configured interval, because a wake whose configured interval
+     * (longer than this bound) has not elapsed since the last round skips the
+     * round. A lengthened interval takes effect at the next wake through the 
same
+     * check, and an interval at or below this bound needs no check at all — 
every
+     * wake runs a round, exactly one per configured period.
+     */
+    private static final long MAX_SLEEP_SLICE_MS = 10_000L;
+
+    private final Supplier<LanceIndexJobManager> jobManagerSupplier;
+
+    /** Wall time of the last executed round, or -1 before the first one. */
+    private long lastRoundMs = -1L;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        this(() -> jobManager);
+    }
+
+    public LanceIndexJobDispatcher(Supplier<LanceIndexJobManager> 
jobManagerSupplier) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManagerSupplier = jobManagerSupplier;
+    }
+
+    /**
+     * Values loaded from fe.conf bypass the config validator (only ADMIN SET 
runs
+     * it), so the positive invariant is re-asserted where a non-positive value
+     * would break the loop: a non-positive interval would kill this thread 
inside
+     * {@code Thread.sleep} or busy-spin it, a non-positive deadline would 
sweep
+     * every dispatched job UNKNOWN on the next round, and a zero cap would 
stall
+     * dispatch forever. The refresh retry interval needs no such defense: a
+     * non-positive value simply disengages the throttle.
+     */
+    private static long dispatchIntervalMs() {
+        return Math.max(1, Config.lance_index_job_dispatch_interval_second) * 
1000L;
+    }
+
+    private static long executeDeadlineMs(long nowMs) {
+        long second = Math.max(1L, 
Config.lance_index_job_execute_deadline_second);
+        return second > (Long.MAX_VALUE - nowMs) / 1000L ? Long.MAX_VALUE : 
nowMs + second * 1000L;
+    }
+
+    @Override
+    protected void runAfterCatalogReady() {
+        if (!Env.getCurrentEnv().isMaster()) {
+            return;
+        }
+        if (Env.isCheckpointThread()) {
+            return;
+        }
+        long configuredMs = dispatchIntervalMs();
+        setInterval(Math.min(configuredMs, MAX_SLEEP_SLICE_MS));
+        if (configuredMs > MAX_SLEEP_SLICE_MS && lastRoundMs >= 0 && nowMs() - 
lastRoundMs < configuredMs) {
+            // A wake inside a long configured interval: the slice elapsed, the
+            // round period has not. Skipping is cheap and writes no journal 
record.
+            return;
+        }
+        lastRoundMs = nowMs();
+        try {
+            runOneRound(jobManagerSupplier.get());
+        } catch (Throwable t) {
+            LOG.warn("Failed to process one round of the lance index job 
dispatcher", t);
+        }
+    }
+
+    /** Clock seam for the round-period check; tests advance it instead of 
sleeping. */
+    protected long nowMs() {
+        return System.currentTimeMillis();
+    }
+
+    private void runOneRound(LanceIndexJobManager jobManager) {
+        long nowMs = System.currentTimeMillis();
+        sweepExpiredRunningJobs(jobManager, nowMs);
+        sweepReplacedProcessEpochs(jobManager);
+        driveRequiredRefreshes(jobManager, nowMs);
+        dispatchPendingJobs(jobManager);
+    }
+
+    /**
+     * Deadline sweep. A RUNNING job past its wait deadline has produced no
+     * complete trusted result, so it converges to UNKNOWN through the same
+     * completeWithResult channel a callback would use. Expiry bounds the wait
+     * only: it never proves termination, so the possible-live slot, the
+     * same-name fence, and the unresolved quota all stay held.
+     */
+    private void sweepExpiredRunningJobs(LanceIndexJobManager jobManager, long 
nowMs) {
+        for (LanceIndexJob job : jobManager.getExpiredRunningJobs(nowMs)) {
+            try {
+                boolean completed = 
jobManager.completeWithResult(job.getJobId(),
+                        dispatchRevisionOf(job), job.getInvocationId(), 
job.getBeProcessEpoch(),
+                        new 
LanceIndexJobResult(LanceIndexJobResultCode.NO_TRUSTED_RESULT,
+                                LanceIndexJobCompletionReason.NONE,
+                                "execute deadline expired without a complete 
trusted result", false));
+                if (completed) {
+                    LOG.info("lance index job {} converged RUNNING -> UNKNOWN 
on deadline expiry",
+                            job.getJobId());
+                } else {
+                    LOG.warn("deadline sweep skipped lance index job {}: 
already converged by a callback or sweep",
+                            job.getJobId());
+                }
+            } catch (Throwable t) {
+                LOG.warn("failed to sweep expired lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    /**
+     * Possible-live sweep. The only slot-release proof this daemon produces is
+     * that the recorded backend process epoch no longer exists: a backend 
entry
+     * reporting a different epoch proves the process that received the 
dispatch
+     * was replaced. A missing backend entry or heartbeat loss proves nothing
+     * (the worker may still be running behind a partition), so such a job 
keeps
+     * its slot until a stronger proof or an operator force release. An epoch
+     * change also proves nothing about the outcome, so the mutation state is
+     * never touched here.
+     */
+    private void sweepReplacedProcessEpochs(LanceIndexJobManager jobManager) {
+        for (LanceIndexJob job : jobManager.getJobsHoldingPossibleLiveSlot()) {
+            try {
+                Backend backend = 
Env.getCurrentSystemInfo().getBackend(job.getBackendId());
+                if (backend == null || backend.getProcessEpoch() == 
job.getBeProcessEpoch()) {
+                    continue;
+                }
+                boolean recorded = 
jobManager.recordTerminationProof(job.getJobId(),
+                        dispatchRevisionOf(job), job.getBackendId(), 
job.getBeProcessEpoch(),
+                        job.getInvocationId(), 
LanceIndexTerminationProof.BE_PROCESS_EPOCH_GONE);
+                if (recorded) {
+                    LOG.info("released possible-live slot of lance index job 
{}: backend process epoch was replaced",
+                            job.getJobId());
+                } else {
+                    LOG.warn("epoch sweep skipped lance index job {}: dispatch 
identity already moved",
+                            job.getJobId());
+                }
+            } catch (Throwable t) {
+                LOG.warn("failed to sweep possible-live slot of lance index 
job " + job.getJobId(), t);
+            }
+        }
+    }
+
+    /**
+     * Refresh driver for terminal jobs with an unfinished refresh obligation.
+     * Completing the refresh is the protocol duty that releases the same-name
+     * fence and the unresolved quota once DONE; it is not a read-visibility
+     * action, because index metadata is never cached. Each job is driven
+     * through markRefreshRunning, the idempotent external-table refresh, then
+     * DONE or FAILED: a FAILED job keeps its fence and is retried, throttled 
to
+     * one attempt per retry interval, while a first REQUIRED refresh is never
+     * delayed. UNKNOWN jobs never appear here; they owe no refresh.
+     */
+    private void driveRequiredRefreshes(LanceIndexJobManager jobManager, long 
nowMs) {
+        for (LanceIndexJob job : jobManager.getJobsNeedingRefresh()) {
+            try {
+                if (job.getRefreshState() == 
LanceIndexJobRefreshState.RUNNING) {
+                    // In flight elsewhere; the master-transfer sweep 
downgrades a stale
+                    // RUNNING back to REQUIRED, so a lost driver cannot 
strand it.
+                    continue;
+                }
+                if (job.getRefreshState() == LanceIndexJobRefreshState.FAILED
+                        && nowMs - job.getUpdateTimeMs()
+                                < Config.lance_index_job_refresh_retry_second 
* 1000L) {
+                    continue;
+                }
+                if (!jobManager.markRefreshRunning(job.getJobId(), 
job.getRevision())) {
+                    // A concurrent driver won the compare-and-set; nothing to 
do here.
+                    continue;
+                }
+                driveOneRefresh(jobManager, job);
+            } catch (Throwable t) {
+                LOG.warn("failed to drive the refresh of lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private void driveOneRefresh(LanceIndexJobManager jobManager, 
LanceIndexJob job) {
+        long refreshRevision = job.getRevision() + 1;
+        CatalogIf catalog = 
Env.getCurrentEnv().getCatalogMgr().getCatalog(job.getCatalogId());
+        if (catalog == null) {
+            // Unreachable while the unresolved-job guard blocks catalog 
drops; kept as a
+            // fail-closed fallback so the job still transitions and retries 
later.
+            LOG.warn("catalog of lance index job {} is gone; marking its 
refresh FAILED", job.getJobId());
+            finishRefreshTransition(jobManager, job.getJobId(), 
refreshRevision, false);
+            return;
+        }
+        try {
+            // A half-orphan target (its db or table already dropped 
externally) is a
+            // silent no-op: nothing is left to invalidate, and DONE is the 
correct end
+            // state for the job.
+            
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(catalog.getName(),
+                    job.getDbName(), job.getTableName(), true);
+        } catch (Throwable t) {
+            // The typed DdlException is the expected failure; an unchecked 
exception out
+            // of the metadata path must still leave the durable refresh 
state, or the
+            // job would strand in refresh RUNNING until the next master 
transfer.
+            LOG.warn("refresh of lance index job {} failed; keeping the fence 
for a retry",
+                    job.getJobId(), t);
+            finishRefreshTransition(jobManager, job.getJobId(), 
refreshRevision, false);
+            return;
+        }
+        finishRefreshTransition(jobManager, job.getJobId(), refreshRevision, 
true);
+    }
+
+    /**
+     * Applies the DONE/FAILED transition with a bounded revision retry. A 
concurrent
+     * termination-proof write can bump the revision after markRefreshRunning 
succeeded,
+     * and silently losing that compare-and-set would leave the refresh 
RUNNING — a
+     * state only the master-transfer sweep downgrades. Re-reading the 
revision and
+     * retrying a few times converges it; a persistent loss is escalated.
+     */
+    private void finishRefreshTransition(LanceIndexJobManager jobManager, long 
jobId, long expectedRevision,
+            boolean done) {
+        long revision = expectedRevision;
+        for (int attempt = 0; attempt < 3; attempt++) {
+            boolean transitioned = done ? jobManager.markRefreshDone(jobId, 
revision)
+                    : jobManager.markRefreshFailed(jobId, revision);
+            if (transitioned) {
+                return;
+            }
+            LanceIndexJob fresh = jobManager.getJob(jobId);
+            if (fresh == null) {
+                break;
+            }
+            revision = fresh.getRevision();
+        }
+        LOG.error("lance index job {} kept its refresh RUNNING: the 
DONE/FAILED transition kept losing the"
+                + " compare-and-set; the master-transfer sweep will downgrade 
it", jobId);
+    }
+
+    /**
+     * PENDING dispatch. Makes at most
+     * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches 
per
+     * round, and only a job this round actually made RUNNING consumes that
+     * budget: skipped jobs (an eligibility gate is closed, or every backend is
+     * at capacity) are scanned past, so a stable subset of permanently
+     * undispatchable jobs can never crowd out later ids. Per backend it never
+     * exceeds {@link Config#lance_index_job_max_inflight_per_backend}
+     * possible-live worker slots, counted from slot ownership (see
+     * {@link LanceIndexJobManager#countPossibleLiveSlotsByBackend()}) plus the
+     * jobs this round already made RUNNING. A job that cannot be dispatched
+     * keeps waiting as PENDING: there is no dispatch-exhaustion terminal state
+     * and no backoff beyond the daemon period.
+     *
+     * <p>{@link Config#lance_index_job_dispatcher_paused} suspends this phase
+     * only — the sweeps and the refresh driver keep running while it is set.
+     * The switch is checked at the phase entry and again before every single
+     * job attempt, which closes the admission race a test or operator cares
+     * about: anyone who sets the switch <em>before</em> admitting a job is
+     * guaranteed the job is never dispatched while paused. A round whose
+     * snapshot was taken before the admission never sees the job at all, and
+     * any round that can see it performs its per-job check after the
+     * admission, hence after the switch was set, and skips it. A skipped job
+     * never consumes the round's dispatch budget.
+     */
+    private void dispatchPendingJobs(LanceIndexJobManager jobManager) {
+        if (Config.lance_index_job_dispatcher_paused) {
+            return;
+        }
+        int maxPerRound = Math.max(1, 
Config.lance_index_job_max_dispatch_per_round);
+        Map<Long, Integer> inflightByBackend = 
jobManager.countPossibleLiveSlotsByBackend();
+        int dispatched = 0;
+        for (LanceIndexJob job : jobManager.getJobsNeedingDispatch()) {
+            if (Config.lance_index_job_dispatcher_paused) {
+                // Flipped mid-round: stop without touching the budget.
+                break;
+            }
+            if (dispatched >= maxPerRound) {
+                break;
+            }
+            try {
+                if (tryDispatch(jobManager, job, inflightByBackend)) {
+                    dispatched++;
+                }
+            } catch (Throwable t) {
+                LOG.warn("failed to dispatch lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    /**
+     * One dispatch attempt for one PENDING job; returns true only when the
+     * attempt made the job durable RUNNING (and so consumes this round's
+     * dispatch budget). Every early return before markRunning leaves the job
+     * PENDING for a later round: the eligibility gates, the backend and
+     * capacity checks, and also the whole request preparation — storage-option
+     * resolution and the wire request build run before the durable boundary,
+     * so an FE-side failure there (for example a catalog id that resolves to
+     * nothing while ALTER CATALOG RENAME has the catalog temporarily removed)
+     * just retries next round instead of stranding the job UNKNOWN without a
+     * single byte sent. Once markRunning succeeds the job is durable RUNNING
+     * and this invocation id gets exactly one send attempt; after that only a
+     * matching callback, the deadline sweep, or the epoch sweep can converge
+     * the job.
+     */
+    private boolean tryDispatch(LanceIndexJobManager jobManager, LanceIndexJob 
job,
+            Map<Long, Integer> inflightByBackend) {
+        boolean localDataset = isLocalFileDataset(job.getNormalizedLocator());
+        if (localDataset && !Config.enable_lance_index_local_file_mutation) {
+            // Operator assertion is off: a local-filesystem mutation stays 
PENDING.
+            return false;
+        }
+        if (localDataset && Env.getCurrentEnv().getFrontends(null).size() != 
1) {
+            // Local files are only shared by a single-node deployment.
+            return false;
+        }
+        SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+        // All schedule-available backends, shuffled by the selection policy: 
the
+        // first one with a free possible-live slot takes the job, so a full
+        // backend defers this attempt only when every selectable backend is at
+        // the cap, never just because the randomly picked one is.
+        List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(
+                new 
BeSelectionPolicy.Builder().needScheduleAvailable().build(), -1);

Review Comment:
   [P2] Include every usable worker in this backend selection. The policy 
defaults to allowOnSameHost=false, so it randomly hides all but one BE per host 
before the slot check; a full BE can hide a free peer and defer the job. It 
also defaults to preferComputeNode=false, which filters every computation-role 
BE. In a compute-only cluster configured to use those nodes for external scans, 
Lance jobs remain PENDING indefinitely. Select all schedule-available instances 
appropriate for Lance work and cover both topologies with the real selector.



##########
regression-test/suites/external_table_p0/lance/test_lance_index_resolve.groovy:
##########
@@ -0,0 +1,254 @@
+// 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.
+
+// PR3E regression scope: the 3D dispatcher is not delivered yet, so a cluster 
cannot
+// produce a genuine UNKNOWN job and every case in this suite is negative or 
static.
+// Deliberately NOT covered here (all wait for 3D / slice 4): FORCE_RELEASE 
happy
+// path e2e, same-name re-admission after FORCE e2e, quota reclaim after FORCE 
e2e,
+// and expiry GC of resolved jobs e2e.
+
+suite("test_lance_index_resolve", "p0,external,nonConcurrent") {
+    // The Lance fixture is preinstalled in the MinIO container of the Iceberg
+    // external environment, so this suite deliberately shares its switch.
+    String enabled = context.config.otherConfigs.get("enableIcebergTest")
+    if (enabled == null || !enabled.equalsIgnoreCase("true")) {
+        logger.info("disable Lance index resolve test because the Iceberg 
MinIO environment is disabled.")
+        return
+    }
+
+    String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+    String minioPort = context.config.otherConfigs.get("iceberg_minio_port")
+    // Admitted jobs are durable and stay PENDING forever in this delivery 
slice:
+    // dispatch, the FORCE happy path and job GC only land in later slices, so 
their
+    // fences and quota charges can never be released here. Every index name 
(and the
+    // filesystem catalog itself, because fence/quota keys include the 
persisted
+    // catalog id) carries this per-run suffix so that rerunning the suite on a
+    // shared pipeline cluster can never collide with a previous run's 
leftovers.
+    String runSuffix = "${System.currentTimeMillis()}"
+    String filesystemCatalog = "test_lance_index_resolve_${runSuffix}"
+    String user = "test_lance_index_resolve_user"
+    String password = "C123_567p"
+    String tableName = "vs_ivf_pq_f32"
+    String resolveIndexName = "idx_resolve_${runSuffix}"
+    // Job ids come from the cluster-wide Env.getNextId() allocator, so this 
literal
+    // can never name a real job.
+    String absentJobId = "9223372036854775806"
+
+    // The filesystem catalog is fresh per run by construction, and once it 
holds
+    // unresolved jobs DROP CATALOG is guarded, so there is deliberately no 
DROP for
+    // it here.
+    try_sql "DROP USER '${user}'@'%'"
+
+    // All three settings are masterOnly. Read them on the master even if the 
suite's
+    // ordinary JDBC connection points at a follower. SHOW uses the 
experimental
+    // display name for the gate, while ADMIN SET accepts its unprefixed alias.
+    def gateRows = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 
'experimental_enable_lance_index_mutation'"""
+    def keepRows = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 
'lance_index_job_keep_max_second'"""
+    def cleanRows = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 
'lance_index_job_clean_interval_second'"""
+    assertEquals(1, gateRows.size())
+    assertEquals(1, keepRows.size())
+    assertEquals(1, cleanRows.size())
+    String originalGate = gateRows[0][1].toString()
+    String originalKeep = keepRows[0][1].toString()
+    String originalClean = cleanRows[0][1].toString()
+    // The retention configs ship with these documented defaults; asserting 
them here
+    // fails loudly if a shared cluster has drifted instead of silently 
restoring a
+    // non-default value afterwards.
+    assertEquals("604800", originalKeep)
+    assertEquals("3600", originalClean)
+    Throwable suiteFailure = null
+
+    // test { ... exception } always runs on the suite's default connection; 
the
+    // masterOnly config rejections below must be asserted on the master 
connection.
+    def expectMasterSqlException = { String stmt, String substring ->
+        String caught = null
+        try {
+            master_sql(stmt)
+        } catch (Throwable t) {
+            caught = t.toString()
+        }
+        assertTrue(caught != null && caught.contains(substring))
+    }
+
+    try {
+        // Rejections reachable with the mutation gate closed. RESOLVE is
+        // deliberately NOT behind enable_lance_index_mutation: it is the 
operator
+        // escape hatch and must stay usable while admission is gated, 
otherwise an
+        // unresolved job would freeze catalog DDL forever. Pin the gate 
closed so
+        // the cases in this block also prove the statement is not answered 
with the
+        // gate rejection.
+        master_sql """ADMIN SET FRONTEND CONFIG ("enable_lance_index_mutation" 
= "false")"""
+
+        // Wrong AS literal: only FORCE_RELEASE is accepted after AS.
+        test {
+            sql """RESOLVE LANCE INDEX JOB ${absentJobId} AS FORCE COMMENT 
'wrong literal'"""
+            exception "mismatched input"
+        }
+
+        // The AS FORCE_RELEASE clause is mandatory.
+        test {
+            sql """RESOLVE LANCE INDEX JOB ${absentJobId} COMMENT 'missing as 
clause'"""
+            exception "mismatched input"
+        }
+
+        // COMMENT is mandatory at grammar level.
+        test {
+            sql """RESOLVE LANCE INDEX JOB ${absentJobId} AS FORCE_RELEASE"""
+            exception "mismatched input"
+        }
+
+        // An empty COMMENT parses but never changes the observable response: 
the
+        // note check runs after the job lookup and the UNKNOWN state gate, so 
with a
+        // missing job the fixed not-found wording answers first. The dedicated
+        // empty-note rejection needs a genuine UNKNOWN job and stays UT-only
+        // (ResolveLanceIndexJobCommandTest) until 3D can produce one.
+        test {
+            sql """RESOLVE LANCE INDEX JOB ${absentJobId} AS FORCE_RELEASE 
COMMENT ''"""
+            exception "Lance index job not found"
+        }
+
+        // A well-formed RESOLVE of a missing job: the fixed non-disclosing 
5103
+        // wording, with no gate rejection even though the gate is closed.
+        test {
+            sql """RESOLVE LANCE INDEX JOB ${absentJobId} AS FORCE_RELEASE 
COMMENT 'no such job'"""
+            exception "Lance index job not found"
+        }
+
+        // Rejections reachable with an admitted PENDING job. masterOnly 
configs set
+        // through ADMIN SET land on the master node locally, which is where
+        // admission reads them; the finally block below restores the gate no 
matter
+        // where the suite fails.
+        master_sql """ADMIN SET FRONTEND CONFIG ("enable_lance_index_mutation" 
= "true")"""
+
+        sql """
+            CREATE CATALOG `${filesystemCatalog}` PROPERTIES (
+                "type" = "lance",
+                "lance.catalog.type" = "filesystem",
+                "warehouse" = "s3://warehouse/lance",
+                "s3.endpoint" = "http://${externalEnvIp}:${minioPort}";,
+                "s3.access_key" = "admin",
+                "s3.secret_key" = "password",
+                "s3.region" = "us-east-1",
+                "use_path_style" = "true"
+            )
+        """
+
+        // CREATE INDEX is admitted and returns a single-column JobId result 
set with
+        // one row; the job stays PENDING because no dispatcher exists in this 
slice.
+        def createRows = sql """CREATE INDEX `${resolveIndexName}` ON 
`${filesystemCatalog}`.`doris`.`${tableName}` (embedding) USING ANN
+                PROPERTIES("index_type"="IVF_PQ", "metric"="l2", 
"num_partitions"="256", "num_sub_vectors"="16")"""
+        assertEquals(1, createRows.size())
+        assertEquals(1, createRows[0].size())
+        String createJobId = createRows[0][0].toString()
+
+        // RESOLVE authorization is checked against the job's target (table 
ALTER, or
+        // global ADMIN for an orphan), and an unauthorized caller gets 
exactly the
+        // same non-disclosing response as a missing job. A fresh user holding 
only
+        // an unrelated SELECT privilege must therefore see "not found".
+        sql """CREATE USER '${user}'@'%' IDENTIFIED BY '${password}'"""
+        sql """GRANT SELECT_PRIV ON regression_test TO '${user}'@'%'"""
+        if (isCloudMode()) {
+            def clusters = sql "SHOW CLUSTERS"
+            assertTrue(!clusters.isEmpty())
+            sql """GRANT USAGE_PRIV ON CLUSTER `${clusters[0][0]}` TO 
'${user}'@'%'"""
+        }
+
+        connect(user, password, context.config.jdbcUrl) {
+            test {
+                sql """RESOLVE LANCE INDEX JOB ${createJobId} AS FORCE_RELEASE 
COMMENT 'unauthorized caller'"""
+                exception "Lance index job not found"
+            }
+        }
+
+        // FORCE_RELEASE only accepts UNKNOWN jobs: the admitted job is 
PENDING, the
+        // one durable state this slice can produce, so an authorized RESOLVE 
is
+        // rejected with the not-in-UNKNOWN wording.
+        test {
+            sql """RESOLVE LANCE INDEX JOB ${createJobId} AS FORCE_RELEASE 
COMMENT 'still pending'"""
+            exception "cannot be resolved: not in UNKNOWN state"
+        }
+
+        // Retention config smoke: both are masterOnly and guarded by the
+        // positive-long callback, whose rejection message carries the field 
name and
+        // the offending value.
+        expectMasterSqlException("""ADMIN SET FRONTEND CONFIG 
("lance_index_job_keep_max_second" = "0")""",
+                "must be a positive long")
+        expectMasterSqlException("""ADMIN SET FRONTEND CONFIG 
("lance_index_job_keep_max_second" = "-1")""",
+                "must be a positive long")
+        expectMasterSqlException("""ADMIN SET FRONTEND CONFIG 
("lance_index_job_clean_interval_second" = "0")""",
+                "must be a positive long")
+        expectMasterSqlException("""ADMIN SET FRONTEND CONFIG 
("lance_index_job_clean_interval_second" = "-1")""",
+                "must be a positive long")
+
+        // A positive value is accepted and visible immediately on the master.
+        master_sql """ADMIN SET FRONTEND CONFIG 
("lance_index_job_clean_interval_second" = "3601")"""
+        def cleanRowsAfterSet = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 
'lance_index_job_clean_interval_second'"""
+        assertEquals("3601", cleanRowsAfterSet[0][1].toString())
+        master_sql """ADMIN SET FRONTEND CONFIG 
("lance_index_job_clean_interval_second" = "${originalClean}")"""
+
+        // Retention GC never touches an unresolved job: with the keep window 
pinned
+        // to one second, any resolved record older than a second is eligible 
for
+        // deletion at the next clean round, while the admitted PENDING job 
must
+        // survive regardless of its age because the GC predicate fails closed 
on
+        // unresolved records. Waiting out a clean interval is impractical 
here, so
+        // this is a static existence assertion; expiry GC e2e waits for 3D / 
slice 4
+        // (see the header comment).
+        try {
+            master_sql """ADMIN SET FRONTEND CONFIG 
("lance_index_job_keep_max_second" = "1")"""
+            def keepRowsAfterSet = master_sql """ADMIN SHOW FRONTEND CONFIG 
LIKE 'lance_index_job_keep_max_second'"""
+            assertEquals("1", keepRowsAfterSet[0][1].toString())
+            def jobsRows = sql_return_maparray """SHOW LANCE INDEX JOBS FROM 
`${filesystemCatalog}`.`doris`
+                    WHERE TableName = "${tableName}" """
+            def pendingJobRow = jobsRows.find { it.IndexName == 
resolveIndexName }
+            assertTrue(pendingJobRow != null)
+            assertEquals(createJobId, pendingJobRow.JobId.toString())
+            assertEquals("PENDING", pendingJobRow.State.toString())

Review Comment:
   [P2] Pause dispatch before relying on PENDING here. This suite admits an 
eligible S3 job but never sets lance_index_job_dispatcher_paused; a 10-second 
daemon round can reach the current BE stub, receive NOT_IMPLEMENTED_ERROR, and 
durably change the job to NOT_COMMITTED while the intervening SQL runs. The 
late PENDING assertion then fails based on scheduling. Set and restore the 
pause switch around admission as the adjacent suites do.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,657 @@
+// 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.datasource.lance.job;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.ClientPool;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.util.MasterDaemon;
+import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.lance.LanceExternalCatalog;
+import org.apache.doris.datasource.lance.storage.LanceStorageOptions;
+import org.apache.doris.persist.gson.GsonUtils;
+import org.apache.doris.system.Backend;
+import org.apache.doris.system.BeSelectionPolicy;
+import org.apache.doris.system.SystemInfoService;
+import org.apache.doris.thrift.BackendService;
+import org.apache.doris.thrift.TLanceIndexJobDispatch;
+import org.apache.doris.thrift.TLanceIndexMutationType;
+import org.apache.doris.thrift.TNetworkAddress;
+import org.apache.doris.thrift.TStatus;
+import org.apache.doris.thrift.TStatusCode;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.apache.thrift.TApplicationException;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+import java.util.function.Supplier;
+
+/**
+ * Master-only daemon that drives the durable Lance index job records through
+ * the lifecycle after admission. Each round runs in a fixed order: converge
+ * expired RUNNING jobs to UNKNOWN, release possible-live slots whose backend
+ * process was replaced, drive the refresh a terminal job still owes, then
+ * dispatch PENDING jobs. Every durable transition goes through
+ * {@link LanceIndexJobManager} under its own lock; the daemon holds no catalog
+ * or manager lock across any call.
+ *
+ * <p>The daemon does not read the admission gate: a job that is already 
durable
+ * must be driven to its terminal state, whatever the gate says now, so the
+ * thread runs unconditionally on the master and simply finds nothing to do
+ * while no jobs exist. An idle round writes no journal record.
+ *
+ * <p>Dispatch follows the durable-before-send boundary: the whole request is
+ * prepared first (so a preparation failure just leaves the job PENDING), then
+ * the markRunning edit log is written and re-read before the first byte of
+ * network I/O, and the invocation id of an attempt that lost the 
compare-and-set
+ * is never reused. After a successful markRunning there is exactly one send;
+ * from that point a job converges only through a matching result callback, the
+ * deadline sweep, or the epoch sweep, never through a resend. A failure that
+ * still proves the dispatch was never enqueued (a clean pre-enqueue error
+ * status, a client-pool borrow failure, or an UNKNOWN_METHOD answer from an
+ * old backend) converges it NOT_COMMITTED through the no-enqueue channel,
+ * which releases the possible-live slot in the same durable transition;
+ * anything ambiguous after the invocation may have started converges UNKNOWN
+ * with the slot retained.
+ *
+ * <p>The manager is resolved from the supplier once per round rather than
+ * captured at construction: {@code Env.loadLanceIndexJobManager} replaces the
+ * Env-owned manager with a brand-new object on every image load, so a cached
+ * reference would keep scanning the abandoned pre-image manager after an FE
+ * restart while replay, admission and SHOW all move on to the restored one.
+ * Every phase of one round shares the single resolved instance.
+ *
+ * <p>The sleep between rounds is sliced at {@link #MAX_SLEEP_SLICE_MS} so a
+ * shortened polling interval takes effect within one slice (see the field
+ * javadoc), and {@link Config#lance_index_job_dispatcher_paused} suspends only
+ * the dispatch phase (see {@link #dispatchPendingJobs}).
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    /**
+     * Upper bound of one sleep slice, equal to the shipped default interval. 
The
+     * daemon never sleeps longer than this, so a shortened
+     * {@link Config#lance_index_job_dispatch_interval_second} takes effect 
within
+     * one slice instead of waiting out a previously adopted long sleep: the
+     * elapsed check in {@link #runAfterCatalogReady} is re-evaluated against 
the
+     * current config at every wake. Slices bound only the sleep; rounds still
+     * honor the configured interval, because a wake whose configured interval
+     * (longer than this bound) has not elapsed since the last round skips the
+     * round. A lengthened interval takes effect at the next wake through the 
same
+     * check, and an interval at or below this bound needs no check at all — 
every
+     * wake runs a round, exactly one per configured period.
+     */
+    private static final long MAX_SLEEP_SLICE_MS = 10_000L;
+
+    private final Supplier<LanceIndexJobManager> jobManagerSupplier;
+
+    /** Wall time of the last executed round, or -1 before the first one. */
+    private long lastRoundMs = -1L;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        this(() -> jobManager);
+    }
+
+    public LanceIndexJobDispatcher(Supplier<LanceIndexJobManager> 
jobManagerSupplier) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManagerSupplier = jobManagerSupplier;
+    }
+
+    /**
+     * Values loaded from fe.conf bypass the config validator (only ADMIN SET 
runs
+     * it), so the positive invariant is re-asserted where a non-positive value
+     * would break the loop: a non-positive interval would kill this thread 
inside
+     * {@code Thread.sleep} or busy-spin it, a non-positive deadline would 
sweep
+     * every dispatched job UNKNOWN on the next round, and a zero cap would 
stall
+     * dispatch forever. The refresh retry interval needs no such defense: a
+     * non-positive value simply disengages the throttle.
+     */
+    private static long dispatchIntervalMs() {
+        return Math.max(1, Config.lance_index_job_dispatch_interval_second) * 
1000L;
+    }
+
+    private static long executeDeadlineMs(long nowMs) {
+        long second = Math.max(1L, 
Config.lance_index_job_execute_deadline_second);
+        return second > (Long.MAX_VALUE - nowMs) / 1000L ? Long.MAX_VALUE : 
nowMs + second * 1000L;
+    }
+
+    @Override
+    protected void runAfterCatalogReady() {
+        if (!Env.getCurrentEnv().isMaster()) {
+            return;
+        }
+        if (Env.isCheckpointThread()) {
+            return;
+        }
+        long configuredMs = dispatchIntervalMs();
+        setInterval(Math.min(configuredMs, MAX_SLEEP_SLICE_MS));
+        if (configuredMs > MAX_SLEEP_SLICE_MS && lastRoundMs >= 0 && nowMs() - 
lastRoundMs < configuredMs) {
+            // A wake inside a long configured interval: the slice elapsed, the
+            // round period has not. Skipping is cheap and writes no journal 
record.
+            return;
+        }
+        lastRoundMs = nowMs();
+        try {
+            runOneRound(jobManagerSupplier.get());
+        } catch (Throwable t) {
+            LOG.warn("Failed to process one round of the lance index job 
dispatcher", t);
+        }
+    }
+
+    /** Clock seam for the round-period check; tests advance it instead of 
sleeping. */
+    protected long nowMs() {
+        return System.currentTimeMillis();
+    }
+
+    private void runOneRound(LanceIndexJobManager jobManager) {
+        long nowMs = System.currentTimeMillis();
+        sweepExpiredRunningJobs(jobManager, nowMs);
+        sweepReplacedProcessEpochs(jobManager);
+        driveRequiredRefreshes(jobManager, nowMs);
+        dispatchPendingJobs(jobManager);
+    }
+
+    /**
+     * Deadline sweep. A RUNNING job past its wait deadline has produced no
+     * complete trusted result, so it converges to UNKNOWN through the same
+     * completeWithResult channel a callback would use. Expiry bounds the wait
+     * only: it never proves termination, so the possible-live slot, the
+     * same-name fence, and the unresolved quota all stay held.
+     */
+    private void sweepExpiredRunningJobs(LanceIndexJobManager jobManager, long 
nowMs) {
+        for (LanceIndexJob job : jobManager.getExpiredRunningJobs(nowMs)) {
+            try {
+                boolean completed = 
jobManager.completeWithResult(job.getJobId(),
+                        dispatchRevisionOf(job), job.getInvocationId(), 
job.getBeProcessEpoch(),
+                        new 
LanceIndexJobResult(LanceIndexJobResultCode.NO_TRUSTED_RESULT,
+                                LanceIndexJobCompletionReason.NONE,
+                                "execute deadline expired without a complete 
trusted result", false));
+                if (completed) {
+                    LOG.info("lance index job {} converged RUNNING -> UNKNOWN 
on deadline expiry",
+                            job.getJobId());
+                } else {
+                    LOG.warn("deadline sweep skipped lance index job {}: 
already converged by a callback or sweep",
+                            job.getJobId());
+                }
+            } catch (Throwable t) {
+                LOG.warn("failed to sweep expired lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    /**
+     * Possible-live sweep. The only slot-release proof this daemon produces is
+     * that the recorded backend process epoch no longer exists: a backend 
entry
+     * reporting a different epoch proves the process that received the 
dispatch
+     * was replaced. A missing backend entry or heartbeat loss proves nothing
+     * (the worker may still be running behind a partition), so such a job 
keeps
+     * its slot until a stronger proof or an operator force release. An epoch
+     * change also proves nothing about the outcome, so the mutation state is
+     * never touched here.
+     */
+    private void sweepReplacedProcessEpochs(LanceIndexJobManager jobManager) {
+        for (LanceIndexJob job : jobManager.getJobsHoldingPossibleLiveSlot()) {
+            try {
+                Backend backend = 
Env.getCurrentSystemInfo().getBackend(job.getBackendId());
+                if (backend == null || backend.getProcessEpoch() == 
job.getBeProcessEpoch()) {
+                    continue;
+                }
+                boolean recorded = 
jobManager.recordTerminationProof(job.getJobId(),
+                        dispatchRevisionOf(job), job.getBackendId(), 
job.getBeProcessEpoch(),
+                        job.getInvocationId(), 
LanceIndexTerminationProof.BE_PROCESS_EPOCH_GONE);
+                if (recorded) {
+                    LOG.info("released possible-live slot of lance index job 
{}: backend process epoch was replaced",
+                            job.getJobId());
+                } else {
+                    LOG.warn("epoch sweep skipped lance index job {}: dispatch 
identity already moved",
+                            job.getJobId());
+                }
+            } catch (Throwable t) {
+                LOG.warn("failed to sweep possible-live slot of lance index 
job " + job.getJobId(), t);
+            }
+        }
+    }
+
+    /**
+     * Refresh driver for terminal jobs with an unfinished refresh obligation.
+     * Completing the refresh is the protocol duty that releases the same-name
+     * fence and the unresolved quota once DONE; it is not a read-visibility
+     * action, because index metadata is never cached. Each job is driven
+     * through markRefreshRunning, the idempotent external-table refresh, then
+     * DONE or FAILED: a FAILED job keeps its fence and is retried, throttled 
to
+     * one attempt per retry interval, while a first REQUIRED refresh is never
+     * delayed. UNKNOWN jobs never appear here; they owe no refresh.
+     */
+    private void driveRequiredRefreshes(LanceIndexJobManager jobManager, long 
nowMs) {
+        for (LanceIndexJob job : jobManager.getJobsNeedingRefresh()) {
+            try {
+                if (job.getRefreshState() == 
LanceIndexJobRefreshState.RUNNING) {
+                    // In flight elsewhere; the master-transfer sweep 
downgrades a stale
+                    // RUNNING back to REQUIRED, so a lost driver cannot 
strand it.
+                    continue;
+                }
+                if (job.getRefreshState() == LanceIndexJobRefreshState.FAILED
+                        && nowMs - job.getUpdateTimeMs()
+                                < Config.lance_index_job_refresh_retry_second 
* 1000L) {
+                    continue;
+                }
+                if (!jobManager.markRefreshRunning(job.getJobId(), 
job.getRevision())) {
+                    // A concurrent driver won the compare-and-set; nothing to 
do here.
+                    continue;
+                }
+                driveOneRefresh(jobManager, job);
+            } catch (Throwable t) {
+                LOG.warn("failed to drive the refresh of lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private void driveOneRefresh(LanceIndexJobManager jobManager, 
LanceIndexJob job) {
+        long refreshRevision = job.getRevision() + 1;
+        CatalogIf catalog = 
Env.getCurrentEnv().getCatalogMgr().getCatalog(job.getCatalogId());
+        if (catalog == null) {
+            // Unreachable while the unresolved-job guard blocks catalog 
drops; kept as a
+            // fail-closed fallback so the job still transitions and retries 
later.
+            LOG.warn("catalog of lance index job {} is gone; marking its 
refresh FAILED", job.getJobId());
+            finishRefreshTransition(jobManager, job.getJobId(), 
refreshRevision, false);
+            return;
+        }
+        try {
+            // A half-orphan target (its db or table already dropped 
externally) is a
+            // silent no-op: nothing is left to invalidate, and DONE is the 
correct end
+            // state for the job.
+            
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(catalog.getName(),

Review Comment:
   [P2] Require a verified refresh before marking this job DONE. On a cold 
external table cache, buildTableForInit can catch a transient remote table-list 
failure and return null. handleRefreshTable(..., true) then returns normally 
before invalidating the table or writing a refresh journal, so this path marks 
DONE and releases the same-name fence/quota even though the dataset still 
exists. Treat uncertain lookup as FAILED for retry; reserve DONE for an actual 
refresh or positively verified absence. This is separate from the existing 
operator FORCE path.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,657 @@
+// 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.datasource.lance.job;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.ClientPool;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.util.MasterDaemon;
+import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.lance.LanceExternalCatalog;
+import org.apache.doris.datasource.lance.storage.LanceStorageOptions;
+import org.apache.doris.persist.gson.GsonUtils;
+import org.apache.doris.system.Backend;
+import org.apache.doris.system.BeSelectionPolicy;
+import org.apache.doris.system.SystemInfoService;
+import org.apache.doris.thrift.BackendService;
+import org.apache.doris.thrift.TLanceIndexJobDispatch;
+import org.apache.doris.thrift.TLanceIndexMutationType;
+import org.apache.doris.thrift.TNetworkAddress;
+import org.apache.doris.thrift.TStatus;
+import org.apache.doris.thrift.TStatusCode;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.apache.thrift.TApplicationException;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+import java.util.function.Supplier;
+
+/**
+ * Master-only daemon that drives the durable Lance index job records through
+ * the lifecycle after admission. Each round runs in a fixed order: converge
+ * expired RUNNING jobs to UNKNOWN, release possible-live slots whose backend
+ * process was replaced, drive the refresh a terminal job still owes, then
+ * dispatch PENDING jobs. Every durable transition goes through
+ * {@link LanceIndexJobManager} under its own lock; the daemon holds no catalog
+ * or manager lock across any call.
+ *
+ * <p>The daemon does not read the admission gate: a job that is already 
durable
+ * must be driven to its terminal state, whatever the gate says now, so the
+ * thread runs unconditionally on the master and simply finds nothing to do
+ * while no jobs exist. An idle round writes no journal record.
+ *
+ * <p>Dispatch follows the durable-before-send boundary: the whole request is
+ * prepared first (so a preparation failure just leaves the job PENDING), then
+ * the markRunning edit log is written and re-read before the first byte of
+ * network I/O, and the invocation id of an attempt that lost the 
compare-and-set
+ * is never reused. After a successful markRunning there is exactly one send;
+ * from that point a job converges only through a matching result callback, the
+ * deadline sweep, or the epoch sweep, never through a resend. A failure that
+ * still proves the dispatch was never enqueued (a clean pre-enqueue error
+ * status, a client-pool borrow failure, or an UNKNOWN_METHOD answer from an
+ * old backend) converges it NOT_COMMITTED through the no-enqueue channel,
+ * which releases the possible-live slot in the same durable transition;
+ * anything ambiguous after the invocation may have started converges UNKNOWN
+ * with the slot retained.
+ *
+ * <p>The manager is resolved from the supplier once per round rather than
+ * captured at construction: {@code Env.loadLanceIndexJobManager} replaces the
+ * Env-owned manager with a brand-new object on every image load, so a cached
+ * reference would keep scanning the abandoned pre-image manager after an FE
+ * restart while replay, admission and SHOW all move on to the restored one.
+ * Every phase of one round shares the single resolved instance.
+ *
+ * <p>The sleep between rounds is sliced at {@link #MAX_SLEEP_SLICE_MS} so a
+ * shortened polling interval takes effect within one slice (see the field
+ * javadoc), and {@link Config#lance_index_job_dispatcher_paused} suspends only
+ * the dispatch phase (see {@link #dispatchPendingJobs}).
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    /**
+     * Upper bound of one sleep slice, equal to the shipped default interval. 
The
+     * daemon never sleeps longer than this, so a shortened
+     * {@link Config#lance_index_job_dispatch_interval_second} takes effect 
within
+     * one slice instead of waiting out a previously adopted long sleep: the
+     * elapsed check in {@link #runAfterCatalogReady} is re-evaluated against 
the
+     * current config at every wake. Slices bound only the sleep; rounds still
+     * honor the configured interval, because a wake whose configured interval
+     * (longer than this bound) has not elapsed since the last round skips the
+     * round. A lengthened interval takes effect at the next wake through the 
same
+     * check, and an interval at or below this bound needs no check at all — 
every
+     * wake runs a round, exactly one per configured period.
+     */
+    private static final long MAX_SLEEP_SLICE_MS = 10_000L;
+
+    private final Supplier<LanceIndexJobManager> jobManagerSupplier;
+
+    /** Wall time of the last executed round, or -1 before the first one. */
+    private long lastRoundMs = -1L;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        this(() -> jobManager);
+    }
+
+    public LanceIndexJobDispatcher(Supplier<LanceIndexJobManager> 
jobManagerSupplier) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManagerSupplier = jobManagerSupplier;
+    }
+
+    /**
+     * Values loaded from fe.conf bypass the config validator (only ADMIN SET 
runs
+     * it), so the positive invariant is re-asserted where a non-positive value
+     * would break the loop: a non-positive interval would kill this thread 
inside
+     * {@code Thread.sleep} or busy-spin it, a non-positive deadline would 
sweep
+     * every dispatched job UNKNOWN on the next round, and a zero cap would 
stall
+     * dispatch forever. The refresh retry interval needs no such defense: a
+     * non-positive value simply disengages the throttle.
+     */
+    private static long dispatchIntervalMs() {
+        return Math.max(1, Config.lance_index_job_dispatch_interval_second) * 
1000L;
+    }
+
+    private static long executeDeadlineMs(long nowMs) {
+        long second = Math.max(1L, 
Config.lance_index_job_execute_deadline_second);
+        return second > (Long.MAX_VALUE - nowMs) / 1000L ? Long.MAX_VALUE : 
nowMs + second * 1000L;
+    }
+
+    @Override
+    protected void runAfterCatalogReady() {
+        if (!Env.getCurrentEnv().isMaster()) {
+            return;
+        }
+        if (Env.isCheckpointThread()) {
+            return;
+        }
+        long configuredMs = dispatchIntervalMs();
+        setInterval(Math.min(configuredMs, MAX_SLEEP_SLICE_MS));
+        if (configuredMs > MAX_SLEEP_SLICE_MS && lastRoundMs >= 0 && nowMs() - 
lastRoundMs < configuredMs) {
+            // A wake inside a long configured interval: the slice elapsed, the
+            // round period has not. Skipping is cheap and writes no journal 
record.
+            return;
+        }
+        lastRoundMs = nowMs();
+        try {
+            runOneRound(jobManagerSupplier.get());
+        } catch (Throwable t) {
+            LOG.warn("Failed to process one round of the lance index job 
dispatcher", t);
+        }
+    }
+
+    /** Clock seam for the round-period check; tests advance it instead of 
sleeping. */
+    protected long nowMs() {
+        return System.currentTimeMillis();
+    }
+
+    private void runOneRound(LanceIndexJobManager jobManager) {
+        long nowMs = System.currentTimeMillis();
+        sweepExpiredRunningJobs(jobManager, nowMs);
+        sweepReplacedProcessEpochs(jobManager);
+        driveRequiredRefreshes(jobManager, nowMs);
+        dispatchPendingJobs(jobManager);
+    }
+
+    /**
+     * Deadline sweep. A RUNNING job past its wait deadline has produced no
+     * complete trusted result, so it converges to UNKNOWN through the same
+     * completeWithResult channel a callback would use. Expiry bounds the wait
+     * only: it never proves termination, so the possible-live slot, the
+     * same-name fence, and the unresolved quota all stay held.
+     */
+    private void sweepExpiredRunningJobs(LanceIndexJobManager jobManager, long 
nowMs) {
+        for (LanceIndexJob job : jobManager.getExpiredRunningJobs(nowMs)) {
+            try {
+                boolean completed = 
jobManager.completeWithResult(job.getJobId(),
+                        dispatchRevisionOf(job), job.getInvocationId(), 
job.getBeProcessEpoch(),
+                        new 
LanceIndexJobResult(LanceIndexJobResultCode.NO_TRUSTED_RESULT,
+                                LanceIndexJobCompletionReason.NONE,
+                                "execute deadline expired without a complete 
trusted result", false));
+                if (completed) {
+                    LOG.info("lance index job {} converged RUNNING -> UNKNOWN 
on deadline expiry",
+                            job.getJobId());
+                } else {
+                    LOG.warn("deadline sweep skipped lance index job {}: 
already converged by a callback or sweep",
+                            job.getJobId());
+                }
+            } catch (Throwable t) {
+                LOG.warn("failed to sweep expired lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    /**
+     * Possible-live sweep. The only slot-release proof this daemon produces is
+     * that the recorded backend process epoch no longer exists: a backend 
entry
+     * reporting a different epoch proves the process that received the 
dispatch
+     * was replaced. A missing backend entry or heartbeat loss proves nothing
+     * (the worker may still be running behind a partition), so such a job 
keeps
+     * its slot until a stronger proof or an operator force release. An epoch
+     * change also proves nothing about the outcome, so the mutation state is
+     * never touched here.
+     */
+    private void sweepReplacedProcessEpochs(LanceIndexJobManager jobManager) {
+        for (LanceIndexJob job : jobManager.getJobsHoldingPossibleLiveSlot()) {
+            try {
+                Backend backend = 
Env.getCurrentSystemInfo().getBackend(job.getBackendId());
+                if (backend == null || backend.getProcessEpoch() == 
job.getBeProcessEpoch()) {
+                    continue;
+                }
+                boolean recorded = 
jobManager.recordTerminationProof(job.getJobId(),
+                        dispatchRevisionOf(job), job.getBackendId(), 
job.getBeProcessEpoch(),
+                        job.getInvocationId(), 
LanceIndexTerminationProof.BE_PROCESS_EPOCH_GONE);
+                if (recorded) {
+                    LOG.info("released possible-live slot of lance index job 
{}: backend process epoch was replaced",
+                            job.getJobId());
+                } else {
+                    LOG.warn("epoch sweep skipped lance index job {}: dispatch 
identity already moved",
+                            job.getJobId());
+                }
+            } catch (Throwable t) {
+                LOG.warn("failed to sweep possible-live slot of lance index 
job " + job.getJobId(), t);
+            }
+        }
+    }
+
+    /**
+     * Refresh driver for terminal jobs with an unfinished refresh obligation.
+     * Completing the refresh is the protocol duty that releases the same-name
+     * fence and the unresolved quota once DONE; it is not a read-visibility
+     * action, because index metadata is never cached. Each job is driven
+     * through markRefreshRunning, the idempotent external-table refresh, then
+     * DONE or FAILED: a FAILED job keeps its fence and is retried, throttled 
to
+     * one attempt per retry interval, while a first REQUIRED refresh is never
+     * delayed. UNKNOWN jobs never appear here; they owe no refresh.
+     */
+    private void driveRequiredRefreshes(LanceIndexJobManager jobManager, long 
nowMs) {
+        for (LanceIndexJob job : jobManager.getJobsNeedingRefresh()) {
+            try {
+                if (job.getRefreshState() == 
LanceIndexJobRefreshState.RUNNING) {
+                    // In flight elsewhere; the master-transfer sweep 
downgrades a stale
+                    // RUNNING back to REQUIRED, so a lost driver cannot 
strand it.
+                    continue;
+                }
+                if (job.getRefreshState() == LanceIndexJobRefreshState.FAILED
+                        && nowMs - job.getUpdateTimeMs()
+                                < Config.lance_index_job_refresh_retry_second 
* 1000L) {
+                    continue;
+                }
+                if (!jobManager.markRefreshRunning(job.getJobId(), 
job.getRevision())) {
+                    // A concurrent driver won the compare-and-set; nothing to 
do here.
+                    continue;
+                }
+                driveOneRefresh(jobManager, job);
+            } catch (Throwable t) {
+                LOG.warn("failed to drive the refresh of lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private void driveOneRefresh(LanceIndexJobManager jobManager, 
LanceIndexJob job) {
+        long refreshRevision = job.getRevision() + 1;
+        CatalogIf catalog = 
Env.getCurrentEnv().getCatalogMgr().getCatalog(job.getCatalogId());
+        if (catalog == null) {
+            // Unreachable while the unresolved-job guard blocks catalog 
drops; kept as a
+            // fail-closed fallback so the job still transitions and retries 
later.
+            LOG.warn("catalog of lance index job {} is gone; marking its 
refresh FAILED", job.getJobId());
+            finishRefreshTransition(jobManager, job.getJobId(), 
refreshRevision, false);
+            return;
+        }
+        try {
+            // A half-orphan target (its db or table already dropped 
externally) is a
+            // silent no-op: nothing is left to invalidate, and DONE is the 
correct end
+            // state for the job.
+            
Env.getCurrentEnv().getRefreshManager().handleRefreshTable(catalog.getName(),
+                    job.getDbName(), job.getTableName(), true);
+        } catch (Throwable t) {
+            // The typed DdlException is the expected failure; an unchecked 
exception out
+            // of the metadata path must still leave the durable refresh 
state, or the
+            // job would strand in refresh RUNNING until the next master 
transfer.
+            LOG.warn("refresh of lance index job {} failed; keeping the fence 
for a retry",
+                    job.getJobId(), t);
+            finishRefreshTransition(jobManager, job.getJobId(), 
refreshRevision, false);
+            return;
+        }
+        finishRefreshTransition(jobManager, job.getJobId(), refreshRevision, 
true);
+    }
+
+    /**
+     * Applies the DONE/FAILED transition with a bounded revision retry. A 
concurrent
+     * termination-proof write can bump the revision after markRefreshRunning 
succeeded,
+     * and silently losing that compare-and-set would leave the refresh 
RUNNING — a
+     * state only the master-transfer sweep downgrades. Re-reading the 
revision and
+     * retrying a few times converges it; a persistent loss is escalated.
+     */
+    private void finishRefreshTransition(LanceIndexJobManager jobManager, long 
jobId, long expectedRevision,
+            boolean done) {
+        long revision = expectedRevision;
+        for (int attempt = 0; attempt < 3; attempt++) {
+            boolean transitioned = done ? jobManager.markRefreshDone(jobId, 
revision)
+                    : jobManager.markRefreshFailed(jobId, revision);
+            if (transitioned) {
+                return;
+            }
+            LanceIndexJob fresh = jobManager.getJob(jobId);
+            if (fresh == null) {
+                break;
+            }
+            revision = fresh.getRevision();
+        }
+        LOG.error("lance index job {} kept its refresh RUNNING: the 
DONE/FAILED transition kept losing the"
+                + " compare-and-set; the master-transfer sweep will downgrade 
it", jobId);
+    }
+
+    /**
+     * PENDING dispatch. Makes at most
+     * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches 
per
+     * round, and only a job this round actually made RUNNING consumes that
+     * budget: skipped jobs (an eligibility gate is closed, or every backend is
+     * at capacity) are scanned past, so a stable subset of permanently
+     * undispatchable jobs can never crowd out later ids. Per backend it never
+     * exceeds {@link Config#lance_index_job_max_inflight_per_backend}
+     * possible-live worker slots, counted from slot ownership (see
+     * {@link LanceIndexJobManager#countPossibleLiveSlotsByBackend()}) plus the
+     * jobs this round already made RUNNING. A job that cannot be dispatched
+     * keeps waiting as PENDING: there is no dispatch-exhaustion terminal state
+     * and no backoff beyond the daemon period.
+     *
+     * <p>{@link Config#lance_index_job_dispatcher_paused} suspends this phase
+     * only — the sweeps and the refresh driver keep running while it is set.
+     * The switch is checked at the phase entry and again before every single
+     * job attempt, which closes the admission race a test or operator cares
+     * about: anyone who sets the switch <em>before</em> admitting a job is
+     * guaranteed the job is never dispatched while paused. A round whose
+     * snapshot was taken before the admission never sees the job at all, and
+     * any round that can see it performs its per-job check after the
+     * admission, hence after the switch was set, and skips it. A skipped job
+     * never consumes the round's dispatch budget.
+     */
+    private void dispatchPendingJobs(LanceIndexJobManager jobManager) {
+        if (Config.lance_index_job_dispatcher_paused) {
+            return;
+        }
+        int maxPerRound = Math.max(1, 
Config.lance_index_job_max_dispatch_per_round);
+        Map<Long, Integer> inflightByBackend = 
jobManager.countPossibleLiveSlotsByBackend();
+        int dispatched = 0;
+        for (LanceIndexJob job : jobManager.getJobsNeedingDispatch()) {
+            if (Config.lance_index_job_dispatcher_paused) {
+                // Flipped mid-round: stop without touching the budget.
+                break;
+            }
+            if (dispatched >= maxPerRound) {
+                break;
+            }
+            try {
+                if (tryDispatch(jobManager, job, inflightByBackend)) {
+                    dispatched++;
+                }
+            } catch (Throwable t) {
+                LOG.warn("failed to dispatch lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    /**
+     * One dispatch attempt for one PENDING job; returns true only when the
+     * attempt made the job durable RUNNING (and so consumes this round's
+     * dispatch budget). Every early return before markRunning leaves the job
+     * PENDING for a later round: the eligibility gates, the backend and
+     * capacity checks, and also the whole request preparation — storage-option
+     * resolution and the wire request build run before the durable boundary,
+     * so an FE-side failure there (for example a catalog id that resolves to
+     * nothing while ALTER CATALOG RENAME has the catalog temporarily removed)
+     * just retries next round instead of stranding the job UNKNOWN without a
+     * single byte sent. Once markRunning succeeds the job is durable RUNNING
+     * and this invocation id gets exactly one send attempt; after that only a
+     * matching callback, the deadline sweep, or the epoch sweep can converge
+     * the job.
+     */
+    private boolean tryDispatch(LanceIndexJobManager jobManager, LanceIndexJob 
job,
+            Map<Long, Integer> inflightByBackend) {
+        boolean localDataset = isLocalFileDataset(job.getNormalizedLocator());
+        if (localDataset && !Config.enable_lance_index_local_file_mutation) {
+            // Operator assertion is off: a local-filesystem mutation stays 
PENDING.
+            return false;
+        }
+        if (localDataset && Env.getCurrentEnv().getFrontends(null).size() != 
1) {
+            // Local files are only shared by a single-node deployment.
+            return false;
+        }
+        SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+        // All schedule-available backends, shuffled by the selection policy: 
the
+        // first one with a free possible-live slot takes the job, so a full
+        // backend defers this attempt only when every selectable backend is at
+        // the cap, never just because the randomly picked one is.
+        List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(
+                new 
BeSelectionPolicy.Builder().needScheduleAvailable().build(), -1);
+        int perBackendCap = Math.max(1, 
Config.lance_index_job_max_inflight_per_backend);
+        Backend backend = null;
+        for (Long backendId : backendIds) {
+            Backend candidate = systemInfo.getBackend(backendId);
+            if (candidate == null) {
+                continue;
+            }
+            if (localDataset && !isOnlyAliveBackend(systemInfo, 
candidate.getId())) {
+                continue;
+            }
+            Integer inflight = inflightByBackend.get(candidate.getId());
+            if (inflight != null && inflight >= perBackendCap) {
+                continue;
+            }
+            backend = candidate;
+            break;
+        }
+        if (backend == null) {
+            return false;
+        }
+        String invocationId = UUID.randomUUID().toString();
+        // The process epoch is captured once, and the same value goes to the
+        // durable record and the wire: a heartbeat landing between the two 
reads
+        // must not split the dispatch identity (the callback matches the 
durable
+        // value, and the epoch sweep releases the slot against it).
+        long beProcessEpoch = backend.getProcessEpoch();
+        long deadlineMs = executeDeadlineMs(System.currentTimeMillis());
+        long expectedDispatchRevision = job.getRevision() + 1;
+        TLanceIndexJobDispatch dispatch;
+        try {
+            dispatch = buildDispatch(job, expectedDispatchRevision, 
invocationId, deadlineMs, beProcessEpoch,
+                    resolveStorageOptions(job));
+        } catch (Exception e) {
+            // Not a trusted worker rejection and not an ambiguity either: 
nothing was
+            // marked and nothing was sent, so the job simply waits for the 
next round.
+            LOG.warn("failed to prepare the dispatch of lance index job {}; 
staying PENDING: {}",
+                    job.getJobId(), e.getMessage());
+            return false;
+        }
+        if (!jobManager.markRunning(job.getJobId(), job.getRevision(), 
backend.getId(),
+                beProcessEpoch, invocationId, deadlineMs)) {
+            // The compare-and-set lost: this attempt's dispatch identity is 
void and its
+            // invocation id is discarded. A fresh identity is built from 
scratch next round.
+            return false;
+        }
+        inflightByBackend.merge(backend.getId(), 1, Integer::sum);

Review Comment:
   [P2] Reclaim this round's local capacity after a proven no-enqueue result. 
markRunning increments inflightByBackend here, but the current BE immediately 
returns NOT_IMPLEMENTED_ERROR; completeProvenNoEnqueue durably frees the slot 
while this map stays at its cap. With one BE and default settings, only two of 
the allowed 16 jobs are processed per round even though no worker is live, so a 
256-job backlog takes about 21 minutes to drain. Decrement after a successful 
no-enqueue transition or refresh the count before selecting another job.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ResolveLanceIndexJobCommand.java:
##########
@@ -0,0 +1,367 @@
+// 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;
+
+import org.apache.doris.catalog.DatabaseIf;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.TableIf;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.DdlException;
+import org.apache.doris.common.ErrorCode;
+import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.CatalogMgr;
+import org.apache.doris.datasource.ExternalDatabase;
+import org.apache.doris.datasource.ExternalTable;
+import org.apache.doris.datasource.lance.LanceExternalCatalog;
+import org.apache.doris.datasource.lance.LanceIndexMutationValidator;
+import org.apache.doris.datasource.lance.job.LanceIndexJob;
+import org.apache.doris.datasource.lance.job.LanceIndexJobManager;
+import org.apache.doris.datasource.lance.job.LanceIndexJobMutationState;
+import org.apache.doris.mysql.privilege.PrivPredicate;
+import org.apache.doris.nereids.trees.plans.PlanType;
+import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.StmtExecutor;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.nio.charset.StandardCharsets;
+
+/**
+ * RESOLVE LANCE INDEX JOB &lt;jobId&gt; AS FORCE_RELEASE COMMENT 
'&lt;note&gt;' — the operator
+ * escape hatch that durably releases a job whose mutation outcome is UNKNOWN 
(design section
+ * 7.1). RESOLVE is deliberately not gated by {@code 
enable_lance_index_mutation}: the gate
+ * controls mutation admission, while FORCE must stay available exactly when 
the gate is off.
+ *
+ * <p>The release protocol keeps the fence, the quota charge and the 
possible-live slot while
+ * it performs one authoritative latest-metadata read and one external-table 
refresh with the
+ * current credentials of the surviving catalog, both outside every 
catalog/manager lock; only
+ * then does the durable release transfer inside the admission critical section
+ * ({@code captureLanceIndexTarget} → lock-free read/refresh → {@code 
withLanceIndexAdmission}
+ * recheck → manager write lock), serialized against DROP CATALOG and identity 
ALTER exactly
+ * like admission. Any failure before the transfer is the typed
+ * {@code ERR_LANCE_INDEX_JOB_RESOLUTION_INCOMPLETE}: nothing is written, 
nothing is released,
+ * and the operator fixes the cause and retries the same statement.
+ *
+ * <p>Target resolution (design section 7.1 step 1) is three-valued. RESOLVED 
means the
+ * persisted names resolve and the catalog's current durable dataset locator 
still matches
+ * the job's — the same revalidation SHOW LANCE INDEX JOBS applies, so a 
repointed dataset
+ * reusing the same names never turns a stale name into table-level 
authorization. MISSING
+ * means the catalog, database or table is verifiably absent, or the locator 
positively
+ * points at a different dataset: that is the orphan family — a full orphan 
(catalog gone)
+ * has no credentials to read with and nothing to invalidate, so it is 
released directly
+ * after global ADMIN authorization, while a half-orphan skips the 
authoritative read and
+ * refreshes with {@code ignoreIfNotExists=true} as a best-effort 
invalidation. FAILED means
+ * a resolution that errors out, or a locator that cannot be resolved right 
now: never an
+ * orphan verdict — after ADMIN authorization the statement fails with the 
typed 5105 so the
+ * fence is kept when "table gone" cannot be told apart from "network down". 
SHOW fails the
+ * same uncertainty closed by hiding the row; RESOLVE fails it closed by not 
releasing.
+ *
+ * <p>Non-disclosure (design section 8): the job is loaded first and 
authorized against its
+ * persisted target — table-level ALTER when the target resolves, global ADMIN 
otherwise — and
+ * a missing job and an unauthorized job share the same fixed 
ERR_LANCE_INDEX_JOB_NOT_FOUND
+ * response naming only the job id. The 5104 state rejection and the 5105 
resolution failure
+ * are only visible to an already authorized caller.
+ *
+ * <p>Success returns an OK packet carrying one warning row with {@link 
#LATE_COMMIT_WARNING},
+ * the same text persisted as the job's durable {@code forceWarning}: the old 
worker may still
+ * overwrite, remove, or reintroduce the index name; the mutation outcome 
remains UNKNOWN.
+ * Retrying FORCE on an already released job is an idempotent success 
returning the existing
+ * release record, never an error.
+ */
+public class ResolveLanceIndexJobCommand extends Command implements 
ForwardWithSync {
+    /**
+     * The late-commit warning (design section 7.1), returned in the OK packet 
and persisted
+     * verbatim as the durable {@code forceWarning}; bounded well under
+     * {@link LanceIndexJob#MAX_FORCE_TEXT_BYTES}.
+     */
+    static final String LATE_COMMIT_WARNING =
+            "the old worker may still overwrite, remove, or reintroduce the 
index name; "
+                    + "the mutation outcome remains UNKNOWN";
+
+    private static final Logger LOG = 
LogManager.getLogger(ResolveLanceIndexJobCommand.class);
+
+    private final long jobId;
+    private final String comment;
+
+    public ResolveLanceIndexJobCommand(long jobId, String comment) {
+        super(PlanType.RESOLVE_LANCE_INDEX_JOB_COMMAND);
+        this.jobId = jobId;
+        this.comment = comment;
+    }
+
+    public long getJobId() {
+        return jobId;
+    }
+
+    public String getComment() {
+        return comment;
+    }
+
+    @Override
+    public void run(ConnectContext ctx, StmtExecutor executor) throws 
Exception {
+        Env env = Env.getCurrentEnv();
+        LanceIndexJobManager manager = env.getLanceIndexJobManager();
+        // 1. Load the job without disclosing any field (design section 7.1 
step 1).
+        LanceIndexJob job = manager.getJob(jobId);
+        if (job == null) {
+            throw notFound();
+        }
+        // 2. Resolve and authorize against the persisted target before any 
state is revealed:
+        //    table-level ALTER when the target resolves, global ADMIN for the 
orphan family
+        //    and for a target whose resolution failed outright.
+        CatalogMgr catalogMgr = env.getCatalogMgr();
+        CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog = 
catalogMgr.getCatalog(job.getCatalogId());
+        TargetResolution resolution = resolveTarget(catalog, job);
+        boolean authorized = resolution == TargetResolution.RESOLVED
+                ? env.getAccessManager().checkTblPriv(ctx, catalog.getName(), 
job.getDbName(), job.getTableName(),
+                        PrivPredicate.ALTER)
+                : env.getAccessManager().checkGlobalPriv(ctx, 
PrivPredicate.ADMIN);
+        if (!authorized) {
+            throw notFound();
+        }
+        // 3. Idempotent replay: a retry returns the existing release record 
(section 7.1).
+        //    This deliberately precedes the resolution-failure rejection: 
once the release
+        //    has landed, a retry during a provider outage is a success, not a 
5105.
+        if (job.isForceReleased()) {
+            ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING);
+            return;
+        }
+        // 4. Only UNKNOWN may be force-released; a null mutation state reads 
as UNKNOWN,
+        //    same as the manager's own gate. The state rejection also 
precedes the
+        //    resolution-failure rejection: for a terminal job the accurate 
answer is 5104,
+        //    not a 5105 claiming the job still holds its fence.
+        if (job.getMutationState() != null && job.getMutationState() != 
LanceIndexJobMutationState.UNKNOWN) {
+            throw new 
AnalysisException(ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN.formatErrorMsg(jobId),
+                    ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN);
+        }
+        if (resolution == TargetResolution.FAILED) {
+            // Never an orphan verdict: "table gone" cannot be told apart from 
"network down",
+            // so nothing is released and nothing beyond the typed error is 
disclosed; the
+            // operator fixes the cause and retries the same statement (design 
7.1 step 4).
+            throw incompleteResolution("the persisted target could not be 
resolved with current catalog"
+                    + " metadata; see fe.log for the cause");
+        }
+        boolean targetResolves = resolution == TargetResolution.RESOLVED;
+        // 5. The grammar makes COMMENT mandatory; here the note must also be 
non-empty after
+        //    trimming and fit the durable force text bound.
+        String note = comment == null ? "" : comment.trim();
+        if (note.isEmpty()) {
+            throw new AnalysisException("force release note must not be empty",
+                    ErrorCode.ERR_LANCE_INDEX_INVALID);
+        }
+        if (note.getBytes(StandardCharsets.UTF_8).length > 
LanceIndexJob.MAX_FORCE_TEXT_BYTES) {
+            throw new AnalysisException("force release note exceeds " + 
LanceIndexJob.MAX_FORCE_TEXT_BYTES
+                    + " UTF-8 bytes", ErrorCode.ERR_LANCE_INDEX_INVALID);
+        }
+        // 6-9. Branch on the orphan state, then the durable release transfer.
+        String actor = ctx.getQualifiedUser();
+        boolean released;
+        if (catalog == null) {
+            // Full orphan: no credentials survive to read with and nothing 
can be
+            // invalidated, so the release goes straight to the manager write 
lock.
+            released = manager.forceRelease(jobId, job.getRevision(), actor, 
note, LATE_COMMIT_WARNING);
+        } else {
+            released = releaseWithLiveCatalog(env, catalogMgr, manager, 
catalog, targetResolves, job, actor, note);
+        }
+        if (!released) {
+            // 10. The expected-revision transfer lost a race. A concurrent 
FORCE_RELEASE
+            //     that already landed makes this an idempotent success; 
anything else means
+            //     the job left UNKNOWN concurrently (UNKNOWN has no other 
outgoing
+            //     transition), so the pinned not-UNKNOWN wording stays 
accurate.
+            LanceIndexJob reread = manager.getJob(jobId);
+            if (reread != null && reread.isForceReleased()) {
+                ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING);
+                return;
+            }
+            throw new 
AnalysisException(ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN.formatErrorMsg(jobId),
+                    ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN);
+        }
+        // 11. The OK packet carries the late-commit warning; it survives the 
forward chain
+        //     byte-identically (proxyExecute serializes the master state).
+        ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING);
+    }
+
+    /**
+     * The live-catalog release path: capture the target identity, then — 
holding no lock —
+     * perform the authoritative read and the refresh, and finally transfer 
the release inside
+     * the admission critical section. Lock order stays CatalogMgr then 
LanceIndexJobManager.
+     */
+    private boolean releaseWithLiveCatalog(Env env, CatalogMgr catalogMgr, 
LanceIndexJobManager manager,
+            CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog, 
boolean targetResolves, LanceIndexJob job,
+            String actor, String note) throws Exception {
+        if (!(catalog instanceof LanceExternalCatalog)
+                || ((LanceExternalCatalog) catalog).isRestCatalogConfigured()) 
{
+            // Defensive: admission never targets a REST catalog and a 
non-Lance catalog
+            // cannot hold a Lance fence, so no UNKNOWN job should ever 
resolve here.
+            LanceIndexMutationValidator.rejectUnsupportedOperation(
+                    "RESOLVE LANCE INDEX JOB AS FORCE_RELEASE", "catalog '" + 
catalog.getName() + "'");
+        }
+        LanceExternalCatalog lanceCatalog = (LanceExternalCatalog) catalog;
+        CatalogMgr.LanceIndexTarget target;
+        try {
+            target = catalogMgr.captureLanceIndexTarget(lanceCatalog);
+        } catch (DdlException e) {
+            throw incompleteResolution(e.getMessage());
+        }
+        if (targetResolves) {
+            // 7. One authoritative latest-metadata read with current 
credentials, proving
+            //    the dataset is reachable before anything is released.
+            authoritativeRead(lanceCatalog, catalog, job);
+        }
+        // 8. Invalidate the external table and broadcast the refresh to every 
FE. A
+        //    half-orphan target is invalidated best-effort (missing db/table 
is a no-op).
+        try {
+            env.getRefreshManager().handleRefreshTable(catalog.getName(), 
job.getDbName(), job.getTableName(),
+                    !targetResolves);
+        } catch (DdlException e) {
+            throw incompleteResolution(e.getMessage());
+        }
+        // 9. The durable transfer rechecks the catalog identity under the 
read lock, then
+        //    runs the revision-checked release under the manager write lock.
+        try {
+            return catalogMgr.withLanceIndexAdmission(lanceCatalog, target,
+                    () -> manager.forceRelease(jobId, job.getRevision(), 
actor, note, LATE_COMMIT_WARNING));
+        } catch (DdlException e) {
+            throw incompleteResolution(e.getMessage());
+        }
+    }
+
+    /**
+     * The authoritative latest-metadata read over the remote names of the 
resolved db/table,
+     * outside every lock (the loader owns its deadline-bound JNI read). Every 
loader failure
+     * already passed through the catalog's sanitized root-cause chain, so its 
message is safe
+     * to echo — locator, credentials and dataset uri are masked there.
+     */
+    private void authoritativeRead(LanceExternalCatalog lanceCatalog,
+            CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog, 
LanceIndexJob job) throws AnalysisException {
+        DatabaseIf<? extends TableIf> db;
+        TableIf table;
+        try {
+            db = catalog.getDbNullable(job.getDbName());
+            table = db == null ? null : 
db.getTableNullable(job.getTableName());
+        } catch (Exception e) {
+            // A remote-metadata blip during resolution is not an orphan 
verdict; the provider
+            // message is not echoed here because it never crossed the 
sanitized chain.
+            LOG.warn("lance index job {}: target re-resolution failed before 
the authoritative read",
+                    job.getJobId(), e);
+            throw incompleteResolution("the persisted target could not be 
resolved with current catalog"
+                    + " metadata; see fe.log for the cause");
+        }
+        if (db == null || table == null) {
+            // The target resolved at authorization time but is gone now: keep 
the fence and
+            // let the retry take the half-orphan branch under global ADMIN.
+            throw incompleteResolution("the persisted target table no longer 
resolves; retry the statement");
+        }
+        if (!(db instanceof ExternalDatabase) || !(table instanceof 
ExternalTable)) {
+            // Boundary guard: a Lance catalog must serve external relations, 
but the command
+            // boundary does not trust that invariant (no raw 
ClassCastException to the user).
+            LOG.warn("lance index job {}: target relation is not external: 
db={}, table={}",
+                    job.getJobId(), db.getClass().getName(), 
table.getClass().getName());
+            throw incompleteResolution("the persisted target is not an 
external relation;"
+                    + " see fe.log for the cause");
+        }
+        String remoteDb = ((ExternalDatabase) db).getRemoteName();
+        String remoteTable = ((ExternalTable) table).getRemoteName();
+        try {
+            lanceCatalog.loadTableIndexAdmissionSnapshot(remoteDb, 
remoteTable);

Review Comment:
   [P2] Check the dataset URI returned by this authoritative read. 
resolveTarget first verifies the job's locator A, but 
loadTableIndexAdmissionSnapshot resolves the namespace again without using that 
result here. If a child-namespace directory table is deregistered and 
re-registered at B between the two calls, this reads B and then releases A's 
UNKNOWN fence under table ALTER; the final catalog check sees no remote change. 
Normalize and compare the snapshot URI with the persisted locator, then retry 
under the orphan/ADMIN path on mismatch.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to