LanceCatalogClient.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;

import org.apache.doris.analysis.TableSnapshot;
import org.apache.doris.common.Config;
import org.apache.doris.common.DdlException;
import org.apache.doris.common.util.TimeUtils;
import org.apache.doris.datasource.lance.index.LanceIndexInspection;
import org.apache.doris.datasource.lance.index.LanceIndexInspectionExecutor;
import org.apache.doris.datasource.lance.index.LancePhysicalIndexEntry;
import org.apache.doris.datasource.lance.index.LanceShowIndexInfo;
import org.apache.doris.datasource.lance.job.LanceIndexDatasetLocator;
import org.apache.doris.datasource.lance.metadata.LanceMetadataLoader;
import org.apache.doris.datasource.lance.metadata.LanceReadOptions;
import org.apache.doris.datasource.lance.metadata.LanceRefSelector;
import org.apache.doris.datasource.lance.metadata.LanceSnapshotResolver;
import org.apache.doris.datasource.lance.metadata.LanceTableAccess;
import org.apache.doris.datasource.lance.metadata.LanceTableMetadata;
import org.apache.doris.datasource.lance.profile.LanceMetadataMetrics;
import org.apache.doris.datasource.lance.profile.LanceMetadataMetrics.Stage;
import org.apache.doris.datasource.property.metastore.AbstractLanceProperties;
import org.apache.doris.datasource.property.storage.StorageProperties;

import com.github.benmanes.caffeine.cache.Ticker;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.vector.types.pojo.Schema;
import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.exception.ExceptionUtils;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.lance.Dataset;
import org.lance.Ref;
import org.lance.Session;
import org.lance.namespace.LanceNamespace;
import org.lance.namespace.errors.TableBranchNotFoundException;
import org.lance.namespace.errors.TableNotFoundException;
import org.lance.namespace.errors.TableVersionNotFoundException;
import org.lance.namespace.model.TableVersion;

import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.NavigableSet;
import java.util.Optional;
import java.util.OptionalLong;
import java.util.TreeSet;
import java.util.function.BiFunction;

/**
 * One catalog generation: Namespace access, snapshot reads and native resource lifetime.
 * Callers must hold a Lease, acquired while selecting the current generation in the catalog.
 */
final class LanceCatalogClient implements AutoCloseable {

    private static final Logger LOG = LogManager.getLogger(LanceCatalogClient.class);
    /** Lance's name for the main chain; {@code @branch(main)} and tags on it mean the main chain. */
    static final String MAIN_BRANCH = "main";
    private static final long METADATA_CACHE_SIZE_BYTES = 64L * 1024 * 1024;
    private static final long INDEX_CACHE_SIZE_BYTES = 128L * 1024 * 1024;

    private final LanceNamespace namespace;
    private final Session session;
    private final LanceNamespaceClient namespaceClient;
    private final Map<String, String> namespaceStorageOptions;
    private final BufferAllocator namespaceAllocator;
    private final List<String> catalogSecrets;
    private int activeOperations;
    private boolean retired;

    static LanceCatalogClient create(AbstractLanceProperties properties,
            List<StorageProperties> storageProperties, Map<String, String> namespaceOptions,
            List<String> catalogSecrets) throws DdlException {
        long limit = Config.lance_catalog_arrow_memory_limit_bytes;
        if (limit <= 0) {
            throw new IllegalArgumentException("lance_catalog_arrow_memory_limit_bytes must be positive");
        }
        BufferAllocator allocator = new RootAllocator(limit);
        LanceNamespace namespace = null;
        Session session = null;
        try {
            List<String> parent = LanceNamespaceName.parseParentNamespace(
                    properties.getNamespaceParent(), properties.getNamespaceDelimiter());
            namespace = properties.createNamespace(allocator, namespaceOptions);
            session = Session.builder().metadataCacheSizeBytes(METADATA_CACHE_SIZE_BYTES)
                    .indexCacheSizeBytes(INDEX_CACHE_SIZE_BYTES).build();
            return new LanceCatalogClient(namespace, allocator, session, properties.getLanceCatalogType(),
                    properties.getRootDatabase(), parent, storageProperties, namespaceOptions, catalogSecrets,
                    properties.getTableAccessCacheTtlSeconds());
        } catch (RuntimeException | Error e) {
            closeResource(namespace);
            closeResource(session);
            closeResource(allocator);
            throw e;
        }
    }

    LanceCatalogClient(LanceNamespace namespace, BufferAllocator allocator, Session session,
            String catalogType, String rootDatabase, List<String> parentNamespace,
            List<StorageProperties> storageProperties, Map<String, String> namespaceStorageOptions,
            List<String> catalogSecrets) {
        this(namespace, allocator, session, catalogType, rootDatabase, parentNamespace,
                storageProperties, namespaceStorageOptions, catalogSecrets,
                AbstractLanceProperties.DEFAULT_TABLE_ACCESS_CACHE_TTL_SECONDS);
    }

