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
84 changes: 52 additions & 32 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -2298,14 +2298,15 @@ def _finish_multipart_upload(
upload_id: str,
futures: list[Future[S3MultipartUploadPart]],
request_kwargs: Mapping[str, Any] | None = None,
abort: bool = True,
) -> S3CompleteMultipartUpload:
"""Collect the uploaded parts and complete the multipart upload.

When any part or the completion fails, or the wait for them is
interrupted, the parts that have not started are cancelled, the
running ones are waited for, and the multipart upload is aborted so
that no incomplete upload or part is left behind. The original error
is then re-raised.
that no incomplete upload or part is left behind, unless ``abort`` is
false. The original error is then re-raised.

Args:
bucket: S3 bucket name.
Expand All @@ -2315,6 +2316,8 @@ def _finish_multipart_upload(
request_kwargs: Parameters of the upload, such as
``RequestPayer`` or the SSE-C parameters; the completion and
the abort receive those that they accept.
abort: Whether to abort the multipart upload on failure. A caller
that keeps the upload to abort it itself passes false.

Returns:
S3CompleteMultipartUpload of the completed upload.
Expand All @@ -2332,6 +2335,8 @@ def _finish_multipart_upload(
**self._get_operation_kwargs("complete_multipart_upload", request_kwargs),
)
except BaseException:
if not abort:
raise
# A part that is still uploading when the upload is aborted may
# be stored after the abort, so wait for the parts that could not
# be cancelled first.
Expand Down Expand Up @@ -3686,20 +3691,23 @@ def _close_without_commit(self) -> None:
Drops the buffered data, so that neither close() nor a deferred
commit() uploads it, and aborts the multipart upload, if any. An
abort failure is logged instead of raised, so it does not mask the
error that the caller is handling. Even if the abort fails or is
interrupted, commit() does not complete the upload afterwards. The
executor is shut down here, as fsspec does not close a closed file
again when it is garbage collected.
error that the caller is handling. If the abort fails or is
interrupted, the upload is kept so that :meth:`discard` can abort it,
and commit() does not complete it. The executor is shut down here, as
fsspec does not close a closed file again when it is garbage
collected.
"""
self.buffer = None
self.closed = True
try:
self.discard()
except Exception:
_logger.exception(f"Failed to abort multipart upload to s3://{self.bucket}/{self.key}.")
# discard() keeps the upload when the abort fails.
upload_id = cast(S3MultipartUpload, self.multipart_upload).upload_id

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, expanded scope (claims, compatibility, operations): FINDINGS (PR text, repaired)

Range ebf798f6dd06f394a84c1cf054a0a7561dd9ef41..56bfa0b4f42b71476008cc81fc13f80b8f7ae07d, full claim pass over the updated PR body, docstrings, comments, and commit message.

  • "Used by pipe_file()/put_file() and the part-limit write": confirmed by grep (s3.py:2266, s3.py:2527, s3_async.py:208, s3_async.py:280, _upload_chunk).
  • "commit() of a failed write does nothing / no longer invalidates the cache": the early return comes before invalidate_cache(). Nothing was written, so no listing changes.
  • "After a transaction commit, only the logged upload ID points to the upload": fsspec Transaction.complete(commit=True) calls commit(), which returns, so the file is not discarded. Rollback calls discard().
  • Docs: docs/filesystem.md:104 ("raises ValueError and aborts its multipart upload") is still true; the upload is only kept when that abort fails. The incomplete-upload section (:214) still applies.
  • Finding: the PR body still described Copy metadata, tags and annotations in multipart copies #1036 as a pending overlap, but Copy metadata, tags and annotations in multipart copies #1036 is merged and in this branch's base. The abort check precedes _abort_multipart_upload() (s3.py:2338). The body is corrected.
  • Operator: both abort-failure logs now carry the upload ID, so a leaked upload can be found or aborted by ID.
  • Limitation: abort failures are mocked; the live run covers the success and rollback paths.

_logger.exception(
f"Failed to abort multipart upload {upload_id} to s3://{self.bucket}/{self.key}."
)
finally:
self.multipart_upload = None
self.multipart_upload_parts = []
self._executor.shutdown()

