diff --git a/paimon-filesystems/paimon-oss-impl/pom.xml b/paimon-filesystems/paimon-oss-impl/pom.xml index b6801595b2e5..6f7e11e28315 100644 --- a/paimon-filesystems/paimon-oss-impl/pom.xml +++ b/paimon-filesystems/paimon-oss-impl/pom.xml @@ -105,6 +105,19 @@ + + org.codehaus.mojo + templating-maven-plugin + + + filter-sources + + filter-sources + + + + + org.apache.maven.plugins diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java-templates/org/apache/paimon/oss/OSSBuildVersions.java b/paimon-filesystems/paimon-oss-impl/src/main/java-templates/org/apache/paimon/oss/OSSBuildVersions.java new file mode 100644 index 000000000000..8887f0c89bf1 --- /dev/null +++ b/paimon-filesystems/paimon-oss-impl/src/main/java-templates/org/apache/paimon/oss/OSSBuildVersions.java @@ -0,0 +1,30 @@ +/* + * 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.paimon.oss; + +/** Versions written in at build time, so reading them needs no class loader lookup. */ +final class OSSBuildVersions { + + static final String PAIMON = "${project.version}"; + + /** The aliyun-sdk-oss version bundled into this plugin. */ + static final String OSS_SDK = "${fs.oss.sdk.version}"; + + private OSSBuildVersions() {} +} diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java index 17842c1f147a..31bc2fba9a3e 100644 --- a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java +++ b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java @@ -164,6 +164,10 @@ public void configure(CatalogContext context) { } } } + // A user-set hadoop-aliyun prefix wins over Paimon's unified User-Agent. + if (!hadoopOptions.containsKey(OSSUserAgent.PREFIX)) { + hadoopOptions.set(OSSUserAgent.PREFIX, OSSUserAgent.prefix(context.options())); + } } @Override diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSUserAgent.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSUserAgent.java new file mode 100644 index 000000000000..2a9ba08d7193 --- /dev/null +++ b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSUserAgent.java @@ -0,0 +1,96 @@ +/* + * 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.paimon.oss; + +import org.apache.paimon.options.Options; +import org.apache.paimon.utils.StringUtils; + +import java.util.Arrays; +import java.util.StringJoiner; + +/** + * Builds Paimon's unified OSS User-Agent {@code module(transport;features) extended}. Each part + * comes from {@code fs.oss.user.agent.*}, falling back to the catalog-wide {@code user-agent.*}. + */ +final class OSSUserAgent { + + static final String MODULE = "fs.oss.user.agent.module"; + static final String FEATURES = "fs.oss.user.agent.features"; + static final String EXTENDED = "fs.oss.user.agent.extended"; + + static final String GENERIC_MODULE = "user-agent.module"; + static final String GENERIC_FEATURES = "user-agent.features"; + static final String GENERIC_EXTENDED = "user-agent.extended"; + + /** hadoop-aliyun sends this followed by {@code ", Hadoop/"}. */ + static final String PREFIX = "fs.oss.user.agent.prefix"; + + static final String DLF_ACCESS_TRACKING_EXTENDED_INFO = "dlf.access-tracking.extended-info"; + + private static final String DEFAULT_MODULE = "Paimon/" + OSSBuildVersions.PAIMON; + + private OSSUserAgent() {} + + static String prefix(Options options) { + String module = part(options, MODULE, GENERIC_MODULE); + StringBuilder builder = new StringBuilder(module == null ? DEFAULT_MODULE : module); + + builder.append("(aliyun-sdk-java/").append(OSSBuildVersions.OSS_SDK); + String features = part(options, FEATURES, GENERIC_FEATURES); + if (features != null) { + for (String feature : features.split("\\s+")) { + builder.append(';').append(feature); + } + } + builder.append(')'); + + // Access tracking info is appended, so a user-set extended value is kept. + String extended = + join( + part(options, EXTENDED, GENERIC_EXTENDED), + options.get(DLF_ACCESS_TRACKING_EXTENDED_INFO)); + if (!extended.isEmpty()) { + builder.append(' ').append(extended); + } + return builder.toString(); + } + + /** + * The OSS-specific value if set, else the catalog-wide one, trimmed; null when both are blank. + */ + private static String part(Options options, String ossKey, String genericKey) { + for (String key : Arrays.asList(ossKey, genericKey)) { + String value = options.get(key); + if (!StringUtils.isNullOrWhitespaceOnly(value)) { + return value.trim(); + } + } + return null; + } + + private static String join(String first, String second) { + StringJoiner joiner = new StringJoiner(" "); + for (String part : Arrays.asList(first, second)) { + if (!StringUtils.isNullOrWhitespaceOnly(part)) { + joiner.add(part.trim()); + } + } + return joiner.toString(); + } +} diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/OSSUserAgentTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/OSSUserAgentTest.java new file mode 100644 index 000000000000..e3ac519b3caf --- /dev/null +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/OSSUserAgentTest.java @@ -0,0 +1,168 @@ +/* + * 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.paimon.oss; + +import org.apache.paimon.catalog.CatalogContext; +import org.apache.paimon.fs.Path; +import org.apache.paimon.options.Options; + +import com.aliyun.oss.common.utils.VersionInfoUtils; +import com.sun.net.httpserver.HttpServer; +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.net.InetAddress; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link OSSUserAgent}. */ +public class OSSUserAgentTest { + + private static final String SDK = "aliyun-sdk-java/" + OSSBuildVersions.OSS_SDK; + + private static final String EMPTY_LISTING = + "\n" + + "bucket1" + + "false0"; + + @Test + public void testDefault() { + // The version written in at build time matches the SDK jar actually bundled. + assertThat(OSSBuildVersions.OSS_SDK).isEqualTo(VersionInfoUtils.getVersion()); + assertThat(OSSBuildVersions.PAIMON).matches("\\d+\\.\\d+\\S*"); + assertThat(OSSUserAgent.prefix(new Options())) + .isEqualTo("Paimon/" + OSSBuildVersions.PAIMON + "(" + SDK + ")"); + } + + @Test + public void testModuleAndFeatures() { + Options options = new Options(); + options.set(OSSUserAgent.MODULE, "MyApp/1.0"); + options.set(OSSUserAgent.FEATURES, " Flink Paimon "); + + assertThat(OSSUserAgent.prefix(options)).isEqualTo("MyApp/1.0(" + SDK + ";Flink;Paimon)"); + } + + @Test + public void testAccessTrackingIsAppendedToUserExtended() { + Options options = new Options(); + options.set(OSSUserAgent.MODULE, "MyApp/1.0"); + options.set(OSSUserAgent.EXTENDED, "vvr"); + options.set(OSSUserAgent.DLF_ACCESS_TRACKING_EXTENDED_INFO, "uid/123 user/alice"); + + assertThat(OSSUserAgent.prefix(options)) + .isEqualTo("MyApp/1.0(" + SDK + ") vvr uid/123 user/alice"); + } + + @Test + public void testAccessTrackingOnly() { + Options options = new Options(); + options.set(OSSUserAgent.MODULE, "MyApp/1.0"); + options.set(OSSUserAgent.EXTENDED, " "); + options.set(OSSUserAgent.DLF_ACCESS_TRACKING_EXTENDED_INFO, "uid/123"); + + assertThat(OSSUserAgent.prefix(options)).isEqualTo("MyApp/1.0(" + SDK + ") uid/123"); + } + + @Test + public void testCatalogWideKeys() { + Options options = new Options(); + options.set(OSSUserAgent.GENERIC_MODULE, "MyApp/1.0"); + options.set(OSSUserAgent.GENERIC_FEATURES, "Flink"); + options.set(OSSUserAgent.GENERIC_EXTENDED, "vvr"); + options.set(OSSUserAgent.DLF_ACCESS_TRACKING_EXTENDED_INFO, "uid/123"); + + assertThat(OSSUserAgent.prefix(options)) + .isEqualTo("MyApp/1.0(" + SDK + ";Flink) vvr uid/123"); + } + + @Test + public void testOssKeysOverrideCatalogWideKeysPerPart() { + Options options = new Options(); + options.set(OSSUserAgent.GENERIC_MODULE, "MyApp/1.0"); + options.set(OSSUserAgent.GENERIC_FEATURES, "Flink"); + options.set(OSSUserAgent.GENERIC_EXTENDED, "vvr"); + options.set(OSSUserAgent.FEATURES, "Spark"); + options.set(OSSUserAgent.EXTENDED, " "); + + assertThat(OSSUserAgent.prefix(options)).isEqualTo("MyApp/1.0(" + SDK + ";Spark) vvr"); + } + + @Test + public void testUserPrefixWins() { + Options options = new Options(); + options.set(OSSUserAgent.PREFIX, "custom"); + options.set(OSSUserAgent.EXTENDED, "vvr"); + OSSFileIO fileIO = new OSSFileIO(); + fileIO.configure(CatalogContext.create(options)); + + assertThat(fileIO.hadoopOptions().get(OSSUserAgent.PREFIX)).isEqualTo("custom"); + } + + @Test + public void testUserAgentOnTheWire() throws IOException { + List userAgents = new CopyOnWriteArrayList<>(); + HttpServer server = + HttpServer.create(new InetSocketAddress(InetAddress.getLoopbackAddress(), 0), 0); + server.createContext( + "/", + exchange -> { + userAgents.add(exchange.getRequestHeaders().getFirst("User-Agent")); + exchange.getResponseHeaders().add("x-oss-request-id", "fake-request-id"); + if ("HEAD".equals(exchange.getRequestMethod())) { + exchange.sendResponseHeaders(404, -1); + } else { + byte[] body = EMPTY_LISTING.getBytes(StandardCharsets.UTF_8); + exchange.getResponseHeaders().add("Content-Type", "application/xml"); + exchange.sendResponseHeaders(200, body.length); + exchange.getResponseBody().write(body); + } + exchange.close(); + }); + server.start(); + try { + Options options = new Options(); + options.set("fs.oss.endpoint", "http://127.0.0.1:" + server.getAddress().getPort()); + options.set("fs.oss.accessKeyId", "ak"); + options.set("fs.oss.accessKeySecret", "sk"); + options.set("fs.oss.attempts.maximum", "1"); + options.set("file-io.allow-cache", "false"); + options.set(OSSUserAgent.GENERIC_FEATURES, "Flink"); + options.set(OSSUserAgent.GENERIC_EXTENDED, "vvr"); + options.set(OSSUserAgent.DLF_ACCESS_TRACKING_EXTENDED_INFO, "uid/123"); + OSSFileIO fileIO = new OSSFileIO(); + fileIO.configure(CatalogContext.create(options)); + + assertThat(fileIO.exists(new Path("oss://bucket/dir/file"))).isFalse(); + assertThat(userAgents) + .isNotEmpty() + .allSatisfy( + ua -> + assertThat(ua) + .startsWith(OSSUserAgent.prefix(options) + ", Hadoop/") + .contains("(" + SDK + ";Flink) vvr uid/123")); + } finally { + server.stop(0); + } + } +} diff --git a/pom.xml b/pom.xml index 0b5512bd31d1..ee9ba053991a 100644 --- a/pom.xml +++ b/pom.xml @@ -980,6 +980,12 @@ under the License. 1.7 + + org.codehaus.mojo + templating-maven-plugin + 3.0.0 + + org.apache.maven.plugins maven-compiler-plugin