    LanceCatalogClient(LanceNamespace namespace, BufferAllocator allocator, Session session,
            String catalogType, String rootDatabase, List<String> parentNamespace,
            List<StorageProperties> storageProperties, Map<String, String> namespaceStorageOptions,
            List<String> catalogSecrets, int tableAccessCacheTtlSeconds) {
        this.catalogSecrets = Collections.unmodifiableList(new ArrayList<>(catalogSecrets));
        this.namespace = namespace;
        this.namespaceAllocator = allocator;
        this.session = session;
        this.namespaceClient = new LanceNamespaceClient(
                namespace, catalogType, rootDatabase, parentNamespace, storageProperties,
                tableAccessCacheTtlSeconds, Ticker.systemTicker());
        this.namespaceStorageOptions = Collections.unmodifiableMap(new HashMap<>(namespaceStorageOptions));
    }

    /** Pins this generation for one operation; the lock does not cover its SDK or JNI calls. */
    synchronized Lease acquire() {
        if (retired) {
            throw new IllegalStateException("Lance catalog resources have been closed");
        }
        activeOperations++;
        return new Lease(this);
    }

    /** Retires this generation immediately; its last active operation performs resource cleanup. */
    @Override
    public void close() {
        boolean release;
        synchronized (this) {
            if (retired) {
                return;
            }
            retired = true;
            release = activeOperations == 0;
        }
        if (release) {
            closeResources();
        }
    }

    private void release() {
        boolean release;
        synchronized (this) {
            activeOperations--;
            release = retired && activeOperations == 0;
        }
        if (release) {
            closeResources();
        }
    }

    private void closeResources() {
        closeResource(namespace);
        closeResource(session);
        closeResource(namespaceAllocator);
    }

    private static void closeResource(Object resource) {
        if (resource instanceof AutoCloseable) {
            try {
                ((AutoCloseable) resource).close();
            } catch (Exception e) {
                // Provider exception messages may contain credentials.
                LOG.warn("Failed to close a Lance catalog resource ({})", resource.getClass().getSimpleName());
            }
        }
    }

    static final class Lease implements AutoCloseable {
        private final LanceCatalogClient client;
        private boolean closed;

        private Lease(LanceCatalogClient client) {
            this.client = client;
        }

        LanceCatalogClient client() {
            return client;
        }

        @Override
        public void close() {
            synchronized (this) {
                if (closed) {
                    return;
                }
                closed = true;
            }
            client.release();
        }
    }

    void invalidateTableAccessCache() {
        namespaceClient.invalidateTableAccessCache();
    }

    List<String> listDatabaseNames() {
        return namespaceClient.listDatabaseNames();
    }

    List<String> listTableNames(String dbName) {
        return namespaceClient.listTableNames(dbName);
    }

    boolean tableExists(String dbName, String tableName) {
        return namespaceClient.tableExists(dbName, tableName);
    }

    public LanceTableMetadata loadTableMetadata(String dbName, String tableName) {
        return loadTableMetadata(dbName, tableName, Optional.empty());
    }

    public LanceTableMetadata loadTableMetadataForSearch(String dbName, String tableName) {
        LanceTableMetadata metadata = loadQueryMetadata(dbName, tableName, Optional.empty(),
                LanceMetadataLoader.MetadataScope.WITH_INDEXES);
        if (!metadata.getIndexMetadataState().canPlanIndexSegments()) {
            throw new IllegalArgumentException("Lance SDK cannot provide field IDs required for search planning");
        }
        return metadata;
    }

    public LanceTableMetadata loadBasicTableMetadata(String dbName, String tableName) {
        return loadQueryMetadata(dbName, tableName, Optional.empty(), LanceMetadataLoader.MetadataScope.BASIC);
    }

    public Schema loadTableSchema(String dbName, String tableName) {
        return readTableSnapshot(dbName, tableName, LanceRefSelector.latest(),
                (dataset, access, metrics) -> metrics.measure(Stage.SCHEMA, dataset::getSchema));
    }

    public LanceTableMetadata loadTableMetadata(String dbName, String tableName,
            Optional<TableSnapshot> tableSnapshot) {
        return loadTableMetadata(dbName, tableName, LanceRefSelector.snapshot(tableSnapshot));
    }

    public LanceTableMetadata loadTableMetadata(String dbName, String tableName, LanceRefSelector selector) {
        return loadQueryMetadata(dbName, tableName, selector, LanceMetadataLoader.MetadataScope.WITH_INDEXES);
    }

