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


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobReportHandler.java:
##########
@@ -0,0 +1,187 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.datasource.lance.job;
+
+import org.apache.doris.thrift.TLanceIndexJobReport;
+import org.apache.doris.thrift.TLanceIndexTerminationProof;
+
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.nio.charset.StandardCharsets;
+import java.security.MessageDigest;
+import java.util.Objects;
+
+/**
+ * Applies one typed result envelope reported by a backend to the durable job
+ * record. The dispatch identity has two halves, checked at two layers:
+ * the journaled identity quad (dispatch revision, invocation id, BE process
+ * epoch, backend) is matched inside {@link LanceIndexJobManager} and rejects
+ * reports that are merely stale, while the per-dispatch invocation secret is
+ * matched here, before the manager is touched at all, and rejects FORGED ones.
+ * The second half exists because the FE thrift server cannot authenticate its
+ * caller and SHOW LANCE INDEX JOB publishes every other identity field, so a
+ * client that can merely reach the port could otherwise assemble a
+ * well-matching envelope; only the secret - handed to the selected BE inside
+ * the dispatch request and never shown or logged - is unforgeable. An
+ * envelope whose secret echo is missing, blank, or wrong, or whose durable
+ * record carries no secret (a legacy record from before the field existed),
+ * is unauthenticated: the whole envelope is dropped, its termination proof
+ * included, because a forged CHILD_REAPED would release the possible-live
+ * slot of a worker that may still be live.
+ *
+ * <p>Beyond authentication this is a thin shim over the manager transitions:
+ * result classification lives in {@link LanceIndexJobManager}, so a stale or
+ * identity-mismatched report only logs a warning and changes nothing. A
+ * malformed envelope (missing result code, a code this FE does not know, or a
+ * sanitized message past the durable bound) has its result dropped rather
+ * than trusted; the job then converges through the dispatcher's deadline
+ * sweep. Only the typed codes are read: message text is never inspected to
+ * infer an outcome.
+ *
+ * <p>A termination proof is validated independently of the result, so a
+ * CHILD_REAPED proof is recorded first and still lands when the result of the
+ * same envelope is malformed: reaping the exact child process proves that
+ * process ended, and dropping that proof together with the result would
+ * strand the possible-live slot until the backend process is replaced.
+ *
+ * <p>The handler runs on the report RPC thread and performs no I/O beyond the
+ * manager's own edit-log write. It starts no refresh: the metadata refresh a
+ * completed job may owe is driven by the dispatcher daemon, not here.
+ */
+public class LanceIndexJobReportHandler {
+
+    private static final Logger LOG = 
LogManager.getLogger(LanceIndexJobReportHandler.class);
+
+    private final LanceIndexJobManager jobManager;
+
+    public LanceIndexJobReportHandler(LanceIndexJobManager jobManager) {
+        this.jobManager = Objects.requireNonNull(jobManager, "jobManager");
+    }
+
+    /**
+     * Handles one report: authentication comes first and gates everything 
else,
+     * then a matched report completes the job with its classified result, and 
a
+     * CHILD_REAPED termination proof additionally releases the possible-live
+     * slot, because reaping the exact child process proves that process ended
+     * (which still says nothing about the outcome). The proof is recorded
+     * before the result is parsed: the two are validated independently, and a
+     * malformed result must not take a valid proof down with it.
+     */
+    public void handle(TLanceIndexJobReport report) {
+        if (report == null) {
+            LOG.warn("dropping null lance index job report");
+            return;
+        }
+        if (!isAuthenticated(report)) {
+            // Deliberately names no secret material, neither the expected nor 
the
+            // presented one: the log only records that the envelope was 
rejected.
+            LOG.warn("dropping unauthenticated lance index job report for job 
{}: invocation secret mismatch",
+                    report.getJobId());
+            return;
+        }
+        if (report.getTerminationProof() == 
TLanceIndexTerminationProof.CHILD_REAPED) {

Review Comment:
   [P2] Apply a valid report and its child proof in one transition. 
`recordChildReaped()` journals a proof while leaving the job RUNNING, then 
`completeWithResult()` journals the result separately. If the deadline sweep 
runs between those calls, it marks the job UNKNOWN and `completeWithResult()` 
rejects this authenticated NATIVE_OK report, leaving a known completed mutation 
fenced until FORCE_RELEASE. Validate the result first and commit the result and 
optional proof atomically; keep the proof-only path for malformed results.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobDispatcher.java:
##########
@@ -0,0 +1,843 @@
+// 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.DatabaseIf;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.TableIf;
+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.LanceIndexDatasetCheck;
+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.commons.codec.binary.Hex;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.apache.thrift.TApplicationException;
+
+import java.security.SecureRandom;
+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. Each attempt also mints a random invocation secret, whose
+ * role completes the dispatch identity rather than duplicating it: the
+ * epoch/invocation pair proves a report is fresh (it matches the current
+ * durable dispatch), while the secret proves its reporter is the BE this
+ * dispatcher actually selected - every other identity field is readable from
+ * SHOW LANCE INDEX JOB, and the FE thrift server cannot authenticate its
+ * caller, so only the never-shown secret can reject a forged report. 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. The blocking time of the send loop and of the
+ * refresh loop is bounded per round (one backend RPC timeout each), because
+ * this thread is also the only thread running the sweeps — see
+ * {@link #dispatchPendingJobs()} and
+ * {@link #driveRequiredRefreshes(LanceIndexJobManager, long)}.
+ *
+ * <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;
+
+    /**
+     * Entropy of one dispatch secret: 16 bytes = 128 bits, well past 
guessability
+     * for a token whose only threat model is a caller forging a report from
+     * outside. Hex-encoded on the wire, so the token is 32 characters.
+     */
+    private static final int INVOCATION_SECRET_BYTES = 16;
+
+    /** Source of the per-dispatch report secret; shared, since SecureRandom 
is thread-safe. */
+    private static final SecureRandom SECURE_RANDOM = new SecureRandom();
+
+    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);
+            }
+        }
+    }
+
+    /**
+     * The clock the FAILED-refresh throttle measures from: the dedicated
+     * refresh-failure timestamp, falling back to the generic update time only 
for
+     * a legacy record replayed before the field existed. Measuring from the
+     * generic time would let an unrelated transition — a CHILD_REAPED proof 
or an
+     * epoch-gone release bumping it — postpone the next retry by a full 
interval
+     * while the fence and quota stay held.
+     */
+    private static long refreshThrottledSinceMs(LanceIndexJob job) {
+        return job.getRefreshFailureTimeMs() == null ? job.getUpdateTimeMs() : 
job.getRefreshFailureTimeMs();
+    }
+
+    /**
+     * 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.
+     *
+     * <p>The blocking time of this loop is bounded per round exactly like the
+     * dispatch send loop ({@link #dispatchPendingJobs()}): one owed refresh 
can
+     * initialize a Lance catalog and list remote databases or tables, so a 
slow
+     * provider can consume a whole metadata timeout, and several owed jobs
+     * could multiply that delay before the next deadline or epoch sweep and
+     * before PENDING jobs dispatch. At most one backend RPC timeout of 
blocking
+     * refresh work runs per round; the remaining owed jobs are logged and
+     * deferred to the next round with their durable state untouched (REQUIRED
+     * or FAILED, never a stranded RUNNING), so the FAILED throttle and the
+     * markRefreshRunning semantics keep their meaning.
+     */
+    private void driveRequiredRefreshes(LanceIndexJobManager jobManager, long 
nowMs) {
+        // fe.conf bypasses the validator, so a non-positive timeout is 
clamped to keep
+        // at least one refresh attempt per round.
+        long blockingBudgetMs = Math.max(1L, Config.backend_rpc_timeout_ms);
+        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 - refreshThrottledSinceMs(job)
+                                < Config.lance_index_job_refresh_retry_second 
* 1000L) {
+                    continue;
+                }
+                if (blockingBudgetMs <= 0) {

Review Comment:
   [P2] Give untouched refresh jobs a turn before retrying a failing prefix. 
`getJobsNeedingRefresh()` can return the same map order each round, and this 
budget break stops at the first slow attempt. With six early FAILED jobs each 
taking about 61 seconds and the default 300-second retry, the first is eligible 
again when the sixth has run; a later REQUIRED job is never reached, so its 
fence and quota stay held indefinitely. Rotate the scan or prioritize 
first-time REQUIRED refreshes ahead of FAILED retries.



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