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
2 changes: 1 addition & 1 deletion docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -156,7 +156,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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round two (documentation reader): the example on the next line uses a trailing slash, so it aborts the uploads under path/to/, as the new comment says. The other examples in this file (list_multipart_uploads on a bucket at line 127, object_version_info on an object path at line 155, ls(..., versions=True)) are unaffected by this change, and no other doc or README references these methods.

fs.clear_multipart_uploads("s3://YOUR_S3_BUCKET/path/to/")
```

Expand Down
57 changes: 43 additions & 14 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,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
Expand Down Expand Up @@ -2100,15 +2101,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] = []
Expand All @@ -2127,7 +2133,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)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round one (implementation behavior) — base aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe, head 3872f3151c34c1032739f9f3aa89590d43dad052. Result: CLEAN.

Covered: pyathena/filesystem/s3.py (list_multipart_uploads, clear_multipart_uploads, object_version_info), pyathena/filesystem/s3_async.py delegates, docs/filesystem.md, tests/pyathena/filesystem/test_s3.py.

  • Multipart filter (this line): prefix is "" for a bucket path, so startswith("") keeps every upload, and u["Key"] == None is never true; the redundant not key guard was dropped before the push. Pagination still follows IsTruncated/markers of the raw response, so a page whose uploads are all filtered out does not stop the loop. ListMultipartUploads always returns Key per upload.
  • Trailing slash (data/): Prefix="data/" is sent unchanged and every returned key starts with data/, so the filter keeps all of them; data itself is not included.
  • clear_multipart_uploads() aborts exactly the filtered list, so the sibling uploads are no longer touched. The aio methods call the sync ones (s3_async.py:484-520), so no separate async path exists.
  • Request count is unchanged: one paginated listing with the same Prefix; existing pagination tests still assert the same request arguments.

)
if not response.get("IsTruncated"):
break
Expand All @@ -2140,7 +2148,15 @@ 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 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
Expand All @@ -2153,6 +2169,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] = []
Expand All @@ -2161,20 +2180,30 @@ 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", [])
)
return versions
# 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"

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed) — reviewer: Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max (different model from the Claude author), session 01a100cc-e27a-78a1-b1ba-1ad4e70f68df. Static review only: codex exec --sandbox read-only on a detached snapshot at 3872f3151c34c1032739f9f3aa89590d43dad052 (no .env), base aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe; prompt had the literal diff and intended behavior, no PR text or prior findings. Snapshot and PR worktree verified unchanged afterwards. Result: FINDINGS.

Covered: exact diff, path parsing, trailing slashes, bucket paths, pagination, delete markers, empty results, result models, async delegation, docs, test fakes, live-resource cleanup.

  1. P2 — explicit URL encoding drops valid versions (pyathena/filesystem/s3.py:1941 at the reviewed head): object_version_info("s3://bucket/a b", EncodingType="url") gets Key="a+b"; botocore decodes only when it set EncodingType itself, so both comparisons reject the key and return [] (old code returned it).
    Verified and fixed in 6a92aa09 (this line): keys are decoded with unquote_plus for matching only, returned as listed. Verified botocore's _decode_list_object condition (encoding_type_auto_set) and live S3: the raw key for a b is a+b, and the fixed call now returns only a b's version, not a b.bak. New offline test test_object_version_info_matches_url_encoded_keys fails without the fix.

Non-actionable observations from the reviewer: the fakes apply S3's plain-prefix matching and distinguish old/new behavior; marker-only precedence differs between delete_markers=False/True (same corner recorded in self-review round one, left as documented); the multipart live test's pre-existing cleanup only covers the sibling in finally (the bucket's AbortIncompleteMultipartUpload lifecycle rule, 1 day, covers a leaked upload).

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up (relayed) — Codex CLI 0.160.0, model gpt-6-astra, session 01a100d5-6dac-70b2-af54-30841cf57a86; static, read-only snapshot at 6a92aa090dc1455185c7a19334683b0f809fac3e, bounded to the repair 3872f3151c34c1032739f9f3aa89590d43dad052..6a92aa090dc1455185c7a19334683b0f809fac3e (single commit, base unchanged). Result: CLEAN — prior finding resolved. Covered: decoding only for explicit EncodingType="url" with defaults unchanged; botocore 1.43.102 _decode_list_object decodes only auto-set encoding via unquote_str = unquote_plus, so the repair matches it; returned objects, encoded keys and order preserved, including delete markers; the new test receives [] before the repair; the aio delegate forwards kwargs unchanged. Snapshot verified clean afterwards.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up (relayed), unified delete-marker rule — Codex CLI 0.160.0, model gpt-6-astra, session 01a1010c-4834-7a60-9845-28dfef5d761f; static, read-only snapshot at 3867c355146304ef5a35b50db367cd028eca2096, bounded to 6a92aa090dc1455185c7a19334683b0f809fac3e..3867c355146304ef5a35b50db367cd028eca2096 (one commit, base unchanged), prompt with the requirement and diff only. Result: CLEAN — inconsistency resolved. Covered: selection independent of delete_markers and docstring accuracy; exact-key/prefix, bucket, trailing-slash, sibling exclusion and EncodingType="url" behavior unchanged; one paginated listing with order preserved; the new test's False case returns [("dir/x", "v1", False)] on the old code; the aio delegate forwards path, flag and kwargs unchanged. Snapshot verified clean afterwards.

keys = [unquote_plus(v.key) if url_encoded else v.key for v in versions]
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.

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:
Expand Down
6 changes: 3 additions & 3 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -570,7 +570,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`.

Expand All @@ -590,7 +590,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.
Expand All @@ -603,7 +603,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)

Expand Down
138 changes: 127 additions & 11 deletions tests/pyathena/filesystem/test_s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -1497,6 +1497,80 @@ 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

@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.
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"):
Expand Down Expand Up @@ -1646,6 +1720,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"):
Expand Down Expand Up @@ -2617,24 +2724,33 @@ 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 = (
f"s3://{ENV.s3_staging_bucket}/{ENV.s3_staging_key}{ENV.schema}/"
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")

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed), finding 2 — P3, not changed: the reviewer noted that the new <uuid>.bak object is never removed. Rejected: the staging bucket's lifecycle rule (cloudformation/github_actions_oidc.yaml:329-338: ExpirationInDays: 1, NoncurrentVersionExpiration 1 day, AbortIncompleteMultipartUpload 1 day) expires it, which is what the existing path object of this test (and the other unique-path filesystem tests) already relies on; an explicit delete would add a request per run without changing the outcome.


versions = fs.object_version_info(path)
assert len(versions) == 1
Expand Down
Loading