    private LanceTableMetadata loadQueryMetadata(String dbName, String tableName,
            Optional<TableSnapshot> tableSnapshot, LanceMetadataLoader.MetadataScope mode) {
        return loadQueryMetadata(dbName, tableName, LanceRefSelector.snapshot(tableSnapshot), mode);
    }

    private LanceTableMetadata loadQueryMetadata(String dbName, String tableName,
            LanceRefSelector selector, LanceMetadataLoader.MetadataScope mode) {
        return readTableSnapshot(dbName, tableName, selector,
                (dataset, access, metrics) -> LanceMetadataLoader.read(dataset, access, mode, metrics));
    }

    /**
     * Pins one resource generation, resolved table access, and the Dataset version for the whole read.
     *
     * <p>Every dataset is read by its URI and a version, as the BE reads it. The latest version of
     * the main chain in storage is opened once as a handle, and the other selectors are checkouts
     * from it, except an explicit version on the main chain and a managed table's branch (below).
     * A tag is resolved first to the chain and version it points at, so a tag created on a branch
     * selects that branch.
     *
     * <p>For a managed table the namespace decides which versions exist. "Latest" is the newest
     * version it records, never the newest manifest in storage, and every version a read selects
     * must be one it records, at the manifest path Doris reads ({@link LanceManifestPaths}). A
     * branch is read from its own directory, as the BE reads it, so it does not depend on the main
     * chain; the main handle only supplies tag files and the manifest listing that FOR TIME AS OF
     * on main takes commit times from.
     */
    private <T> T readTableSnapshot(String dbName, String tableName, LanceRefSelector selector,
            SnapshotReader<T> reader) {
        ReadState state = new ReadState(selector, dbName + "." + tableName);
        LanceMetadataMetrics metrics = LanceMetadataMetrics.startMetadataRead();
        try {
            T result;
            try (BufferAllocator allocator = namespaceAllocator.newChildAllocator(
                    "lance-metadata-read", 0, namespaceAllocator.getLimit())) {
                state.access = metrics.measure(Stage.TABLE_ACCESS,
                        () -> namespaceClient.resolveTableAccess(dbName, tableName));
                OptionalLong direct = directMainVersion(state, metrics);
                if (state.access.isManagedVersioning() && state.branch.isPresent() && !selector.getTag().isPresent()) {
                    result = readManagedBranch(allocator, state, reader, metrics);
                } else if (direct.isPresent() || isLatestMain(selector)) {
                    state.version = direct;
                    try (Dataset dataset = openDataset(allocator, state.access, direct, metrics)) {
                        result = reader.read(dataset, state.access, metrics);
                    }
                } else {
                    try (Dataset main = openDataset(allocator, state.access, OptionalLong.empty(), metrics)) {
                        result = readFromLatest(main, state, reader, metrics);
                    }
                }
            }
            metrics.succeeded();
            return result;
        } catch (LanceUserFacingException e) {
            throw new RuntimeException(e.getMessage(), e);
        } catch (Exception e) {
            LanceTableAccess access = state.access;
            String uri = access == null ? null : access.getDatasetUri();
            Map<String, String> options = access == null ? namespaceStorageOptions : access.getStorageOptions();
            String what = state.displayName();
            // Lance's Directory namespace reports a branch it lacks as a missing table. The table
            // was described in this read, so a missing table from a branch's version request
            // means the branch.
            boolean namespaceLacksBranch = access != null && access.isManagedVersioning()
                    && ExceptionUtils.indexOfType(e, TableNotFoundException.class) >= 0;
            if (state.branch.isPresent() && !state.branchExists
                    && (namespaceLacksBranch || isBranchNotFound(e, state.branch.get()))) {
                throw new RuntimeException("Lance branch '" + state.branch.get() + "' of " + state.tableName
                        + state.selector.getTag().map(tag -> " (tag '" + tag + "')").orElse("")
                        + " was not found" + (namespaceLacksBranch || isNamespaceMiss(e) ? " in the namespace" : ""),
                        sanitizedCause(e, uri, options));
            }
            if (isVersionNotFound(e) && state.pinned != null
                    && state.pinned.manifest == LanceManifestPaths.Recorded.STAGED) {
                throw new RuntimeException(unreadableStaged(state.pinned, state), sanitizedCause(e, uri, options));
            }
            if (state.version.isPresent() && isVersionNotFound(e)) {
                throw new RuntimeException("Lance version " + state.version.getAsLong() + " of " + what
                        + state.selector.getTag().map(tag -> " (tag '" + tag + "')").orElse("")
                        + " was not found" + (isNamespaceMiss(e) ? " in the namespace" : ""),
                        sanitizedCause(e, uri, options));
            }
            throw LanceErrorMessages.failure("Failed to load Lance table metadata for " + what, e, uri, options,
                    catalogSecrets);
        } finally {
            metrics.close();
        }
    }

