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
104 changes: 47 additions & 57 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -1879,49 +1879,6 @@ def _check_multipart_upload_size(self, path: str, size: int, block_size: int) ->
f"{min_block_size} bytes."
)

@staticmethod
def _write_and_close(f: S3File, value: bytes | bytearray | memoryview) -> None:
"""Write the whole value to a file opened for writing and close it.

Unlike a ``with`` block, a failed write closes the file without
committing it, so the existing object is left unchanged.

Args:
f: The file to write to.
value: The bytes to write.
"""
try:
if isinstance(value, memoryview) and not value.c_contiguous:
# The buffer of the file cannot write a non-contiguous memoryview.
value = value.tobytes()
f.write(value)
except BaseException:
f._close_without_commit()
raise
f.close()

@staticmethod
def _write_file_and_close(f: S3File, local: BinaryIO, callback: Callback) -> None:
"""Write the rest of a local file to a file opened for writing and close it.

Unlike a ``with`` block, a failed read, write, or progress update
closes the file without committing it, so the existing object is
left unchanged.

Args:
f: The file to write to.
local: The local file to read from.
callback: Progress callback, updated with the size of each block.
"""
try:
while data := local.read(f.blocksize):
f.write(data)
callback.relative_update(len(data))
except BaseException:
f._close_without_commit()
raise
f.close()

def pipe_file(
self, path: str, value: bytes | bytearray | memoryview, mode: str = "overwrite", **kwargs
) -> None:
Expand Down Expand Up @@ -1975,9 +1932,7 @@ def pipe_file(
# Defer to the buffered open() path, which keeps the
# deferred-commit semantics of fsspec transactions and uploads
# large data as a parallel multipart upload.
self._write_and_close(
self.open(path, "xb" if mode == "create" else "wb", **kwargs), value
)
self.open(path, "xb" if mode == "create" else "wb", **kwargs)._write_and_close(value)

@laughingman7743 laughingman7743 Oct 3, 2026 •

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

Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only, model reported gpt-6-astra), session 01a1022e-eefe-7122-9d90-c62474c23aaf. Static review only. Base fe21250a3e228e3d38e61898334c5baacc1285c2, head 8011487869ccf0cde9a48ce9b129b59bfce2e74c, detached snapshot (unchanged afterwards). The prompt held the diff and conventions, without the PR number, PR text, or prior findings.

Covered: sync and aio pipe_file/put_file in and outside transactions incl. mode="create"; open/_open, compression wrappers, deferred commit, discard, exception cleanup, executor shutdown; AbstractBufferedFile attributes and AioS3File inheritance (no name collisions or shadowing); changed tests; docstrings and conventions.

P2 — Regression: compressed buffered writes now fail. fsspec's open implementation (spec.py:1427) can return a compression wrapper instead of S3File. For example, with either filesystem: with fs.transaction: fs.pipe_file("s3://bucket/data.gz", b"payload", compression="gzip"). Previously, the static helper called the wrapper's write() and close(), allowing the transaction to commit the compressed payload. Now method lookup raises AttributeError because GzipFile has no _write_and_close. This also affects buffered writes outside transactions and exclusive creation of an absent object. (Also s3_async.py:201.)

The four edited mock-based tests lose execution coverage of the read/write/callback/close loop because the new helper becomes another MagicMock. Their parameter assertions remain meaningful; they are not vacuous for their stated purpose.

Author verification (offline, mocked requests):

  • Confirmed. pipe_file(..., compression="gzip") on the buffered path (in a transaction, or larger than the block size) uploads a valid gzip body on master fe21250a and before Write pipe_file() data without committing a failed write #1003 (5a60dd19); on 80114878 it raises AttributeError: 'GzipFile' object has no attribute '_write_and_close'. put_file() does not pass compression to open() and is unaffected.
  • Pre-existing, found while verifying: on master, a failed compressed write on that path raises AttributeError: 'GzipFile' object has no attribute '_close_without_commit' (masking the original error) and still uploads a 10-byte object holding only the gzip header, so Write pipe_file() data without committing a failed write #1003's guarantee does not hold with compression. The single-request path (small values outside a transaction) passes compression to PutObject, which botocore parameter validation rejects (ParamValidationError: Unknown parameter in input: "compression").

The repair changes how pipe_file() handles compression, which is a design decision; this PR stays Draft until it is decided.

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.

