Move the write-and-close helpers from S3FileSystem to S3File - #1035
Conversation
| self.multipart_upload_parts = [] | ||
| self._executor.shutdown() | ||
|
|
||
| def _write_and_close(self, value: bytes | bytearray | memoryview) -> None: |
There was a problem hiding this comment.
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 anS3FileandAioS3FileSystem.open()anAioS3File, which inherits both methods. The static helpers already required anS3File(they called its_close_without_commit()).- No name collision: neither fsspec 2026.9.0
AbstractBufferedFilenorAioS3Filedefines 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 removedblocksizeassignments were unused; the tests assert only theopen()arguments. - Validation:
just lintpassed;tests/pyathena/filesystem/530 passed against live S3.
| 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) |
There was a problem hiding this comment.
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.
| 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) |
There was a problem hiding this comment.
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 ofS3File. 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'swrite()andclose(), allowing the transaction to commit the compressed payload. Now method lookup raisesAttributeErrorbecauseGzipFilehas no_write_and_close. This also affects buffered writes outside transactions and exclusive creation of an absent object. (Alsos3_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 masterfe21250aand before Write pipe_file() data without committing a failed write #1003 (5a60dd19); on80114878it raisesAttributeError: 'GzipFile' object has no attribute '_write_and_close'.put_file()does not passcompressiontoopen()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) passescompressionto 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.
There was a problem hiding this comment.
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:
- Round 1 (behavior): every caller of
_write_and_close()/_write_file_and_close()now receives the file fromopen()in binary write mode withoutcompression(pipe_file()pops it;put_file()and_put_file_in_transaction()pass onlyblock_size,max_workers, ands3_additional_kwargs), so the object is always anS3File/AioS3File. - Round 2 (claims): PR body updated (dependency on Compress pipe_file() values up front, and route and key them consistently #1038, tested commit, results).
- Validation on
880ba8a0:just lintpassed;tests/pyathena/filesystem/616 passed against live S3, including the 23 compression tests of Compress pipe_file() values up front, and route and key them consistently #1038 (compressed buffered writes in and outside transactions, multipart).
An independent re-review of the rebased PR follows.
There was a problem hiding this comment.
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.0open(). Their arguments produceS3File/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 norAioS3Fileshadows 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.
_write_and_close() and _write_file_and_close() were static methods of S3FileSystem that used no filesystem state: they wrote to the S3File they received and called its private _close_without_commit(). Make them instance methods of S3File, next to _close_without_commit(), and call them on the opened file. AioS3FileSystem no longer reaches them through its internal filesystem, as AioS3File inherits them. Drop the blocksize of the mocked files in the tests that mock open(), whose _write_file_and_close() is now a mock too. No behavior change. Closes #1034 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
8011487 to
880ba8a
Compare
WHAT
Move
_write_and_close()(added forpipe_file()in #1003) and_write_file_and_close()(added forput_file()in #1017) from static methods ofS3FileSystemto instance methods ofS3File, next to_close_without_commit().S3FileSystem.pipe_file()/put_file()andAioS3FileSystem._pipe_file_in_transaction()/_put_file_in_transaction()call them on the opened file:self.open(...)._write_and_close(value)andself.open(...)._write_file_and_close(local, callback).AioS3FileSystemno longer calls them through its internal filesystem (self._sync_fs), asAioS3Fileinherits them fromS3File.open()no longer setblocksizeon the mocked file:_write_file_and_close()of the mocked file is a mock too, so nothing reads it.Private methods only; no behavior change.
WHY
Depends on #1038 (merged):
pipe_file()used to pass fsspec'scompressionargument toopen(), which then returned a compression wrapper (e.g.GzipFile) without these methods; the review of this PR found that. Since #1038,pipe_file()compresses the value itself andopen()always returns theS3File.Closes #1034.
Neither helper used the filesystem: both operated only on the
S3Filethey received and called its private_close_without_commit()from outside the class.TEST
Tested commit: 880ba8a (rebased onto
abc99b09, #1038)just lint: passed.uv run --env-file .env pytest -n 8 tests/pyathena/filesystem/: 616 passed (live S3). This includes the regression tests of Write pipe_file() data without committing a failed write #1003 and Leave the existing object unchanged when put_file() fails #1017 that exercise the moved methods: a failingS3File.write/AioS3File.write(test_pipe_file_failed_write,test_transaction_pipe_file_write,test_transaction_put_file_failed_write[RuntimeError]), a failing callback or unreadable local file (test_put_file_failed_write,test_transaction_put_file_failed_write[PermissionError]), a failing part submission or callback after a multipart upload started (test_pipe_file_failed_write_aborts_multipart_upload,test_put_file_failed_write_aborts_multipart_upload), and a non-contiguous memoryview (test_pipe_file_buffered_non_contiguous_memoryview). The compression tests of Compress pipe_file() values up front, and route and key them consistently #1038 (test_pipe_file_compression*,TestCompressedBuffer, aiotest_pipe_file_compression; 23 passed) cover the compressed buffered writes that this change broke before Compress pipe_file() values up front, and route and key them consistently #1038.🤖 Generated with Claude Code