    /** What a read has resolved so far; the catch block reports errors against it. */
    private static final class ReadState {
        private final LanceRefSelector selector;
        private final String tableName;
        /** The table's access; a branch's access is derived from it with {@code onBranch}. */
        private LanceTableAccess access;
        private Optional<String> branch;
        /**
         * Set once the branch is known to exist: the namespace recorded versions for it, or its
         * latest version was checked out. Later failures are not reported as a missing branch.
         */
        private boolean branchExists;
        private OptionalLong version = OptionalLong.empty();
        /** The managed version this read opens next, as the namespace records it. */
        private Recorded pinned;
        /** The namespace's version list of the chain a FOR TIME AS OF reads ("" is main), fetched once. */
        private final Map<String, List<TableVersion>> namespaceVersions = new HashMap<>();

        private ReadState(LanceRefSelector selector, String tableName) {
            this.selector = selector;
            this.tableName = tableName;
            this.branch = selector.getBranch();
        }

        private String displayName() {
            return tableName + branch.map(name -> "@" + name).orElse("");
        }
    }

    /** A version of a managed chain the namespace records, and how it records its manifest. */
    private static final class Recorded {
        private final Optional<String> branch;
        private final long version;
        private final LanceManifestPaths.Recorded manifest;

        private Recorded(Optional<String> branch, long version, LanceManifestPaths.Recorded manifest) {
            this.branch = branch;
            this.version = version;
            this.manifest = manifest;
        }
    }

    private static boolean isLatestMain(LanceRefSelector selector) {
        return !selector.getTag().isPresent() && !selector.getBranch().isPresent()
                && !selector.getSnapshot().isPresent();
    }

    /**
     * The main-chain version a selector names without looking at the latest manifest: an explicit
     * version, or the latest version of a managed table, which the namespace records.
     */
    private OptionalLong directMainVersion(ReadState state, LanceMetadataMetrics metrics) {
        LanceRefSelector selector = state.selector;
        if (selector.getTag().isPresent() || selector.getBranch().isPresent()) {
            return OptionalLong.empty();
        }
        if (!selector.getSnapshot().isPresent()) {
            return state.access.isManagedVersioning()
                    ? OptionalLong.of(recordedHead(state, Optional.empty(), metrics))
                    : OptionalLong.empty();
        }
        TableSnapshot snapshot = selector.getSnapshot().get();
        if (snapshot.getType() != TableSnapshot.VersionType.VERSION) {
            return OptionalLong.empty();
        }
        state.version = OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
        requireRecorded(state, Optional.empty(), state.version.getAsLong(), metrics);
        return state.version;
    }

    /** Resolves the selector against the open latest main chain and reads the selected snapshot. */
    private <T> T readFromLatest(Dataset main, ReadState state, SnapshotReader<T> reader, LanceMetadataMetrics metrics)
            throws Exception {
        LanceRefSelector selector = state.selector;
        if (selector.getTag().isPresent()) {
            // Only this tag's file is read, however many tags the table has. The SDK checks the tag
            // out on the branch of the version it points at.
            String tag = selector.getTag().get();
            state.version = OptionalLong.of(metrics.measure(Stage.VERSION_RESOLVE, () -> tagVersion(main, tag, state)));
            try (Dataset target = checkout(main, Ref.ofTag(tag), metrics)) {
                // The checkout reads the tag file again; the version it read is the one to check.
                state.version = OptionalLong.of(target.version());
                state.branch = branchOf(target.uri(), state.access.getDatasetUri());
                requireRecorded(state, state.branch, target.version(), metrics);
                return reader.read(target, accessOf(target, state), metrics);
            }
        }
        if (state.branch.isPresent()) {
            String branch = state.branch.get();
            // Check out the branch's latest version first even when a version is already known, so
            // a missing branch and a missing version inside an existing branch are told apart.
            try (Dataset latest = checkout(main, Ref.ofBranch(branch), metrics)) {
                state.branchExists = true;
                LanceTableAccess branchAccess = accessOf(latest, state);
                if (!state.version.isPresent() && selector.getSnapshot().isPresent()) {
                    state.version = resolveSnapshotVersion(latest, branchAccess, selector.getSnapshot().get(), state,
                            metrics);
                }
                if (!state.version.isPresent()) {
                    return reader.read(latest, branchAccess, metrics);
                }
                try (Dataset dataset = checkout(latest, Ref.ofBranch(branch, state.version.getAsLong()), metrics)) {
                    return reader.read(dataset, branchAccess, metrics);
                }
            }
        }
        // FOR TIME AS OF on the main chain, resolved from manifest commit times.
        state.version = resolveSnapshotVersion(main, state.access, selector.getSnapshot().get(), state, metrics);
        try (Dataset dataset = checkout(main, Ref.ofMain(state.version.getAsLong()), metrics)) {
            return reader.read(dataset, state.access, metrics);
        }
    }