Resolved by #1038 (merged as abc99b0984f9e19755580e5be494003b725e61cb, maintainer's choice "compress up front"): pipe_file() and AioS3FileSystem._pipe_file_in_transaction() now pop compression, compress the value with CompressedBuffer.compress(), and never pass compression to open(), so open() returns the S3File. #1038 also fixed the pre-existing failure path noted above.

Rebased onto abc99b09 as 880ba8a018353cec2fdcd025151ce06176056ff9 (--force-with-lease from 80114878); no conflicts, and git range-diff fe21250a..80114878 abc99b09..880ba8a0 shows the commit unchanged.

Self-review after the rebase:

An independent re-review of the rebased PR follows.

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 re-review of the rebased PR (relayed): CLEAN

Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only, model reported gpt-6-astra), session 01a102b6-5215-7f12-874c-824af00ceb4f. Static review only. Scope git diff abc99b09..880ba8a0 (full PR on the new base); snapshot at 880ba8a018353cec2fdcd025151ce06176056ff9, unchanged afterwards. The prompt held the diff and conventions, without the PR number, PR text, or prior findings, and limited reading to the snapshot and the fsspec source.

All four changed callers through PyAthena _open() and fsspec 2026.9.0 open(). Their arguments produce S3File/AioS3File: binary modes exclude text wrappers, pipe callers consume compression beforehand, and put callers pass no wrapper options. [...] Inheritance and name collisions: neither fsspec's base class nor AioS3File shadows the moved methods. Changed docstrings remain accurate for these callers. The four edited tests [...]: the edited mocks now bypass helper execution, losing incidental loop execution, but their parameter assertions remain meaningful. Other tests still exercise actual file objects and verify bytes, callbacks, transaction outcomes, compression, and failure cleanup.

CLEAN — No actionable regressions found. No separate pre-existing findings to report.

