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


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceSdkNamespace.java:
##########
@@ -0,0 +1,375 @@
+// 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.datasource.lance.storage.LanceStorageOptions;
+
+import com.google.common.collect.ImmutableSet;
+import com.google.common.hash.Hasher;
+import com.google.common.hash.Hashing;
+import org.apache.arrow.memory.BufferAllocator;
+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.namespace.LanceNamespace;
+import org.lance.namespace.model.DescribeTableRequest;
+import org.lance.namespace.model.DescribeTableResponse;
+import org.lance.namespace.model.DescribeTableVersionRequest;
+import org.lance.namespace.model.DescribeTableVersionResponse;
+import org.lance.namespace.model.ListTableVersionsRequest;
+import org.lance.namespace.model.ListTableVersionsResponse;
+import org.lance.namespace.model.TableVersion;
+
+import java.io.ByteArrayOutputStream;
+import java.math.BigInteger;
+import java.nio.charset.StandardCharsets;
+import java.security.SecureRandom;
+import java.util.HashMap;
+import java.util.Locale;
+import java.util.Map;
+import java.util.Set;
+import java.util.TreeMap;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Supplier;
+import java.util.regex.Pattern;
+
+/**
+ * The namespace the Lance SDK is handed to open one namespace-managed 
dataset: the catalog's
+ * namespace, with three things the SDK gets wrong on its own.
+ *
+ * <p>The SDK opens with the options it is handed plus whatever its own 
describe vends, spelled as
+ * the namespace spells them. Lance then adds the process environment for any 
option whose
+ * canonical key is missing, so a vended {@code endpoint} next to an {@code 
AWS_ENDPOINT} in the
+ * FE environment leaves the FE on whichever endpoint object_store folds last, 
while the BE, handed
+ * the canonical {@code aws_endpoint}, keeps the vended one. {@link 
#describeTable} therefore
+ * returns the vended options in the vocabulary Doris uses for everything else.
+ *
+ * <p>The SDK also caches the object store of a namespace-opened dataset in 
the catalog Session by
+ * the namespace's id and the table id alone, ignoring the options: a read 
overlapping another one
+ * that still holds a store for the same table reuses that store, whatever 
endpoint it was built
+ * for (lance-io {@code StorageOptionsAccessor::accessor_id}, {@code 
ObjectStoreRegistry::get_store}).
+ * {@link #namespaceId} therefore also identifies the table location and the 
options the store is
+ * built with, less credentials the store refreshes from the namespace anyway.
+ *
+ * <p>Lance opens a finalized manifest a namespace records wherever it is, 
while the BE opens a
+ * version by the dataset URI and number, at the canonical path. {@link 
#describeTableVersion} and
+ * {@link #listTableVersions} therefore reject a finalized manifest anywhere 
else.
+ *
+ * <p>A namespace Lance does not implement natively is called back through 
JNI, which reports an
+ * exception thrown by the callback only as "Java exception was thrown". The 
last one is kept for
+ * {@link #unwrapCallbackFailure}, so a missing version or branch is still 
reported as such.
+ *
+ * <p>One instance is created per open. The datasets checked out from it 
resolve versions through
+ * it, and the store the open builds refreshes credentials through it, also 
for later reads that
+ * share that store.
+ */
+final class LanceSdkNamespace implements LanceNamespace {
+    private static final Logger LOG = 
LogManager.getLogger(LanceSdkNamespace.class);
+
+    /** What jni-rs reports for a Java exception a callback threw. */
+    private static final String CALLBACK_FAILURE = "Java exception was thrown";
+
+    private static final String EXPIRES_AT_MILLIS = "expires_at_millis";
+
+    /** The values lance-core's {@code str_is_truthy} accepts, lower case. */
+    private static final Set<String> TRUTHY = ImmutableSet.of("1", "true", 
"on", "yes", "y");
+
+    private static final String MANIFEST_EXTENSION = ".manifest";
+
+    /** A URL scheme; lance-io takes a single letter before the colon for a 
Windows drive instead. */
+    private static final Pattern URL_SCHEME = 
Pattern.compile("^[A-Za-z][A-Za-z0-9+.-]+:");
+
+    /**
+     * The credentials a store takes from its credential provider rather than 
fixing them when it
+     * is built: every spelling lance-io's {@code DynamicCredentials} 
conversions read for AWS,
+     * Azure and GCS, the OSS keys its dynamic OpenDAL store re-reads, and the 
refresh deadline.
+     * The provider refreshes them from the namespace only when they carry
+     * {@value #EXPIRES_AT_MILLIS} and the store is not an OpenDAL one; see 
{@link #storeIdentity}.
+     */
+    private static final Set<String> CREDENTIAL_OPTIONS = ImmutableSet.of(
+            "aws_access_key_id", "access_key_id", "aws_secret_access_key", 
"secret_access_key",
+            "aws_session_token", "aws_token", "aws_security_token", 
"session_token", "token",
+            "azure_storage_sas_token", "azure_storage_sas_key", "sas_token", 
"sas_key",
+            "azure_storage_token", "bearer_token", 
"azure_storage_account_key", "azure_storage_access_key",
+            "azure_storage_master_key", "access_key", "master_key", 
"account_key",
+            "google_storage_token",
+            "oss_access_key_id", "oss_secret_access_key", "oss_security_token",
+            EXPIRES_AT_MILLIS);
+
+    /** Keys the store digest, so an id in a log cannot be checked against 
guessed credentials. */
+    private static final byte[] IDENTITY_KEY = new byte[32];
+
+    static {
+        new SecureRandom().nextBytes(IDENTITY_KEY);
+    }
+
+    private final LanceNamespace catalogNamespace;
+    private final Map<String, String> sdkStorageOptions;
+    private final AtomicReference<RuntimeException> callbackFailure = new 
AtomicReference<>();
+    /** The thread that opens the dataset and issues the SDK's describe; every 
other call is a JNI callback. */
+    private final Thread openingThread = Thread.currentThread();
+    /** Set by the SDK's own describe while it opens the dataset. */
+    private volatile String storeIdentity;
+    /** The table location the SDK's own describe returned; set with {@link 
#storeIdentity}. */
+    private volatile String tableLocation;
+
+    /**
+     * @param sdkStorageOptions the options the SDK is handed in its read 
options, which it opens
+     *     with under what its describe vends
+     */
+    LanceSdkNamespace(LanceNamespace catalogNamespace, Map<String, String> 
sdkStorageOptions) {
+        this.catalogNamespace = catalogNamespace;
+        this.sdkStorageOptions = sdkStorageOptions;
+    }
+
+    @Override
+    public void initialize(Map<String, String> configProperties, 
BufferAllocator allocator) {
+        throw new UnsupportedOperationException("A Lance SDK namespace wraps 
an initialized catalog namespace");
+    }
+
+    /**
+     * Read by the SDK once, when it opens the dataset: Lance 12 describes the 
table in
+     * {@code OpenDatasetBuilder.buildFromNamespaceClient} first and reads the 
id when the JNI
+     * wraps this namespace. It keys the store cache, through the credential 
provider the SDK
+     * builds from this namespace.
+     */
+    @Override
+    public String namespaceId() {
+        String identity = storeIdentity;
+        if (identity == null) {
+            throw new IllegalStateException("The Lance SDK read the namespace 
id before describing the table");
+        }
+        return "DorisSdkNamespace[" + catalogNamespace.namespaceId() + ", 
store=" + identity + "]";
+    }
+
+    /**
+     * The catalog namespace's describe, with the vended options normalized. 
The first call is the
+     * SDK's own describe while it opens the dataset, whose options the store 
is built with; later
+     * ones refresh credentials.
+     */
+    @Override
+    public DescribeTableResponse describeTable(DescribeTableRequest request) {
+        return record(() -> {
+            DescribeTableResponse response = 
catalogNamespace.describeTable(request);
+            Map<String, String> vended = 
LanceStorageOptions.normalizeVendedStorageOptions(
+                    response.getLocation(), response.getStorageOptions());
+            // Left null when nothing was vended: a credential refresh then 
keeps the options it has.
+            if (response.getStorageOptions() != null) {
+                response.setStorageOptions(vended);
+            }
+            if (storeIdentity == null) {
+                Map<String, String> opened = new HashMap<>(sdkStorageOptions);
+                opened.putAll(vended);
+                tableLocation = response.getLocation();
+                storeIdentity = storeIdentity(tableLocation, opened);
+            }
+            return response;
+        });
+    }
+
+    /**
+     * The catalog namespace's version list. Lance takes the newest entry's 
manifest as the head of
+     * a chain, so every entry is checked with {@link #checkManifestPath}.
+     */
+    @Override
+    public ListTableVersionsResponse 
listTableVersions(ListTableVersionsRequest request) {
+        return record(() -> {
+            ListTableVersionsResponse response = 
catalogNamespace.listTableVersions(request);
+            if (response.getVersions() != null) {
+                response.getVersions().forEach(version -> 
checkManifestPath(request.getBranch(), version));
+            }
+            return response;
+        });
+    }
+
+    /** The catalog namespace's describe of one version, whose manifest Lance 
opens; see {@link #checkManifestPath}. */
+    @Override
+    public DescribeTableVersionResponse 
describeTableVersion(DescribeTableVersionRequest request) {
+        return record(() -> {
+            DescribeTableVersionResponse response = 
catalogNamespace.describeTableVersion(request);
+            checkManifestPath(request.getBranch(), response.getVersion());
+            return response;
+        });
+    }
+
+    /**
+     * Rejects a finalized manifest the BE would not open. The BE opens a 
version by the dataset URI
+     * and number, at {@code <chain>/_versions/<u64::MAX - v>.manifest} (or 
{@code <v>.manifest} for
+     * the V1 naming scheme). Lance opens a manifest path ending in {@code 
.manifest} as recorded,
+     * and copies any other (staged) one to that canonical path first
+     * ({@code ExternalManifestCommitHandler::resolve_version_location}), so 
only a finalized path
+     * elsewhere can leave the FE and the BE reading different manifests.
+     */
+    private void checkManifestPath(String branch, TableVersion version) {
+        String path = version == null ? null : version.getManifestPath();
+        // Lance parses the path first, which drops surrounding slashes.
+        String recorded = path == null ? "" : StringUtils.strip(path, "/");
+        if (version == null || version.getVersion() == null || 
!recorded.endsWith(MANIFEST_EXTENSION)) {
+            return;
+        }
+        if (storeIdentity == null) {
+            throw new IllegalStateException("The Lance SDK resolved a version 
before describing the table");
+        }
+        String chain = objectStorePath(tableLocation);
+        if (branch != null) {
+            chain = (chain.isEmpty() ? "" : chain + "/") + "tree/" + branch;
+        }
+        String versions = (chain.isEmpty() ? "" : chain + "/") + "_versions/";
+        long number = version.getVersion();
+        String canonical = versions + String.format("%020d",

Review Comment:
   [P2] Format the canonical V2 manifest name with ASCII digits. 
`String.format("%020d", invertedVersion)` uses the FE JVM's default FORMAT 
locale; under a locale with non-ASCII digits it builds a different filename 
from Lance's ASCII `_versions/18446744073709551612.manifest`, so this check 
rejects a valid namespace-managed version before the FE can open it. Use 
`Locale.ROOT` (or another locale-independent ASCII conversion) and cover a 
non-Latin FORMAT locale. This differs from the existing noncanonical-manifest 
thread: the namespace path is canonical here.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceCatalogClient.java:
##########
@@ -235,68 +252,581 @@ public LanceTableMetadata loadBasicTableMetadata(String 
dbName, String tableName
     }
 
     public Schema loadTableSchema(String dbName, String tableName) {
-        return readTableSnapshot(dbName, tableName, Optional.empty(),
+        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 loadQueryMetadata(dbName, tableName, tableSnapshot, 
LanceMetadataLoader.MetadataScope.WITH_INDEXES);
+        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 readTableSnapshot(dbName, tableName, tableSnapshot,
+        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. */
-    private <T> T readTableSnapshot(String dbName, String tableName, 
Optional<TableSnapshot> tableSnapshot,
+    /**
+     * Pins one resource generation, resolved table access, and the Dataset 
version for the whole read.
+     *
+     * <p>The latest version of the main chain is opened once and every other 
selector is a
+     * checkout from that handle, so the SDK resolves the ref with the same 
commit handler
+     * (the namespace's, for a managed table). A tag is resolved first to the 
chain and version it
+     * points at, so a tag created on a branch selects that branch. The two 
shortcuts that skip the
+     * latest open are an explicit version on the main chain, and the latest 
version of a managed
+     * table. For a managed table, "latest" is always the newest version the 
namespace records,
+     * never the newest manifest in storage.
+     */
+    private <T> T readTableSnapshot(String dbName, String tableName, 
LanceRefSelector selector,
             SnapshotReader<T> reader) {
-        LanceTableAccess tableAccess = null;
+        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())) {
-                tableAccess = metrics.measure(Stage.TABLE_ACCESS,
+                state.access = metrics.measure(Stage.TABLE_ACCESS,
                         () -> namespaceClient.resolveTableAccess(dbName, 
tableName));
-                OptionalLong version = OptionalLong.empty();
-                if (tableSnapshot.isPresent()) {
-                    TableSnapshot snapshot = tableSnapshot.get();
-                    if (snapshot.getType() == 
TableSnapshot.VersionType.VERSION) {
-                        version = 
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
-                    } else {
-                        long timestamp = 
TimeUtils.timeStringToLong(snapshot.getValue(), TimeUtils.getTimeZone());
-                        if (timestamp < 0) {
-                            throw new IllegalArgumentException(
-                                    "Cannot parse Lance FOR TIME AS OF value 
'" + snapshot.getValue() + "'");
-                        }
-                        try (Dataset latest = openDataset(allocator, 
tableAccess, OptionalLong.empty(), metrics)) {
-                            version = 
OptionalLong.of(metrics.measure(Stage.VERSION_RESOLVE,
-                                    () -> 
LanceSnapshotResolver.getVersionAtOrBefore(latest, timestamp)));
-                        }
+                OptionalLong direct = directMainVersion(state, metrics);
+                if (direct.isPresent() || isLatestMain(selector)) {
+                    state.version = direct;
+                    try (Dataset dataset = openDataset(allocator, state, 
direct, isLatestMain(selector), metrics)) {
+                        result = reader.read(dataset, state.access, metrics);
+                    }
+                } else {
+                    OptionalLong mainVersion = 
state.access.isManagedVersioning()
+                            ? OptionalLong.of(recordedLatestVersion(state, 
Optional.empty(), metrics))
+                            : OptionalLong.empty();
+                    try (Dataset main = openDataset(allocator, state, 
mainVersion, true, metrics)) {
+                        result = readFromLatest(main, state, reader, metrics);
                     }
-                }
-                try (Dataset dataset = openDataset(allocator, tableAccess, 
version, metrics)) {
-                    result = reader.read(dataset, tableAccess, metrics);
                 }
             }
             metrics.succeeded();
             return result;
-        } catch (Exception e) {
-            throw LanceErrorMessages.failure("Failed to load Lance table 
metadata for " + dbName + "." + tableName, e,
-                    tableAccess == null ? null : tableAccess.getDatasetUri(),
-                    tableAccess == null ? namespaceStorageOptions : 
tableAccess.getStorageOptions(), catalogSecrets);
+        } catch (LanceUserFacingException e) {
+            throw new RuntimeException(e.getMessage(), e);
+        } catch (Exception sdkError) {
+            Exception e = unwrapCallbackFailure(state, sdkError);
+            LanceTableAccess access = state.access;
+            String uri = access == null ? null : access.getDatasetUri();
+            Map<String, String> options = access == null ? 
namespaceStorageOptions : access.getStorageOptions();
+            String what = state.displayName();
+            if (state.branch.isPresent() && !state.branchExists && 
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" + (isNamespaceMiss(e, "table branch 
not found") ? " in the namespace" : ""),
+                        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, "table 
version not found") ? " in the namespace" : ""),
+                        sanitizedCause(e, uri, options));
+            }
+            String hint = access != null && access.isManagedVersioning() && 
isAccessDenied(e)
+                    ? " (reading a namespace-managed Lance table may need 
write access to finalize a staged manifest)"
+                    : "";
+            throw LanceErrorMessages.failure("Failed to load Lance table 
metadata for " + what + hint, e, uri, options,
+                    catalogSecrets);
         } finally {
             metrics.close();
         }
     }
 
-    private Dataset openDataset(BufferAllocator allocator, LanceTableAccess 
access, OptionalLong version,
+    /** 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;
+        private LanceTableAccess access;
+        /** The namespace the SDK opened a managed table through, which keeps 
its callbacks' failures. */
+        private LanceSdkNamespace sdkNamespace;
+        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 namespace's version list per chain ("" is main), fetched at 
most once per read. */
+        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("");
+        }
+    }
+
+    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(recordedLatestVersion(state, 
Optional.empty(), metrics))
+                    : OptionalLong.empty();
+        }
+        TableSnapshot snapshot = selector.getSnapshot().get();
+        return snapshot.getType() == TableSnapshot.VersionType.VERSION
+                ? 
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()))
+                : OptionalLong.empty();
+    }
+
+    /** 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; for a managed 
table that is an
+            // explicit version the namespace resolves, never a storage 
fallback.
+            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)) {
+                state.branch = branchOf(target.uri(), 
state.access.getDatasetUri());
+                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.
+            Ref branchHead = Ref.ofBranch(branch);
+            if (state.access.isManagedVersioning()) {

Review Comment:
   [P1] Recheck the table location before checking out this managed branch. 
`main` was opened and reconciled at A, but the branch head is listed afterward. 
If the namespace repoints the table to B between those steps, `main.checkout` 
still uses A's branch root while the namespace returns B's version path. A 
finalized B manifest fails the path check; a staged B manifest in the same 
bucket can be copied to A's canonical `tree/dev/_versions/` path, overwriting 
A's same-number manifest. Restart from a fresh main handle or fail with a retry 
before checkout when the table access changed, and cover a move after the main 
open. This is a later race than the existing head-before-open relocation thread.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to