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


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobManager.java:
##########
@@ -596,6 +633,32 @@ private static boolean hasFenceIdentity(LanceIndexJob job) 
{
                 && job.getNormalizedIndexName() != null;
     }
 
+    /**
+     * True when the durable record carries the complete target identity the
+     * dispatcher needs to address a job: catalog, database, table, normalized
+     * locator, and the index/mutation identity. This is the eligibility of a
+     * PENDING job for dispatch; the dispatch quad (backend id, BE process
+     * epoch, dispatch revision, invocation id) does not exist before
+     * {@code markRunning} and is not required here. Corrupt identity-less
+     * records are never dispatchable.
+     */
+    private static boolean hasDispatchTarget(LanceIndexJob job) {
+        return job.getProvider() != null && job.getNormalizedLocator() != null
+                && job.getNormalizedIndexName() != null && 
job.getDisplayIndexName() != null
+                && job.getDbName() != null && job.getTableName() != null
+                && job.getMutationType() != null && job.getCatalogId() > 0;
+    }
+
+    /**
+     * True when the record additionally carries the full dispatch quad 
recorded
+     * by {@code markRunning}. The possible-live sweep requires the quad: only
+     * with it can {@code recordTerminationProof} address and release the slot.
+     */
+    private static boolean hasDispatchIdentity(LanceIndexJob job) {
+        return hasDispatchTarget(job) && job.getBackendId() != null && 
job.getBeProcessEpoch() != null

Review Comment:
   [P2] Sweep slot holders using only the identity the proof writer needs. A 
replayed RUNNING or UNKNOWN job can keep its backend ID, epoch, and invocation 
while lacking an unrelated target field or the legacy `dispatchRevision` (which 
`dispatchRevisionOf` already falls back from). 
`countPossibleLiveSlotsByBackend` still charges its slot, but this helper 
excludes it from the epoch sweep, so even a replaced BE process cannot free 
capacity. Align this filter with `recordTerminationProof`'s matching fields.



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

Review Comment:
   [P2] Publish both mutable dispatcher switches to the daemon thread. ADMIN 
SET writes this plain static pause flag under `ConfigBase`'s monitor, while the 
dispatcher reads it without that monitor or a volatile read; the regression 
suites use separate JDBC connections for SET and admission, so their claimed 
pre-admission barrier has no guaranteed visibility. The local-file mutation 
switch below has the same issue: a stale `true` can permit local dispatch after 
an operator disables it. Use volatile fields or shared synchronization.



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

Review Comment:
   [P3] Base refresh retries on the last refresh attempt time. A FAILED job may 
still own a possible-live slot; a later CHILD_REAPED or BE epoch proof updates 
the job's generic `updateTimeMs` without attempting refresh. If refresh fails 
at t=0 and a proof arrives at t=299s, this check postpones the next 300s retry 
to t=599s while its fence and quota remain held. Keep a separate refresh 
failure timestamp for this throttle.



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

Review Comment:
   [P2] Preserve the backend rejection code in the job result. In this build 
the BE always returns `NOT_IMPLEMENTED_ERROR`, but this branch persists only 
`PRE_INVOCATION_RESOURCE_REJECTED` with a generic message, so SHOW and logs 
cannot tell an unavailable worker from resource or other clean rejections. 
Include at least the bounded `TStatusCode` name in the persisted reason and 
log; keep raw BE messages out unless sanitized.



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

Review Comment:
   [P2] Keep blocking sends off the sole sweep and refresh thread. This loop 
waits synchronously for each `submitLanceIndexJob` reply, while deadline, BE 
epoch, and refresh work runs only before the loop. The backend pool's default 
RPC timeout is 60 seconds and the round can send 16 jobs, so heartbeating BEs 
with stalled job RPCs can delay all other jobs' lifecycle work for about 16 
minutes despite a 10-second interval. Bound total send time per round or 
dispatch asynchronously while preserving one durable send per identity.



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

Review Comment:
   [P3] Bind the refresh to the persisted catalog ID. This reads catalog A by 
job ID but passes A's mutable name into `handleRefreshTable`, which resolves 
that name again. If A is renamed and catalog B takes the old name between those 
steps, this job refreshes B and records A as DONE or FAILED based on B's 
metadata path; a B-side failure delays A's fence release until another retry. 
Refresh by ID or verify the resolved catalog ID at the refresh boundary.



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