u70b3 commented on code in PR #67754: URL: https://github.com/apache/doris/pull/67754#discussion_r4141836590
########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobCleaner.java: ########## @@ -0,0 +1,65 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.datasource.lance.job; + +import org.apache.doris.catalog.Env; +import org.apache.doris.common.Config; +import org.apache.doris.common.util.MasterDaemon; + +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.List; + +/** + * Master-only retention sweeper of resolved Lance index job records. Resolved + * records serve audit only; once one has been resolved for longer than + * {@code Config.lance_index_job_keep_max_second} it is removed on all FEs + * through one batch edit-log record per round, so every FE serves the same + * SHOW LANCE INDEX JOBS view. Unresolved records are never removed, no matter + * their age (fail-closed). + * + * <p>MasterDaemon itself performs no master check: the master-only semantics + * come entirely from being started in {@code Env.startMasterOnlyDaemonThreads}. + * Both retention configs are mutable and re-read every round. + */ +public class LanceIndexJobCleaner extends MasterDaemon { + private static final Logger LOG = LogManager.getLogger(LanceIndexJobCleaner.class); + + /** Upper bound of jobs removed per round, mirroring MAX_REMOVE_TXN_PER_ROUND. */ + private static final int MAX_REMOVE_PER_ROUND = 1024; + + public LanceIndexJobCleaner() { + super("LanceIndexJobCleaner", Config.lance_index_job_clean_interval_second * 1000L); + } + + @Override + protected void runAfterCatalogReady() { + // Both configs are mutable; re-read them every round. + setInterval(Config.lance_index_job_clean_interval_second * 1000L); Review Comment: done: both conversions are clamped at consumption — the interval to at least one second (a non-positive period would kill the thread or spin it) and the keep window to a saturated Long.MAX_VALUE milliseconds, never a wrapped negative window; fe.conf-loaded values are covered because the clamp does not depend on the validator (keepWindowConversionSaturatesInsteadOfOverflowing, e2a8747744). ########## fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ResolveLanceIndexJobCommand.java: ########## @@ -0,0 +1,367 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.nereids.trees.plans.commands; + +import org.apache.doris.catalog.DatabaseIf; +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.TableIf; +import org.apache.doris.common.AnalysisException; +import org.apache.doris.common.DdlException; +import org.apache.doris.common.ErrorCode; +import org.apache.doris.datasource.CatalogIf; +import org.apache.doris.datasource.CatalogMgr; +import org.apache.doris.datasource.ExternalDatabase; +import org.apache.doris.datasource.ExternalTable; +import org.apache.doris.datasource.lance.LanceExternalCatalog; +import org.apache.doris.datasource.lance.LanceIndexMutationValidator; +import org.apache.doris.datasource.lance.job.LanceIndexJob; +import org.apache.doris.datasource.lance.job.LanceIndexJobManager; +import org.apache.doris.datasource.lance.job.LanceIndexJobMutationState; +import org.apache.doris.mysql.privilege.PrivPredicate; +import org.apache.doris.nereids.trees.plans.PlanType; +import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor; +import org.apache.doris.qe.ConnectContext; +import org.apache.doris.qe.StmtExecutor; + +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.nio.charset.StandardCharsets; + +/** + * RESOLVE LANCE INDEX JOB <jobId> AS FORCE_RELEASE COMMENT '<note>' — the operator + * escape hatch that durably releases a job whose mutation outcome is UNKNOWN (design section + * 7.1). RESOLVE is deliberately not gated by {@code enable_lance_index_mutation}: the gate + * controls mutation admission, while FORCE must stay available exactly when the gate is off. + * + * <p>The release protocol keeps the fence, the quota charge and the possible-live slot while + * it performs one authoritative latest-metadata read and one external-table refresh with the + * current credentials of the surviving catalog, both outside every catalog/manager lock; only + * then does the durable release transfer inside the admission critical section + * ({@code captureLanceIndexTarget} → lock-free read/refresh → {@code withLanceIndexAdmission} + * recheck → manager write lock), serialized against DROP CATALOG and identity ALTER exactly + * like admission. Any failure before the transfer is the typed + * {@code ERR_LANCE_INDEX_JOB_RESOLUTION_INCOMPLETE}: nothing is written, nothing is released, + * and the operator fixes the cause and retries the same statement. + * + * <p>Target resolution (design section 7.1 step 1) is three-valued. RESOLVED means the + * persisted names resolve and the catalog's current durable dataset locator still matches + * the job's — the same revalidation SHOW LANCE INDEX JOBS applies, so a repointed dataset + * reusing the same names never turns a stale name into table-level authorization. MISSING + * means the catalog, database or table is verifiably absent, or the locator positively + * points at a different dataset: that is the orphan family — a full orphan (catalog gone) + * has no credentials to read with and nothing to invalidate, so it is released directly + * after global ADMIN authorization, while a half-orphan skips the authoritative read and + * refreshes with {@code ignoreIfNotExists=true} as a best-effort invalidation. FAILED means + * a resolution that errors out, or a locator that cannot be resolved right now: never an + * orphan verdict — after ADMIN authorization the statement fails with the typed 5105 so the + * fence is kept when "table gone" cannot be told apart from "network down". SHOW fails the + * same uncertainty closed by hiding the row; RESOLVE fails it closed by not releasing. + * + * <p>Non-disclosure (design section 8): the job is loaded first and authorized against its + * persisted target — table-level ALTER when the target resolves, global ADMIN otherwise — and + * a missing job and an unauthorized job share the same fixed ERR_LANCE_INDEX_JOB_NOT_FOUND + * response naming only the job id. The 5104 state rejection and the 5105 resolution failure + * are only visible to an already authorized caller. + * + * <p>Success returns an OK packet carrying one warning row with {@link #LATE_COMMIT_WARNING}, + * the same text persisted as the job's durable {@code forceWarning}: the old worker may still + * overwrite, remove, or reintroduce the index name; the mutation outcome remains UNKNOWN. + * Retrying FORCE on an already released job is an idempotent success returning the existing + * release record, never an error. + */ +public class ResolveLanceIndexJobCommand extends Command implements ForwardWithSync { + /** + * The late-commit warning (design section 7.1), returned in the OK packet and persisted + * verbatim as the durable {@code forceWarning}; bounded well under + * {@link LanceIndexJob#MAX_FORCE_TEXT_BYTES}. + */ + static final String LATE_COMMIT_WARNING = + "the old worker may still overwrite, remove, or reintroduce the index name; " + + "the mutation outcome remains UNKNOWN"; + + private static final Logger LOG = LogManager.getLogger(ResolveLanceIndexJobCommand.class); + + private final long jobId; + private final String comment; + + public ResolveLanceIndexJobCommand(long jobId, String comment) { + super(PlanType.RESOLVE_LANCE_INDEX_JOB_COMMAND); + this.jobId = jobId; + this.comment = comment; + } + + public long getJobId() { + return jobId; + } + + public String getComment() { + return comment; + } + + @Override + public void run(ConnectContext ctx, StmtExecutor executor) throws Exception { + Env env = Env.getCurrentEnv(); + LanceIndexJobManager manager = env.getLanceIndexJobManager(); + // 1. Load the job without disclosing any field (design section 7.1 step 1). + LanceIndexJob job = manager.getJob(jobId); + if (job == null) { + throw notFound(); + } + // 2. Resolve and authorize against the persisted target before any state is revealed: + // table-level ALTER when the target resolves, global ADMIN for the orphan family + // and for a target whose resolution failed outright. + CatalogMgr catalogMgr = env.getCatalogMgr(); + CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog = catalogMgr.getCatalog(job.getCatalogId()); + TargetResolution resolution = resolveTarget(catalog, job); + boolean authorized = resolution == TargetResolution.RESOLVED + ? env.getAccessManager().checkTblPriv(ctx, catalog.getName(), job.getDbName(), job.getTableName(), + PrivPredicate.ALTER) + : env.getAccessManager().checkGlobalPriv(ctx, PrivPredicate.ADMIN); + if (!authorized) { + throw notFound(); + } + // 3. Idempotent replay: a retry returns the existing release record (section 7.1). + // This deliberately precedes the resolution-failure rejection: once the release + // has landed, a retry during a provider outage is a success, not a 5105. + if (job.isForceReleased()) { + ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING); + return; + } + // 4. Only UNKNOWN may be force-released; a null mutation state reads as UNKNOWN, + // same as the manager's own gate. The state rejection also precedes the + // resolution-failure rejection: for a terminal job the accurate answer is 5104, + // not a 5105 claiming the job still holds its fence. + if (job.getMutationState() != null && job.getMutationState() != LanceIndexJobMutationState.UNKNOWN) { + throw new AnalysisException(ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN.formatErrorMsg(jobId), + ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN); + } + if (resolution == TargetResolution.FAILED) { + // Never an orphan verdict: "table gone" cannot be told apart from "network down", + // so nothing is released and nothing beyond the typed error is disclosed; the + // operator fixes the cause and retries the same statement (design 7.1 step 4). + throw incompleteResolution("the persisted target could not be resolved with current catalog" + + " metadata; see fe.log for the cause"); + } + boolean targetResolves = resolution == TargetResolution.RESOLVED; + // 5. The grammar makes COMMENT mandatory; here the note must also be non-empty after + // trimming and fit the durable force text bound. + String note = comment == null ? "" : comment.trim(); + if (note.isEmpty()) { + throw new AnalysisException("force release note must not be empty", + ErrorCode.ERR_LANCE_INDEX_INVALID); + } + if (note.getBytes(StandardCharsets.UTF_8).length > LanceIndexJob.MAX_FORCE_TEXT_BYTES) { + throw new AnalysisException("force release note exceeds " + LanceIndexJob.MAX_FORCE_TEXT_BYTES + + " UTF-8 bytes", ErrorCode.ERR_LANCE_INDEX_INVALID); + } + // 6-9. Branch on the orphan state, then the durable release transfer. + String actor = ctx.getQualifiedUser(); + boolean released; + if (catalog == null) { + // Full orphan: no credentials survive to read with and nothing can be + // invalidated, so the release goes straight to the manager write lock. + released = manager.forceRelease(jobId, job.getRevision(), actor, note, LATE_COMMIT_WARNING); + } else { + released = releaseWithLiveCatalog(env, catalogMgr, manager, catalog, targetResolves, job, actor, note); + } + if (!released) { + // 10. The expected-revision transfer lost a race. A concurrent FORCE_RELEASE + // that already landed makes this an idempotent success; anything else means + // the job left UNKNOWN concurrently (UNKNOWN has no other outgoing + // transition), so the pinned not-UNKNOWN wording stays accurate. + LanceIndexJob reread = manager.getJob(jobId); + if (reread != null && reread.isForceReleased()) { + ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING); + return; + } + throw new AnalysisException(ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN.formatErrorMsg(jobId), + ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN); + } + // 11. The OK packet carries the late-commit warning; it survives the forward chain + // byte-identically (proxyExecute serializes the master state). + ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING); + } + + /** + * The live-catalog release path: capture the target identity, then — holding no lock — + * perform the authoritative read and the refresh, and finally transfer the release inside + * the admission critical section. Lock order stays CatalogMgr then LanceIndexJobManager. + */ + private boolean releaseWithLiveCatalog(Env env, CatalogMgr catalogMgr, LanceIndexJobManager manager, + CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog, boolean targetResolves, LanceIndexJob job, + String actor, String note) throws Exception { + if (!(catalog instanceof LanceExternalCatalog) + || ((LanceExternalCatalog) catalog).isRestCatalogConfigured()) { + // Defensive: admission never targets a REST catalog and a non-Lance catalog + // cannot hold a Lance fence, so no UNKNOWN job should ever resolve here. + LanceIndexMutationValidator.rejectUnsupportedOperation( + "RESOLVE LANCE INDEX JOB AS FORCE_RELEASE", "catalog '" + catalog.getName() + "'"); + } + LanceExternalCatalog lanceCatalog = (LanceExternalCatalog) catalog; + CatalogMgr.LanceIndexTarget target; + try { + target = catalogMgr.captureLanceIndexTarget(lanceCatalog); + } catch (DdlException e) { + throw incompleteResolution(e.getMessage()); + } + if (targetResolves) { + // 7. One authoritative latest-metadata read with current credentials, proving + // the dataset is reachable before anything is released. + authoritativeRead(lanceCatalog, catalog, job); + } + // 8. Invalidate the external table and broadcast the refresh to every FE. A + // half-orphan target is invalidated best-effort (missing db/table is a no-op). + try { + env.getRefreshManager().handleRefreshTable(catalog.getName(), job.getDbName(), job.getTableName(), + !targetResolves); + } catch (DdlException e) { + throw incompleteResolution(e.getMessage()); + } + // 9. The durable transfer rechecks the catalog identity under the read lock, then + // runs the revision-checked release under the manager write lock. + try { + return catalogMgr.withLanceIndexAdmission(lanceCatalog, target, + () -> manager.forceRelease(jobId, job.getRevision(), actor, note, LATE_COMMIT_WARNING)); + } catch (DdlException e) { + throw incompleteResolution(e.getMessage()); + } + } + + /** + * The authoritative latest-metadata read over the remote names of the resolved db/table, + * outside every lock (the loader owns its deadline-bound JNI read). Every loader failure + * already passed through the catalog's sanitized root-cause chain, so its message is safe + * to echo — locator, credentials and dataset uri are masked there. + */ + private void authoritativeRead(LanceExternalCatalog lanceCatalog, + CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog, LanceIndexJob job) throws AnalysisException { + DatabaseIf<? extends TableIf> db; + TableIf table; + try { + db = catalog.getDbNullable(job.getDbName()); + table = db == null ? null : db.getTableNullable(job.getTableName()); + } catch (Exception e) { + // A remote-metadata blip during resolution is not an orphan verdict; the provider + // message is not echoed here because it never crossed the sanitized chain. + LOG.warn("lance index job {}: target re-resolution failed before the authoritative read", + job.getJobId(), e); + throw incompleteResolution("the persisted target could not be resolved with current catalog" + + " metadata; see fe.log for the cause"); + } + if (db == null || table == null) { + // The target resolved at authorization time but is gone now: keep the fence and + // let the retry take the half-orphan branch under global ADMIN. + throw incompleteResolution("the persisted target table no longer resolves; retry the statement"); + } + if (!(db instanceof ExternalDatabase) || !(table instanceof ExternalTable)) { + // Boundary guard: a Lance catalog must serve external relations, but the command + // boundary does not trust that invariant (no raw ClassCastException to the user). + LOG.warn("lance index job {}: target relation is not external: db={}, table={}", + job.getJobId(), db.getClass().getName(), table.getClass().getName()); + throw incompleteResolution("the persisted target is not an external relation;" + + " see fe.log for the cause"); + } + String remoteDb = ((ExternalDatabase) db).getRemoteName(); + String remoteTable = ((ExternalTable) table).getRemoteName(); + try { + lanceCatalog.loadTableIndexAdmissionSnapshot(remoteDb, remoteTable); + } catch (Exception e) { + throw incompleteResolution(e.getMessage()); + } + } + + /** + * The three-way verdict on the job's persisted target. Deliberately not SHOW's + * {@code targetResolves}: SHOW must keep listing through provider outages, so it folds + * every failed resolution into its orphan rule; RESOLVE takes a durable action on the + * verdict and only releases on positive evidence, so a failed resolution stays its own + * outcome here. + */ + private enum TargetResolution { + /** Names resolve and the catalog's current durable locator still matches the job's. */ + RESOLVED, + /** Catalog, database or table verifiably absent, or the locator positively repointed. */ + MISSING, + /** Resolution errored out, or the locator cannot be resolved right now. */ + FAILED + } + + /** + * Resolves the persisted target once, up front, distinguishing "verifiably gone" (the + * orphan family, releasable under global ADMIN) from "could not tell" (fail with 5105, + * keep the fence). The locator leg mirrors {@link ShowLanceIndexJobsCommand}: a null + * current locator means the provider is unreachable or the names no longer resolve + * remotely, which is absence of evidence either way — it fails closed here instead of + * granting the half-orphan release path. + */ + static TargetResolution resolveTarget(CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog, + LanceIndexJob job) { + if (catalog == null) { + return TargetResolution.MISSING; + } + DatabaseIf<? extends TableIf> db; + try { + db = catalog.getDbNullable(job.getDbName()); + } catch (Exception e) { + LOG.warn("lance index job {}: target database resolution failed", job.getJobId(), e); + return TargetResolution.FAILED; + } + if (db == null) { + return TargetResolution.MISSING; + } + TableIf table; + try { + table = db.getTableNullable(job.getTableName()); + } catch (Exception e) { + LOG.warn("lance index job {}: target table resolution failed", job.getJobId(), e); + return TargetResolution.FAILED; + } + if (table == null) { + return TargetResolution.MISSING; + } + if (!(catalog instanceof LanceExternalCatalog)) { + return TargetResolution.RESOLVED; + } + String currentLocator = ((LanceExternalCatalog) catalog).resolveCurrentIndexJobLocator( + job.getDbName(), job.getTableName()); + if (currentLocator == null) { + // The catalog folds provider outages and unresolvable names into null (and logs Review Comment: done: the cause is preserved and logged at the lookup boundary — resolveCurrentIndexJobLocator logs the failure it folds into null, and the new checkIndexJobDataset logs every UNRESOLVED cause (the sanitized client chain already masks locators and credentials), so fe.log now carries the namespace timeout / credential story behind a 5105 (59d120c268, boundary warn added in the dispatcher PR 089a93147d which this branch includes). ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobManager.java: ########## @@ -498,6 +543,122 @@ private void downgradeRunningRefresh(LanceIndexJob sweepCandidate) { } } + /** + * Durable FORCE_RELEASE of an UNKNOWN job, the operator escape hatch for a + * mutation whose outcome cannot be established. Validation order is fixed: + * unknown job -> idempotent short-circuit on an already released record + * (deliberately ahead of the revision CAS, so an operator retry carrying the + * pre-release revision still observes the existing release record, unlike + * the stale-callback convention) -> revision CAS -> UNKNOWN state gate + * (a null mutation state is treated as UNKNOWN, the safe direction, mirroring + * {@link LanceIndexJob#isUnresolved()}). The short-circuit only honors a + * released record in a null/UNKNOWN state: a record that is force-released + * on a non-UNKNOWN state can only come from a corrupt journal or image, so + * it falls through to the CAS and the state gate instead (fail-closed). + * + * <p>The staged copy sets only the five FORCE audit fields and bumps + * revision and update time; {@code possibleLiveOwned} and + * {@code terminationProof} keep their values for audit. Swapping the record + * in flips {@link LanceIndexJob#isUnresolved()} and + * {@link LanceIndexJob#holdsPossibleLiveSlot()} to false, so + * {@link #applyToMemory(LanceIndexJob)} releases the fence, the quota + * charge, and the corrupt-admission blocker with no explicit teardown. + * + * @return false (with a warning) on an unknown job id, a revision mismatch, + * or a non-UNKNOWN state; true when the job is released or already was + * @throws DdlException when a bounded FORCE text field exceeds its limit + */ + public boolean forceRelease(long jobId, long expectedRevision, String actor, String note, String warning) + throws DdlException { + writeLock(); + try { + LanceIndexJob current = jobs.get(jobId); + if (current == null) { + LOG.warn("reject force release of unknown lance index job {}", jobId); + return false; + } + if (current.isForceReleased() && (current.getMutationState() == null + || current.getMutationState() == LanceIndexJobMutationState.UNKNOWN)) { + return true; + } + if (current.getRevision() != expectedRevision) { + LOG.warn("reject force release of lance index job {}: expected revision {}, current {}", + jobId, expectedRevision, current); + return false; + } + if (current.getMutationState() != null + && current.getMutationState() != LanceIndexJobMutationState.UNKNOWN) { + LOG.warn("reject force release of lance index job {} in mutation state {}: only UNKNOWN may be" + + " force-released", jobId, current.getMutationState()); + return false; + } + long now = System.currentTimeMillis(); + LanceIndexJob updated = new LanceIndexJob(current); + try { + updated.setForceReleased(true); + updated.setForceActor(actor); + updated.setForceTimeMs(now); + updated.setForceNote(note); + updated.setForceWarning(warning); + } catch (IllegalArgumentException e) { + throw new DdlException("invalid lance index job force release: " + e.getMessage(), e); + } + updated.setRevision(current.getRevision() + 1); + updated.setUpdateTimeMs(now); + writeEditLog(updated); + applyToMemory(updated); + return true; + } finally { + writeUnlock(); + } + } + + /** + * Retention GC of resolved job records, master-only: every record that + * stayed resolved for longer than {@code keepMs} (measured on the durable + * update time, bumped by every durable transition including the force + * release) is removed on all FEs through one batch edit-log record per + * round, so every FE serves the same SHOW LANCE INDEX JOBS view. An + * unresolved record is never removed, no matter its age (fail-closed). At + * most {@code maxPerRound} jobs are removed per round, oldest first. + * + * @return the ids actually removed, in oldest-first order + */ + public List<Long> removeResolvedJobsOlderThan(long keepMs, int maxPerRound) { + writeLock(); + try { + long now = System.currentTimeMillis(); + List<LanceIndexJob> expired = new ArrayList<>(); + for (LanceIndexJob job : jobs.values()) { + if (job != null && !job.isUnresolved() && now - job.getUpdateTimeMs() > keepMs) { Review Comment: done: the GC predicate also requires !holdsPossibleLiveSlot(), and the keep window starts at the slot release — the termination proof is the final transition and restarts the aging, so a slot-owning resolved record survives no matter its age (slotOwningResolvedJobsSurviveUntilTheSlotIsProvenGone plus the reworked family test, e2a8747744). ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/job/LanceIndexJobCleaner.java: ########## @@ -0,0 +1,65 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.datasource.lance.job; + +import org.apache.doris.catalog.Env; +import org.apache.doris.common.Config; +import org.apache.doris.common.util.MasterDaemon; + +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.List; + +/** + * Master-only retention sweeper of resolved Lance index job records. Resolved + * records serve audit only; once one has been resolved for longer than + * {@code Config.lance_index_job_keep_max_second} it is removed on all FEs + * through one batch edit-log record per round, so every FE serves the same + * SHOW LANCE INDEX JOBS view. Unresolved records are never removed, no matter + * their age (fail-closed). + * + * <p>MasterDaemon itself performs no master check: the master-only semantics + * come entirely from being started in {@code Env.startMasterOnlyDaemonThreads}. + * Both retention configs are mutable and re-read every round. Review Comment: done: each round drains the expired backlog in repeated bounded batches (full batch repeats, capped at 16 batches = 16,384 records per round, an order of magnitude beyond the sustained resolution throughput) instead of one capped batch per hour (cleanRoundDrainsAFullBacklogInRepeatedBatches, e2a8747744). ########## regression-test/suites/external_table_p0/lance/test_lance_index_resolve.groovy: ########## @@ -0,0 +1,254 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +// PR3E regression scope: the 3D dispatcher is not delivered yet, so a cluster cannot +// produce a genuine UNKNOWN job and every case in this suite is negative or static. +// Deliberately NOT covered here (all wait for 3D / slice 4): FORCE_RELEASE happy +// path e2e, same-name re-admission after FORCE e2e, quota reclaim after FORCE e2e, +// and expiry GC of resolved jobs e2e. + +suite("test_lance_index_resolve", "p0,external,nonConcurrent") { + // The Lance fixture is preinstalled in the MinIO container of the Iceberg + // external environment, so this suite deliberately shares its switch. + String enabled = context.config.otherConfigs.get("enableIcebergTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("disable Lance index resolve test because the Iceberg MinIO environment is disabled.") + return + } + + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + // Admitted jobs are durable and stay PENDING forever in this delivery slice: + // dispatch, the FORCE happy path and job GC only land in later slices, so their + // fences and quota charges can never be released here. Every index name (and the + // filesystem catalog itself, because fence/quota keys include the persisted + // catalog id) carries this per-run suffix so that rerunning the suite on a + // shared pipeline cluster can never collide with a previous run's leftovers. + String runSuffix = "${System.currentTimeMillis()}" + String filesystemCatalog = "test_lance_index_resolve_${runSuffix}" + String user = "test_lance_index_resolve_user" + String password = "C123_567p" + String tableName = "vs_ivf_pq_f32" + String resolveIndexName = "idx_resolve_${runSuffix}" + // Job ids come from the cluster-wide Env.getNextId() allocator, so this literal + // can never name a real job. + String absentJobId = "9223372036854775806" + + // The filesystem catalog is fresh per run by construction, and once it holds + // unresolved jobs DROP CATALOG is guarded, so there is deliberately no DROP for + // it here. + try_sql "DROP USER '${user}'@'%'" + + // All three settings are masterOnly. Read them on the master even if the suite's + // ordinary JDBC connection points at a follower. SHOW uses the experimental + // display name for the gate, while ADMIN SET accepts its unprefixed alias. + def gateRows = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 'experimental_enable_lance_index_mutation'""" + def keepRows = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 'lance_index_job_keep_max_second'""" + def cleanRows = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 'lance_index_job_clean_interval_second'""" + assertEquals(1, gateRows.size()) + assertEquals(1, keepRows.size()) + assertEquals(1, cleanRows.size()) + String originalGate = gateRows[0][1].toString() + String originalKeep = keepRows[0][1].toString() + String originalClean = cleanRows[0][1].toString() + // The retention configs ship with these documented defaults; asserting them here + // fails loudly if a shared cluster has drifted instead of silently restoring a + // non-default value afterwards. + assertEquals("604800", originalKeep) + assertEquals("3600", originalClean) + Throwable suiteFailure = null + + // test { ... exception } always runs on the suite's default connection; the + // masterOnly config rejections below must be asserted on the master connection. + def expectMasterSqlException = { String stmt, String substring -> + String caught = null + try { + master_sql(stmt) + } catch (Throwable t) { + caught = t.toString() + } + assertTrue(caught != null && caught.contains(substring)) + } + + try { + // Rejections reachable with the mutation gate closed. RESOLVE is + // deliberately NOT behind enable_lance_index_mutation: it is the operator + // escape hatch and must stay usable while admission is gated, otherwise an + // unresolved job would freeze catalog DDL forever. Pin the gate closed so + // the cases in this block also prove the statement is not answered with the + // gate rejection. + master_sql """ADMIN SET FRONTEND CONFIG ("enable_lance_index_mutation" = "false")""" + + // Wrong AS literal: only FORCE_RELEASE is accepted after AS. + test { + sql """RESOLVE LANCE INDEX JOB ${absentJobId} AS FORCE COMMENT 'wrong literal'""" + exception "mismatched input" + } + + // The AS FORCE_RELEASE clause is mandatory. + test { + sql """RESOLVE LANCE INDEX JOB ${absentJobId} COMMENT 'missing as clause'""" + exception "mismatched input" + } + + // COMMENT is mandatory at grammar level. + test { + sql """RESOLVE LANCE INDEX JOB ${absentJobId} AS FORCE_RELEASE""" + exception "mismatched input" + } + + // An empty COMMENT parses but never changes the observable response: the + // note check runs after the job lookup and the UNKNOWN state gate, so with a + // missing job the fixed not-found wording answers first. The dedicated + // empty-note rejection needs a genuine UNKNOWN job and stays UT-only + // (ResolveLanceIndexJobCommandTest) until 3D can produce one. + test { + sql """RESOLVE LANCE INDEX JOB ${absentJobId} AS FORCE_RELEASE COMMENT ''""" + exception "Lance index job not found" + } + + // A well-formed RESOLVE of a missing job: the fixed non-disclosing 5103 + // wording, with no gate rejection even though the gate is closed. + test { + sql """RESOLVE LANCE INDEX JOB ${absentJobId} AS FORCE_RELEASE COMMENT 'no such job'""" + exception "Lance index job not found" + } + + // Rejections reachable with an admitted PENDING job. masterOnly configs set + // through ADMIN SET land on the master node locally, which is where + // admission reads them; the finally block below restores the gate no matter + // where the suite fails. + master_sql """ADMIN SET FRONTEND CONFIG ("enable_lance_index_mutation" = "true")""" + + sql """ + CREATE CATALOG `${filesystemCatalog}` PROPERTIES ( + "type" = "lance", + "lance.catalog.type" = "filesystem", + "warehouse" = "s3://warehouse/lance", + "s3.endpoint" = "http://${externalEnvIp}:${minioPort}", + "s3.access_key" = "admin", + "s3.secret_key" = "password", + "s3.region" = "us-east-1", + "use_path_style" = "true" + ) + """ + + // CREATE INDEX is admitted and returns a single-column JobId result set with + // one row; the job stays PENDING because no dispatcher exists in this slice. + def createRows = sql """CREATE INDEX `${resolveIndexName}` ON `${filesystemCatalog}`.`doris`.`${tableName}` (embedding) USING ANN + PROPERTIES("index_type"="IVF_PQ", "metric"="l2", "num_partitions"="256", "num_sub_vectors"="16")""" + assertEquals(1, createRows.size()) + assertEquals(1, createRows[0].size()) + String createJobId = createRows[0][0].toString() + + // RESOLVE authorization is checked against the job's target (table ALTER, or + // global ADMIN for an orphan), and an unauthorized caller gets exactly the + // same non-disclosing response as a missing job. A fresh user holding only + // an unrelated SELECT privilege must therefore see "not found". + sql """CREATE USER '${user}'@'%' IDENTIFIED BY '${password}'""" + sql """GRANT SELECT_PRIV ON regression_test TO '${user}'@'%'""" + if (isCloudMode()) { + def clusters = sql "SHOW CLUSTERS" + assertTrue(!clusters.isEmpty()) + sql """GRANT USAGE_PRIV ON CLUSTER `${clusters[0][0]}` TO '${user}'@'%'""" + } + + connect(user, password, context.config.jdbcUrl) { + test { + sql """RESOLVE LANCE INDEX JOB ${createJobId} AS FORCE_RELEASE COMMENT 'unauthorized caller'""" + exception "Lance index job not found" + } + } + + // FORCE_RELEASE only accepts UNKNOWN jobs: the admitted job is PENDING, the + // one durable state this slice can produce, so an authorized RESOLVE is + // rejected with the not-in-UNKNOWN wording. + test { + sql """RESOLVE LANCE INDEX JOB ${createJobId} AS FORCE_RELEASE COMMENT 'still pending'""" + exception "cannot be resolved: not in UNKNOWN state" + } + + // Retention config smoke: both are masterOnly and guarded by the + // positive-long callback, whose rejection message carries the field name and + // the offending value. + expectMasterSqlException("""ADMIN SET FRONTEND CONFIG ("lance_index_job_keep_max_second" = "0")""", + "must be a positive long") + expectMasterSqlException("""ADMIN SET FRONTEND CONFIG ("lance_index_job_keep_max_second" = "-1")""", + "must be a positive long") + expectMasterSqlException("""ADMIN SET FRONTEND CONFIG ("lance_index_job_clean_interval_second" = "0")""", + "must be a positive long") + expectMasterSqlException("""ADMIN SET FRONTEND CONFIG ("lance_index_job_clean_interval_second" = "-1")""", + "must be a positive long") + + // A positive value is accepted and visible immediately on the master. + master_sql """ADMIN SET FRONTEND CONFIG ("lance_index_job_clean_interval_second" = "3601")""" + def cleanRowsAfterSet = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 'lance_index_job_clean_interval_second'""" + assertEquals("3601", cleanRowsAfterSet[0][1].toString()) + master_sql """ADMIN SET FRONTEND CONFIG ("lance_index_job_clean_interval_second" = "${originalClean}")""" + + // Retention GC never touches an unresolved job: with the keep window pinned + // to one second, any resolved record older than a second is eligible for + // deletion at the next clean round, while the admitted PENDING job must + // survive regardless of its age because the GC predicate fails closed on + // unresolved records. Waiting out a clean interval is impractical here, so + // this is a static existence assertion; expiry GC e2e waits for 3D / slice 4 + // (see the header comment). + try { + master_sql """ADMIN SET FRONTEND CONFIG ("lance_index_job_keep_max_second" = "1")""" + def keepRowsAfterSet = master_sql """ADMIN SHOW FRONTEND CONFIG LIKE 'lance_index_job_keep_max_second'""" + assertEquals("1", keepRowsAfterSet[0][1].toString()) + def jobsRows = sql_return_maparray """SHOW LANCE INDEX JOBS FROM `${filesystemCatalog}`.`doris` + WHERE TableName = "${tableName}" """ + def pendingJobRow = jobsRows.find { it.IndexName == resolveIndexName } + assertTrue(pendingJobRow != null) + assertEquals(createJobId, pendingJobRow.JobId.toString()) + assertEquals("PENDING", pendingJobRow.State.toString()) Review Comment: done: the suite now saves, sets, and restores lance_index_job_dispatcher_paused around the admission like the adjacent suites, and unpauses afterwards to converge the job (see the reply on the quota-leak thread), so the PENDING assertions are barriered instead of racing the 10s daemon round (59d120c268). ########## fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ResolveLanceIndexJobCommand.java: ########## @@ -0,0 +1,367 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.nereids.trees.plans.commands; + +import org.apache.doris.catalog.DatabaseIf; +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.TableIf; +import org.apache.doris.common.AnalysisException; +import org.apache.doris.common.DdlException; +import org.apache.doris.common.ErrorCode; +import org.apache.doris.datasource.CatalogIf; +import org.apache.doris.datasource.CatalogMgr; +import org.apache.doris.datasource.ExternalDatabase; +import org.apache.doris.datasource.ExternalTable; +import org.apache.doris.datasource.lance.LanceExternalCatalog; +import org.apache.doris.datasource.lance.LanceIndexMutationValidator; +import org.apache.doris.datasource.lance.job.LanceIndexJob; +import org.apache.doris.datasource.lance.job.LanceIndexJobManager; +import org.apache.doris.datasource.lance.job.LanceIndexJobMutationState; +import org.apache.doris.mysql.privilege.PrivPredicate; +import org.apache.doris.nereids.trees.plans.PlanType; +import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor; +import org.apache.doris.qe.ConnectContext; +import org.apache.doris.qe.StmtExecutor; + +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.nio.charset.StandardCharsets; + +/** + * RESOLVE LANCE INDEX JOB <jobId> AS FORCE_RELEASE COMMENT '<note>' — the operator + * escape hatch that durably releases a job whose mutation outcome is UNKNOWN (design section + * 7.1). RESOLVE is deliberately not gated by {@code enable_lance_index_mutation}: the gate + * controls mutation admission, while FORCE must stay available exactly when the gate is off. + * + * <p>The release protocol keeps the fence, the quota charge and the possible-live slot while + * it performs one authoritative latest-metadata read and one external-table refresh with the + * current credentials of the surviving catalog, both outside every catalog/manager lock; only + * then does the durable release transfer inside the admission critical section + * ({@code captureLanceIndexTarget} → lock-free read/refresh → {@code withLanceIndexAdmission} + * recheck → manager write lock), serialized against DROP CATALOG and identity ALTER exactly + * like admission. Any failure before the transfer is the typed + * {@code ERR_LANCE_INDEX_JOB_RESOLUTION_INCOMPLETE}: nothing is written, nothing is released, + * and the operator fixes the cause and retries the same statement. + * + * <p>Target resolution (design section 7.1 step 1) is three-valued. RESOLVED means the + * persisted names resolve and the catalog's current durable dataset locator still matches + * the job's — the same revalidation SHOW LANCE INDEX JOBS applies, so a repointed dataset + * reusing the same names never turns a stale name into table-level authorization. MISSING + * means the catalog, database or table is verifiably absent, or the locator positively + * points at a different dataset: that is the orphan family — a full orphan (catalog gone) + * has no credentials to read with and nothing to invalidate, so it is released directly + * after global ADMIN authorization, while a half-orphan skips the authoritative read and + * refreshes with {@code ignoreIfNotExists=true} as a best-effort invalidation. FAILED means + * a resolution that errors out, or a locator that cannot be resolved right now: never an + * orphan verdict — after ADMIN authorization the statement fails with the typed 5105 so the + * fence is kept when "table gone" cannot be told apart from "network down". SHOW fails the + * same uncertainty closed by hiding the row; RESOLVE fails it closed by not releasing. + * + * <p>Non-disclosure (design section 8): the job is loaded first and authorized against its + * persisted target — table-level ALTER when the target resolves, global ADMIN otherwise — and + * a missing job and an unauthorized job share the same fixed ERR_LANCE_INDEX_JOB_NOT_FOUND + * response naming only the job id. The 5104 state rejection and the 5105 resolution failure + * are only visible to an already authorized caller. + * + * <p>Success returns an OK packet carrying one warning row with {@link #LATE_COMMIT_WARNING}, + * the same text persisted as the job's durable {@code forceWarning}: the old worker may still + * overwrite, remove, or reintroduce the index name; the mutation outcome remains UNKNOWN. + * Retrying FORCE on an already released job is an idempotent success returning the existing + * release record, never an error. + */ +public class ResolveLanceIndexJobCommand extends Command implements ForwardWithSync { + /** + * The late-commit warning (design section 7.1), returned in the OK packet and persisted + * verbatim as the durable {@code forceWarning}; bounded well under + * {@link LanceIndexJob#MAX_FORCE_TEXT_BYTES}. + */ + static final String LATE_COMMIT_WARNING = + "the old worker may still overwrite, remove, or reintroduce the index name; " + + "the mutation outcome remains UNKNOWN"; + + private static final Logger LOG = LogManager.getLogger(ResolveLanceIndexJobCommand.class); + + private final long jobId; + private final String comment; + + public ResolveLanceIndexJobCommand(long jobId, String comment) { + super(PlanType.RESOLVE_LANCE_INDEX_JOB_COMMAND); + this.jobId = jobId; + this.comment = comment; + } + + public long getJobId() { + return jobId; + } + + public String getComment() { + return comment; + } + + @Override + public void run(ConnectContext ctx, StmtExecutor executor) throws Exception { + Env env = Env.getCurrentEnv(); + LanceIndexJobManager manager = env.getLanceIndexJobManager(); + // 1. Load the job without disclosing any field (design section 7.1 step 1). + LanceIndexJob job = manager.getJob(jobId); + if (job == null) { + throw notFound(); + } + // 2. Resolve and authorize against the persisted target before any state is revealed: + // table-level ALTER when the target resolves, global ADMIN for the orphan family + // and for a target whose resolution failed outright. + CatalogMgr catalogMgr = env.getCatalogMgr(); + CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog = catalogMgr.getCatalog(job.getCatalogId()); + TargetResolution resolution = resolveTarget(catalog, job); + boolean authorized = resolution == TargetResolution.RESOLVED + ? env.getAccessManager().checkTblPriv(ctx, catalog.getName(), job.getDbName(), job.getTableName(), + PrivPredicate.ALTER) + : env.getAccessManager().checkGlobalPriv(ctx, PrivPredicate.ADMIN); + if (!authorized) { + throw notFound(); + } + // 3. Idempotent replay: a retry returns the existing release record (section 7.1). + // This deliberately precedes the resolution-failure rejection: once the release + // has landed, a retry during a provider outage is a success, not a 5105. + if (job.isForceReleased()) { + ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING); + return; + } + // 4. Only UNKNOWN may be force-released; a null mutation state reads as UNKNOWN, + // same as the manager's own gate. The state rejection also precedes the + // resolution-failure rejection: for a terminal job the accurate answer is 5104, + // not a 5105 claiming the job still holds its fence. + if (job.getMutationState() != null && job.getMutationState() != LanceIndexJobMutationState.UNKNOWN) { + throw new AnalysisException(ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN.formatErrorMsg(jobId), + ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN); + } + if (resolution == TargetResolution.FAILED) { + // Never an orphan verdict: "table gone" cannot be told apart from "network down", + // so nothing is released and nothing beyond the typed error is disclosed; the + // operator fixes the cause and retries the same statement (design 7.1 step 4). + throw incompleteResolution("the persisted target could not be resolved with current catalog" + + " metadata; see fe.log for the cause"); + } + boolean targetResolves = resolution == TargetResolution.RESOLVED; + // 5. The grammar makes COMMENT mandatory; here the note must also be non-empty after + // trimming and fit the durable force text bound. + String note = comment == null ? "" : comment.trim(); + if (note.isEmpty()) { + throw new AnalysisException("force release note must not be empty", + ErrorCode.ERR_LANCE_INDEX_INVALID); + } + if (note.getBytes(StandardCharsets.UTF_8).length > LanceIndexJob.MAX_FORCE_TEXT_BYTES) { + throw new AnalysisException("force release note exceeds " + LanceIndexJob.MAX_FORCE_TEXT_BYTES + + " UTF-8 bytes", ErrorCode.ERR_LANCE_INDEX_INVALID); + } + // 6-9. Branch on the orphan state, then the durable release transfer. + String actor = ctx.getQualifiedUser(); + boolean released; + if (catalog == null) { + // Full orphan: no credentials survive to read with and nothing can be + // invalidated, so the release goes straight to the manager write lock. + released = manager.forceRelease(jobId, job.getRevision(), actor, note, LATE_COMMIT_WARNING); + } else { + released = releaseWithLiveCatalog(env, catalogMgr, manager, catalog, targetResolves, job, actor, note); + } + if (!released) { + // 10. The expected-revision transfer lost a race. A concurrent FORCE_RELEASE + // that already landed makes this an idempotent success; anything else means + // the job left UNKNOWN concurrently (UNKNOWN has no other outgoing + // transition), so the pinned not-UNKNOWN wording stays accurate. + LanceIndexJob reread = manager.getJob(jobId); + if (reread != null && reread.isForceReleased()) { + ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING); + return; + } + throw new AnalysisException(ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN.formatErrorMsg(jobId), + ErrorCode.ERR_LANCE_INDEX_JOB_NOT_UNKNOWN); + } + // 11. The OK packet carries the late-commit warning; it survives the forward chain + // byte-identically (proxyExecute serializes the master state). + ctx.getState().setOk(0, 1, LATE_COMMIT_WARNING); + } + + /** + * The live-catalog release path: capture the target identity, then — holding no lock — + * perform the authoritative read and the refresh, and finally transfer the release inside + * the admission critical section. Lock order stays CatalogMgr then LanceIndexJobManager. + */ + private boolean releaseWithLiveCatalog(Env env, CatalogMgr catalogMgr, LanceIndexJobManager manager, + CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog, boolean targetResolves, LanceIndexJob job, + String actor, String note) throws Exception { + if (!(catalog instanceof LanceExternalCatalog) + || ((LanceExternalCatalog) catalog).isRestCatalogConfigured()) { + // Defensive: admission never targets a REST catalog and a non-Lance catalog + // cannot hold a Lance fence, so no UNKNOWN job should ever resolve here. + LanceIndexMutationValidator.rejectUnsupportedOperation( + "RESOLVE LANCE INDEX JOB AS FORCE_RELEASE", "catalog '" + catalog.getName() + "'"); + } + LanceExternalCatalog lanceCatalog = (LanceExternalCatalog) catalog; + CatalogMgr.LanceIndexTarget target; + try { + target = catalogMgr.captureLanceIndexTarget(lanceCatalog); + } catch (DdlException e) { + throw incompleteResolution(e.getMessage()); + } + if (targetResolves) { + // 7. One authoritative latest-metadata read with current credentials, proving + // the dataset is reachable before anything is released. + authoritativeRead(lanceCatalog, catalog, job); + } + // 8. Invalidate the external table and broadcast the refresh to every FE. A + // half-orphan target is invalidated best-effort (missing db/table is a no-op). + try { + env.getRefreshManager().handleRefreshTable(catalog.getName(), job.getDbName(), job.getTableName(), + !targetResolves); + } catch (DdlException e) { + throw incompleteResolution(e.getMessage()); + } + // 9. The durable transfer rechecks the catalog identity under the read lock, then + // runs the revision-checked release under the manager write lock. + try { + return catalogMgr.withLanceIndexAdmission(lanceCatalog, target, + () -> manager.forceRelease(jobId, job.getRevision(), actor, note, LATE_COMMIT_WARNING)); + } catch (DdlException e) { + throw incompleteResolution(e.getMessage()); + } + } + + /** + * The authoritative latest-metadata read over the remote names of the resolved db/table, + * outside every lock (the loader owns its deadline-bound JNI read). Every loader failure + * already passed through the catalog's sanitized root-cause chain, so its message is safe + * to echo — locator, credentials and dataset uri are masked there. + */ + private void authoritativeRead(LanceExternalCatalog lanceCatalog, + CatalogIf<? extends DatabaseIf<? extends TableIf>> catalog, LanceIndexJob job) throws AnalysisException { + DatabaseIf<? extends TableIf> db; + TableIf table; + try { + db = catalog.getDbNullable(job.getDbName()); + table = db == null ? null : db.getTableNullable(job.getTableName()); + } catch (Exception e) { + // A remote-metadata blip during resolution is not an orphan verdict; the provider + // message is not echoed here because it never crossed the sanitized chain. + LOG.warn("lance index job {}: target re-resolution failed before the authoritative read", + job.getJobId(), e); + throw incompleteResolution("the persisted target could not be resolved with current catalog" + + " metadata; see fe.log for the cause"); + } + if (db == null || table == null) { + // The target resolved at authorization time but is gone now: keep the fence and + // let the retry take the half-orphan branch under global ADMIN. + throw incompleteResolution("the persisted target table no longer resolves; retry the statement"); + } + if (!(db instanceof ExternalDatabase) || !(table instanceof ExternalTable)) { + // Boundary guard: a Lance catalog must serve external relations, but the command + // boundary does not trust that invariant (no raw ClassCastException to the user). + LOG.warn("lance index job {}: target relation is not external: db={}, table={}", + job.getJobId(), db.getClass().getName(), table.getClass().getName()); + throw incompleteResolution("the persisted target is not an external relation;" + + " see fe.log for the cause"); + } + String remoteDb = ((ExternalDatabase) db).getRemoteName(); + String remoteTable = ((ExternalTable) table).getRemoteName(); + try { + lanceCatalog.loadTableIndexAdmissionSnapshot(remoteDb, remoteTable); Review Comment: done: the authoritative read now normalizes the returned snapshot dataset URI and compares it with the persisted locator; a re-registered dataset under the same names answers a retryable resolution-incomplete and the retry takes the orphan path where the verdict comes from the namespace (testSnapshotDatasetMismatchKeepsTheFence, 59d120c268). -- 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]
