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
17 changes: 15 additions & 2 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -187,7 +187,12 @@ def _pipe_file_in_transaction(
Raises:
FileExistsError: If the mode is "create" and the path already
exists.
ValueError: If the data takes more than
``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
"""
block_size = kwargs.get("block_size") or self._sync_fs.default_block_size
# The size in bytes; the length of a memoryview counts its items.
self._sync_fs._check_multipart_upload_size(path, memoryview(value).nbytes, block_size)

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 (behavior and implementation) — base 92c9e3e67180eb3d52ef7b85579a24540a30fd64, head bab08235f0050aebdc4cf2d85a8fcba08b6dc59b. Result: CLEAN.

Covered: _pipe_file_in_transaction() and _put_file_in_transaction() against S3FileSystem.pipe_file()/put_file() (same check, same order: before the mode="create" existence check, before callback.set_size(), before open()); the block_size that pipe still passes to open() through kwargs; the error propagating through asyncio.to_thread and the generated sync wrappers; no file opened, so nothing joins transaction.files; non-transaction paths unchanged (they delegate to the sync methods). Tests: 3 new cases fail on the base and pass here; offline tests/pyathena/filesystem/ 262 passed vs 259 on the base, same 110 credential-less failures. Out of scope: the non-contiguous memoryview failure in the same with self.open(...) pattern (#997); max_workers routing (#969).

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, sandbox read-only, session 01a10115-dd19-7f23-90dc-25c82e2a475f. Base 92c9e3e67180eb3d52ef7b85579a24540a30fd64, head bab08235f0050aebdc4cf2d85a8fcba08b6dc59b, detached snapshot; given the diff, the intended behavior, and the fsspec source only. Static review only. Result: CLEAN.

Covered (reviewer): equivalence with the sync pipe_file()/put_file() (memoryview(value).nbytes, the default block size, the check before the mode="create" existence check and open(), block_size removed from the S3 parameters and passed to open(); directory/bucket early returns, MIME inference, progress, cache invalidation); errors through asyncio.to_thread, fsspec's wrapper generation, _runner, and sync, reaching async and generated sync callers; rejection before either file opens (fsspec registers transaction files only after _open() returns); docstrings and comments; the new tests fail against the base and cannot pass merely because open() is mocked.

Coverage limits noted by the reviewer (not defects): the forwarding test checks argument routing only (its 8-byte block works only with the mocked open()), and the tests omit multi-byte memoryviews, oversized data with mode="create", and an explicit block_size different from the default; the implementation handles these correctly. These paths share _check_multipart_upload_size(), which tests/pyathena/filesystem/test_s3.py covers for the sync methods, including the multi-byte memoryview.

if mode == "create" and self._sync_fs.exists(path):
raise FileExistsError(path)
with self.open(path, "wb", **kwargs) as f:
Expand All @@ -211,21 +216,29 @@ def _put_file_in_transaction(self, lpath: str, rpath: str, callback, **kwargs) -
rpath: S3 destination path (s3://bucket/key).
callback: Progress callback for tracking upload progress.
**kwargs: Additional S3 parameters (e.g., ContentType, StorageClass).
The ``block_size`` parameter of ``open()`` is also accepted.

Raises:
ValueError: If the file takes more than
``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
"""
if os.path.isdir(lpath):
return
_, key, _ = self.parse_path(rpath)
if not key:
return

callback.set_size(os.path.getsize(lpath))
size = os.path.getsize(lpath)
block_size = kwargs.pop("block_size", None) or self._sync_fs.default_block_size

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) — base 92c9e3e67180eb3d52ef7b85579a24540a30fd64, head bab08235f0050aebdc4cf2d85a8fcba08b6dc59b. Result: CLEAN.

Claims checked: #988's branch head does not contain #968 (git merge-base --is-ancestor ae2d2a04 <#988 head> is false), so its transaction paths mirror the pre-#968 sync methods, as the PR body says; inside a transaction the data previously went through open() and failed at part 10,001 (S3File._upload_chunk guard), and put(..., block_size=...) placed block_size in s3_additional_kwargs. docs/filesystem.md ("put and pipe check the size before uploading anything") now holds for AioS3FileSystem transactions too, so no docs change. Callers: put(..., block_size=...) in an async transaction previously failed at the S3 API, so accepting it breaks nobody; no new API requests.

self._sync_fs._check_multipart_upload_size(rpath, size, block_size)
callback.set_size(size)
if "ContentType" not in kwargs:
content_type, _ = mimetypes.guess_type(lpath)
if content_type is not None:
kwargs["ContentType"] = content_type

with (
self.open(rpath, "wb", s3_additional_kwargs=kwargs) as remote,
self.open(rpath, "wb", block_size=block_size, s3_additional_kwargs=kwargs) as remote,
open(lpath, "rb") as local,
):
while data := local.read(remote.blocksize):
Expand Down
37 changes: 37 additions & 0 deletions tests/pyathena/filesystem/test_s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,43 @@ def test_transaction_pipe_file_create_existing(self):
fs.pipe_file("s3://bucket/key", b"data", mode="create")
fs._sync_fs.exists.assert_called_once_with("s3://bucket/key")

@pytest.mark.parametrize("kwargs", [{"block_size": 4}, {}])
def test_transaction_pipe_put_file_exceeding_max_parts(self, tmp_path, kwargs):
# GH-953: in a transaction, as outside one, pipe_file() and put_file()
# reject data that does not fit in the maximum number of parts before
# opening the file.
fs = AioS3FileSystem(connection=mock.MagicMock(), skip_instance_cache=True)
fs._sync_fs.MULTIPART_UPLOAD_MAX_PARTS = 3
fs._sync_fs.default_block_size = 4
fs._sync_fs._call = mock.MagicMock()
fs.open = mock.MagicMock()
local = tmp_path / "local"
local.write_bytes(b"a" * 13)

with fs.transaction:
with pytest.raises(ValueError, match="block_size"):
fs.pipe_file("s3://bucket/k1", b"a" * 13, **kwargs)
with pytest.raises(ValueError, match="block_size"):
fs.put_file(str(local), "s3://bucket/k2", **kwargs)
fs.open.assert_not_called()
fs._sync_fs._call.assert_not_called()

def test_transaction_put_file_block_size(self, tmp_path):
# In a transaction, put_file() passes block_size to open() instead of
# the S3 API, as outside one.
fs = AioS3FileSystem(connection=mock.MagicMock(), skip_instance_cache=True)
fs.open = mock.MagicMock()
fs.open.return_value.__enter__.return_value.blocksize = 8
local = tmp_path / "local"
local.write_bytes(b"a" * 13)

with fs.transaction:
fs.put_file(str(local), "s3://bucket/key", block_size=8)

fs.open.assert_called_once_with(
"s3://bucket/key", "wb", block_size=8, s3_additional_kwargs={}
)

def test_touch_sync_wrapper(self):
# GH-977: touch() used to be fsspec's open()-based default, which
# dropped the PutObject parameters and returned None.
Expand Down
Loading