return
bucket, key, version_id = self.parse_path(path)
if version_id:
Expand Down Expand Up @@ -2210,17 +2165,13 @@ def put_file(
# The local file is opened first, so that an unreadable one fails
# before the remote file is opened.
with open(lpath, "rb") as local:
self._write_file_and_close(
self.open(
rpath,
"xb" if mode == "create" else "wb",
block_size=block_size,
max_workers=max_workers,
s3_additional_kwargs=s3_additional_kwargs,
),
local,
callback,
)
self.open(
rpath,
"xb" if mode == "create" else "wb",
block_size=block_size,
max_workers=max_workers,
s3_additional_kwargs=s3_additional_kwargs,
)._write_file_and_close(local, callback)

self.invalidate_cache(rpath)

Expand Down Expand Up @@ -3403,6 +3354,45 @@ def _close_without_commit(self) -> None:
self.multipart_upload_parts = []
self._executor.shutdown()

def _write_and_close(self, value: bytes | bytearray | memoryview) -> None:

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 1 (behavior and implementation): CLEAN

Base fe21250a3e228e3d38e61898334c5baacc1285c2 (merge-base with master), head 8011487869ccf0cde9a48ce9b129b59bfce2e74c.

Covered: the moved S3File._write_and_close() / _write_file_and_close(), their four callers (S3FileSystem.pipe_file(), put_file(), AioS3FileSystem._pipe_file_in_transaction(), _put_file_in_transaction()), and the four tests that mock open().

  • The method bodies are unchanged apart from f → self; the callers pass the same file, opened with the same mode and arguments, so the behavior is identical.
  • S3FileSystem.open() returns an S3File and AioS3FileSystem.open() an AioS3File, which inherits both methods. The static helpers already required an S3File (they called its _close_without_commit()).
  • No name collision: neither fsspec 2026.9.0 AbstractBufferedFile nor AioS3File defines these names. TestS3File._write_and_close(f, writes) is an unrelated static helper of the test class.
  • In the four tests that mock open(), _write_file_and_close() of the mocked file is a mock, so the removed blocksize assignments were unused; the tests assert only the open() arguments.
  • Validation: just lint passed; tests/pyathena/filesystem/ 530 passed against live S3.

"""Write the whole value and close the file.

Unlike a ``with`` block, a failed write closes the file without
committing it, so the existing object is left unchanged.

Args:
value: The bytes to write.
"""
try:
if isinstance(value, memoryview) and not value.c_contiguous:
# The buffer of the file cannot write a non-contiguous memoryview.
value = value.tobytes()
self.write(value)
except BaseException:
self._close_without_commit()
raise
self.close()

def _write_file_and_close(self, local: BinaryIO, callback: Callback) -> None:
"""Write the rest of a local file and close the file.

Unlike a ``with`` block, a failed read, write, or progress update
closes the file without committing it, so the existing object is
left unchanged.

Args:
local: The local file to read from.
callback: Progress callback, updated with the size of each block.
"""
try:
while data := local.read(self.blocksize):
self.write(data)
callback.relative_update(len(data))
except BaseException:
self._close_without_commit()
raise
self.close()

def _initiate_upload(self) -> None:
if not self.append_block and self.tell() < self.blocksize:
# Files smaller than block size in size cannot be multipart uploaded.
Expand Down
22 changes: 8 additions & 14 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -205,9 +205,7 @@ def _pipe_file_in_transaction(
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)
self._sync_fs._write_and_close(
self.open(path, "xb" if mode == "create" else "wb", **kwargs), value
)
self.open(path, "xb" if mode == "create" else "wb", **kwargs)._write_and_close(value)

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 2 (claims, callers, operations): FINDINGS (PR text only, repaired)

Base fe21250a3e228e3d38e61898334c5baacc1285c2, head 8011487869ccf0cde9a48ce9b129b59bfce2e74c.

Claims checked: "No behavior change" (the bodies match the removed static methods line for line apart from f → self). "Neither helper used the filesystem" (neither body referenced the filesystem). "AioS3File inherits them" (class AioS3File(S3File) defines neither). The commit message and docstrings agree. Existing callers: the methods are private and had no other callers in pyathena/, tests/, or docs/. No AWS request changes.

Finding (PR text, repaired): the TEST section said that all the listed regression tests patch S3File.write/AioS3File.write. Only some do; the put_file() ones fail through the callback, an unreadable local file, or the executor. The PR body now lists each test with the failure it injects.


async def _put_file(
self,
Expand Down Expand Up @@ -273,17 +271,13 @@ def _put_file_in_transaction(

# See S3FileSystem.put_file.
with open(lpath, "rb") as local:
self._sync_fs._write_file_and_close(
self.open(
rpath,
"xb" if mode == "create" else "wb",
block_size=block_size,
max_workers=max_workers,
s3_additional_kwargs=s3_additional_kwargs,
),
local,
callback,
)
self.open(
rpath,
"xb" if mode == "create" else "wb",
block_size=block_size,
max_workers=max_workers,
s3_additional_kwargs=s3_additional_kwargs,
)._write_file_and_close(local, callback)
self.invalidate_cache(rpath)

async def _get_file(self, rpath: str, lpath: str, callback=_DEFAULT_CALLBACK, **kwargs) -> None:
Expand Down
2 changes: 0 additions & 2 deletions tests/pyathena/filesystem/test_s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -2079,7 +2079,6 @@ def test_put_file_block_size(self, tmp_path):
# API.
fs = self._make_fs()
fs.open = mock.MagicMock()
fs.open.return_value.blocksize = 8
lpath = tmp_path / "data"
lpath.write_bytes(b"a" * 13)

Expand Down Expand Up @@ -2107,7 +2106,6 @@ def test_put_file_content_type(self, tmp_path, filesystem_kwargs, kwargs, expect
fs = self._make_fs()
fs.s3_additional_kwargs = filesystem_kwargs
fs.open = mock.MagicMock()
fs.open.return_value.blocksize = 8
lpath = tmp_path / "data.csv"
lpath.write_bytes(b"a")

Expand Down
2 changes: 0 additions & 2 deletions tests/pyathena/filesystem/test_s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -368,7 +368,6 @@ def test_transaction_put_file_block_size(self, tmp_path):
# 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.blocksize = 8
local = tmp_path / "local"
local.write_bytes(b"a" * 13)

Expand Down Expand Up @@ -517,7 +516,6 @@ def test_put_file_in_transaction_open_parameters(self, tmp_path, mode, open_mode
# GH-972: fsspec's mode argument selects the mode of the file.
fs = AioS3FileSystem(connection=mock.MagicMock(), skip_instance_cache=True)
fs.open = mock.MagicMock()
fs.open.return_value.blocksize = 4
lpath = tmp_path / "data.csv"
lpath.write_bytes(b"a")

Expand Down
Loading