def _write_and_close(self, value: bytes | bytearray | memoryview) -> None:
Expand Down Expand Up @@ -3869,49 +3877,61 @@ def commit(self) -> None:
Creates an empty object if nothing was written, uploads the buffered
data with PutObject if no multipart upload part was submitted, and
otherwise completes the multipart upload, which is aborted if the
completion fails or is interrupted. Invalidates the cache of the path
afterwards.
completion fails or is interrupted. If the abort also fails, the
upload is kept so that :meth:`discard` can abort it. Invalidates the
cache of the path afterwards. Does nothing for a file whose failed
write dropped the written data.

Raises:
FileExistsError: If an object was created at the path after the
file was opened in exclusive-create mode.
RuntimeError: If parts were submitted but no multipart upload is
initialized.
"""
if self.buffer is 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 one, expanded scope (implementation behavior): CLEAN

Range ebf798f6dd06f394a84c1cf054a0a7561dd9ef41..56bfa0b4f42b71476008cc81fc13f80b8f7ae07d (expanded contract, so this is a full pass over _close_without_commit()/commit() and their callers). Branch base is ec5323ea after the rebase; git range-diff showed the earlier four commits unchanged.

Covered: _close_without_commit() and its three callers (_write_and_close() for sync/aio pipe_file, _write_file_and_close() for sync/aio put_file, and _upload_chunk() at the part limit), commit(), discard(), AioS3File (no overrides), S3AioExecutor.shutdown() (no-op), and fsspec close()/flush().

  • buffer is None happens only after _close_without_commit(): fsspec only assigns io.BytesIO() to the buffer of a write-mode file. So the guard affects only failed writes. The two removed inner checks were for the same state.
  • The cast in the except is safe: discard() raises only from _call(), inside if self.multipart_upload, and it does not clear on failure.
  • Retry: discard() on the closed file uses neither the buffer nor the executor; it waits for uncancellable parts and aborts again. close() on the closed file returns early.
  • An interrupted abort propagates through finally (executor shut down) with the upload kept. The callers re-raise in any case.
  • Design choice: commit() of a failed write stays a no-op rather than retrying the abort, so a transaction commit does not newly raise abort errors. Rollback (discard()) retries.
  • Tests: the two part-limit abort tests now assert the kept upload, the logged upload ID, and an identical second abort request from discard(), followed by cleared state. They fail on ebf798f6. Live tests/pyathena/filesystem/: 643 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.

Correction: after the rebase, the branch base (merge base with master) is 8c9d676b77e86a7069dc1a47cd7a54d82a09c802 (#1035 merge), not ec5323ea. The review conclusions are unchanged.

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, expanded scope (relayed): Codex CLI 0.160.0, model gpt-6-astra (OpenAI), sandbox read-only, --ephemeral, session 01a1044e-d11e-7791-839c-4b62ee10550d. Static review only.

Range 8c9d676b77e86a7069dc1a47cd7a54d82a09c802..56bfa0b4f42b71476008cc81fc13f80b8f7ae07d, full pass, on a detached snapshot that stayed unchanged. The prompt had no PR number, PR text, commit messages, or prior findings. The reviewer was pointed to the installed fsspec source to check the buffer lifecycle.

Result: FINDINGS, one pre-existing item and no regressions:

Reviewer output (verbatim)

Covered the named upload/cleanup methods, sync and async callers, executor cancellation/waiting, event-loop-thread discard, fsspec transactions, autocommit close, xb, abort parameters, logging, and test assertions against the base revision. Inspected the installed fsspec buffer lifecycle: normal write-mode files retain a BytesIO, so the new buffer is None guard does not skip legitimate commits.

FINDINGS

  • P2 — Pre-existing mismatch with the stated error-preservation contract: s3.py:3930, also s3.py:3704. If completion raises RuntimeError and abort then raises KeyboardInterrupt, except Exception misses the abort interruption. The upload state survives, but KeyboardInterrupt replaces the original propagated error and no abort-failure message with the upload ID is logged. The same applies to cleanup after a failed write. This behavior already exists in the base; the new test_s3.py:5612 explicitly expects it. Preserving the original error across abort interruptions requires changing these handlers and that expectation.

No other actionable defects or introduced regressions found. The retention/retry assertions would fail against the original implementation by source inspection; the interrupted-abort commit test preserves existing behavior. The changes remain straightforward.

Static review only; no builds, tests, or GitHub access.

Author disposition: this is the same pre-existing interrupt behavior the first independent review raised. It is kept on purpose: an interrupt during the abort must propagate rather than be swallowed or re-ordered behind the completion error, the original error remains its __context__, and the upload is kept for discard(). No change. The reviewer confirmed that the buffer is None guard does not skip legitimate commits (fsspec keeps a BytesIO for normal write-mode files).

# _close_without_commit() dropped the written data. A multipart
# upload that it failed to abort is kept for discard(), not
# completed.
return
if self.tell() == 0:
if self.buffer is not None:
self.discard()
self.fs.touch(self.path, **self._get_request_kwargs("put_object"))
self.discard()
self.fs.touch(self.path, **self._get_request_kwargs("put_object"))
elif not self.multipart_upload_parts:
if self.buffer is not None:
# Upload files smaller than block size.
self.buffer.seek(0)
data = self.buffer.read()
self.fs._put_object(
bucket=self.bucket,
key=self.key,
body=data,
**self._get_request_kwargs("put_object"),
)
# Upload files smaller than block size.
self.buffer.seek(0)
data = self.buffer.read()
self.fs._put_object(
bucket=self.bucket,
key=self.key,
body=data,
**self._get_request_kwargs("put_object"),
)
else:
if not self.multipart_upload:
raise RuntimeError("Multipart upload is not initialized.")

upload_id = cast(str, self.multipart_upload.upload_id)
try:
self.fs._finish_multipart_upload(
bucket=self.bucket,
key=self.key,
upload_id=cast(str, self.multipart_upload.upload_id),
upload_id=upload_id,
futures=self.multipart_upload_parts,
request_kwargs=self.s3_additional_kwargs,
abort=False,

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) — FINDINGS (repaired)

Base 28ec68d9f70b5bc2fb447f6a4fb745ee0a467683, reviewed head f9eedeae41d7864839ea2eec75f21d8074486100, repair head 819d7d1adc6f4bb6d5734d82307b9286e1f543ee.

Covered: _finish_multipart_upload() (new abort flag) and both callers (cp_file multipart copy keeps the default; S3File.commit() passes abort=False), S3File.commit()/discard()/_close_without_commit(), AioS3File (inherits commit()/discard(), no override), fsspec Transaction.complete() (2026.9.0), and the changed tests.

  • Failure path equivalence: discard() cancels the pending parts, waits for the running ones, and sends abort_multipart_upload with _get_operation_kwargs("abort_multipart_upload", s3_additional_kwargs). That matches what the helper sent with request_kwargs=self.s3_additional_kwargs, including the event-loop-thread case (cancelled parts are not waited for).
  • The bare raise after the inner except Exception re-raises the original completion error, which match="complete failed" asserts.
  • An abort interrupted inside discard() propagates before the state is cleared, so the upload is kept, as before.
  • Finding (repaired in 819d7d1): the comment said a transaction calls discard() after a failed commit(). fsspec is unpinned, and only recent Transaction.complete() re-queues the failed file for discard(), so the comment now says "such as a transaction rollback".
  • Not changed (pre-existing, intentional): _close_without_commit() clears the upload after a failed abort so that a deferred commit() cannot complete a failed write (s3.py:3337).
  • Tests: the 4 test_commit_failure_and_discard cases and test_finish_multipart_upload_without_abort fail with the source change reverted. Offline run of the affected tests: 42 passed after the repair.

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.

Round one on the repair (5e09af9b..721d24a1, tests only): CLEAN. test_commit_failure_and_discard now also asserts that the part futures are kept after a failed abort. The new test_commit_failure_and_interrupted_abort drives the real _finish_multipart_upload() with abort=False, makes the first abort_multipart_upload raise KeyboardInterrupt, and asserts that the interrupt propagates, the upload and parts are kept, and discard() retries the abort and clears the state. No source changed. Offline: 69 passed.

)
except Exception:
# The multipart upload has been aborted by the helper;
# prevent discard() from aborting it again. An interrupt may
# have stopped the helper before the abort, so the upload is
# kept for discard() then.
self.multipart_upload = None
self.multipart_upload_parts = []
except BaseException:
# discard() keeps the upload if the abort fails or is
# interrupted, so that a later discard(), such as a
# transaction rollback, retries the abort. An abort failure
# is logged so that it does not mask the original error.
try:
self.discard()
except Exception:

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): Codex CLI 0.160.0, model gpt-6-astra (OpenAI), sandbox read-only, --ephemeral, session 01a102a7-5eca-7790-96c6-8aeb2d1b75b6. Static review only.

Base 28ec68d9f70b5bc2fb447f6a4fb745ee0a467683, head 5e09af9bbac5504612108c3bbecb7718420c6200, on a detached snapshot worktree. The prompt had no PR number, PR text, commit messages, or prior findings, and asked for no edits, tests, or GitHub access. Both the snapshot and the PR worktree were unchanged afterwards.

Result: FINDINGS, both pre-existing, not introduced by this diff:

Reviewer output (verbatim)

Covered commit(), discard(), _close_without_commit(), autocommit close, both helper callers, async delegation/inheritance, executor cancellation, repository transaction integration, request filtering, logging, subclasses, and regression tests.

FINDINGS — two pre-existing gaps remain; neither is introduced by this diff.

  1. P2 — An abort-time interrupt replaces the original completion error. s3.py:3533
    CompleteMultipartUpload raises OSError, then discard() raises KeyboardInterrupt while waiting or aborting → commit() propagates KeyboardInterrupt, bypassing the original-error re-raise and abort-failure logging. The upload state survives, but the stated original-error guarantee does not. The previous helper had the same limitation. Preserving that guarantee requires handling cleanup BaseException separately before re-raising the completion error.

  2. P2 — _close_without_commit() still destroys failed-abort retry state. s3.py:3355
    A multipart write exceeds the part limit, and its abort raises PermissionError → the finally block clears the upload ID and futures anyway. A later discard() makes no abort request, leaving the upload behind. This unchanged path falls outside the repaired completion branch. Retain cleanup state while separately preventing the failed write from being committed.

The changed completion path otherwise retains both fields after ordinary abort failures and clears them after successful abort. Abort request filtering and field precedence are unchanged, including RequestPayer, ExpectedBucketOwner, and SSE-C filtering. The default helper behavior remains compatible with the sync copy caller; AioS3File inherits the repair.

By inspection, the ordinary completion-error/abort-failure test would fail on the original code’s upload-state assertion. It checks retry requests and logging, but does not assert retained futures before retry or cover interruption during abort. Reusing discard() with an explicit helper opt-out is a simple design consistent with the surrounding code.

Static review only: no edits, builds, tests, or GitHub access. Transaction assessment used tracked integration code and tests; fsspec internals are not tracked here.

Author disposition (verified against the code):

  1. Interrupt during the abort replaces the original error (s3.py:3533): kept as is. An interrupt must propagate, not be swallowed. The completion error stays as __context__, and the upload is kept for discard(). The old helper behaved the same way.
  2. _close_without_commit() still clears the upload after a failed abort (s3.py:3337): deferred as out of scope. The clearing is what stops a deferred commit() from completing a failed write. Keeping the upload ID while blocking the commit needs a separate design and will be proposed as a follow-up issue.
  3. Test gaps the reviewer noted (part futures kept before retry; interrupted abort not covered): repaired in 721d24a1. test_commit_failure_and_interrupted_abort passes on the original source too, so it guards existing behavior rather than proving a regression.

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): Codex CLI 0.160.0, model gpt-6-astra, sandbox read-only, --ephemeral, session 01a102b1-a9af-7113-864f-ecb713634492. Static review of git diff 5e09af9bbac5504612108c3bbecb7718420c6200..721d24a1fa1aca199a38224b6a95d9e425c0d2e6 (tests only), on a detached snapshot that stayed unchanged.

Result: CLEAN. Reviewer summary: the new assertion checks that the part futures are kept before the retry, across both completion errors and successful and failed aborts. The new test drives completion failure, then an interrupted abort, then the propagated KeyboardInterrupt, the kept state, a successful discard() retry, and the cleared state. _finish_multipart_upload delegates to the production code, and commit()/discard() are real, so premature clearing, a swallowed interrupt, or a missing retry would be caught. No spurious-pass issue was found.

_logger.exception(
f"Failed to abort multipart upload {upload_id} "

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, compatibility, operations) — FINDINGS (repaired)

Base 28ec68d9f70b5bc2fb447f6a4fb745ee0a467683, reviewed head 819d7d1adc6f4bb6d5734d82307b9286e1f543ee, repair head 5e09af9bbac5504612108c3bbecb7718420c6200.

Claims checked:

  • "With abort=False the helper re-raises without cancelling or aborting": if not abort: raise comes before the cancel/wait/abort; asserted by test_finish_multipart_upload_without_abort.
  • "The multipart copy keeps the default": cp_file's call (s3.py:1777) passes no abort.
  • "discard() clears only after a successful abort": the clearing follows _call() with no finally, so a raising abort skips it.
  • "Previously, an interrupted completion was aborted twice": the old helper aborted on BaseException, and old commit() caught only Exception, so it kept the state for a second abort in discard(). The PR body does not claim what S3 returns for that second abort.
  • "_close_without_commit() clears on purpose": its docstring states that commit() must not complete a failed write afterwards. Keeping the upload there would let commit() complete it.
  • Docs: docs/filesystem.md (transactions and incomplete-upload management) is still accurate.
  • Existing caller: _finish_multipart_upload() is private, and the new keyword has a compatible default. AioS3File inherits the methods.

Findings (repaired in 5e09af9):

  1. Operator: the new commit() log line dropped the upload ID that the helper logged ("Failed to abort multipart upload {upload_id} to s3://…"). On a non-transactional close(), nothing calls discard() again, so the log is the only pointer for a manual cleanup. The ID is restored here, and the test asserts it.
  2. The PR body claimed a conflict with Copy metadata, tags and annotations in multipart copies #1036. git merge-tree merges the two heads without conflicts, so the body is corrected.

Limitation: no live S3 run; abort failures are simulated with mocks. AWS CI comes after Ready.

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.

Round two on the repair (5e09af9b..721d24a1): CLEAN after a PR text update. The TEST section now names 721d24a1 and 69 passing tests, and states that the interrupted-abort test also passes on the original source, so it is a guard rather than regression evidence. The commit() comment's claim that the upload is kept if the abort is interrupted is now backed by that test. No other claims changed.

f"to s3://{self.bucket}/{self.key}."
)
raise

self.fs.invalidate_cache(self.path)
Expand Down
99 changes: 88 additions & 11 deletions tests/pyathena/filesystem/test_s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -3106,6 +3106,26 @@ def test_finish_multipart_upload_aborts_on_failure(self, error):
UploadId="uploadid",
)

def test_finish_multipart_upload_without_abort(self):
# A caller that aborts the upload itself, as S3File.commit() does,
# gets the original error with the parts and the upload left alone.
fs = self._make_fs()
fs._complete_multipart_upload = mock.MagicMock()
failed: Future[SimpleNamespace] = Future()
failed.set_exception(RuntimeError("upload failed"))
pending: Future[SimpleNamespace] = Future()

with pytest.raises(RuntimeError, match="upload failed"):
fs._finish_multipart_upload(
bucket="bucket",
key="key",
upload_id="uploadid",
futures=[failed, pending],
abort=False,
)
assert not pending.cancelled()
fs._call.assert_not_called()

def test_finish_multipart_upload_abort_failure_does_not_mask_the_original_error(self):
fs = self._make_fs()
fs._complete_multipart_upload = mock.MagicMock()
Expand Down Expand Up @@ -5247,10 +5267,11 @@ def test_write_exceeding_max_parts(self, existing, mode, writes, autocommit):
fs._put_object.assert_not_called()

@pytest.mark.parametrize("autocommit", [True, False])
def test_write_exceeding_max_parts_abort_failure(self, autocommit):
def test_write_exceeding_max_parts_abort_failure(self, caplog, autocommit):
# An abort failure is logged; the part limit error propagates, and
# neither closing the file nor committing a deferred write retries
# the upload or completes it.
# the upload or completes it. GH-945: the upload is kept, so that
# discard() retries the abort.
fs = self._make_append_fs(b"")
fs.MULTIPART_UPLOAD_MAX_PARTS = 3
fs._call.side_effect = PermissionError("abort failed")
Expand All @@ -5275,10 +5296,20 @@ def test_write_exceeding_max_parts_abort_failure(self, autocommit):
fs._call.assert_called_once()
fs._finish_multipart_upload.assert_not_called()
fs._put_object.assert_not_called()
assert "Failed to abort multipart upload uploadid to s3://bucket/key.txt." in caplog.text

assert f.multipart_upload is not None
fs._call.side_effect = None
f.discard()
assert fs._call.call_count == 2
assert fs._call.call_args_list[1] == fs._call.call_args_list[0]
assert f.multipart_upload is None
assert f.multipart_upload_parts == []

def test_write_exceeding_max_parts_abort_interrupted(self):
# GH-997: an interrupted abort propagates, and a deferred commit
# still does not complete the upload; the executor is shut down.
# GH-945: the upload is kept, so that discard() retries the abort.
fs = self._make_append_fs(b"")
fs.MULTIPART_UPLOAD_MAX_PARTS = 3
fs._call.side_effect = KeyboardInterrupt
Expand All @@ -5302,6 +5333,14 @@ def test_write_exceeding_max_parts_abort_interrupted(self):
fs._finish_multipart_upload.assert_not_called()
fs._put_object.assert_not_called()

assert f.multipart_upload is not None
fs._call.side_effect = None
f.discard()
assert fs._call.call_count == 2
assert fs._call.call_args_list[1] == fs._call.call_args_list[0]
assert f.multipart_upload is None
assert f.multipart_upload_parts == []

def test_write_exceeding_max_parts_without_close(self):
# The executor of the closed file is shut down, as fsspec does not
# close it again when it is garbage collected.
Expand Down Expand Up @@ -5526,21 +5565,59 @@ def wait_parts(futures):
waited.assert_called_once_with([running])
assert pending.cancelled()

@pytest.mark.parametrize(("error", "aborts"), [(RuntimeError, 0), (KeyboardInterrupt, 1)])
def test_commit_failure_and_discard(self, error, aborts):
# GH-1014: an error from _finish_multipart_upload() follows its
# abort, so a later discard(), as a transaction calls after a failed
# commit(), does not abort the upload again. An interrupt may have
# stopped it before the abort, so the upload is kept for discard().
@pytest.mark.parametrize("abort_fails", [False, True])
@pytest.mark.parametrize("error", [RuntimeError, KeyboardInterrupt])
def test_commit_failure_and_discard(self, caplog, error, abort_fails):
# GH-1014: a failed or interrupted completion is aborted by commit(),
# so a later discard(), such as a transaction rollback, does not
# abort the upload again. GH-945: if the abort also fails, the upload
# is kept so that discard() retries the abort.
file = self._make_multipart_write_file(b"x" * 16, autocommit=False)
file._upload_chunk(final=True)
file.fs._finish_multipart_upload.side_effect = functools.partial(
S3FileSystem._finish_multipart_upload, file.fs
)
file.fs._complete_multipart_upload.side_effect = error("complete failed")
if abort_fails:
file.fs._call.side_effect = [PermissionError("abort failed"), None]

# The abort failure is logged, and the original error propagates.
with pytest.raises(error, match="complete failed"):
file.commit()
assert (file.multipart_upload is not None) is abort_fails
assert bool(file.multipart_upload_parts) is abort_fails
assert (
"Failed to abort multipart upload uploadid to s3://bucket/key.txt." in caplog.text
) is abort_fails
file.discard()

assert file.fs._call.call_args_list == [
mock.call("abort_multipart_upload", Bucket="bucket", Key="key.txt", UploadId="uploadid")
] * (2 if abort_fails else 1)
assert file.multipart_upload is None
assert file.multipart_upload_parts == []

def test_commit_failure_and_interrupted_abort(self):
# GH-945: if the abort after a failed completion is interrupted, the
# interrupt propagates and the upload is kept so that discard()
# retries the abort.
file = self._make_multipart_write_file(b"x" * 16, autocommit=False)
file._upload_chunk(final=True)
file.fs._finish_multipart_upload.side_effect = error("failed")
file.fs._finish_multipart_upload.side_effect = functools.partial(
S3FileSystem._finish_multipart_upload, file.fs
)
file.fs._complete_multipart_upload.side_effect = RuntimeError("complete failed")
file.fs._call.side_effect = [KeyboardInterrupt, None]

with pytest.raises(error):
with pytest.raises(KeyboardInterrupt):
file.commit()
assert file.multipart_upload is not None
assert file.multipart_upload_parts
file.discard()

assert file.fs._call.call_count == aborts
assert file.fs._call.call_count == 2
assert file.multipart_upload is None
assert file.multipart_upload_parts == []

def test_discard_on_event_loop_thread(self):
# GH-976: the parts that have not started are cancelled and not
Expand Down
Loading