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]
