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
3 changes: 3 additions & 0 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -401,6 +401,9 @@ upload_url = core.generate_presigned_url(path, "put_object", ContentType="text/c
`create_bucket()` and `delete_bucket()` create and delete a bucket when called. The
`allow_bucket_creation` and `allow_bucket_deletion` options apply only to the
filesystem's `mkdir`/`makedirs` and `rmdir`, not to calls through the core.
`get_bucket_versioning()` returns the versioning state of a bucket, `Enabled` or
`Suspended`, or `None` if versioning has never been enabled on it.
`S3Path.is_directory_bucket` tells whether the bucket of a path is a directory bucket.

`create_bucket()` sends a `LocationConstraint` for the `region_name` argument, or for
the client's region by default, except in `us-east-1`. The fields of a
Expand Down
9 changes: 2 additions & 7 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -1438,15 +1438,10 @@ def _move_pairs(
buckets = {
p.bucket
for p in paths.values()
if p.version_id == "null"
and p.name in unversioned
and not self.core._is_directory_bucket(p.bucket)
if p.version_id == "null" and p.name in unversioned and not p.is_directory_bucket
}
versioned_buckets = {
bucket
for bucket in buckets
if self._call(self._client.get_bucket_versioning, Bucket=bucket).get("Status")
== "Enabled"
bucket for bucket in buckets if self.core.get_bucket_versioning(bucket) == "Enabled"

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): CLEAN

  • Base 47fb6f31749f505a723c314b03ebf1ca3c5b8114, head e12d5d584c0f5b6f6ff7944585ca66941be18e1a.
  • Requests: _move_pairs() sends the same GetBucketVersioning, once per distinct candidate bucket.
    • Before: self._call(self._client.get_bucket_versioning, Bucket=...), which delegated to S3Core.call().
    • Now: S3Core.get_bucket_versioning(), which sends self.call(self._client.get_bucket_versioning, Bucket=..., **params).
    • The client, the retry policy, the filtering of inherited request_kwargs and the error translation are therefore unchanged.
  • Result: the comparison with "Enabled" is unchanged; None and "Suspended" still mean the bucket is not versioning-enabled.
  • Directory buckets: S3Path.is_directory_bucket applies the same --x-s3 suffix rule to the same bucket name at all three former call sites. These are the adapter's candidate filter and the two checks in plan_multipart_copy().
  • Tests: the aio mv tests now set _sync_fs._call = _sync_fs._core.call, as the other tests do.
    • With the request sent by the core, the old wiring made test_mv_null_version_onto_key fail.
    • It also made the directory-bucket test pass without observing the request.
  • Callers: no other code used S3Core._is_directory_bucket() (repository grep).
  • Limitation: offline only at this head. A live run is in progress.

}
missing = {
source
Expand Down
35 changes: 19 additions & 16 deletions pyathena/filesystem/s3_core.py
Original file line number Diff line number Diff line change
Expand Up @@ -835,6 +835,23 @@ def delete_bucket(self, bucket: str, **params) -> None:
_logger.debug(f"Delete bucket: s3://{bucket}")
self.call(self._client.delete_bucket, Bucket=bucket, **params)

def get_bucket_versioning(self, bucket: str, **params) -> str | None:
"""Get the versioning state of a bucket with GetBucketVersioning.

Args:
bucket: The name of the bucket.
**params: Additional request parameters, sent as given.

Returns:
``Enabled`` or ``Suspended``, or None if versioning has never
been enabled on the bucket.

Raises:
FileNotFoundError: If the bucket does not exist.
"""
response = self.call(self._client.get_bucket_versioning, Bucket=bucket, **params)
return cast(str | None, response.get("Status"))

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 (claims, callers, operations): CLEAN

  • Base 47fb6f31749f505a723c314b03ebf1ca3c5b8114, head e12d5d584c0f5b6f6ff7944585ca66941be18e1a.
  • I checked these claims:
    • "None for a bucket whose versioning has never been enabled": GetBucketVersioning omits Status in that state, according to the botocore model docs and the S3 API reference. Stubber pins the empty response.
    • "A missing bucket raises FileNotFoundError": S3ClientError maps NoSuchBucket to it, and Stubber pins this.
    • "ExpectedBucketOwner is passed through and RequestPayer is filtered": the operation's input shape has only Bucket and ExpectedBucketOwner, measured with botocore 1.43.102.
    • "No change of requests": the existing sync mv tests assert fs._call.assert_called_once_with(fs._client.get_bucket_versioning, Bucket="bucket") through the shared core mock, and they pass unchanged.
  • The docs sentence matches the docstring.
  • Breaking-change check: the private staticmethod is removed. It was underscore-private and only the core used it, so no release note is needed.


def delete_object(self, path: S3Path, **params) -> None:
"""Delete an object, or a version of it, with DeleteObject.

Expand Down Expand Up @@ -1266,7 +1283,7 @@ def plan_multipart_copy(
request.pop("Tagging", None)
# Directory buckets do not support GetObjectTagging, and their
# objects have no tags.
if not self._is_directory_bucket(source.bucket):
if not source.is_directory_bucket:
tags = self.get_object_tagging(
source, **self.operation_params("get_object_tagging", source_params)
)
Expand All @@ -1290,7 +1307,7 @@ def plan_multipart_copy(
)
if annotation_directive == "COPY"
and "CopySourceSSECustomerAlgorithm" not in params
and not self._is_directory_bucket(source.bucket)
and not source.is_directory_bucket
else ()
)
return S3MultipartCopyPlan(
Expand Down Expand Up @@ -1604,20 +1621,6 @@ def generate_presigned_url(
),
)

@staticmethod
def _is_directory_bucket(bucket: str) -> bool:
"""Return whether the bucket is a directory bucket (S3 Express One Zone).

Directory bucket names end with ``--x-s3``.

Args:
bucket: S3 bucket name.

Returns:
True if the bucket is a directory bucket.
"""
return bucket.endswith("--x-s3")

@staticmethod
def _copy_source_params(params: Mapping[str, Any]) -> dict[str, Any]:
"""Map the parameters of a copy to those of the requests that read its source.
Expand Down
8 changes: 8 additions & 0 deletions pyathena/filesystem/s3_path.py
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,14 @@ def is_bucket(self) -> bool:
"""Whether the path names the bucket: it has no key, or a key of only slashes."""
return not self.key or not self.key.strip("/")

@property
def is_directory_bucket(self) -> bool:

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 and /code-review at e12d5d584c0f5b6f6ff7944585ca66941be18e1a (base 47fb6f31749f505a723c314b03ebf1ca3c5b8114, detached snapshot, static only)

  • Codex (codex exec -s read-only): CLEAN. It covered the versioning request parameters, the status returns and error translation, the sync/aio move behavior and lookup counts, directory-bucket detection, the multipart-copy guards, the tests and the docs.
  • claude-fable-5-1 (max profile, effort high): CLEAN. It confirmed that the request, the request count, the retry policy, the error translation and the == "Enabled" result are unchanged. It also confirmed that the four aio tests are the only ones that mocked _call alone.
  • /code-review: no regression. Nine design, test-precision and docs notes, none changed:
    • The name is_directory_bucket is kept. "Directory bucket" is the AWS term, and the docstring already says it concerns the bucket of the path. The maintainer decided this.
    • Docstring and docs additions about directory-bucket support and paragraph placement were judged unnecessary and reverted. Other core operations do not list S3 feature support either.
    • The str | None return was agreed on Finish the public S3Core API extraction from filesystem adapters #1086.
    • Not changed because they are pre-existing behavior or out of scope: sequential GetBucketVersioning per bucket, --xa-s3 access point aliases, and test assertions through the shared _call alias.
  • Live AWS at e12d5d584c0f5b6f6ff7944585ca66941be18e1a: pytest -n 8 tests/pyathena/filesystem/ tests/pyathena/s3fs/ tests/pyathena/aio/s3fs/ → 1235 passed, 10 skipped.

"""Whether the bucket is a directory bucket (S3 Express One Zone).

The names of directory buckets end with ``--x-s3``.
"""
return self.bucket.endswith("--x-s3")

@property
def name(self) -> str:
"""The path without a scheme or version, in ``bucket/key`` form, or the bucket."""
Expand Down
10 changes: 6 additions & 4 deletions tests/pyathena/filesystem/test_s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -1036,7 +1036,7 @@ async def copy_file(path1, path2, **kwargs):
@pytest.mark.asyncio
async def test_mv_null_version_onto_key(self, status):
fs = AioS3FileSystem(connection=mock.MagicMock(), skip_instance_cache=True)
fs._sync_fs._call = mock.MagicMock(
fs._sync_fs._call = fs._sync_fs._core.call = mock.MagicMock(
return_value={} if status is None else {"Status": status}
)
fs._copy_file = mock.AsyncMock(return_value=True)
Expand All @@ -1063,7 +1063,7 @@ async def test_mv_null_version_onto_key(self, status):
@pytest.mark.asyncio
async def test_mv_null_version_conflicts_depend_on_bucket_state(self, status):
fs = AioS3FileSystem(connection=mock.MagicMock(), skip_instance_cache=True)
fs._sync_fs._call = mock.MagicMock(
fs._sync_fs._call = fs._sync_fs._core.call = mock.MagicMock(
return_value={} if status is None else {"Status": status}
)
fs._copy_file = mock.AsyncMock(return_value=True)
Expand All @@ -1090,7 +1090,7 @@ async def test_mv_null_version_conflicts_depend_on_bucket_state(self, status):
@pytest.mark.asyncio
async def test_mv_null_version_directory_bucket_does_not_read_bucket_state(self):
fs = AioS3FileSystem(connection=mock.MagicMock(), skip_instance_cache=True)
fs._sync_fs._call = mock.MagicMock()
fs._sync_fs._call = fs._sync_fs._core.call = mock.MagicMock()
fs._copy_file = mock.AsyncMock(return_value=True)
fs._delete_objects = mock.AsyncMock()
key = "s3://example--usw2-az1--x-s3/key"
Expand All @@ -1105,7 +1105,9 @@ async def test_mv_null_version_directory_bucket_does_not_read_bucket_state(self)
@pytest.mark.asyncio
async def test_mv_null_version_failure_does_not_delete(self, stage):
fs = AioS3FileSystem(connection=mock.MagicMock(), skip_instance_cache=True)
fs._sync_fs._call = mock.MagicMock(return_value={"Status": "Enabled"})
fs._sync_fs._call = fs._sync_fs._core.call = mock.MagicMock(
return_value={"Status": "Enabled"}
)
fs._copy_file = mock.AsyncMock(return_value=True)
fs._delete_objects = mock.AsyncMock()
if stage == "lookup":
Expand Down
30 changes: 30 additions & 0 deletions tests/pyathena/filesystem/test_s3_core.py
Original file line number Diff line number Diff line change
Expand Up @@ -425,6 +425,36 @@ def test_delete_bucket(self):
core.delete_bucket("bucket")
stubber.assert_no_pending_responses()

@pytest.mark.parametrize(
("response", "expected"),
[
({"Status": "Enabled"}, "Enabled"),
({"Status": "Suspended", "MFADelete": "Disabled"}, "Suspended"),
# A bucket whose versioning has never been enabled.
({}, None),
],
)
def test_get_bucket_versioning(self, response, expected):
core, stubber = _make_core(request_kwargs={"RequestPayer": "requester"})
# RequestPayer, which GetBucketVersioning does not accept, is not sent.
stubber.add_response(
"get_bucket_versioning",
response,
{"Bucket": "bucket", "ExpectedBucketOwner": "123456789012"},
)
with stubber:
status = core.get_bucket_versioning("bucket", ExpectedBucketOwner="123456789012")
stubber.assert_no_pending_responses()
assert status == expected

def test_get_bucket_versioning_translates_errors(self):
core, stubber = _make_core()
stubber.add_client_error(
"get_bucket_versioning", service_error_code="NoSuchBucket", http_status_code=404
)
with stubber, pytest.raises(FileNotFoundError):
core.get_bucket_versioning("bucket")

def test_list_objects(self):
core, stubber = _make_core()
stubber.add_response(
Expand Down
12 changes: 12 additions & 0 deletions tests/pyathena/filesystem/test_s3_path.py
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,18 @@ def test_parse_invalid(self, path):
def test_is_bucket(self, path, expected):
assert path.is_bucket is expected

@pytest.mark.parametrize(
("path", "expected"),
[
(S3Path("bucket--usw2-az1--x-s3", "key"), True),
(S3Path("bucket--usw2-az1--x-s3"), True),
(S3Path("bucket", "key--x-s3"), False),
(S3Path("bucket--x-s3-other"), False),
],
)
def test_is_directory_bucket(self, path, expected):
assert path.is_directory_bucket is expected

@pytest.mark.parametrize(
("path", "name", "string", "uri"),
[
Expand Down
Loading