    /**
     * Reads a branch of a managed table from the branch's directory, by URI and version as the BE
     * reads it. The namespace answers for the branch first, so a branch it lacks is reported as
     * missing, and the main chain need not be readable. FOR TIME AS OF lists the branch's
     * manifests from its newest version in storage and, as on main, selects among the versions the
     * namespace records.
     */
    private <T> T readManagedBranch(BufferAllocator allocator, ReadState state, SnapshotReader<T> reader,
            LanceMetadataMetrics metrics) throws Exception {
        String branch = state.branch.get();
        LanceTableAccess branchAccess = state.access.onBranch(branch,
                branchUri(state.access.getDatasetUri(), branch));
        Optional<TableSnapshot> snapshot = state.selector.getSnapshot();
        if (!snapshot.isPresent()) {
            state.version = OptionalLong.of(recordedHead(state, state.branch, metrics));
        } else if (snapshot.get().getType() == TableSnapshot.VersionType.VERSION) {
            state.version = OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.get().getValue()));
            requireRecorded(state, state.branch, state.version.getAsLong(), metrics);
        } else {
            namespaceVersions(state, state.branch, metrics);
            state.branchExists = true;
            try (Dataset latest = openDataset(allocator, branchAccess, OptionalLong.empty(), metrics)) {
                state.version = resolveSnapshotVersion(latest, branchAccess, snapshot.get(), state, metrics);
            }
        }
        state.branchExists = true;
        try (Dataset dataset = openDataset(allocator, branchAccess, state.version, metrics)) {
            return reader.read(dataset, branchAccess, metrics);
        }
    }

    /**
     * The URI of a branch's directory, joined as Lance's {@code BranchLocation} joins it:
     * {@code tree/<branch>} under the table root, before the URI's query string.
     */
    static String branchUri(String tableUri, String branch) {
        int query = tableUri.indexOf('?');
        String path = query < 0 ? tableUri : tableUri.substring(0, query);
        String joined = path + (path.endsWith("/") ? "" : "/") + "tree/" + StringUtils.stripStart(branch, "/");
        return query < 0 ? joined : joined + tableUri.substring(query);
    }

    private static long tagVersion(Dataset main, String tag, ReadState state) {
        try {
            return main.tags().getVersion(tag);
        } catch (RuntimeException e) {
            String rootMessage = ExceptionUtils.getRootCauseMessage(e);
            if (rootMessage != null && rootMessage.contains("tag " + tag + " does not exist")) {
                throw new LanceUserFacingException("Lance tag '" + tag + "' of " + state.tableName + " was not found");
            }
            throw e;
        }
    }

    /**
     * A dataset URI without its query and trailing slash. The query may carry credentials, which a
     * namespace can vend anew on every describe.
     */
    private static String location(String uri) {
        return StringUtils.removeEnd(StringUtils.substringBefore(uri, "?"), "/");
    }

    /**
     * The branch a dataset checked out from the table root is on, from its root directory: the
     * table root for main, {@code <root>/tree/<branch>} otherwise. Lance inserts the branch path
     * before a URI's query string, so the query is compared apart. A URI that is neither is an
     * error rather than main, which would hand the BE the wrong chain.
     */
    static Optional<String> branchOf(String checkedOutUri, String tableUri) {
        String root = location(tableUri);
        String uri = location(checkedOutUri);
        if (uri.equals(root)) {
            return Optional.empty();
        }
        String branchRoot = root + "/tree/";
        if (!uri.startsWith(branchRoot) || uri.length() == branchRoot.length()) {
            // The URIs may carry credentials in their query, so they stay out of the message.
            throw new IllegalStateException("Cannot tell which branch a Lance tag was checked out on");
        }
        return Optional.of(uri.substring(branchRoot.length()));
    }

    /**
     * The access for a dataset checked out from the table: the main chain keeps the table access,
     * and a branch takes the directory the SDK checked out, which is what the BE opens by URI.
     */
    private static LanceTableAccess accessOf(Dataset dataset, ReadState state) {
        return state.branch.isPresent() ? state.access.onBranch(state.branch.get(), dataset.uri()) : state.access;
    }

    /** A selector error whose message is user-facing as is, such as a tag that does not exist. */
    private static final class LanceUserFacingException extends RuntimeException {
        private LanceUserFacingException(String message) {
            super(message);
        }
    }

    /**
     * The newest version the namespace records for a managed chain, which the read then opens.
     * Doris asks for it itself: opening "latest" by URI would read the newest manifest in storage,
     * which the namespace may not have published.
     */
    private long recordedHead(ReadState state, Optional<String> branch, LanceMetadataMetrics metrics) {
        state.pinned = null;
        Optional<TableVersion> head = metrics.measure(Stage.VERSION_RESOLVE,
                () -> namespaceClient.latestManagedVersion(state.access, branch));
        if (!head.isPresent()) {
            throw new LanceUserFacingException("Lance namespace lists no versions for " + state.tableName
                    + branch.map(name -> "@" + name).orElse(""));
        }
        long version = head.get().getVersion();
        state.pinned = new Recorded(branch, version, LanceManifestPaths.check(state.access.getDatasetUri(), branch,
                version, head.get().getManifestPath(), state.tableName));
        return version;
    }

    /**
     * Requires the namespace of a managed table to record {@code version} of the chain on
     * {@code branch}, at the manifest path Doris reads; nothing for a storage-versioned table.
     */
    private void requireRecorded(ReadState state, Optional<String> branch, long version,
            LanceMetadataMetrics metrics) {
        if (!state.access.isManagedVersioning()) {
            return;
        }
        state.pinned = null;
        TableVersion recorded = metrics.measure(Stage.VERSION_RESOLVE,
                () -> namespaceClient.describeManagedVersion(state.access, branch, version));
        state.pinned = new Recorded(branch, version, LanceManifestPaths.check(state.access.getDatasetUri(), branch,
                version, recorded.getManifestPath(), state.tableName));
    }

    /**
     * The error for a version the namespace records at a staged manifest while its canonical
     * manifest, which Doris reads, does not exist. Either the commit reserved the version and was
     * not finalized, which a reader that uses the namespace would finish and Doris does not, or
     * the version was finalized and cleanup later removed it; the namespace's record does not tell
     * the two apart.
     */
    private static String unreadableStaged(Recorded pinned, ReadState state) {
        return "Lance version " + pinned.version + " of " + state.tableName
                + pinned.branch.map(name -> "@" + name).orElse("") + " cannot be read: " + stagedOnly();
    }

    private static String stagedOnly() {
        return "the namespace records it at a staged manifest, and its canonical manifest, which Doris reads,"
                + " does not exist (its commit was not finalized, or cleanup removed it)";
    }

    private RuntimeException sanitizedCause(Throwable error, String uri, Map<String, String> options) {
        return new RuntimeException(LanceErrorMessages.sanitize(error, uri, options, catalogSecrets));
    }

    /** Checks out a ref of an already open dataset; the SDK resolves the ref from the dataset directory. */
    private static Dataset checkout(Dataset dataset, Ref ref, LanceMetadataMetrics metrics) {
        return metrics.measure(Stage.VERSION_RESOLVE, () -> dataset.checkout(ref));
    }

    /**
     * Resolves a {@code FOR VERSION AS OF} / {@code FOR TIME AS OF} snapshot against the chain
     * {@code latest} is checked out on: the main chain, or a branch when {@code access} is a
     * branch access.
     */
    private OptionalLong resolveSnapshotVersion(Dataset latest, LanceTableAccess access, TableSnapshot snapshot,
            ReadState state, LanceMetadataMetrics metrics) {
        if (snapshot.getType() == TableSnapshot.VersionType.VERSION) {
            state.version = OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
            requireRecorded(state, access.getBranch(), state.version.getAsLong(), metrics);
            return state.version;
        }
        long timestamp = parseTimeTravelTimestamp(snapshot.getValue());
        try {
            return OptionalLong.of(resolveVersionAtOrBefore(latest, access, timestamp, snapshot.getValue(), state,
                    metrics));
        } catch (LanceSnapshotResolver.NoVersionAtOrBeforeException e) {
            if (!access.getBranch().isPresent()) {
                throw new LanceUserFacingException("Lance table " + state.tableName + " has no version at or before '"
                        + snapshot.getValue() + "'");
            }
            // A branch's chain starts at the version it was created from and carries its own
            // commit times, so an earlier timestamp has nothing to select on the branch.
            throw new LanceUserFacingException("Lance branch '" + access.getBranch().get() + "' of "
                    + state.tableName + " has no version at or before '" + snapshot.getValue()
                    + "'; a branch only holds the versions from its creation on");
        }
    }

    /**
     * Whether a failure means the branch does not exist. A checkout reports "branch <name> does
     * not exist", or a missing manifest under the branch directory when nothing was ever
     * committed there; a namespace that reports the branch itself throws
     * {@link TableBranchNotFoundException}.
     */
    private static boolean isBranchNotFound(Throwable throwable, String branch) {
        if (ExceptionUtils.indexOfType(throwable, TableBranchNotFoundException.class) >= 0) {
            return true;
        }
        String rootMessage = ExceptionUtils.getRootCauseMessage(throwable);
        if (rootMessage == null) {
            return false;
        }
        String lower = rootMessage.toLowerCase(Locale.ROOT);
        String name = branch.toLowerCase(Locale.ROOT);
        return lower.contains("branch " + name + " does not exist")
                || (lower.contains("not found") && lower.contains("tree/" + name + "/"));
    }

    /** Whether a not-found came from the namespace client rather than storage. */
    private static boolean isNamespaceMiss(Throwable throwable) {
        return ExceptionUtils.indexOfType(throwable, TableVersionNotFoundException.class) >= 0
                || ExceptionUtils.indexOfType(throwable, TableBranchNotFoundException.class) >= 0;
    }

    /**
     * Every version the namespace records for the chain on {@code branch}, listed once per read.
     * The whole list is needed: FOR TIME AS OF selects among it, and neither the order a namespace
     * returns nor monotonic commit times can be relied on to stop early.
     */
    private List<TableVersion> namespaceVersions(ReadState state, Optional<String> branch,
            LanceMetadataMetrics metrics) {
        return state.namespaceVersions.computeIfAbsent(branch.orElse(""), chain -> {
            List<TableVersion> versions = metrics.measure(Stage.VERSION_RESOLVE,
                    () -> namespaceClient.listManagedVersions(state.access, branch));
            if (versions.isEmpty()) {
                throw new LanceUserFacingException("Lance namespace lists no versions for "
                        + state.tableName + (chain.isEmpty() ? "" : "@" + chain));
            }
            return versions;
        });
    }

    /**
     * Resolves {@code FOR TIME AS OF} to a version on the chain {@code latest} is checked out on,
     * from the commit times the manifests in storage record, over the history
     * {@link LanceSnapshotResolver} describes. A managed table only selects among the versions its
     * namespace records. A version it no longer records between recorded ones cuts the history
     * like a removed one, and so does a recorded version storage lacks, since its commit time is
     * unknown.
     */
    private long resolveVersionAtOrBefore(Dataset latest, LanceTableAccess access, long timestamp,
            String requestedText, ReadState state, LanceMetadataMetrics metrics) {
        state.pinned = null;
        Map<Long, TableVersion> records = null;
        if (access.isManagedVersioning()) {
            records = new HashMap<>();
            for (TableVersion recorded : namespaceVersions(state, access.getBranch(), metrics)) {
                if (recorded.getVersion() != null) {
                    // Selection compares the commit times of the manifests Doris reads, so each
                    // recorded version must be at one of them, not only the selected one.
                    LanceManifestPaths.check(state.access.getDatasetUri(), access.getBranch(), recorded.getVersion(),
                            recorded.getManifestPath(), state.tableName);
                    records.put(recorded.getVersion(), recorded);
                }
            }
        }
        Map<Long, TableVersion> recordedById = records;
        NavigableSet<Long> recorded = records == null ? null : new TreeSet<>(records.keySet());
        long version = metrics.measure(Stage.VERSION_RESOLVE, () -> {
            try {
                return LanceSnapshotResolver.versionAtOrBefore(latest.listVersions(), recorded, timestamp,
                        requestedText);
            } catch (LanceSnapshotResolver.HistoryRemovedException e) {
                // A removed version the namespace still records may only be staged; one it no
                // longer records is gone.
                TableVersion removed = recordedById == null ? null : recordedById.get(e.getVersion());
                boolean staged = removed != null && LanceManifestPaths.check(state.access.getDatasetUri(),
                        access.getBranch(), e.getVersion(), removed.getManifestPath(), state.tableName)
                        == LanceManifestPaths.Recorded.STAGED;
                throw historyRemoved(e.getVersion(), staged, requestedText, state);
            }
        });
        if (records != null) {
            state.pinned = new Recorded(access.getBranch(), version, LanceManifestPaths.check(
                    state.access.getDatasetUri(), access.getBranch(), version, records.get(version).getManifestPath(),
                    state.tableName));
        }
        LOG.debug("Resolved Lance FOR TIME AS OF '{}' to version {} from manifest commit times", requestedText,
                version);
        return version;
    }

    private static LanceUserFacingException historyRemoved(long version, boolean staged, String requestedText,
            ReadState state) {
        return new LanceUserFacingException("Lance cannot resolve FOR TIME AS OF '" + requestedText + "' on "
                + state.displayName() + ": version " + version + ", which may hold the state at that time, "
                + (staged ? "cannot be read: " + stagedOnly() : "no longer exists"));
    }

    /**
     * Parses a {@code FOR TIME AS OF} value in the session time zone. Second and millisecond
     * precision are accepted. Commit times are compared in full, so a commit later within the
     * requested millisecond is not selected.
     */
    private static long parseTimeTravelTimestamp(String value) {
        long timestamp = TimeUtils.timeStringToLong(value, TimeUtils.getTimeZone());
        if (timestamp < 0) {
            timestamp = TimeUtils.msTimeStringToLong(value, TimeUtils.getTimeZone());
        }
        if (timestamp < 0) {
            throw new IllegalArgumentException("Cannot parse Lance FOR TIME AS OF value '" + value
                    + "', expected 'yyyy-MM-dd HH:mm:ss' or 'yyyy-MM-dd HH:mm:ss.SSS'");
        }
        return timestamp;
    }

    /**
     * Whether a failed open of an explicitly requested version means that version does not exist.
     * A namespace reports it through {@link TableVersionNotFoundException}. The storage reader
     * reports it as a missing manifest under {@code _versions/} or as Lance's own version-not-found
     * error; a missing dataset or an unreachable store fails differently and keeps its message.
     */
    private static boolean isVersionNotFound(Throwable throwable) {
        if (ExceptionUtils.indexOfType(throwable, TableVersionNotFoundException.class) >= 0) {
            return true;
        }
        String rootMessage = ExceptionUtils.getRootCauseMessage(throwable);
        if (rootMessage == null) {
            return false;
        }
        String lower = rootMessage.toLowerCase(Locale.ROOT);
        return lower.contains("version not found")
                || (lower.contains("not found") && lower.contains("_versions/"));
    }

    private Dataset openDataset(BufferAllocator allocator, LanceTableAccess access, OptionalLong version,
            LanceMetadataMetrics metrics) {
        return metrics.measure(Stage.DATASET_OPEN, () -> Dataset.open().allocator(allocator).uri(access.getDatasetUri())
                .readOptions(LanceReadOptions.forSharedSession(access.getStorageOptions(), version, session)).build());
    }

    @FunctionalInterface
    private interface SnapshotReader<T> {
        T read(Dataset dataset, LanceTableAccess access, LanceMetadataMetrics metrics);
    }

    public List<LanceShowIndexInfo> loadTableIndexesForShow(
            String dbName, String tableName) {
        return inspectTableIndexes(dbName, tableName,
                (dataset, uri) -> LanceIndexInspection.readIndexesForShow(dataset));
    }

    public List<LancePhysicalIndexEntry> loadTableIndexEntries(
            String dbName, String tableName) {
        return inspectTableIndexes(dbName, tableName,
                (dataset, uri) -> LanceIndexInspection.readPhysicalEntries(dataset));
    }

    String resolveCurrentIndexJobLocator(String dbName, String tableName) {
        return LanceIndexDatasetLocator.normalize(
                namespaceClient.resolveTableAccessUncached(dbName, tableName).getDatasetUri());
    }

    public LanceIndexAdmissionSnapshot loadTableIndexAdmissionSnapshot(String dbName, String tableName) {
        return inspectTableIndexes(dbName, tableName, LanceIndexInspection::readAdmissionSnapshot);
    }

    private <T> T inspectTableIndexes(String dbName, String tableName, BiFunction<Dataset, String, T> inspection) {
        LanceTableAccess tableAccess = null;
        try {
            // The worker owns its Dataset and Session even if the caller releases its lease on timeout.
            // Index admission must verify the current target even while query access is cached.
            tableAccess = namespaceClient.resolveTableAccessUncached(dbName, tableName);
            // Index paths open the dataset by URI at storage-latest. They are only reachable for
            // filesystem catalogs, whose Directory namespace never manages versions; a managed
            // table here would bypass the namespace.
            if (tableAccess.isManagedVersioning()) {
                throw new IllegalStateException("Lance index inspection does not support namespace-managed tables");
            }
            LanceTableAccess access = tableAccess;
            return LanceIndexInspectionExecutor.execute(() -> {
                // The caller can time out while JNI is running; the worker must own resource cleanup.
                try (BufferAllocator allocator = new RootAllocator(LanceMetadataLoader.READ_ALLOCATOR_LIMIT);
                        Dataset dataset = Dataset.open().allocator(allocator).uri(access.getDatasetUri())
                                .readOptions(LanceReadOptions.forIndependentRead(
                                        access.getStorageOptions(), OptionalLong.empty()))
                                .build()) {
                    return inspection.apply(dataset, access.getDatasetUri());
                }
            });
        } catch (Exception e) {
            throw LanceErrorMessages.failure("Failed to load Lance index metadata for " + dbName + "." + tableName, e,
                    tableAccess == null ? null : tableAccess.getDatasetUri(),
                    tableAccess == null ? namespaceStorageOptions : tableAccess.getStorageOptions(), catalogSecrets);
        }
    }

}