LanceIndexJobDispatcher.java

// 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,
                    "backend returned a clean error status before enqueueing the dispatch");
        }
        // OK: enqueued exactly once. The result arrives through the report callback;
        // nothing more is done here, and the deadline sweep bounds the wait.
        return true;
    }

    /**
     * A send failure that proves the dispatch never reached a worker: the
     * client could not be borrowed, so no connection was ever established, or
     * an old backend answered UNKNOWN_METHOD for the new RPC during a rolling
     * upgrade, so it provably never enqueued the dispatch. Every other failure
     * after the invocation may have started (a broken write, a read timeout)
     * stays ambiguous and converges UNKNOWN instead.
     */
    public static class PreInvocationSendException extends Exception {
        public PreInvocationSendException(String message, Throwable cause) {
            super(message, cause);
        }
    }

    /**
     * Sends one dispatch to the backend's thrift service and returns its status.
     * The connection is borrowed per send, returned only when the call completed,
     * and invalidated after a failed call. Two failure families are wrapped into
     * {@link PreInvocationSendException} because they prove the dispatch never
     * reached a worker: a borrow failure means no connection was ever
     * established, and an UNKNOWN_METHOD answer means an old backend (a rolling
     * upgrade not yet serving this RPC) provably never enqueued the dispatch —
     * with no BE capability bit to gate on, classifying that answer is what lets
     * a dispatch retry on another, already-upgraded backend. Everything thrown
     * later propagates unwrapped as ambiguous. Test seam: subclasses override
     * this method to record the request or inject faults without a live client
     * pool.
     */
    protected TStatus sendExecuteRequest(Backend backend, TLanceIndexJobDispatch dispatch) throws Exception {
        TNetworkAddress address = new TNetworkAddress(backend.getHost(), backend.getBePort());
        BackendService.Client client = null;
        boolean callCompleted = false;
        try {
            try {
                client = ClientPool.backendPool.borrowObject(address);
            } catch (Exception e) {
                throw new PreInvocationSendException(
                        "no backend client could be borrowed; the dispatch was never sent", e);
            }
            TStatus status;
            try {
                status = client.submitLanceIndexJob(dispatch);
            } catch (TApplicationException e) {
                if (e.getType() == TApplicationException.UNKNOWN_METHOD) {
                    throw new PreInvocationSendException(
                            "backend does not serve submitLanceIndexJob (rolling upgrade); not enqueued", e);
                }
                throw e;
            }
            callCompleted = true;
            return status;
        } finally {
            if (client != null) {
                if (callCompleted) {
                    ClientPool.backendPool.returnObject(address, client);
                } else {
                    ClientPool.backendPool.invalidateObject(address, client);
                }
            }
        }
    }

    /**
     * Builds the wire request from the job record and the dispatch identity
     * that markRunning is about to make durable (the dispatch revision is the
     * pre-computed {@code job.revision + 1}; the pre-send recheck pins that the
     * durable record landed with exactly this identity). Definition fields a
     * DROP never carries travel as the empty string: the wire marks them
     * required, and the worker only reads them for CREATE and REPLACE.
     */
    private TLanceIndexJobDispatch buildDispatch(LanceIndexJob job, long dispatchRevision, String invocationId,
            long deadlineMs, long beProcessEpoch, Map<String, String> storageOptions) {
        TLanceIndexJobDispatch dispatch = new TLanceIndexJobDispatch();
        dispatch.setJobId(job.getJobId());
        dispatch.setDispatchRevision(dispatchRevision);
        dispatch.setInvocationId(invocationId);
        dispatch.setBeProcessEpoch(beProcessEpoch);
        dispatch.setDeadlineMs(deadlineMs);
        dispatch.setMutationType(TLanceIndexMutationType.valueOf(job.getMutationType().name()));
        dispatch.setIndexName(job.getDisplayIndexName());
        dispatch.setColumnName(job.getColumnName() == null ? "" : job.getColumnName());
        dispatch.setIndexType(job.getIndexType() == null ? "" : job.getIndexType());
        if (job.getPropertiesJson() != null) {
            dispatch.setPropertiesJson(job.getPropertiesJson());
        }
        dispatch.setIfNotExists(job.isIfNotExists());
        dispatch.setIfExists(job.isIfExists());
        dispatch.setDatasetUri(job.getNormalizedLocator());
        dispatch.setAdmittedDatasetVersion(job.getAdmittedDatasetVersion());
        dispatch.setSchemaContractJson(job.getSchemaContract() == null ? ""
                : GsonUtils.GSON.toJson(job.getSchemaContract()));
        if (!storageOptions.isEmpty()) {
            dispatch.setStorageOptions(storageOptions);
        }
        return dispatch;
    }

    /**
     * Resolves the storage options of one dataset at send time from the
     * catalog's current storage properties, in the vocabulary of the provider
     * the dataset URI routes to. The result is used for this dispatch only: it
     * is never persisted in the job record and never logged, so a rotated
     * credential takes effect on the next dispatch without any journal record.
     */
    private Map<String, String> resolveStorageOptions(LanceIndexJob job) {
        CatalogIf catalog = Env.getCurrentEnv().getCatalogMgr().getCatalog(job.getCatalogId());
        if (!(catalog instanceof LanceExternalCatalog)) {
            throw new IllegalStateException(
                    "catalog of lance index job " + job.getJobId() + " does not resolve to a Lance catalog");
        }
        return LanceStorageOptions.fromDorisStorageProperties(job.getNormalizedLocator(),
                ((LanceExternalCatalog) catalog).getCatalogProperty().getOrderedStoragePropertiesList());
    }

    private void completeNoTrusted(LanceIndexJobManager jobManager, LanceIndexJob job, String reason) {
        boolean completed = jobManager.completeWithResult(job.getJobId(),
                dispatchRevisionOf(job), job.getInvocationId(), job.getBeProcessEpoch(),
                new LanceIndexJobResult(LanceIndexJobResultCode.NO_TRUSTED_RESULT,
                        LanceIndexJobCompletionReason.NONE, reason, false));
        if (!completed) {
            LOG.warn("no-trusted-result convergence skipped for lance index job {}: already converged by a callback"
                    + " or sweep", job.getJobId());
        }
    }

    private void completePreInvocationRejected(LanceIndexJobManager jobManager, LanceIndexJob job, String reason) {
        boolean completed = jobManager.completeProvenNoEnqueue(job.getJobId(),
                dispatchRevisionOf(job), job.getInvocationId(), job.getBeProcessEpoch(),
                new LanceIndexJobResult(LanceIndexJobResultCode.PRE_INVOCATION_RESOURCE_REJECTED,
                        LanceIndexJobCompletionReason.NONE, reason, false));
        if (!completed) {
            LOG.warn("rejection convergence skipped for lance index job {}: already converged by a callback or sweep",
                    job.getJobId());
        }
    }

    /**
     * True for datasets on the local filesystem: a scheme-less absolute path or
     * a {@code file://} URI, matching the provider routing of the dataset URL.
     */
    private static boolean isLocalFileDataset(String normalizedLocator) {
        int separator = normalizedLocator.indexOf("://");
        if (separator < 0) {
            return true;
        }
        return "file".equals(normalizedLocator.substring(0, separator).toLowerCase(Locale.ROOT));
    }

    private static boolean isOnlyAliveBackend(SystemInfoService systemInfo, long backendId) {
        List<Long> aliveBackendIds = systemInfo.getAllBackendIds(true);
        return aliveBackendIds.size() == 1 && aliveBackendIds.get(0) == backendId;
    }

    private static long dispatchRevisionOf(LanceIndexJob job) {
        return job.getDispatchRevision() == null ? job.getRevision() : job.getDispatchRevision();
    }
}