Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions paimon-python/pypaimon/build_info.py
Original file line number Diff line number Diff line change
Expand Up @@ -84,3 +84,11 @@ def _load_full_version():
def full_version():
"""Return ``<pypaimon-version>-<commit-id>`` 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
10 changes: 9 additions & 1 deletion paimon-python/pypaimon/filesystem/jindo_file_system_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.")
Expand Down Expand Up @@ -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/<version> 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


Expand Down
77 changes: 77 additions & 0 deletions paimon-python/pypaimon/filesystem/oss_user_agent.py
Original file line number Diff line number Diff line change
@@ -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/<version>`` 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 []
3 changes: 3 additions & 0 deletions paimon-python/pypaimon/filesystem/pyarrow_file_io.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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/<version> 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),
Expand Down
55 changes: 54 additions & 1 deletion paimon-python/pypaimon/tests/jindo_file_system_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
150 changes: 150 additions & 0 deletions paimon-python/pypaimon/tests/oss_user_agent_test.py
Original file line number Diff line number Diff line change
@@ -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()
Loading