Skip to content

Write pipe_file() data without committing a failed write - #1003

Merged
laughingman7743 merged 2 commits into
masterfrom
fix/997-pipe-file-memoryview
Oct 3, 2026
Merged

laughingman7743 merged 2 commits into
masterfrom
fix/997-pipe-file-memoryview

Conversation

@laughingman7743

@laughingman7743 laughingman7743 commented Oct 3, 2026 •

Copy link
Copy Markdown
Member

WHAT

Fix the buffered path of pipe_file() (data larger than the block size, and every write inside an fsspec transaction) for non-contiguous memoryviews and failed writes.

  • S3FileSystem.pipe_file() and AioS3FileSystem._pipe_file_in_transaction() write through a new S3FileSystem._write_and_close(f, value) instead of fsspec's with self.open(...) as f: f.write(value).
    • A non-contiguous memoryview is copied to bytes first, as the single-request path already does with bytes(value). A contiguous value (bytes, bytearray, contiguous memoryview) is written as before, without an extra copy.
    • When write() raises (any BaseException, so also KeyboardInterrupt), the file is closed without being committed: its buffer is dropped, its multipart upload, if any, is aborted, and the original exception is raised. Before, the with block committed the file, so a write that failed before buffering anything replaced the object with an empty one through touch(). Inside a transaction, the empty object was written when the transaction committed, if the caller had caught the error.
  • The close-without-commit steps of the part-limit error in S3File._upload_chunk() (Keep multipart uploads within the 10,000-part limit #968) move into S3File._close_without_commit(), which both paths use. It now clears the multipart state and shuts the executor down in a finally, so an abort that is interrupted (e.g. KeyboardInterrupt while waiting for running parts) still leaves nothing for a deferred commit() to complete; before, a caller that caught the interrupt inside a transaction had the submitted parts completed at commit. An abort that fails with an Exception is logged, as before.
  • With mode="create", the buffered path still opens the file in xb mode (Implement exclusive create for open("xb"), put_file() and pipe_file() #1009); a failed write is closed without its conditional (IfNoneMatch) commit.

Release-note items:

  • pipe() accepts a non-contiguous memoryview larger than the block size, and any non-contiguous memoryview inside a transaction (sync and async); before, it raised BufferError.
  • A pipe() whose write fails no longer replaces the existing object with an empty one, and aborts its multipart upload, if one was started.

WHY

Closes #997.

The buffered path was fsspec's generic pipe_file(), whose with block commits the file on any exception. A file whose write failed was therefore treated like a completed one.

Not changed here:

  • Writes through open() in the caller's own with block keep fsspec's semantics: an exception inside the block still closes and commits the file. Changing that would affect every open() user.
  • S3File.write() of a non-contiguous memoryview still raises BufferError, as Python's own files do; only pipe() converts it, because its signature accepts memoryviews.
  • An interrupt while close() waits for the parts in the final completion (_finish_multipart_upload() handles only Exception) still leaves the multipart upload behind. This applies to every write, not only pipe().
  • pipe(..., compression=...) is not supported (the single-request path sends compression to PutObject, which rejects it); on the buffered path, a failed write of such a call raises AttributeError from the compressor wrapper instead of being cleaned up.
  • pipe_file() still chooses between the two paths by len(value), which counts the items of a memoryview, as pinned by test_pipe_file_non_contiguous_memoryview.

TEST

Tested commit: 5206acc (rebased onto 65a305b, which adds #1009, #989 and #985; the conflict with #1009's open(path, "xb" if mode == "create" else "wb") was resolved by passing that file to _write_and_close(), dropping this PR's separate exists() check).

  • just format, just lint: pass.
  • New offline unit tests (dummy environment, --noconftest):
    • TestS3FileSystem::test_pipe_file_buffered_non_contiguous_memoryview: a 5 MiB + 1 byte non-contiguous memoryview is uploaded as a multipart upload with the expected bytes, and nothing is touched.
    • TestS3FileSystem::test_pipe_file_failed_write[False/True]: a failing write() on the buffered path sends no PutObject, outside a transaction and inside one whose caller catches the error.
    • TestS3FileSystem::test_pipe_file_failed_write_aborts_multipart_upload: a write that fails after its multipart upload started aborts the upload, does not complete it, and shuts the executor down.
    • TestAioS3FileSystem::test_transaction_pipe_file_write: in an async transaction, a small non-contiguous memoryview is written and a failed write leaves no PutObject at commit.
    • TestS3File::test_write_exceeding_max_parts_abort_interrupted: when the abort after the part-limit error raises KeyboardInterrupt, the interrupt propagates, the executor is shut down, and a deferred commit() does not complete the upload. It fails on 9a05411 (the upload is completed).
    • With the source changes reverted, the first five fail (BufferError, PutObject from touch(), or the upload not aborted).
  • The whole tests/pyathena/filesystem/ directory offline: no failure beyond the base's (master 65a305b); the remaining failures are AWS integration tests without credentials.
  • AWS integration tests: uv run --env-file .env pytest -n 4 -p no:cacheprovider tests/pyathena/filesystem/ (local, against the CI account) at 5206acc: 474 passed (441 at 2ff6a9b and 395 at 70d0887 before the rebases).
  • Not covered against real S3: the failed-write paths, which need an injected failure.

🤖 Generated with Claude Code

Comment thread pyathena/filesystem/s3.py
try:
f.write(value)
except BaseException:
f._close_without_commit()

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 (implementation behavior) — CLEAN

Scope: git diff 49509962cb3e81f819c59985d9f6e2ece576ffe0..9a05411b566d227ad33649e3501c3421c95731df (base = merge-base with master, head = published head). Covered: _write_and_close(), both callers (S3FileSystem.pipe_file() buffered path, AioS3FileSystem._pipe_file_in_transaction()), S3File._close_without_commit() and its two callers (this one and the part-limit error in _upload_chunk()), the mode="create" check, fsspec's Transaction commit/rollback of a file closed this way, and the five new tests.

Checked without findings:

  • Deferred commit after a failed write in a transaction: commit() sees buffer is None and, after discard(), no parts, so it neither touches nor puts nor completes; a rollback's discard() is a no-op.
  • Failure in _initiate_upload() (CreateMultipartUpload): fsspec's flush() already sets closed; the second close here only shuts the executor down, which close() used to skip for a closed file.
  • Part-limit ValueError: _close_without_commit() runs twice (in _upload_chunk() and here); discard() is then a no-op and a repeated shutdown() is harmless for both executors (ThreadPoolExecutor.shutdown, S3AioExecutor.shutdown is a no-op).
  • BaseException: an interrupted write no longer completes from a half-consumed buffer (the old with block's final flush re-read the buffer from offset 0 after some parts were submitted).
  • Tests: the five new tests fail on the base (BufferError, touch()'s PutObject, upload not aborted) and pass at the head; offline tests/pyathena/filesystem/ differs from the base only by them; live tests/pyathena/filesystem/: 394 passed.

Reasoned non-change: pipe(..., compression=...) would hand a compressor wrapper to _write_and_close(), whose failure path would raise AttributeError instead of the write error. compression is not a supported pipe parameter (the single-request path sends it to PutObject, which rejects it), so this is not addressed.

Comment thread pyathena/filesystem/s3.py Outdated
super().pipe_file(path, value, mode=mode, **kwargs)
if mode == "create" and self.exists(path):
raise FileExistsError(path)
self._write_and_close(self.open(path, "wb", **kwargs), 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 description corrected; no code change)

Scope: same base/head; the PR body, commit message, the new and changed docstrings/comments, and docs/filesystem.md.

Claims checked:

  • The single-request path converts with bytes(value) (s3.py, unchanged lines after this call): confirmed.
  • A failed write before buffering replaced the object through touch(), also later at a transaction's commit if the caller caught the error: confirmed by test_pipe_file_failed_write[False/True] failing on the base with _put_object(body=None); without catching, the transaction rolls back and nothing was written (stated as such).
  • fsspec's pipe_file() raises FileExistsError without arguments (fsspec 2026.9.0 source): confirmed, so the path in the message is new.
  • Python's own files reject a non-contiguous memoryview: tempfile.TemporaryFile().write(memoryview(b"ab"*4)[::2]) raises BufferError.
  • The caller's own with fs.open(...) keeps committing on exception: AbstractBufferedFile.__exit__ calls close() unconditionally.
  • docs/filesystem.md: no statement is made obsolete (it describes the paths, not failure handling).
  • Operations: no request is added on success; a failed write after a multipart upload started now sends one AbortMultipartUpload instead of re-submitting parts on close.

Corrections to the PR body:

  • "written as before, without a copy" → "without an extra copy" (BytesIO.write() copies into the buffer anyway).
  • TEST now states that the local runs used the head's tree before a final, comment-only test edit.

Not covered against real S3: the failed-write paths (need an injected failure); stated in TEST.

Comment thread pyathena/filesystem/s3.py Outdated
"""
if isinstance(value, memoryview) and not value.c_contiguous:
# The buffer of the file cannot write a non-contiguous memoryview.
value = value.tobytes()

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: OpenAI Codex CLI 0.160.0, model gpt-6-astra (local Codex config), reasoning effort high, sandbox read-only; session 01a10140-0e8c-7871-be1e-ad9638bfe80f. Static review only: no tests, builds, edits, or network access. Snapshot: detached worktree at 9a05411b566d227ad33649e3501c3421c95731df, diff git diff 49509962cb3e81f819c59985d9f6e2ece576ffe0..9a05411b566d227ad33649e3501c3421c95731df; the prompt omitted the PR number, description, commit messages, and prior findings. The snapshot and the PR worktree were unchanged afterwards.

Covered (reviewer): sync/async buffered writes, transactions, bytes-like inputs, compression wrappers; fsspec 2026.9.0 write/flush/close/destruction and transaction commit/discard; executor shutdown, multipart cancellation/abort/completion, direct open(), the 10,000-part guard; tests and docstring claims. The new tests were judged to detect the original failures (the transaction cases suppress the error before the transaction exits, so its commit exposes an upload).

Introduced:

  1. P2 — s3.py:1524: the tobytes() conversion runs outside the try; a MemoryError there leaves the opened file untouched, so a later transaction commit (or garbage collection) touches an empty object.
  2. P2 — s3.py:1528: with pipe_file(..., compression="gzip"), open() returns a compressor wrapper without _close_without_commit(), so a failed write raises AttributeError, masking the error and skipping cleanup.

Pre-existing:
3. P2 — s3.py:2637: a KeyboardInterrupt during discard()'s wait or abort bypasses except Exception; the multipart state stays populated and the executor is not shut down, so a deferred commit can complete the submitted prefix. The extracted helper now also serves ordinary failed writes.
4. P2 — s3.py:1628: an interrupt while close() waits in _finish_multipart_upload() bypasses its Exception handler, so the upload is never aborted.

Author verification: 1 confirmed. 2 declined: compression is not a supported pipe parameter (the single-request path sends it to PutObject, which rejects it); listed in "Not changed here". 3 confirmed (reproduced as a test that fails on 9a05411: the deferred commit calls _finish_multipart_upload()); a contained fix, folded in. 4 pre-existing for every write, not pipe()-specific; listed in "Not changed here".

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.

Repair 70d0887 for findings 1 and 3, with both self-review perspectives on the repair

Range: git range-diff 49509962cb3e81f819c59985d9f6e2ece576ffe0..9a05411b566d227ad33649e3501c3421c95731df 49509962cb3e81f819c59985d9f6e2ece576ffe0..70d088752828b13b34129f01b153cc7221baa1fa (same base; one new commit).

Repair:

  1. _write_and_close() converts a non-contiguous memoryview inside the try, so a failed conversion also closes the file without committing it.
  2. _close_without_commit() clears multipart_upload/multipart_upload_parts and shuts the executor down in a finally; an Exception from the abort is still logged, not raised.

Behavior (round-1 perspective): on success, discard() already cleared the same state, so the finally changes nothing; on an Exception it does what the old except did; on an interrupt, the abort may not have been sent (an upload can remain, unavoidable without completing the wait), but commit() finds no parts and does nothing, and the interrupt propagates. Both callers (part-limit error, _write_and_close()) get this. A second invocation (part-limit error inside write() reaching _write_and_close()'s handler) stays a no-op apart from a repeated, harmless shutdown().

Claims (round-2 perspective): the docstring now says commit() does not complete the upload even if the abort fails or is interrupted, which the new test checks; the PR body describes the finally, lists findings 2 and 4 under "Not changed here", and updates TEST.

Tests: new TestS3File::test_write_exceeding_max_parts_abort_interrupted fails on 9a05411 (_finish_multipart_upload() called by the deferred commit) and passes on 70d0887. Finding 1 has no dedicated test (a MemoryError from memoryview.tobytes() cannot be injected without patching the builtin); the change is the position of the conversion only. just lint passes; offline tests/pyathena/filesystem/ failures equal the base's; live tests/pyathena/filesystem/ at 70d0887: 395 passed.

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) on 70d0887 — CLEAN

Reviewer: OpenAI Codex CLI 0.160.0, model gpt-6-astra, effort high, sandbox read-only; session 01a10149-0197-7cf3-9c8f-277f6af8ee82. Static review of git diff 9a05411b..70d08875 (same base) at a detached snapshot of 70d0887, which stayed unchanged, as did the PR worktree.

Covered (reviewer): both buffered pipe_file() callers, _write_and_close(), _close_without_commit(), the multipart part-limit path, commit()/discard()/close(), both executors, fsspec 2026.9.0 transactions, and the new test.

Reviewer result: both reported issues are resolved; the success path is unchanged; an abort Exception is still logged and suppressed while a BaseException propagates after cleanup; repeated cleanup sends no second abort and both executors tolerate repeated shutdown; the new test fails without the repair through its deferred-commit assertion (its shutdown assertion alone would not isolate the repair, because the test helper also calls close()); the docstring is accurate. No actionable introduced or pre-existing defects within scope.

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.

Rebase onto 28826be (#1000, #991, #995) — 2ff6a9b

git range-diff 49509962cb3e81f819c59985d9f6e2ece576ffe0..70d088752828b13b34129f01b153cc7221baa1fa 28826be7d37983cb54050e849b32a4fe2331ea7d..2ff6a9b7f4cb8e778b4611cdb45398808ad290d4: the repair commit is identical; the first commit differs only by the test module's import contextlib, which master already has.

Upstream contracts checked against this change: #1000 rewrote S3File.discard() (abort parameters from _get_request_kwargs(), own fields last) without changing its state handling, which _close_without_commit() relies on; open() keyword parameters now reach the file's parameters, so the buffered pipe_file() path through _write_and_close(self.open(path, "wb", **kwargs), value) keeps them (#1000's test_pipe_file_buffered_s3_parameters passes on 2ff6a9b); _make_fs now uses the real service model, and the new tests pass with it.

Validation at 2ff6a9b: just lint passes; offline tests/pyathena/filesystem/ failures equal master 28826be's (only AWS integration tests); live tests/pyathena/filesystem/: 441 passed. Pushed with --force-with-lease against 70d0887; current CI pending.

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.

Rebase onto 65a305b (#1009, #989, #985) with a conflict resolution — 5206acc

git range-diff 28826be7d37983cb54050e849b32a4fe2331ea7d..2ff6a9b7f4cb8e778b4611cdb45398808ad290d4 65a305b18e8c7e79a24a60c12b913e0ad71fbe6a..5206accedbb2c417b588d9e5fc32b28911d06ffb: the repair commit is identical. The first commit conflicted with #1009, which made the buffered path of pipe_file() and _pipe_file_in_transaction() open the file with "xb" if mode == "create" else "wb" inside a with block. Resolution: self._write_and_close(self.open(path, "xb" if mode == "create" else "wb", **kwargs), value) (sync and async); this PR's separate exists() check is dropped, because open() in xb mode raises FileExistsError before the file is created.

Both self-review perspectives on the resolution: in xb mode the file is a write-mode S3File with IfNoneMatch="*" in its parameters; a failed write closes it through _close_without_commit(), so neither the conditional PutObject/CompleteMultipartUpload nor touch() is sent; a successful write keeps #1009's conditional commit. #1009's tests (test_exclusive_create offline cases, the async transaction mode test that replaced test_transaction_pipe_file_create_existing) pass on 5206acc. The PR body's former bullet about moving the create check is replaced accordingly.

Validation at 5206acc: just lint passes; offline tests/pyathena/filesystem/ failures equal master 65a305b's (only AWS integration tests); live tests/pyathena/filesystem/: 474 passed. Pushed with --force-with-lease against 2ff6a9b; current CI pending; an independent follow-up of the resolution is running.

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) on the conflict resolution at 5206acc — CLEAN

Reviewer: OpenAI Codex CLI 0.160.0, model gpt-6-astra, effort high, sandbox read-only; session 01a101c9-fa7d-79a0-9339-c9980ed361bb. Static review of git range-diff 28826be7..2ff6a9b7 65a305b1..5206acce and the resulting code at a detached snapshot of 5206acc, which stayed unchanged, as did the PR worktree.

Covered (reviewer): both resolved call sites, _write_and_close(), _close_without_commit(), S3File init/upload/commit/discard in xb mode, async transaction ownership, fsspec 2026.9.0 transactions and buffered files, conditional request filtering and error translation.

Reviewer result: exclusive create stays enforced at open and at commit (IfNoneMatch="*" for empty, small and multipart commits), so dropping the separate exists() check is safe; conversion or write failures clear the buffer and multipart state and close the file even if the abort fails or is interrupted, so a later transaction commit cannot publish the failed write; both paths keep deferred transaction commits, including registration on the async filesystem's transaction. No introduced or pre-existing actionable findings.

laughingman7743 and others added 2 commits October 3, 2026 21:39
pipe_file() larger than the block size, and every pipe_file() inside a
transaction, wrote through fsspec's `with open(...)` block. A
non-contiguous memoryview raised BufferError in the buffer of the file,
and the with block then committed the file, which replaced the existing
object with an empty one (also later, when the transaction committed).

The buffered path of S3FileSystem.pipe_file() and of
AioS3FileSystem._pipe_file_in_transaction() now writes through
_write_and_close(), which copies a non-contiguous memoryview to bytes and
closes a file whose write fails without committing it, aborting its
multipart upload. The close-without-commit steps of the part-limit error
move to S3File._close_without_commit() and are shared.

Closes #997

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Convert a non-contiguous memoryview inside the protected region of
_write_and_close(), so a failed conversion closes the file without
committing it too. _close_without_commit() now clears the multipart
state and shuts the executor down even when the abort is interrupted
(e.g. KeyboardInterrupt), so a deferred commit() cannot complete the
parts submitted before the failure.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@laughingman7743
laughingman7743 force-pushed the fix/997-pipe-file-memoryview branch from 2ff6a9b to 5206acc Compare October 3, 2026 12:42
@laughingman7743
laughingman7743 merged commit 16aef64 into master Oct 3, 2026
7 checks passed
@laughingman7743
laughingman7743 deleted the fix/997-pipe-file-memoryview branch October 3, 2026 12:52
laughingman7743 added a commit that referenced this pull request Oct 3, 2026
put_file() wrote the local file inside a with block that held the remote
file, so a failure inside the loop committed the blocks written so far,
and an unreadable local file, opened after the remote one, replaced the
object with an empty one. The same applied to put_file() inside a
transaction of AioS3FileSystem when the transaction committed.

Open the local file first, and write it through a new
_write_file_and_close() helper that closes the remote file with
_close_without_commit() on any failure, as pipe_file() does since #1003.

_finish_multipart_upload() now also aborts the multipart upload when the
wait for the parts or the completion is interrupted, instead of leaving
it behind.

Closes #1014

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

pipe_file() fails on non-contiguous memoryviews and replaces the object with an empty one

1 participant