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


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobManager.java:
##########
@@ -693,6 +721,75 @@ && hasFenceIdentity(job)
         }
     }
 
+    /**
+     * PENDING jobs eligible for the dispatcher, in job id order (FIFO 
fairness),
+     * at most {@code limit} of them. Target identity must be complete: corrupt
+     * identity-less records are never dispatchable (their only exit is the
+     * force-release transition by job id) and are skipped here. The dispatch
+     * quad is deliberately not required: it is written by {@code markRunning},
+     * which is the step this query feeds. All matches are collected and 
ordered
+     * before truncating, so a stable subset of permanently undispatchable jobs
+     * can never crowd out later ids.
+     */
+    public List<LanceIndexJob> getJobsNeedingDispatch(int limit) {
+        readLock();
+        try {
+            List<LanceIndexJob> result = new ArrayList<>();
+            for (LanceIndexJob job : jobs.values()) {
+                if (job != null && job.getMutationState() == 
LanceIndexJobMutationState.PENDING
+                        && hasDispatchTarget(job)) {
+                    result.add(new LanceIndexJob(job));
+                }

Review Comment:
   [P1] Scan beyond ineligible pending jobs before applying the round cap. With 
the default cap of 16, 16 admitted file jobs stay PENDING while local mutation 
is disabled; this subList always returns those same IDs, so a later S3 job is 
never visited. Count actual dispatches or advance a fair cursor after skipped 
jobs.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobReportHandler.java:
##########
@@ -0,0 +1,129 @@
+// 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.thrift.TLanceIndexJobReport;
+import org.apache.doris.thrift.TLanceIndexTerminationProof;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.Objects;
+
+/**
+ * Applies one typed result envelope reported by a backend to the durable job
+ * record. This is a thin shim over the manager transitions: dispatch-identity
+ * checking and result classification all live in {@link LanceIndexJobManager},
+ * so a stale or identity-mismatched report only logs a warning and changes
+ * nothing. A malformed envelope (missing result code, a code this FE does not
+ * know, or a sanitized message past the durable bound) is dropped rather than
+ * trusted; the job then converges through the dispatcher's deadline sweep.
+ * Only the typed codes are read: message text is never inspected to infer an
+ * outcome.
+ *
+ * <p>The handler runs on the report RPC thread and performs no I/O beyond the
+ * manager's own edit-log write. It starts no refresh: the metadata refresh a
+ * completed job may owe is driven by the dispatcher daemon, not here.
+ */
+public class LanceIndexJobReportHandler {
+
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobReportHandler.class);
+
+    private final LanceIndexJobManager jobManager;
+
+    public LanceIndexJobReportHandler(LanceIndexJobManager jobManager) {
+        this.jobManager = Objects.requireNonNull(jobManager, "jobManager");
+    }
+
+    /**
+     * Handles one report: a matched report completes the job with its
+     * classified result, and a CHILD_REAPED termination proof additionally
+     * releases the possible-live slot, because reaping the exact child process
+     * proves that process ended (which still says nothing about the outcome).
+     */
+    public void handle(TLanceIndexJobReport report) {
+        if (report == null) {
+            LOG.warn("dropping null lance index job report");
+            return;
+        }
+        LanceIndexJobResult result;
+        try {
+            result = toResult(report);
+        } catch (IllegalArgumentException e) {

Review Comment:
   [P2] Process an identity-matched CHILD_REAPED proof even when the result is 
malformed. If a reaped child reports an overlong message or unknown result 
code, toResult throws and this return drops its valid termination proof. The 
deadline then leaves a false possible-live slot on a stable BE; 
recordTerminationProof can validate the proof independently while the result 
remains untrusted.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    private final LanceIndexJobManager jobManager;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManager = jobManager;
+    }
+
+    /**
+     * 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;
+        }
+        setInterval(dispatchIntervalMs());
+        try {
+            runOneRound();
+        } catch (Throwable t) {
+            LOG.warn("Failed to process one round of the lance index job 
dispatcher", t);
+        }
+    }
+
+    private void runOneRound() {
+        long nowMs = System.currentTimeMillis();
+        sweepExpiredRunningJobs(nowMs);
+        sweepReplacedProcessEpochs();
+        driveRequiredRefreshes(nowMs);
+        dispatchPendingJobs();
+    }
+
+    /**
+     * 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(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() {
+        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(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(job);
+            } catch (Throwable t) {
+                LOG.warn("failed to drive the refresh of lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private void driveOneRefresh(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(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(job.getJobId(), refreshRevision, false);
+            return;
+        }
+        finishRefreshTransition(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(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. Attempts at most
+     * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches 
per
+     * round, and never more than {@link 
Config#lance_index_job_max_inflight_per_backend}
+     * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+     */
+    private void dispatchPendingJobs() {
+        int maxPerRound = Math.max(1, 
Config.lance_index_job_max_dispatch_per_round);
+        Map<Long, Integer> inflightByBackend = countInflightByBackend();
+        int attempted = 0;
+        for (LanceIndexJob job : 
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+            if (++attempted > maxPerRound) {
+                break;
+            }
+            try {
+                tryDispatch(job, inflightByBackend);
+            } catch (Throwable t) {
+                LOG.warn("failed to dispatch lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private Map<Long, Integer> countInflightByBackend() {
+        Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+        for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+            if (job.getMutationState() == LanceIndexJobMutationState.RUNNING 
&& job.getBackendId() != null) {
+                inflightByBackend.merge(job.getBackendId(), 1, Integer::sum);
+            }
+        }
+        return inflightByBackend;
+    }
+
+    /**
+     * One dispatch attempt for one PENDING job. Every early return before
+     * markRunning leaves the job PENDING for a later round. 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 void tryDispatch(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;
+        }
+        if (localDataset && Env.getCurrentEnv().getFrontends(null).size() != 
1) {
+            // Local files are only shared by a single-node deployment.
+            return;
+        }
+        SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+        List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(
+                new 
BeSelectionPolicy.Builder().needScheduleAvailable().build(), 1);
+        if (backendIds.isEmpty()) {
+            return;
+        }
+        Backend backend = systemInfo.getBackend(backendIds.get(0));
+        if (backend == null) {
+            return;
+        }
+        if (localDataset && !isOnlyAliveBackend(systemInfo, backend.getId())) {
+            return;
+        }
+        Integer inflight = inflightByBackend.get(backend.getId());
+        if (inflight != null && inflight >= Math.max(1, 
Config.lance_index_job_max_inflight_per_backend)) {
+            return;
+        }
+        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());
+        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;
+        }
+        inflightByBackend.merge(backend.getId(), 1, Integer::sum);
+        long expectedDispatchRevision = job.getRevision() + 1;
+        LanceIndexJob fresh = jobManager.getJob(job.getJobId());
+        if (!Env.getCurrentEnv().isMaster() || fresh == null
+                || fresh.getMutationState() != 
LanceIndexJobMutationState.RUNNING
+                || fresh.getDispatchRevision() == null
+                || fresh.getDispatchRevision() != expectedDispatchRevision
+                || !invocationId.equals(fresh.getInvocationId())) {
+            // The recheck failed right before the send: no send, and no 
resend either.
+            // The job is durable RUNNING, so the deadline sweep or a matching 
callback
+            // converges it.
+            LOG.warn("lance index job {} did not survive the pre-send recheck; 
not sending", job.getJobId());
+            return;
+        }
+        TLanceIndexJobDispatch dispatch;
+        try {
+            dispatch = buildDispatch(fresh, invocationId, deadlineMs, 
beProcessEpoch,
+                    resolveStorageOptions(fresh));
+        } catch (Exception e) {
+            // An FE-side resolution failure is not a trusted worker 
rejection, so it must
+            // not fabricate NOT_COMMITTED. The job is already RUNNING without 
a send, and
+            // the send may never happen, so converge it to UNKNOWN 
fail-closed.
+            LOG.warn("failed to prepare the dispatch of lance index job {}: 
{}", job.getJobId(), e.getMessage());
+            completeNoTrusted(fresh, "dispatch preparation failed before 
send");
+            return;
+        }
+        TStatus status;
+        try {
+            status = sendExecuteRequest(backend, dispatch);
+        } catch (Exception e) {
+            // The request may have reached the backend, so its outcome cannot 
be trusted.
+            LOG.warn("dispatch send of lance index job {} failed: {}", 
job.getJobId(), e.getMessage());
+            completeNoTrusted(fresh, "dispatch send failed; the result cannot 
be trusted");
+            return;
+        }
+        if (status == null || status.getStatusCode() == null) {
+            // Absence of a status is the absence of a trusted answer, not a 
clean
+            // rejection; only a complete error status proves the dispatch was 
not
+            // enqueued.
+            LOG.warn("dispatch send of lance index job {} returned no status", 
job.getJobId());
+            completeNoTrusted(fresh, "dispatch send returned no status");
+            return;
+        }
+        if (status.getStatusCode() != TStatusCode.OK) {

Review Comment:
   [P2] Release the possible-live slot on a clean pre-enqueue rejection. The 
shipped BE handler always returns NOT_IMPLEMENTED_ERROR without enqueuing, yet 
completePreInvocationRejected only changes the mutation state: 
possibleLiveOwned remains true and SHOW reports a worker as possibly live until 
that BE restarts. Record a durable no-enqueue proof for the matched dispatch 
identity.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    private final LanceIndexJobManager jobManager;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManager = jobManager;
+    }
+
+    /**
+     * 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;
+        }
+        setInterval(dispatchIntervalMs());
+        try {
+            runOneRound();
+        } catch (Throwable t) {
+            LOG.warn("Failed to process one round of the lance index job 
dispatcher", t);
+        }
+    }
+
+    private void runOneRound() {
+        long nowMs = System.currentTimeMillis();
+        sweepExpiredRunningJobs(nowMs);
+        sweepReplacedProcessEpochs();
+        driveRequiredRefreshes(nowMs);
+        dispatchPendingJobs();
+    }
+
+    /**
+     * 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(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() {
+        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(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(job);
+            } catch (Throwable t) {
+                LOG.warn("failed to drive the refresh of lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private void driveOneRefresh(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(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(job.getJobId(), refreshRevision, false);
+            return;
+        }
+        finishRefreshTransition(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(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. Attempts at most
+     * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches 
per
+     * round, and never more than {@link 
Config#lance_index_job_max_inflight_per_backend}
+     * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+     */
+    private void dispatchPendingJobs() {
+        int maxPerRound = Math.max(1, 
Config.lance_index_job_max_dispatch_per_round);
+        Map<Long, Integer> inflightByBackend = countInflightByBackend();
+        int attempted = 0;
+        for (LanceIndexJob job : 
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+            if (++attempted > maxPerRound) {
+                break;
+            }
+            try {
+                tryDispatch(job, inflightByBackend);
+            } catch (Throwable t) {
+                LOG.warn("failed to dispatch lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private Map<Long, Integer> countInflightByBackend() {
+        Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+        for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+            if (job.getMutationState() == LanceIndexJobMutationState.RUNNING 
&& job.getBackendId() != null) {
+                inflightByBackend.merge(job.getBackendId(), 1, Integer::sum);
+            }
+        }
+        return inflightByBackend;
+    }
+
+    /**
+     * One dispatch attempt for one PENDING job. Every early return before
+     * markRunning leaves the job PENDING for a later round. 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 void tryDispatch(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;
+        }
+        if (localDataset && Env.getCurrentEnv().getFrontends(null).size() != 
1) {
+            // Local files are only shared by a single-node deployment.
+            return;
+        }
+        SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+        List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(
+                new 
BeSelectionPolicy.Builder().needScheduleAvailable().build(), 1);
+        if (backendIds.isEmpty()) {
+            return;
+        }
+        Backend backend = systemInfo.getBackend(backendIds.get(0));
+        if (backend == null) {
+            return;
+        }
+        if (localDataset && !isOnlyAliveBackend(systemInfo, backend.getId())) {
+            return;
+        }
+        Integer inflight = inflightByBackend.get(backend.getId());
+        if (inflight != null && inflight >= Math.max(1, 
Config.lance_index_job_max_inflight_per_backend)) {
+            return;
+        }
+        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());
+        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;
+        }
+        inflightByBackend.merge(backend.getId(), 1, Integer::sum);
+        long expectedDispatchRevision = job.getRevision() + 1;
+        LanceIndexJob fresh = jobManager.getJob(job.getJobId());
+        if (!Env.getCurrentEnv().isMaster() || fresh == null
+                || fresh.getMutationState() != 
LanceIndexJobMutationState.RUNNING
+                || fresh.getDispatchRevision() == null
+                || fresh.getDispatchRevision() != expectedDispatchRevision
+                || !invocationId.equals(fresh.getInvocationId())) {
+            // The recheck failed right before the send: no send, and no 
resend either.
+            // The job is durable RUNNING, so the deadline sweep or a matching 
callback
+            // converges it.
+            LOG.warn("lance index job {} did not survive the pre-send recheck; 
not sending", job.getJobId());
+            return;
+        }
+        TLanceIndexJobDispatch dispatch;
+        try {
+            dispatch = buildDispatch(fresh, invocationId, deadlineMs, 
beProcessEpoch,
+                    resolveStorageOptions(fresh));
+        } catch (Exception e) {

Review Comment:
   [P1] Keep a failed pre-send preparation out of UNKNOWN. During ALTER CATALOG 
RENAME, CatalogMgr temporarily removes the catalog ID; a job can be marked 
RUNNING, then resolveStorageOptions throws before any RPC, and this catch 
retains its fence and quota indefinitely as UNKNOWN. Prepare the request before 
markRunning, or durably record a proven no-send result and release the slot for 
this identity.



##########
fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java:
##########
@@ -860,6 +863,7 @@ public Env(boolean isCheckpointCatalog) {
         this.eventProcessor = new EventProcessor(mtmvService);
         this.insertOverwriteManager = new InsertOverwriteManager();

Review Comment:
   [P1] Rebind the dispatcher after loading the job-manager image. The 
constructor passes its initial manager into a final dispatcher field, but 
loadLanceIndexJobManager later replaces Env.lanceIndexJobManager with a new 
object. After restart, replay, admission, SHOW, and the master-transfer sweep 
use the restored manager while every dispatcher round scans the abandoned empty 
one, so no restored or newly admitted job advances. Resolve Env's current 
manager each round or construct the dispatcher after image load, and cover an 
image restart in the wiring test.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    private final LanceIndexJobManager jobManager;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManager = jobManager;
+    }
+
+    /**
+     * 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;
+        }
+        setInterval(dispatchIntervalMs());
+        try {
+            runOneRound();
+        } catch (Throwable t) {
+            LOG.warn("Failed to process one round of the lance index job 
dispatcher", t);
+        }
+    }
+
+    private void runOneRound() {
+        long nowMs = System.currentTimeMillis();
+        sweepExpiredRunningJobs(nowMs);
+        sweepReplacedProcessEpochs();
+        driveRequiredRefreshes(nowMs);
+        dispatchPendingJobs();
+    }
+
+    /**
+     * 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(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() {
+        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(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(job);
+            } catch (Throwable t) {
+                LOG.warn("failed to drive the refresh of lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private void driveOneRefresh(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(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(job.getJobId(), refreshRevision, false);
+            return;
+        }
+        finishRefreshTransition(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(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. Attempts at most
+     * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches 
per
+     * round, and never more than {@link 
Config#lance_index_job_max_inflight_per_backend}
+     * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+     */
+    private void dispatchPendingJobs() {
+        int maxPerRound = Math.max(1, 
Config.lance_index_job_max_dispatch_per_round);
+        Map<Long, Integer> inflightByBackend = countInflightByBackend();
+        int attempted = 0;
+        for (LanceIndexJob job : 
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+            if (++attempted > maxPerRound) {
+                break;
+            }
+            try {
+                tryDispatch(job, inflightByBackend);
+            } catch (Throwable t) {
+                LOG.warn("failed to dispatch lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private Map<Long, Integer> countInflightByBackend() {
+        Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+        for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+            if (job.getMutationState() == LanceIndexJobMutationState.RUNNING 
&& job.getBackendId() != null) {
+                inflightByBackend.merge(job.getBackendId(), 1, Integer::sum);
+            }
+        }
+        return inflightByBackend;
+    }
+
+    /**
+     * One dispatch attempt for one PENDING job. Every early return before
+     * markRunning leaves the job PENDING for a later round. 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 void tryDispatch(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;
+        }
+        if (localDataset && Env.getCurrentEnv().getFrontends(null).size() != 
1) {
+            // Local files are only shared by a single-node deployment.
+            return;
+        }
+        SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+        List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(

Review Comment:
   [P2] Choose a backend with a free job slot before consuming this attempt. 
selectBackendIdsByPolicy(..., 1) returns a random BE; when it is at the per-BE 
cap, tryDispatch returns even if another schedule-available BE is idle. With 
one full and one idle BE, half the jobs are deferred for a whole polling round 
on average. Filter/select across all available backends.



##########
regression-test/suites/external_table_p0/lance/test_lance_index_admission.groovy:
##########
@@ -64,13 +65,28 @@ suite("test_lance_index_admission", 
"p0,external,nonConcurrent") {
     assertEquals(1, quotaRows.size())
     String originalGate = gateRows[0][1].toString()
     String originalQuota = quotaRows[0][1].toString()
+    def dispatchIntervalRows = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 
'lance_index_job_dispatch_interval_second'"""
+    assertEquals(1, dispatchIntervalRows.size())
+    String originalDispatchInterval = dispatchIntervalRows[0][1].toString()
+    // The interval pin below outlives the daemon's in-flight wait on this 
shipped
+    // default, so the wait is sized correctly.
+    assertEquals("10", originalDispatchInterval)
     // The main scenario admits two jobs on one table, independently of the
     // cluster's original quota. The dedicated quota case temporarily lowers 
it.
     String suiteQuota = Math.max(2L, originalQuota.toLong()).toString()
     Throwable suiteFailure = null
 
     try {
         master_sql """ADMIN SET FRONTEND CONFIG 
("lance_index_job_max_unresolved_per_table" = "${suiteQuota}")"""
+        // Pin the dispatcher's polling interval to one hour so no dispatch 
round can
+        // fire between admission and the PENDING assertions below: this 
slice's
+        // backends answer submit_lance_index_job with a clean not-implemented 
error,
+        // which converges a dispatched job to NOT_COMMITTED and would break 
this
+        // suite's PENDING premise. A cycle already sleeping on the shipped 
interval

Review Comment:
   [P2] Synchronize with a completed dispatcher cycle before asserting PENDING. 
Setting the interval to 3600 and sleeping 11 seconds does not stop a round 
already running: a slow refresh before dispatch can reach the newly admitted 
job, or a round blocked in a backend RPC can finish and sleep its old 10-second 
interval before the next round picks up the job. The shipped BE stub then moves 
it to NOT_COMMITTED and makes these PENDING assertions flaky. Use a daemon 
barrier or scoped pause in both regression suites.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    private final LanceIndexJobManager jobManager;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManager = jobManager;
+    }
+
+    /**
+     * 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;
+        }
+        setInterval(dispatchIntervalMs());
+        try {
+            runOneRound();
+        } catch (Throwable t) {
+            LOG.warn("Failed to process one round of the lance index job 
dispatcher", t);
+        }
+    }
+
+    private void runOneRound() {
+        long nowMs = System.currentTimeMillis();
+        sweepExpiredRunningJobs(nowMs);
+        sweepReplacedProcessEpochs();
+        driveRequiredRefreshes(nowMs);
+        dispatchPendingJobs();
+    }
+
+    /**
+     * 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(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() {
+        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(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(job);
+            } catch (Throwable t) {
+                LOG.warn("failed to drive the refresh of lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private void driveOneRefresh(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(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(job.getJobId(), refreshRevision, false);
+            return;
+        }
+        finishRefreshTransition(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(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. Attempts at most
+     * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches 
per
+     * round, and never more than {@link 
Config#lance_index_job_max_inflight_per_backend}
+     * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+     */
+    private void dispatchPendingJobs() {
+        int maxPerRound = Math.max(1, 
Config.lance_index_job_max_dispatch_per_round);
+        Map<Long, Integer> inflightByBackend = countInflightByBackend();
+        int attempted = 0;
+        for (LanceIndexJob job : 
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+            if (++attempted > maxPerRound) {
+                break;
+            }
+            try {
+                tryDispatch(job, inflightByBackend);
+            } catch (Throwable t) {
+                LOG.warn("failed to dispatch lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private Map<Long, Integer> countInflightByBackend() {
+        Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+        for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+            if (job.getMutationState() == LanceIndexJobMutationState.RUNNING 
&& job.getBackendId() != null) {

Review Comment:
   [P1] Base backend capacity on possible-live ownership rather than RUNNING 
state. After an epoch change, the sweep proves old workers dead and clears 
their slots but their jobs stay RUNNING until deadline, so they block the 
restarted BE for up to an hour. Conversely deadline expiry makes a potentially 
live child UNKNOWN, so it stops counting and another job can exceed the 
intended worker cap. Both transitions run before this count in the same round; 
clear proven no-enqueue slots as well.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    private final LanceIndexJobManager jobManager;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManager = jobManager;
+    }
+
+    /**
+     * 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;
+        }
+        setInterval(dispatchIntervalMs());

Review Comment:
   [P2] Make a shorter mutable dispatch interval take effect while the daemon 
sleeps. Both changed regression suites set this to 3600 seconds, wait for the 
daemon to adopt it, then restore 10 seconds; Daemon.run is still in its 
one-hour Thread.sleep, so subsequent jobs can sit PENDING for nearly an hour. 
Use a bounded wake cadence or wake the daemon on config change.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    private final LanceIndexJobManager jobManager;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManager = jobManager;
+    }
+
+    /**
+     * 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;
+        }
+        setInterval(dispatchIntervalMs());
+        try {
+            runOneRound();
+        } catch (Throwable t) {
+            LOG.warn("Failed to process one round of the lance index job 
dispatcher", t);
+        }
+    }
+
+    private void runOneRound() {
+        long nowMs = System.currentTimeMillis();
+        sweepExpiredRunningJobs(nowMs);
+        sweepReplacedProcessEpochs();
+        driveRequiredRefreshes(nowMs);
+        dispatchPendingJobs();
+    }
+
+    /**
+     * 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(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() {
+        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(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(job);
+            } catch (Throwable t) {
+                LOG.warn("failed to drive the refresh of lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private void driveOneRefresh(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(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(job.getJobId(), refreshRevision, false);
+            return;
+        }
+        finishRefreshTransition(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(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. Attempts at most
+     * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches 
per
+     * round, and never more than {@link 
Config#lance_index_job_max_inflight_per_backend}
+     * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+     */
+    private void dispatchPendingJobs() {
+        int maxPerRound = Math.max(1, 
Config.lance_index_job_max_dispatch_per_round);
+        Map<Long, Integer> inflightByBackend = countInflightByBackend();
+        int attempted = 0;
+        for (LanceIndexJob job : 
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+            if (++attempted > maxPerRound) {
+                break;
+            }
+            try {
+                tryDispatch(job, inflightByBackend);
+            } catch (Throwable t) {
+                LOG.warn("failed to dispatch lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private Map<Long, Integer> countInflightByBackend() {
+        Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+        for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {
+            if (job.getMutationState() == LanceIndexJobMutationState.RUNNING 
&& job.getBackendId() != null) {
+                inflightByBackend.merge(job.getBackendId(), 1, Integer::sum);
+            }
+        }
+        return inflightByBackend;
+    }
+
+    /**
+     * One dispatch attempt for one PENDING job. Every early return before
+     * markRunning leaves the job PENDING for a later round. 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 void tryDispatch(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;
+        }
+        if (localDataset && Env.getCurrentEnv().getFrontends(null).size() != 
1) {
+            // Local files are only shared by a single-node deployment.
+            return;
+        }
+        SystemInfoService systemInfo = Env.getCurrentSystemInfo();
+        List<Long> backendIds = systemInfo.selectBackendIdsByPolicy(
+                new 
BeSelectionPolicy.Builder().needScheduleAvailable().build(), 1);
+        if (backendIds.isEmpty()) {
+            return;
+        }
+        Backend backend = systemInfo.getBackend(backendIds.get(0));
+        if (backend == null) {
+            return;
+        }
+        if (localDataset && !isOnlyAliveBackend(systemInfo, backend.getId())) {
+            return;
+        }
+        Integer inflight = inflightByBackend.get(backend.getId());
+        if (inflight != null && inflight >= Math.max(1, 
Config.lance_index_job_max_inflight_per_backend)) {
+            return;
+        }
+        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());
+        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;
+        }
+        inflightByBackend.merge(backend.getId(), 1, Integer::sum);
+        long expectedDispatchRevision = job.getRevision() + 1;
+        LanceIndexJob fresh = jobManager.getJob(job.getJobId());
+        if (!Env.getCurrentEnv().isMaster() || fresh == null
+                || fresh.getMutationState() != 
LanceIndexJobMutationState.RUNNING
+                || fresh.getDispatchRevision() == null
+                || fresh.getDispatchRevision() != expectedDispatchRevision
+                || !invocationId.equals(fresh.getInvocationId())) {
+            // The recheck failed right before the send: no send, and no 
resend either.
+            // The job is durable RUNNING, so the deadline sweep or a matching 
callback
+            // converges it.
+            LOG.warn("lance index job {} did not survive the pre-send recheck; 
not sending", job.getJobId());
+            return;
+        }
+        TLanceIndexJobDispatch dispatch;
+        try {
+            dispatch = buildDispatch(fresh, invocationId, deadlineMs, 
beProcessEpoch,
+                    resolveStorageOptions(fresh));
+        } catch (Exception e) {
+            // An FE-side resolution failure is not a trusted worker 
rejection, so it must
+            // not fabricate NOT_COMMITTED. The job is already RUNNING without 
a send, and
+            // the send may never happen, so converge it to UNKNOWN 
fail-closed.
+            LOG.warn("failed to prepare the dispatch of lance index job {}: 
{}", job.getJobId(), e.getMessage());
+            completeNoTrusted(fresh, "dispatch preparation failed before 
send");
+            return;
+        }
+        TStatus status;
+        try {
+            status = sendExecuteRequest(backend, dispatch);
+        } catch (Exception e) {

Review Comment:
   [P1] Separate known no-enqueue failures from ambiguous RPC failures here. An 
old BE returns Thrift UNKNOWN_METHOD for this new RPC during a rolling upgrade; 
a client-pool borrow can also fail while opening the socket before 
submitLanceIndexJob is called. Neither case can execute a worker, yet this 
catch makes the job UNKNOWN and holds its fence and quota indefinitely. Gate 
old BEs and classify proven pre-call failures separately, while retaining 
UNKNOWN for failures after invocation starts.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,522 @@
+// 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 com.google.common.collect.Maps;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 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 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.
+ */
+public class LanceIndexJobDispatcher extends MasterDaemon {
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobDispatcher.class);
+
+    private final LanceIndexJobManager jobManager;
+
+    public LanceIndexJobDispatcher(LanceIndexJobManager jobManager) {
+        super("lance index job dispatcher", dispatchIntervalMs());
+        this.jobManager = jobManager;
+    }
+
+    /**
+     * 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;
+        }
+        setInterval(dispatchIntervalMs());
+        try {
+            runOneRound();
+        } catch (Throwable t) {
+            LOG.warn("Failed to process one round of the lance index job 
dispatcher", t);
+        }
+    }
+
+    private void runOneRound() {
+        long nowMs = System.currentTimeMillis();
+        sweepExpiredRunningJobs(nowMs);
+        sweepReplacedProcessEpochs();
+        driveRequiredRefreshes(nowMs);
+        dispatchPendingJobs();
+    }
+
+    /**
+     * 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(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() {
+        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(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(job);
+            } catch (Throwable t) {
+                LOG.warn("failed to drive the refresh of lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private void driveOneRefresh(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(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(job.getJobId(), refreshRevision, false);
+            return;
+        }
+        finishRefreshTransition(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(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. Attempts at most
+     * {@link Config#lance_index_job_max_dispatch_per_round} fresh dispatches 
per
+     * round, and never more than {@link 
Config#lance_index_job_max_inflight_per_backend}
+     * in-flight jobs per backend, counted from the RUNNING snapshot 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.
+     */
+    private void dispatchPendingJobs() {
+        int maxPerRound = Math.max(1, 
Config.lance_index_job_max_dispatch_per_round);
+        Map<Long, Integer> inflightByBackend = countInflightByBackend();
+        int attempted = 0;
+        for (LanceIndexJob job : 
jobManager.getJobsNeedingDispatch(maxPerRound)) {
+            if (++attempted > maxPerRound) {
+                break;
+            }
+            try {
+                tryDispatch(job, inflightByBackend);
+            } catch (Throwable t) {
+                LOG.warn("failed to dispatch lance index job " + 
job.getJobId(), t);
+            }
+        }
+    }
+
+    private Map<Long, Integer> countInflightByBackend() {
+        Map<Long, Integer> inflightByBackend = Maps.newHashMap();
+        for (LanceIndexJob job : jobManager.getAllJobsSnapshot()) {

Review Comment:
   [P2] Avoid copying and sorting all historical jobs on every dispatcher tick. 
getAllJobsSnapshot clones every durable job and sorts by ID under the manager 
read lock, but this count only needs active backend slots. Resolved jobs are 
never removed in this slice and quotas bound only unresolved jobs, so an idle 
master pays unbounded O(history log history) work and allocation every 10 
seconds. Query or maintain just the active slot holders without sorting 
terminal history.



-- 
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