zy-kkk commented on code in PR #68453:
URL: https://github.com/apache/doris/pull/68453#discussion_r4123551105


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceSdkNamespace.java:
##########
@@ -0,0 +1,217 @@
+// 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.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 java.nio.charset.StandardCharsets;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Set;
+import java.util.TreeMap;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Supplier;
+
+/**
+ * The namespace the Lance SDK is handed to open one namespace-managed 
dataset: the catalog's
+ * namespace, with two 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 options the store is 
built with, less the
+ * credentials the SDK re-reads through its credential provider anyway.
+ *
+ * <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";
+
+    /**
+     * The options the SDK serves to an object store through its credential 
provider, refreshing
+     * them from the namespace, rather than fixing them when the store is 
built: every spelling
+     * lance-io's {@code DynamicCredentials} conversions read for AWS, Azure 
and GCS, and the
+     * refresh deadline. OSS credentials are not refreshed, but are left out 
as well: a store is
+     * shared across credentials as it was before, and a namespace that vends 
new credentials on
+     * every describe would otherwise leave one registry entry behind per read.
+     */
+    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");
+
+    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;
+
+    /**
+     * @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);
+                storeIdentity = storeIdentity(opened);
+            }
+            return response;
+        });
+    }
+
+    @Override
+    public ListTableVersionsResponse 
listTableVersions(ListTableVersionsRequest request) {
+        return record(() -> catalogNamespace.listTableVersions(request));
+    }
+
+    @Override
+    public DescribeTableVersionResponse 
describeTableVersion(DescribeTableVersionRequest request) {
+        return record(() -> catalogNamespace.describeTableVersion(request));
+    }
+
+    /**
+     * The exception a failed SDK call reported as a failed callback, in place 
of that report, or
+     * {@code sdkError} itself. Each kept exception is handed out once.
+     */
+    Exception unwrapCallbackFailure(Exception sdkError) {
+        if (!isCallbackFailure(sdkError)) {
+            return sdkError;
+        }
+        RuntimeException failure = callbackFailure.getAndSet(null);
+        if (failure == null) {
+            return sdkError;
+        }
+        failure.addSuppressed(sdkError);
+        return failure;
+    }
+
+    private static boolean isCallbackFailure(Throwable error) {
+        return ExceptionUtils.getThrowableList(error).stream()
+                .anyMatch(cause -> cause.getMessage() != null && 
cause.getMessage().contains(CALLBACK_FAILURE));
+    }
+
+    private <T> T record(Supplier<T> call) {
+        try {
+            return call.get();
+        } catch (RuntimeException e) {
+            callbackFailure.set(e);
+            if (Thread.currentThread() != openingThread) {
+                // The JNI leaves the exception pending on a thread it 
attached only for this call,
+                // and the JVM reports it as uncaught when that thread 
detaches. It is not: it
+                // reaches the read through unwrapCallbackFailure.
+                Thread.currentThread().setUncaughtExceptionHandler(
+                        (thread, error) -> LOG.debug("Lance namespace callback 
failed", error));
+            }
+            throw e;
+        }
+    }
+
+    /**
+     * A digest of every option that fixes where and how a store connects. The 
options can name
+     * endpoints and account names, so only the digest reaches the id, which 
Lance logs.
+     */
+    static String storeIdentity(Map<String, String> options) {
+        Hasher hasher = Hashing.sha256().newHasher();
+        new TreeMap<>(options).forEach((key, value) -> {
+            if (!CREDENTIAL_OPTIONS.contains(key)) {

Review Comment:
   Confirmed: without `expires_at_millis`, 
`StorageOptionsAccessor::needs_refresh` never refreshes, so a store shared 
across a key rotation keeps the key it was built with. Lance keys a store 
opened without a namespace by all of its options, credentials included (the 
static `accessor_id`), so the digest was looser than Lance's own key.
   
   Fixed in 78cb959f7b7. Credentials now count toward the store identity unless 
the options carry `expires_at_millis`. With an expiry, the store refreshes them 
from the namespace before they expire (AWS, Azure and GCS through 
`DynamicCredentials`, OSS through its dynamic OpenDAL store), so they stay out 
of the identity. Otherwise a namespace that vends new credentials on every 
describe would leave one registry entry behind per read: the registry drops a 
dead entry only when the same key is looked up again, and each entry holds the 
JNI reference to its namespace. Since the digest now covers secrets and Lance 
logs the id, it is an HMAC under a random per-process key.
   
   Test: 
`LanceManagedS3StoreTest.testOverlappingReadsKeepTheirOwnCredentialsWithoutAnExpiry`.
 The S3 stub rejects a revoked access key. Q1 holds a store built with `ak-1`; 
the namespace then vends `ak-2` without an expiry, and `ak-1` is revoked. Q2 
must plan with `ak-2`. With the previous rule, Q2 fails with 403 on its first 
HEAD, sent with `ak-1` through the store Q1 built. 
`LanceSdkNamespaceTest.testStoreIdentityCountsCredentialsTheStoreDoesNotRefresh`
 covers the rule without JNI.
   
   With an expiry, an overlapping read uses the store's current credentials 
until the store refreshes them, as a single long read does. lance-io does not 
handle a vended credential revoked before its expiry in either case.
   



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceSdkNamespace.java:
##########
@@ -0,0 +1,217 @@
+// 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.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 java.nio.charset.StandardCharsets;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Set;
+import java.util.TreeMap;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Supplier;
+
+/**
+ * The namespace the Lance SDK is handed to open one namespace-managed 
dataset: the catalog's
+ * namespace, with two 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 options the store is 
built with, less the
+ * credentials the SDK re-reads through its credential provider anyway.
+ *
+ * <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";
+
+    /**
+     * The options the SDK serves to an object store through its credential 
provider, refreshing
+     * them from the namespace, rather than fixing them when the store is 
built: every spelling
+     * lance-io's {@code DynamicCredentials} conversions read for AWS, Azure 
and GCS, and the
+     * refresh deadline. OSS credentials are not refreshed, but are left out 
as well: a store is
+     * shared across credentials as it was before, and a namespace that vends 
new credentials on
+     * every describe would otherwise leave one registry entry behind per read.
+     */
+    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");
+
+    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;
+
+    /**
+     * @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(

Review Comment:
   The fold is real: `with_env_azure` only checks `azure_storage_endpoint`, and 
object_store 0.14.1 parses `endpoint` and `azure_endpoint` as the same config 
key, so which one wins depends on map order.
   
   It is not specific to managed tables or to this PR, though. Azure options go 
through `LancePassThroughStorageProvider`, which keeps them in the namespace's 
spelling for every Lance table. A non-managed `az://` table with a vended 
`endpoint` and `AZURE_STORAGE_ENDPOINT` in the FE environment folds the same 
way on the FE, before and after this PR, and the BE merges its own environment 
independently. The same applies to the other Azure keys and to GCS. The S3 case 
in the earlier thread was different: the previous round's `forManagedSdkOpen` 
replaced the canonical key Doris already produced with the raw alias, so this 
PR had introduced it.
   
   Fixing this means giving the pass-through provider an Azure vocabulary for 
both the FE and the BE. The provider's class comment explains why it stays 
inert until there is an Azure backend to test such a translation against, and I 
don't have one for this PR. As with `AWS_ENDPOINT_URL_S3` in the earlier 
thread, I'd leave this interaction with the deployment environment to a 
follow-up that comes with Azure coverage.
   



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceCatalogClient.java:
##########
@@ -235,68 +252,508 @@ 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, metrics)) {
+                        result = reader.read(dataset, state.access, metrics);
+                    }
+                } else {
+                    OptionalLong mainVersion = 
state.access.isManagedVersioning()

Review Comment:
   Right: every selector other than an explicit version on `main` is checked 
out from an open of `main`'s latest recorded version, as the PR description 
says. The Java SDK only opens a dataset on `main`; Lance 12 has no Java open at 
a branch or tag, and `Dataset.checkout` needs an open dataset.
   
   The state in the example is not one Lance produces. A managed table records 
`main` version 1 at its first commit, a branch is created from a version that 
already exists, and nothing in `rust/lance` removes a version record: 
`batch_delete_table_versions` has no caller there, and cleanup deletes files, 
not namespace records. So `main` has no recorded version only when the table 
has no data at all, which is reported as such, or when a client deleted 
`main`'s records through the namespace API and kept a branch's. The stub's 
`versions=[]` models the empty table.
   
   I'd keep the rejection of an empty `main` history. Once the Java SDK can 
open a branch directly, Doris can open the selected branch without going 
through `main`, and this dependency goes away.
   



##########
fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LanceManagedS3StoreTest.java:
##########
@@ -0,0 +1,412 @@
+// 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.common.util.JsonUtil;
+import org.apache.doris.datasource.lance.metadata.LanceTableMetadata;
+
+import com.google.common.io.ByteStreams;
+import com.sun.jna.Library;
+import com.sun.jna.Native;
+import com.sun.net.httpserver.HttpExchange;
+import com.sun.net.httpserver.HttpServer;
+import org.apache.arrow.memory.BufferAllocator;
+import org.apache.arrow.memory.RootAllocator;
+import org.apache.arrow.vector.IntVector;
+import org.apache.arrow.vector.VectorSchemaRoot;
+import org.apache.arrow.vector.types.pojo.ArrowType;
+import org.apache.arrow.vector.types.pojo.Field;
+import org.apache.arrow.vector.types.pojo.FieldType;
+import org.apache.arrow.vector.types.pojo.Schema;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Disabled;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.TestInstance;
+import org.lance.Dataset;
+import org.lance.Fragment;
+import org.lance.FragmentMetadata;
+import org.lance.FragmentOperation;
+import org.lance.WriteParams;
+
+import java.io.IOException;
+import java.io.OutputStream;
+import java.math.BigInteger;
+import java.net.InetSocketAddress;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.time.Instant;
+import java.time.ZoneOffset;
+import java.time.format.DateTimeFormatter;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+import java.util.stream.Stream;
+
+/**
+ * Reads a namespace-managed table on S3 through two S3 stubs that serve 
different data under one
+ * key, standing in for two endpoints of a table. The namespace vends the 
endpoint in its own
+ * spelling ({@code endpoint}), as a REST catalog may.
+ *
+ * <p>The FE must plan from the endpoint the BE is handed. That breaks if the 
SDK reuses the object
+ * store another read still holds for the same table, or if the FE 
environment's
+ * {@code AWS_ENDPOINT} takes the place of the vended endpoint.
+ */
+@Disabled("Re-enable after fixing Arrow C Data JNI compatibility: CI libstdc++ 
lacks CXXABI_1.3.9")

Review Comment:
   These classes need the Arrow C Data JNI library, which does not load on the 
FE UT hosts: their libstdc++ lacks CXXABI_1.3.9. A read-only case would not 
avoid it either, since `Dataset.getSchema()` imports the schema through Arrow C 
Data. That is why the existing dataset-backed Lance tests in fe-core 
(`LanceSharedSessionTest`, `LanceSchemaContractBuilderRealDatasetTest`) carry 
the same `@Disabled` reason; these follow them, so all of them can be 
re-enabled together once the runner is fixed. Changing the runner image is 
outside this PR.
   
   What runs on CI today: `LanceSdkNamespaceTest` (normalization, store 
identity, callback failures), `LanceSdkOpenedAccessTest` and 
`LanceSnapshotTest` without JNI, and `test_lance_rest_time_travel`, which reads 
the managed fixtures (`time_travel_managed`, `_partial`, `_lagging`, 
`_untimed`, `_unprefixed`) through the REST stub and MinIO with version, time, 
tag and branch selectors. For this PR I run every Lance test class in fe-core 
locally with 
`-Djunit.jupiter.conditions.deactivate='org.junit.*DisabledCondition'`, which 
lifts the `@Disabled`s; on 78cb959f7b7 that is 49 classes and 630 cases, all 
passing.
   



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