diff --git a/paimon-python/pypaimon/build_info.py b/paimon-python/pypaimon/build_info.py index e034f356fd4e..22649caf7789 100644 --- a/paimon-python/pypaimon/build_info.py +++ b/paimon-python/pypaimon/build_info.py @@ -84,3 +84,11 @@ def _load_full_version(): def full_version(): """Return ``-`` for snapshot provenance.""" return _FULL_VERSION + + +def version(): + """Return the pypaimon version embedded at build time, or None when unknown.""" + prefix = "python-" + if not _FULL_VERSION.startswith(prefix): + return None + return _FULL_VERSION[len(prefix):].rsplit("-", 1)[0] or None diff --git a/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py b/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py index eaad4a89be29..bade450d9b7b 100644 --- a/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py +++ b/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py @@ -44,6 +44,7 @@ from pypaimon.common.options import Options from pypaimon.common.options.config import OssOptions +from pypaimon.filesystem import oss_user_agent _JINDO_CONFIG_PREFIXES = ("fs.", "logger.") @@ -100,7 +101,14 @@ def build_jindo_config(catalog_options: Options): config.set("fs.oss.endpoint", endpoint_clean) if region: config.set("fs.oss.region", region) - config.set("fs.oss.user.agent.features", "pypaimon") + # This backend fills the module itself, so pypaimon/ leads the features. + user_agent_module = oss_user_agent.module(catalog_options) + if user_agent_module: + config.set(oss_user_agent.USER_AGENT_MODULE, user_agent_module) + config.set(oss_user_agent.USER_AGENT_FEATURES, oss_user_agent.features(catalog_options)) + user_agent_extended = oss_user_agent.extended(catalog_options) + if user_agent_extended: + config.set(oss_user_agent.USER_AGENT_EXTENDED, user_agent_extended) return config diff --git a/paimon-python/pypaimon/filesystem/oss_user_agent.py b/paimon-python/pypaimon/filesystem/oss_user_agent.py new file mode 100644 index 000000000000..844499c3a53d --- /dev/null +++ b/paimon-python/pypaimon/filesystem/oss_user_agent.py @@ -0,0 +1,77 @@ +# 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. + +"""Parts of Paimon's unified OSS User-Agent ``module(transport;features) extended``.""" + +import functools +from typing import List, Optional + +from pypaimon import build_info +from pypaimon.common.options import Options + +USER_AGENT_MODULE = "fs.oss.user.agent.module" +USER_AGENT_FEATURES = "fs.oss.user.agent.features" +USER_AGENT_EXTENDED = "fs.oss.user.agent.extended" +# Catalog-wide keys shared with the REST client; the fs.oss keys above take precedence. +COMMON_USER_AGENT_MODULE = "user-agent.module" +COMMON_USER_AGENT_FEATURES = "user-agent.features" +COMMON_USER_AGENT_EXTENDED = "user-agent.extended" +DLF_ACCESS_TRACKING_EXTENDED_INFO = "dlf.access-tracking.extended-info" + +_NAME = "pypaimon" + + +@functools.lru_cache(maxsize=None) +def identity() -> str: + """``pypaimon/`` with the version embedded at build time, or bare ``pypaimon``.""" + version = build_info.version() + return "{}/{}".format(_NAME, version) if version else _NAME + + +def module(options: Options) -> Optional[str]: + """The configured module, or None to keep the backend's own.""" + return _effective(options, USER_AGENT_MODULE, COMMON_USER_AGENT_MODULE) + + +def features(options: Options) -> str: + """The identity followed by the configured features, space separated.""" + user_features = _split(_effective(options, USER_AGENT_FEATURES, COMMON_USER_AGENT_FEATURES)) + if any(f == _NAME or f.startswith(_NAME + "/") for f in user_features): + return " ".join(user_features) + return " ".join([identity()] + user_features) + + +def extended(options: Options) -> Optional[str]: + """The configured extended info with the DLF access-tracking info appended, or None.""" + parts = [_effective(options, USER_AGENT_EXTENDED, COMMON_USER_AGENT_EXTENDED), + _effective(options, DLF_ACCESS_TRACKING_EXTENDED_INFO)] + parts = [p for p in parts if p] + return " ".join(parts) if parts else None + + +def _effective(options: Options, *keys: str) -> Optional[str]: + """The first non-blank value among ``keys``, stripped.""" + data = options.to_map() + for key in keys: + value = data.get(key) + if value is not None and str(value).strip(): + return str(value).strip() + return None + + +def _split(value) -> List[str]: + return str(value).split() if value is not None else [] diff --git a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py index 362994fb65ef..ad5d3c8a4ea9 100644 --- a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py +++ b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py @@ -35,6 +35,7 @@ from pypaimon.common.options.config import OssOptions, S3Options, SecurityOptions from pypaimon.common.options.options_utils import OptionsUtils from pypaimon.common.uri_reader import UriReaderFactory +from pypaimon.filesystem import oss_user_agent from pypaimon.filesystem.jindo_file_system_handler import JindoFileSystemHandler, JINDO_AVAILABLE from pypaimon.schema.data_types import (AtomicType, DataField, PyarrowFieldParser) @@ -198,6 +199,8 @@ def _initialize_oss_fs(self, path) -> FileSystem: # Uses setdefault so that an explicit user setting is never overridden. # Note: this is process-wide and affects all AWS SDK clients. os.environ.setdefault("AWS_EC2_METADATA_DISABLED", "true") + # S3FileSystem takes no User-Agent; this adds app/pypaimon/ to it, process-wide. + os.environ.setdefault("AWS_SDK_UA_APP_ID", oss_user_agent.identity()) client_kwargs = { "access_key": self.properties.get(OssOptions.OSS_ACCESS_KEY_ID), diff --git a/paimon-python/pypaimon/tests/jindo_file_system_test.py b/paimon-python/pypaimon/tests/jindo_file_system_test.py index 8c62113ca98a..4cfae7303b48 100644 --- a/paimon-python/pypaimon/tests/jindo_file_system_test.py +++ b/paimon-python/pypaimon/tests/jindo_file_system_test.py @@ -27,6 +27,7 @@ from pypaimon.common.options import Options from pypaimon.common.options.config import OssOptions from pypaimon.filesystem import jindo_file_system_handler as jindo_module +from pypaimon.filesystem import oss_user_agent from pypaimon.filesystem.jindo_file_system_handler import ( JindoFileSystemHandler, JindoInputFile, @@ -141,7 +142,10 @@ def test_forwards_native_options_to_connect(self): self.assertEqual(config.values["logger.dir"], "/tmp/jindo-log") self.assertEqual(config.values["logger.verbose"], "3") self.assertEqual(config.values["logger.console.log.enable"], "false") - self.assertEqual(config.values["fs.oss.user.agent.features"], "pypaimon") + self.assertEqual( + config.values["fs.oss.user.agent.features"], oss_user_agent.identity()) + self.assertNotIn("fs.oss.user.agent.module", config.values) + self.assertNotIn("fs.oss.user.agent.extended", config.values) self.assertNotIn(OssOptions.OSS_IMPL.key(), config.values) self.assertNotIn("fs.oss.unset.option", config.values) self.assertNotIn("metastore", config.values) @@ -172,6 +176,55 @@ def test_does_not_load_external_config(self): self.assertNotIn("fs.oss.provider.endpoint", config.values) self.assertNotIn("fs.oss.provider.format", config.values) + def test_user_agent_keeps_user_values_and_appends_access_tracking(self): + fake_jutil = types.SimpleNamespace(Config=_RecordingConfig) + options = Options({ + "fs.oss.user.agent.module": "MyModule/1.0", + "fs.oss.user.agent.features": "Flink", + "fs.oss.user.agent.extended": "user/ext", + "dlf.access-tracking.extended-info": "acs/xxx k/v", + }) + + with mock.patch.object(jindo_module, "JINDO_AVAILABLE", True), \ + mock.patch.object(jindo_module, "jutil", fake_jutil), \ + mock.patch.object(oss_user_agent, "identity", return_value="pypaimon/2.2.dev0"): + config = jindo_module.build_jindo_config(options) + + self.assertEqual(config.values["fs.oss.user.agent.module"], "MyModule/1.0") + self.assertEqual( + config.values["fs.oss.user.agent.features"], "pypaimon/2.2.dev0 Flink") + self.assertEqual( + config.values["fs.oss.user.agent.extended"], "user/ext acs/xxx k/v") + self.assertNotIn("dlf.access-tracking.extended-info", config.values) + + def test_user_agent_from_common_keys(self): + fake_jutil = types.SimpleNamespace(Config=_RecordingConfig) + options = Options({ + "user-agent.module": "MyApp/1.0", + "user-agent.features": "Flink", + "user-agent.extended": "vvr", + "dlf.access-tracking.extended-info": "uid/123", + }) + + with mock.patch.object(jindo_module, "JINDO_AVAILABLE", True), \ + mock.patch.object(jindo_module, "jutil", fake_jutil), \ + mock.patch.object(oss_user_agent, "identity", return_value="pypaimon/2.2.dev0"): + config = jindo_module.build_jindo_config(options) + + self.assertEqual(config.values["fs.oss.user.agent.module"], "MyApp/1.0") + self.assertEqual(config.values["fs.oss.user.agent.features"], "pypaimon/2.2.dev0 Flink") + self.assertEqual(config.values["fs.oss.user.agent.extended"], "vvr uid/123") + + def test_user_agent_extended_from_access_tracking_only(self): + fake_jutil = types.SimpleNamespace(Config=_RecordingConfig) + options = Options({"dlf.access-tracking.extended-info": "acs/xxx"}) + + with mock.patch.object(jindo_module, "JINDO_AVAILABLE", True), \ + mock.patch.object(jindo_module, "jutil", fake_jutil): + config = jindo_module.build_jindo_config(options) + + self.assertEqual(config.values["fs.oss.user.agent.extended"], "acs/xxx") + def test_forwards_native_options_to_jindo_oss_filesystem(self): created_config = _RecordingConfig() config_factory = mock.Mock(return_value=created_config) diff --git a/paimon-python/pypaimon/tests/oss_user_agent_test.py b/paimon-python/pypaimon/tests/oss_user_agent_test.py new file mode 100644 index 000000000000..dee7fd822ce3 --- /dev/null +++ b/paimon-python/pypaimon/tests/oss_user_agent_test.py @@ -0,0 +1,150 @@ +# 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. + +import os +import unittest +from unittest import mock + +from pypaimon import build_info +from pypaimon.common.options import Options +from pypaimon.common.options.config import OssOptions +from pypaimon.filesystem import oss_user_agent +from pypaimon.filesystem import pyarrow_file_io +from pypaimon.filesystem.pyarrow_file_io import PyArrowFileIO + +IDENTITY = "pypaimon/2.2.dev0" + + +class OssUserAgentIdentityTest(unittest.TestCase): + + def setUp(self): + oss_user_agent.identity.cache_clear() + self.addCleanup(oss_user_agent.identity.cache_clear) + + def test_identity_uses_build_version(self): + with mock.patch.object(build_info, "version", return_value="2.2.dev"): + self.assertEqual("pypaimon/2.2.dev", oss_user_agent.identity()) + + def test_identity_without_build_version(self): + with mock.patch.object(build_info, "version", return_value=None): + self.assertEqual("pypaimon", oss_user_agent.identity()) + + def test_build_version_from_full_version(self): + for full_version, expected in (("python-2.2.dev-abc123", "2.2.dev"), + ("python-2.2.0-UNKNOWN", "2.2.0"), + ("UNKNOWN", None)): + with mock.patch.object(build_info, "_FULL_VERSION", full_version): + self.assertEqual(expected, build_info.version()) + + +class OssUserAgentTest(unittest.TestCase): + + def setUp(self): + patcher = mock.patch.object(oss_user_agent, "identity", return_value=IDENTITY) + patcher.start() + self.addCleanup(patcher.stop) + + def test_features_default_to_identity(self): + self.assertEqual(IDENTITY, oss_user_agent.features(Options({}))) + + def test_user_features_follow_identity(self): + options = Options({oss_user_agent.USER_AGENT_FEATURES: " Flink Spark "}) + self.assertEqual(IDENTITY + " Flink Spark", oss_user_agent.features(options)) + + def test_identity_is_not_duplicated(self): + for user_features in ("pypaimon Flink", "Flink pypaimon/1.0"): + options = Options({oss_user_agent.USER_AGENT_FEATURES: user_features}) + self.assertEqual(user_features, oss_user_agent.features(options)) + + def test_access_tracking_is_appended_to_user_extended(self): + options = Options({ + oss_user_agent.USER_AGENT_EXTENDED: "user/ext", + oss_user_agent.DLF_ACCESS_TRACKING_EXTENDED_INFO: "acs/xxx k/v", + }) + self.assertEqual("user/ext acs/xxx k/v", oss_user_agent.extended(options)) + + def test_extended_from_either_side_alone(self): + self.assertEqual("acs/xxx", oss_user_agent.extended( + Options({oss_user_agent.DLF_ACCESS_TRACKING_EXTENDED_INFO: "acs/xxx"}))) + self.assertEqual("user/ext", oss_user_agent.extended( + Options({oss_user_agent.USER_AGENT_EXTENDED: "user/ext"}))) + + def test_common_keys_are_used_without_oss_keys(self): + options = Options({ + oss_user_agent.COMMON_USER_AGENT_MODULE: "MyApp/1.0", + oss_user_agent.COMMON_USER_AGENT_FEATURES: "Flink", + oss_user_agent.COMMON_USER_AGENT_EXTENDED: "vvr", + oss_user_agent.DLF_ACCESS_TRACKING_EXTENDED_INFO: "uid/123", + }) + self.assertEqual("MyApp/1.0", oss_user_agent.module(options)) + self.assertEqual(IDENTITY + " Flink", oss_user_agent.features(options)) + self.assertEqual("vvr uid/123", oss_user_agent.extended(options)) + + def test_oss_keys_override_common_keys_per_part(self): + options = Options({ + oss_user_agent.COMMON_USER_AGENT_MODULE: "MyApp/1.0", + oss_user_agent.COMMON_USER_AGENT_FEATURES: "Flink", + oss_user_agent.COMMON_USER_AGENT_EXTENDED: "vvr", + oss_user_agent.USER_AGENT_FEATURES: "Spark", + oss_user_agent.USER_AGENT_EXTENDED: " ", + }) + self.assertEqual("MyApp/1.0", oss_user_agent.module(options)) + self.assertEqual(IDENTITY + " Spark", oss_user_agent.features(options)) + self.assertEqual("vvr", oss_user_agent.extended(options)) + options = Options({oss_user_agent.COMMON_USER_AGENT_EXTENDED: "vvr", + oss_user_agent.USER_AGENT_EXTENDED: "oss/ext"}) + self.assertEqual("oss/ext", oss_user_agent.extended(options)) + self.assertIsNone(oss_user_agent.module(Options({}))) + + def test_blank_extended_is_ignored(self): + options = Options({ + oss_user_agent.USER_AGENT_EXTENDED: " ", + oss_user_agent.DLF_ACCESS_TRACKING_EXTENDED_INFO: "", + }) + self.assertIsNone(oss_user_agent.extended(options)) + self.assertIsNone(oss_user_agent.extended(Options({}))) + + +class LegacyOssUserAgentTest(unittest.TestCase): + + def _new_legacy_file_io(self): + options = Options({ + OssOptions.OSS_ACCESS_KEY_ID.key(): "ak", + OssOptions.OSS_ACCESS_KEY_SECRET.key(): "sk", + OssOptions.OSS_ENDPOINT.key(): "oss-cn-test.example.com", + OssOptions.OSS_REGION.key(): "cn-test", + OssOptions.OSS_IMPL.key(): "legacy", + }) + with mock.patch.object(pyarrow_file_io.pafs, "S3FileSystem") as s3_filesystem, \ + mock.patch.object(oss_user_agent, "identity", return_value=IDENTITY): + PyArrowFileIO("oss://test-bucket/", options) + s3_filesystem.assert_called_once() + + def test_sets_aws_app_id_when_absent(self): + with mock.patch.dict(os.environ): + os.environ.pop("AWS_SDK_UA_APP_ID", None) + self._new_legacy_file_io() + self.assertEqual(IDENTITY, os.environ["AWS_SDK_UA_APP_ID"]) + + def test_keeps_user_aws_app_id(self): + with mock.patch.dict(os.environ, {"AWS_SDK_UA_APP_ID": "my-app"}): + self._new_legacy_file_io() + self.assertEqual("my-app", os.environ["AWS_SDK_UA_APP_ID"]) + + +if __name__ == "__main__": + unittest.main()