From 3872f3151c34c1032739f9f3aa89590d43dad052 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 17:02:52 +0900 Subject: [PATCH 1/3] Limit multipart uploads and object versions to the key and keys under it list_multipart_uploads(), clear_multipart_uploads() and object_version_info() sent the path key to S3 as a plain string Prefix, so they also listed, aborted or returned sibling keys that merely start with the same characters (data2/..., data.csv, a.csv.bak). The multipart methods now keep only the uploads to the key itself and to the keys under key/. object_version_info() returns the versions of the key itself when it has any, and otherwise those under key/; a path with a trailing slash always selects the keys under it. Closes #981 Co-Authored-By: Claude Opus 5.5 --- docs/filesystem.md | 2 +- pyathena/filesystem/s3.py | 40 ++++++++--- pyathena/filesystem/s3_async.py | 6 +- tests/pyathena/filesystem/test_s3.py | 102 ++++++++++++++++++++++++--- 4 files changed, 126 insertions(+), 24 deletions(-) diff --git a/docs/filesystem.md b/docs/filesystem.md index b19f93d43..88293969b 100644 --- a/docs/filesystem.md +++ b/docs/filesystem.md @@ -128,7 +128,7 @@ uploads = fs.list_multipart_uploads("s3://YOUR_S3_BUCKET") for upload in uploads: print(upload.key, upload.upload_id, upload.initiated) -# Abort all incomplete uploads under a bucket or key prefix. +# Abort all incomplete uploads to a key and the keys under it. fs.clear_multipart_uploads("s3://YOUR_S3_BUCKET/path/to/") ``` diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 3d039832c..dc74bd2f6 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -1854,15 +1854,20 @@ def list_multipart_uploads(self, path: str) -> list[S3MultipartUpload]: to abort all of them. Args: - path: S3 bucket or prefix path (e.g., "bucket", "s3://bucket" or - "s3://bucket/prefix"). If the path contains a key prefix, - only the uploads under that prefix are listed. + path: S3 bucket or key path (e.g., "bucket", "s3://bucket" or + "s3://bucket/prefix"). If the path contains a key, only the + uploads to that key and to the keys under ``key/`` are + listed, not those to sibling keys that merely start with the + same characters (e.g., ``prefix2/a``). Returns: List of S3MultipartUpload instances describing the in-progress multipart uploads. """ bucket, key, _ = self.parse_path(path) + # S3 matches Prefix as a plain string, so the uploads are filtered to + # the key itself and the keys under it. + prefix = f"{key.rstrip('/')}/" if key else "" _logger.debug(f"List multipart uploads: s3://{bucket}/{key}") uploads: list[S3MultipartUpload] = [] @@ -1881,7 +1886,9 @@ def list_multipart_uploads(self, path: str) -> list[S3MultipartUpload]: **request, ) uploads.extend( - S3MultipartUpload({**u, "Bucket": bucket}) for u in response.get("Uploads", []) + S3MultipartUpload({**u, "Bucket": bucket}) + for u in response.get("Uploads", []) + if u["Key"] == key or u["Key"].startswith(prefix) ) if not response.get("IsTruncated"): break @@ -1894,7 +1901,14 @@ def list_multipart_uploads(self, path: str) -> list[S3MultipartUpload]: def object_version_info( self, path: str, delete_markers: bool = False, **kwargs ) -> list[S3ObjectVersion]: - """List the versions of the objects under the path. + """List the versions of the object or of the objects under the path. + + A key path without a trailing slash returns the versions of that key + if it has any, and otherwise the versions of the keys under ``key/``. + A key path with a trailing slash returns the versions of the keys + under it, and a bucket path returns the versions of all the keys in + the bucket. Sibling keys that merely start with the same characters + (e.g., ``key.bak``) are never included. Args: path: S3 path (s3://bucket/key or a key prefix) to list the @@ -1907,6 +1921,9 @@ def object_version_info( List of S3ObjectVersion instances describing the versions. """ bucket, key, _ = self.parse_path(path) + # S3 matches Prefix as a plain string, so the versions are filtered to + # the key itself or the keys under it. + prefix = f"{key.rstrip('/')}/" if key else "" _logger.debug(f"List object versions: s3://{bucket}/{key}") versions: list[S3ObjectVersion] = [] @@ -1920,15 +1937,20 @@ def object_version_info( S3ObjectVersion(bucket=bucket, is_delete_marker=True, response=m) for m in response.get("DeleteMarkers", []) ) - return versions + if key and not key.endswith("/"): + object_versions = [v for v in versions if v.key == key] + if object_versions: + return object_versions + return [v for v in versions if v.key.startswith(prefix)] def clear_multipart_uploads(self, path: str) -> None: """Abort any incomplete multipart uploads in the bucket. Args: - path: S3 bucket or prefix path (e.g., "bucket", "s3://bucket" or - "s3://bucket/prefix"). If the path contains a key prefix, - only the uploads under that prefix are aborted. + path: S3 bucket or key path (e.g., "bucket", "s3://bucket" or + "s3://bucket/prefix"). If the path contains a key, only the + uploads to that key and to the keys under ``key/`` are + aborted, as listed by :meth:`list_multipart_uploads`. """ uploads = self.list_multipart_uploads(path) if not uploads: diff --git a/pyathena/filesystem/s3_async.py b/pyathena/filesystem/s3_async.py index 77deccc2a..1742bf1f3 100644 --- a/pyathena/filesystem/s3_async.py +++ b/pyathena/filesystem/s3_async.py @@ -484,7 +484,7 @@ def chmod(self, path: str, acl: str, recursive: bool = False, **kwargs) -> None: def object_version_info( self, path: str, delete_markers: bool = False, **kwargs ) -> list[S3ObjectVersion]: - """List the versions of the objects under the path. + """List the versions of the object or of the objects under the path. See :meth:`S3FileSystem.object_version_info`. @@ -504,7 +504,7 @@ def list_multipart_uploads(self, path: str) -> list[S3MultipartUpload]: See :meth:`S3FileSystem.list_multipart_uploads`. Args: - path: S3 bucket or prefix path (e.g., "s3://bucket" or "s3://bucket/prefix"). + path: S3 bucket or key path (e.g., "s3://bucket" or "s3://bucket/prefix"). Returns: List of S3MultipartUpload instances describing the uploads. @@ -517,7 +517,7 @@ def clear_multipart_uploads(self, path: str) -> None: See :meth:`S3FileSystem.clear_multipart_uploads`. Args: - path: S3 bucket or prefix path (e.g., "s3://bucket" or "s3://bucket/prefix"). + path: S3 bucket or key path (e.g., "s3://bucket" or "s3://bucket/prefix"). """ self._sync_fs.clear_multipart_uploads(path) diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 2131eb1e3..9e8d0c300 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -786,6 +786,44 @@ def test_object_version_info_with_delete_markers(self): ("m1", True), ] + @pytest.mark.parametrize( + ("path", "expected"), + [ + # The key itself wins over the keys under it. + ("s3://bucket/a.csv", ["a.csv"]), + # Without an object, the keys under the path are returned. + ("s3://bucket/dir", ["dir/", "dir/x"]), + ("s3://bucket/dir/", ["dir/", "dir/x"]), + # A trailing slash selects the keys under the path even if the + # key without it exists. + ("s3://bucket/a.csv/", ["a.csv/x"]), + ("s3://bucket", ["a.csv", "a.csv.bak", "a.csv/x", "dir/", "dir/x", "dir2/y"]), + ], + ) + def test_object_version_info_excludes_sibling_keys(self, path, expected): + fs = self._make_fs() + keys = ["a.csv", "a.csv.bak", "a.csv/x", "dir/", "dir/x", "dir2/y"] + # S3 matches Prefix as a plain string prefix. + fs._call.side_effect = lambda _, **request: { + "Versions": [ + {"Key": k, "VersionId": f"v-{k}", "IsLatest": True} + for k in keys + if k.startswith(request["Prefix"]) + ], + "DeleteMarkers": [ + {"Key": k, "VersionId": f"m-{k}", "IsLatest": False} + for k in keys + if k.startswith(request["Prefix"]) + ], + "IsTruncated": False, + } + + actual = fs.object_version_info(path) + assert [v.key for v in actual] == expected + actual = fs.object_version_info(path, delete_markers=True) + assert sorted(v.key for v in actual if not v.is_delete_marker) == expected + assert sorted(v.key for v in actual if v.is_delete_marker) == expected + def test_ls_versions_requires_version_aware(self): fs = self._make_fs() with pytest.raises(ValueError, match="version aware"): @@ -908,6 +946,39 @@ def test_list_multipart_uploads_paginates(self): UploadIdMarker="upload1", ) + @pytest.mark.parametrize( + ("path", "expected"), + [ + ("s3://bucket/data", ["data", "data/part.csv"]), + ("s3://bucket/data/", ["data/part.csv"]), + ("s3://bucket", ["data", "data.csv", "data/part.csv", "data2/other.csv"]), + ], + ) + def test_list_and_clear_multipart_uploads_exclude_sibling_keys(self, path, expected): + fs = self._make_fs() + keys = ["data", "data.csv", "data/part.csv", "data2/other.csv"] + aborted = [] + + def call(method, **request): + if method is fs._client.abort_multipart_upload: + aborted.append(request["Key"]) + return {} + # S3 matches Prefix as a plain string prefix. + return { + "Uploads": [ + {"Key": k, "UploadId": f"u-{k}"} + for k in keys + if k.startswith(request.get("Prefix", "")) + ], + "IsTruncated": False, + } + + fs._call.side_effect = call + + assert [u.key for u in fs.list_multipart_uploads(path)] == expected + fs.clear_multipart_uploads(path) + assert sorted(aborted) == expected + @pytest.fixture(scope="class") def fs(self, request): if not hasattr(request, "param"): @@ -1879,17 +1950,24 @@ def test_list_and_clear_multipart_uploads(self, fs): prefix_path = f"s3://{bucket}/{prefix}" key = f"{prefix}/file" upload = fs._create_multipart_upload(bucket=bucket, key=key) - - uploads = fs.list_multipart_uploads(prefix_path) - listed = next((u for u in uploads if u.upload_id == upload.upload_id), None) - assert listed - assert listed.bucket == bucket - assert listed.key == key - assert listed.initiated - - fs.clear_multipart_uploads(prefix_path) - uploads = fs.list_multipart_uploads(prefix_path) - assert not any(u.upload_id == upload.upload_id for u in uploads) + # A sibling key that starts with the same characters as the prefix. + sibling = fs._create_multipart_upload(bucket=bucket, key=f"{prefix}2/file") + try: + uploads = fs.list_multipart_uploads(prefix_path) + listed = next((u for u in uploads if u.upload_id == upload.upload_id), None) + assert listed + assert listed.bucket == bucket + assert listed.key == key + assert listed.initiated + assert not any(u.upload_id == sibling.upload_id for u in uploads) + + fs.clear_multipart_uploads(prefix_path) + uploads = fs.list_multipart_uploads(prefix_path) + assert not any(u.upload_id == upload.upload_id for u in uploads) + uploads = fs.list_multipart_uploads(f"{prefix_path}2") + assert any(u.upload_id == sibling.upload_id for u in uploads) + finally: + fs.clear_multipart_uploads(f"{prefix_path}2") def test_object_version_info(self, fs): path = ( @@ -1897,6 +1975,8 @@ def test_object_version_info(self, fs): f"filesystem/test_object_version_info/{uuid.uuid4()}" ) fs.pipe(path, b"data") + # A sibling key that starts with the same characters as the path. + fs.pipe(f"{path}.bak", b"backup") versions = fs.object_version_info(path) assert len(versions) == 1 From 6a92aa090dc1455185c7a19334683b0f809fac3e Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 17:14:29 +0900 Subject: [PATCH 2/3] Match URL-encoded keys in object_version_info() With an explicit EncodingType="url", botocore leaves the listed keys URL-encoded, so the key filter compared "a+b" with "a b" and dropped the versions of the requested key. Decode the keys for matching while returning them as listed. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 9 +++++++-- tests/pyathena/filesystem/test_s3.py | 14 ++++++++++++++ 2 files changed, 21 insertions(+), 2 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index dc74bd2f6..bb94708e2 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -13,6 +13,7 @@ from multiprocessing import cpu_count from re import Pattern from typing import Any, cast +from urllib.parse import unquote_plus import botocore.exceptions from boto3 import Session @@ -1937,11 +1938,15 @@ def object_version_info( S3ObjectVersion(bucket=bucket, is_delete_marker=True, response=m) for m in response.get("DeleteMarkers", []) ) + # botocore decodes the keys only when it sets EncodingType itself, so + # the keys of an explicit EncodingType="url" are decoded for matching. + url_encoded = kwargs.get("EncodingType") == "url" + keys = [unquote_plus(v.key) if url_encoded else v.key for v in versions] if key and not key.endswith("/"): - object_versions = [v for v in versions if v.key == key] + object_versions = [v for v, k in zip(versions, keys, strict=True) if k == key] if object_versions: return object_versions - return [v for v in versions if v.key.startswith(prefix)] + return [v for v, k in zip(versions, keys, strict=True) if k.startswith(prefix)] def clear_multipart_uploads(self, path: str) -> None: """Abort any incomplete multipart uploads in the bucket. diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 9e8d0c300..968c26867 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -824,6 +824,20 @@ def test_object_version_info_excludes_sibling_keys(self, path, expected): assert sorted(v.key for v in actual if not v.is_delete_marker) == expected assert sorted(v.key for v in actual if v.is_delete_marker) == expected + def test_object_version_info_matches_url_encoded_keys(self): + fs = self._make_fs() + # With an explicit EncodingType="url", botocore leaves the keys encoded. + fs._call.return_value = { + "Versions": [ + {"Key": "a+b", "VersionId": "v1", "IsLatest": True}, + {"Key": "a+b.bak", "VersionId": "v2", "IsLatest": True}, + ], + "IsTruncated": False, + } + + actual = fs.object_version_info("s3://bucket/a b", EncodingType="url") + assert [(v.key, v.version_id) for v in actual] == [("a+b", "v1")] + def test_ls_versions_requires_version_aware(self): fs = self._make_fs() with pytest.raises(ValueError, match="version aware"): From 3867c355146304ef5a35b50db367cd028eca2096 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 18:15:02 +0900 Subject: [PATCH 3/3] Choose the key in object_version_info() regardless of delete_markers A key that had only delete markers and was also a folder returned the markers with delete_markers=True but the versions under key/ without them. Decide on the key itself whenever it has any versions or delete markers, and drop the delete markers only after that choice, so both modes select the same key from the same single listing. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 34 +++++++++++++++------------- tests/pyathena/filesystem/test_s3.py | 22 ++++++++++++++++++ 2 files changed, 40 insertions(+), 16 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index bb94708e2..5d8596b2d 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -1904,12 +1904,13 @@ def object_version_info( ) -> list[S3ObjectVersion]: """List the versions of the object or of the objects under the path. - A key path without a trailing slash returns the versions of that key - if it has any, and otherwise the versions of the keys under ``key/``. - A key path with a trailing slash returns the versions of the keys - under it, and a bucket path returns the versions of all the keys in - the bucket. Sibling keys that merely start with the same characters - (e.g., ``key.bak``) are never included. + A key path without a trailing slash selects that key if it has any + versions or delete markers, and otherwise the keys under ``key/``. + The choice does not depend on ``delete_markers``, so a key that has + only delete markers yields no versions without them. A key path with + a trailing slash selects the keys under it, and a bucket path selects + all the keys in the bucket. Sibling keys that merely start with the + same characters (e.g., ``key.bak``) are never included. Args: path: S3 path (s3://bucket/key or a key prefix) to list the @@ -1933,20 +1934,21 @@ def object_version_info( S3ObjectVersion(bucket=bucket, is_delete_marker=False, response=v) for v in response.get("Versions", []) ) - if delete_markers: - versions.extend( - S3ObjectVersion(bucket=bucket, is_delete_marker=True, response=m) - for m in response.get("DeleteMarkers", []) - ) + # Delete markers are kept until the key is chosen, so that the + # choice is the same with and without them. + versions.extend( + S3ObjectVersion(bucket=bucket, is_delete_marker=True, response=m) + for m in response.get("DeleteMarkers", []) + ) # botocore decodes the keys only when it sets EncodingType itself, so # the keys of an explicit EncodingType="url" are decoded for matching. url_encoded = kwargs.get("EncodingType") == "url" keys = [unquote_plus(v.key) if url_encoded else v.key for v in versions] - if key and not key.endswith("/"): - object_versions = [v for v, k in zip(versions, keys, strict=True) if k == key] - if object_versions: - return object_versions - return [v for v, k in zip(versions, keys, strict=True) if k.startswith(prefix)] + if key and not key.endswith("/") and key in keys: + selected = [v for v, k in zip(versions, keys, strict=True) if k == key] + else: + selected = [v for v, k in zip(versions, keys, strict=True) if k.startswith(prefix)] + return [v for v in selected if delete_markers or not v.is_delete_marker] def clear_multipart_uploads(self, path: str) -> None: """Abort any incomplete multipart uploads in the bucket. diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 968c26867..49dbc4ae4 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -824,6 +824,28 @@ def test_object_version_info_excludes_sibling_keys(self, path, expected): assert sorted(v.key for v in actual if not v.is_delete_marker) == expected assert sorted(v.key for v in actual if v.is_delete_marker) == expected + @pytest.mark.parametrize( + ("delete_markers", "expected"), + [ + (False, []), + (True, [("dir", "m1", True)]), + ], + ) + def test_object_version_info_chooses_key_with_only_delete_markers( + self, delete_markers, expected + ): + fs = self._make_fs() + # "dir" is both a deleted object, with only a delete marker left, and + # a folder. + fs._call.return_value = { + "Versions": [{"Key": "dir/x", "VersionId": "v1", "IsLatest": True}], + "DeleteMarkers": [{"Key": "dir", "VersionId": "m1", "IsLatest": True}], + "IsTruncated": False, + } + + actual = fs.object_version_info("s3://bucket/dir", delete_markers=delete_markers) + assert [(v.key, v.version_id, v.is_delete_marker) for v in actual] == expected + def test_object_version_info_matches_url_encoded_keys(self): fs = self._make_fs() # With an explicit EncodingType="url", botocore leaves the keys encoded.