-
Notifications
You must be signed in to change notification settings - Fork 116
Keep a multipart upload whose abort fails so that discard() retries it #1047
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
baf1f8b
f39c21c
0deeb8a
ebf798f
56bfa0b
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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. | ||
|
|
@@ -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. | ||
|
|
@@ -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. | ||
|
|
@@ -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 | ||
| _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: | ||
|
|
@@ -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: | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review round one, expanded scope (implementation behavior): CLEAN Range Covered:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Correction: after the rebase, the branch base (merge base with master) is
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review, expanded scope (relayed): Codex CLI 0.160.0, model Range 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, FINDINGS
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 |
||
| # _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, | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review round one (implementation behavior) — FINDINGS (repaired) Base Covered:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Round one on the repair ( |
||
| ) | ||
| 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: | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review (relayed): Codex CLI 0.160.0, model Base Result: FINDINGS, both pre-existing, not introduced by this diff: Reviewer output (verbatim)Covered FINDINGS — two pre-existing gaps remain; neither is introduced by this diff.
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 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 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):
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent follow-up (relayed): Codex CLI 0.160.0, model 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 |
||
| _logger.exception( | ||
| f"Failed to abort multipart upload {upload_id} " | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review round two (claims, compatibility, operations) — FINDINGS (repaired) Base Claims checked:
Findings (repaired in 5e09af9):
Limitation: no live S3 run; abort failures are simulated with mocks. AWS CI comes after Ready.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Round two on the repair ( |
||
| f"to s3://{self.bucket}/{self.key}." | ||
| ) | ||
| raise | ||
|
|
||
| self.fs.invalidate_cache(self.path) | ||
|
|
||
There was a problem hiding this comment.
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.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 beforeinvalidate_cache(). Nothing was written, so no listing changes.Transaction.complete(commit=True)callscommit(), which returns, so the file is not discarded. Rollback callsdiscard().docs/filesystem.md:104("raisesValueErrorand aborts its multipart upload") is still true; the upload is only kept when that abort fails. The incomplete-upload section (:214) still applies.abortcheck precedes_abort_multipart_upload()(s3.py:2338). The body is corrected.