-
Notifications
You must be signed in to change notification settings - Fork 116
Abort a cancelled async multipart copy after its running parts #1056
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
09b8aa6
e35aa6b
50fe7cb
28f0b6a
cd0c4a4
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 |
|---|---|---|
|
|
@@ -23,6 +23,7 @@ | |
| from pyathena.filesystem.s3 import CompressedBuffer, S3File, S3FileSystem | ||
| from pyathena.filesystem.s3_executor import S3AioExecutor, S3Executor, S3ThreadPoolExecutor | ||
| from pyathena.filesystem.s3_object import ( | ||
| S3CompleteMultipartUpload, | ||
| S3Metadata, | ||
| S3MultipartUpload, | ||
| S3Object, | ||
|
|
@@ -37,6 +38,8 @@ | |
| from pyathena.connection import Connection | ||
|
|
||
| _logger = logging.getLogger(__name__) | ||
| # The running cleanups of cancelled multipart copies. | ||
| _cleanup_tasks: set[asyncio.Future[None]] = set() | ||
|
|
||
|
|
||
| class AioS3FileSystem(AsyncFileSystem): | ||
|
|
@@ -485,8 +488,12 @@ async def _copy_object_with_multipart_upload( | |
| """Copy an object with a multipart upload of its byte ranges. | ||
|
|
||
| See :meth:`S3FileSystem._copy_object_with_multipart_upload`. The part | ||
| and annotation copies run in parallel with ``asyncio.gather`` and | ||
| ``asyncio.to_thread``. | ||
| and annotation copies run in parallel as asyncio tasks with | ||
| ``asyncio.to_thread``. On a cancellation after the upload is | ||
| created, the running part copies and the completion are waited for, | ||
|
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 review (relayed), static review only. Reviewer: Codex CLI 0.160.0, model Finding 4 (P3): the docstring overstated the guarantee. A cancellation during the unshielded CreateMultipartUpload await can let the worker create an upload whose ID is never kept, and no cleanup runs. That leak is pre-existing, but the blanket claim was new. Repaired in 28f0b6a: the docstring now says "On a cancellation after the upload is created". The PR description already lists the creation window under Not handled. Self-review of these repairs: they change one test and one docstring, and add no new behavior. Both rounds: CLEAN. The new assertions observe the documented wait, and the docstring and the PR text now state the same scope. |
||
| the upload is aborted unless it has completed, and the cancellation | ||
| is re-raised. A repeated cancellation returns without stopping this | ||
| cleanup. | ||
|
|
||
| Args: | ||
| bucket1: Source S3 bucket name. | ||
|
|
@@ -589,25 +596,55 @@ async def _upload_part(i: int, range_: tuple[int, int]) -> dict[str, Any] | None | |
| } | ||
|
|
||
| tasks = [asyncio.ensure_future(_upload_part(i, r)) for i, r in enumerate(ranges)] | ||
| try: | ||
| # gather keeps the part-number order of the tasks. | ||
| parts = await asyncio.gather(*tasks) | ||
| completed = await asyncio.to_thread( | ||
| self._sync_fs._complete_multipart_upload, | ||
| bucket=bucket2, | ||
| key=key2, | ||
| upload_id=upload_id, | ||
| parts=cast(list[dict[str, Any]], parts), | ||
| **self._sync_fs._get_operation_kwargs("complete_multipart_upload", kwargs), | ||
| ) | ||
| except Exception: | ||
| failed = True | ||
| completion: asyncio.Task[S3CompleteMultipartUpload] | None = None | ||
|
|
||
| async def _abort() -> 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 of the repair (both perspectives), base Round one (behavior): CLEAN.
Round two (claims and operations): the docstring claims hold within the stated limits, and the PR description is rewritten for the final behavior. A cancelled completion that succeeded leaves the object without its annotations; the description now states this. The cancellation can take as long as the slowest running part or the completion. The loop-shutdown and CreateMultipartUpload limits are listed. Evidence: |
||
| # A part that is still copying when the upload is aborted may be | ||
| # stored after the abort, so wait for the running parts first. | ||
| await asyncio.gather(*tasks, return_exceptions=True) | ||
| if completion is not None: | ||
| await asyncio.wait([completion]) | ||
| if not completion.cancelled() and completion.exception() is None: | ||
| # The upload completed despite the cancellation, so | ||
| # there is nothing to abort. | ||
| return | ||
| await asyncio.to_thread( | ||
| self._sync_fs._abort_multipart_upload, bucket2, key2, upload_id, kwargs | ||
| ) | ||
|
|
||
| try: | ||
| # Unlike gather, wait does not cancel the parts when this task is | ||
| # cancelled; their threads would keep copying, so they are waited | ||
| # for in _abort(). | ||
| done, _ = await asyncio.wait(tasks, return_when=asyncio.FIRST_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. Self-review round one, checked: CLEAN.
|
||
| for task in done: | ||
| if (error := task.exception()) is not None: | ||
| raise error | ||
| # The tasks are in part-number order. | ||
| parts = [task.result() for task in tasks] | ||
| completion = asyncio.ensure_future( | ||
| asyncio.to_thread( | ||
| self._sync_fs._complete_multipart_upload, | ||
| bucket=bucket2, | ||
| key=key2, | ||
| upload_id=upload_id, | ||
| parts=cast(list[dict[str, Any]], parts), | ||
| **self._sync_fs._get_operation_kwargs("complete_multipart_upload", kwargs), | ||
| ) | ||
| ) | ||
| # shield keeps a cancellation from cancelling the completion, whose | ||
| # thread would keep running, so that _abort() can wait for it. | ||
| completed = await asyncio.shield(completion) | ||
|
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), static review only. Reviewer: Codex CLI 0.160.0, model Finding 2 (P2): the completion-cancellation guarantee was too strong. Reviewer's scenario: cancel while Verified. The 185115d version fails the new |
||
| except BaseException: | ||
|
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, callers, and operations): FINDINGS, repaired in the PR description. Scope: full pass over base Finding: the description listed "a cancellation during CompleteMultipartUpload sends the abort while the completion may still run" under Unchanged. That is false. Before this PR, Also added to the description: a cancelled copy now returns only after the running part copies finish, which can take up to the slowest running part copy ( Claims checked and held:
Evidence limits, as stated in the PR: the AWS test run belongs to 43b3fe4 (185115d changes only the docstring). No real-S3 cancellation and no test for cancellation during completion were run. |
||
| # Also on cancellation, as S3FileSystem._finish_multipart_upload | ||
| # does on an interrupt. | ||
| failed = True | ||
| cleanup = asyncio.ensure_future(_abort()) | ||
| # A repeated cancellation of this task returns without stopping | ||
| # the cleanup; the event loop keeps only weak references to tasks. | ||
| _cleanup_tasks.add(cleanup) | ||
|
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 review (relayed), static review only. Reviewer: Codex CLI 0.160.0, model Finding 1 (P2): shutdown can abort before the worker threads finish. At event-loop shutdown, the part and completion tasks are cancelled directly, and Deferred, documented. On master, shutdown during a multipart copy never aborts, so the whole upload and all its parts are left behind. With this PR, the cleanup that runs aborts the upload, and only a part still copying at that moment can remain. That is no worse than master. Tracking the worker threads apart from their asyncio wrappers would need a different executor design, which this cancellation fix does not call for. The PR description's Not handled list now describes both shutdown outcomes. Finding 2 (P2): exceptions of a detached cleanup are never retrieved. After a repeated cancellation detaches the cleanup, a manually managed loop may shut down its default executor. The abort's Deferred, documented. This needs a loop that shuts down its default executor while a detached cleanup still runs. Even then, the error is not silent: asyncio logs it as "Task exception was never retrieved" with its traceback. The PR description lists it under Not handled. |
||
| cleanup.add_done_callback(_cleanup_tasks.discard) | ||
| await asyncio.shield(cleanup) | ||
|
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), static review only. Reviewer: Codex CLI 0.160.0, model Finding 1 (P2): a repeated cancellation can skip the abort. Reviewer's scenario: two part copies are running. Cancel the copy, then cancel it again while the cleanup awaits Verified. The 185115d version fails the new |
||
| raise | ||
|
|
||
| failed = False | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -356,6 +356,151 @@ def upload_part_copy(**kw): | |
| assert events[2:] == ["end 1", "abort"] | ||
| sync_fs._complete_multipart_upload.assert_not_called() | ||
|
|
||
| @pytest.mark.parametrize("cancellations", [1, 2]) | ||
| @pytest.mark.asyncio | ||
| async def test_copy_object_with_multipart_upload_cancelled(self, cancellations): | ||
|
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), static review only. Reviewer: Codex CLI 0.160.0, model Finding 3 (P2): the new test could fail with correct code. The semaphore proved that two parts had started, but a 200 ms sleep was all that kept them running until the cancellation. Under a scheduling delay, the parts could finish, part 3 could start, or the copy could complete before Verified. Repaired in 50fe7cb. The parts block on a |
||
| # GH-1046: a cancellation waits for the part copies that are running, | ||
| # aborts the upload, and is re-raised, as S3FileSystem does on an | ||
| # interrupt. A repeated cancellation returns without stopping the | ||
| # cleanup. | ||
| fs = AioS3FileSystem(connection=mock.MagicMock(), max_workers=2, skip_instance_cache=True) | ||
| sync_fs = fs._sync_fs | ||
| sync_fs._create_multipart_upload = mock.MagicMock( | ||
| return_value=SimpleNamespace(upload_id="uploadid") | ||
| ) | ||
| events = [] | ||
| lock = threading.Lock() | ||
| started = threading.Semaphore(0) | ||
| release = threading.Event() | ||
| aborted = threading.Event() | ||
|
|
||
| def upload_part_copy(**kw): | ||
| part_number = kw["part_number"] | ||
| with lock: | ||
| events.append(f"start {part_number}") | ||
| started.release() | ||
| # The finally blocks of the test always release it. | ||
| release.wait() | ||
| with lock: | ||
| events.append(f"end {part_number}") | ||
| return SimpleNamespace(etag='"e"', part_number=part_number) | ||
|
|
||
| def abort_multipart_upload(*args): | ||
| events.append("abort") | ||
| aborted.set() | ||
|
|
||
| sync_fs._upload_part_copy = mock.MagicMock(side_effect=upload_part_copy) | ||
| sync_fs._complete_multipart_upload = mock.MagicMock() | ||
| # The HeadObject of the source, for its version. | ||
| sync_fs._call = mock.MagicMock(return_value={}) | ||
| sync_fs._abort_multipart_upload = mock.MagicMock(side_effect=abort_multipart_upload) | ||
|
|
||
| task = asyncio.ensure_future( | ||
| fs._copy_object_with_multipart_upload( | ||
| bucket1="bucket", | ||
| key1="src", | ||
| size1=3 * S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE, | ||
| bucket2="bucket", | ||
| key2="dst", | ||
| block_size=S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE, | ||
| MetadataDirective="REPLACE", | ||
| TaggingDirective="REPLACE", | ||
| AnnotationDirective="EXCLUDE", | ||
| ) | ||
| ) | ||
| try: | ||
| for _ in range(2): | ||
| assert await asyncio.to_thread(started.acquire, timeout=5) | ||
| # The running parts are held until the copy has been cancelled. | ||
| for _ in range(cancellations): | ||
| task.cancel() | ||
| # Lets the copy enter its cleanup. | ||
| await asyncio.sleep(0) | ||
| if cancellations > 1: | ||
| # The repeated cancellation returns while the parts still run. | ||
| assert task.done() | ||
| assert events[2:] == [] | ||
| else: | ||
| assert not task.done() | ||
| except BaseException: | ||
| task.cancel() | ||
| raise | ||
| finally: | ||
| release.set() | ||
| with pytest.raises(asyncio.CancelledError): | ||
| await task | ||
| assert await asyncio.to_thread(aborted.wait, 5) | ||
|
|
||
| # Part 3 waits for a worker and is not started after the cancellation. | ||
| assert sorted(events[:2]) == ["start 1", "start 2"] | ||
| assert sorted(events[2:4]) == ["end 1", "end 2"] | ||
| assert events[4:] == ["abort"] | ||
| sync_fs._complete_multipart_upload.assert_not_called() | ||
|
|
||
| @pytest.mark.parametrize("completion_fails", [False, True]) | ||
| @pytest.mark.asyncio | ||
| async def test_copy_object_with_multipart_upload_cancelled_completion(self, completion_fails): | ||
| # GH-1046: a cancellation during CompleteMultipartUpload waits for it, | ||
| # aborts the upload only if it failed, and is re-raised. | ||
| fs = AioS3FileSystem(connection=mock.MagicMock(), skip_instance_cache=True) | ||
| sync_fs = fs._sync_fs | ||
| sync_fs._create_multipart_upload = mock.MagicMock( | ||
| return_value=SimpleNamespace(upload_id="uploadid") | ||
| ) | ||
| events = [] | ||
| started = threading.Event() | ||
| release = threading.Event() | ||
|
|
||
| def complete_multipart_upload(**kw): | ||
| started.set() | ||
| # The finally blocks of the test always release it. | ||
| release.wait() | ||
|
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 narrow follow-up (relayed), static review only. Reviewer: Codex CLI 0.160.0, model Finding (P3): the hold depended on elapsed time. The mock ignored the result of Repaired in cd0c4a4, in both cancellation tests. The mocks call |
||
| events.append("complete") | ||
| if completion_fails: | ||
| raise OSError("completion failed") | ||
| return SimpleNamespace() | ||
|
|
||
| sync_fs._upload_part_copy = mock.MagicMock( | ||
| side_effect=lambda **kw: SimpleNamespace(etag='"e"', part_number=kw["part_number"]) | ||
| ) | ||
| sync_fs._complete_multipart_upload = mock.MagicMock(side_effect=complete_multipart_upload) | ||
| # The HeadObject of the source, for its version. | ||
| sync_fs._call = mock.MagicMock(return_value={}) | ||
| sync_fs._abort_multipart_upload = mock.MagicMock( | ||
| side_effect=lambda *args: events.append("abort") | ||
| ) | ||
|
|
||
| task = asyncio.ensure_future( | ||
| fs._copy_object_with_multipart_upload( | ||
| bucket1="bucket", | ||
| key1="src", | ||
| size1=2 * S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE, | ||
| bucket2="bucket", | ||
| key2="dst", | ||
| block_size=S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE, | ||
| MetadataDirective="REPLACE", | ||
| TaggingDirective="REPLACE", | ||
| AnnotationDirective="EXCLUDE", | ||
| ) | ||
| ) | ||
| try: | ||
| assert await asyncio.to_thread(started.wait, 5) | ||
| task.cancel() | ||
| # Gives the cleanup time to abort early or to return, which it | ||
| # must not do while the completion is held. | ||
| await asyncio.sleep(0.1) | ||
|
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 narrow follow-up (relayed), static review only. Reviewer: Codex CLI 0.160.0, model Finding (P2): the cleanup was not synchronized before the completion was released. A single Repaired in cd0c4a4. The completion stays held for
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 narrow follow-up of the repair (relayed), static review only: CLEAN. Reviewer: Codex CLI 0.160.0, model Reviewer's notes: both |
||
| assert not task.done() | ||
|
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 review (relayed), static review only. Reviewer: Codex CLI 0.160.0, model Finding 3 (P3): the successful-completion case did not prove that the copy waits. The test released the completion right after Repaired in 28f0b6a. The completion stays held while the test cancels the copy, yields once so that the cancellation is processed, and asserts |
||
| assert events == [] | ||
| except BaseException: | ||
| task.cancel() | ||
| raise | ||
| finally: | ||
| release.set() | ||
| with pytest.raises(asyncio.CancelledError): | ||
| await task | ||
|
|
||
| assert events == (["complete", "abort"] if completion_fails else ["complete"]) | ||
|
|
||
| @pytest.mark.parametrize( | ||
| "block_size", | ||
| [ | ||
|
|
||
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 one (behavior and implementation): FINDINGS, repaired.
Scope: base
e49d576ce1e8cb60d74e1fde56de9c72ff3cd8e1(merge-base with master), head43b3fe471fdcf4ccf657a2383ea6df284e18e8c4, both changed files:pyathena/filesystem/s3_async.pyandtests/pyathena/filesystem/test_s3_async.py.Finding: this docstring still said that the part copies run in parallel with
asyncio.gather. After this PR they run as tasks that are awaited withasyncio.wait, so the docstring no longer matched the code. It also did not mention the new behavior on cancellation.Repair in 185115d: the docstring now says the copies run as asyncio tasks with
asyncio.to_thread. It also states that a cancellation while the parts are copied or the upload is completed waits for the running part copies, aborts the upload, and is re-raised. The wording leaves out a cancellation during CreateMultipartUpload, which this PR does not handle (see the PR description).just lintpassed.