Wait for running parts before aborting multipart uploads, and validate S3File before init - #992
Conversation
A failed multipart upload or copy cancelled the remaining parts and aborted the upload at once. Future.cancel() does not stop a running part, so an UploadPart or UploadPartCopy still in flight could be stored after the abort. _finish_multipart_upload() and S3File.discard() now wait for the running parts before the abort. The futures of S3AioExecutor (AioS3File) could be cancelled while their function kept running, so wait() returned at once. The future is now started and resolved by the function's thread, so that, as with ThreadPoolExecutor, it cannot be cancelled once running. S3File.__init__ checked the write block size after the base class initializer, so a failed open() left a half-initialized file whose garbage collection closed it and raised AttributeError. The arguments are now validated, and the existing object of an append looked up, before the base class initializer. The write block size must also be at most 5 GiB, the maximum part size, and the block size error messages state the accepted range. Closes #976, closes #926, closes #952. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| executor = S3AioExecutor(loop=loop) | ||
| running = executor.submit(work) | ||
| pending = executor.submit(events.append, "pending finished") | ||
| while not started.is_set(): |
There was a problem hiding this comment.
Self-review round 1 (implementation behavior) — FINDINGS (1, repaired)
Scope: git diff aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe..b2c0fe7ece79402433e1e39ebbad37ae175a50d4 (base = merge-base with master, head = published head). Covered: _finish_multipart_upload() and its callers (S3File.commit(), sync multipart copy), S3File.discard() (transaction rollback, #968's part-limit abort calls it too), S3AioExecutor.submit() and its only caller AioS3File (_fetch_range via as_completed, part uploads), S3File.__init__ read/write/append paths, the copy block size checks (sync and async), docs, and the new tests.
Finding: test_loop_shutdown returned from main() after await asyncio.sleep(0), which does not guarantee that the first function had started. If asyncio.run() cancelled its task before it started, running would be cancelled as well and the test would fail intermittently. Repaired in 0191c1f: main() waits for a threading.Event set by the function. 5 repeated offline runs pass.
Checked without findings:
wait()aftercancel(): queued parts are cancelled (never run), running parts are waited for; a completion failure has all futures done already.S3AioExecutor: the returned future is resolved by the worker thread, not by an event loop callback, sowait()/result()no longer depend on the loop running; a task cancelled or failed before the function starts settles the future (covered bytest_loop_shutdownandtest_executor_shut_down).settlecannot see a RUNNING future with a task exception, becauseruncatchesBaseException.S3File.__init__: the append lookups usepathinstead ofself.pathand noversion_id; a version is rejected for writing earlier, so the requests are the same as before.self.locis still overwritten with the existing size afterwrite().- Live:
uv run --env-file .env pytest -n 4 tests/pyathena/filesystem/at b2c0fe7: 298 passed.
| # Carry the version in the path, as with the ?versionId= suffix, | ||
| # so that a reopened (e.g., unpickled) file reads the same version. | ||
| path = f"{path}?versionId={self.version_id}" | ||
| if "r" not in mode and not ( |
There was a problem hiding this comment.
Self-review round 2 (claims, callers, operations) — FINDINGS (PR description corrected; no code change)
Scope: git diff aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe..0191c1f683ff6ee0283d4cbdf431ecb8f876b87b, the PR body, commit messages, changed docstrings/comments, and docs/filesystem.md.
Claims checked:
- AbortMultipartUpload: "parts currently in progress might or might not succeed … it might be necessary to abort a given multipart upload multiple times" — confirmed in the botocore S3 service model shipped in
uv.lock. - 5 MiB to 5 GiB part size: confirmed in the S3 user guide's multipart upload limits table; PutObject single-operation limit of 5 GB confirmed in the "Uploading objects" guide. The PR body now says the >5 GiB part rejection is documented, not reproduced.
- Issue premise "missing-key check raises after the base initializer": true at e0e85da (issue base), but PR Accept version_id in S3FileSystem.open() and cat_file() #958 (bbf3759) moved it before
super().__init__. Reproduced on aa0fc91: the issue script prints one "Exception ignored" line (the block size file), not two. Corrected the WHY section; the missing-key test cases pass on the base and are kept for the no-request contract. - Caller impact of the new upper bound:
open()/put_file()withdefault_block_size> 5 GiB now fail atopen()even for small writes that used to succeed with a single PutObject;pipe_file()of data that fits in one PutObject never opens a file and is unaffected (s3.pypipe_fileroutes onmin(block_size, MULTIPART_UPLOAD_MAX_PART_SIZE)). Added this to the release-note items. Read-mode block sizes (e.g.S3FSResultSet'sdefault_block_size) are not checked. - Operator: the added wait is bounded per part by botocore timeouts and
_call's retry configuration, as the part request itself already was; stated in the release note. No new requests are added. docs/filesystem.md: no other statement about block size limits exists indocs/or the README.- Evidence: offline results are local; the live filesystem run (298 passed) is at b2c0fe7, and 0191c1f changes only an offline test. The in-flight abort is not exercised against real S3 (stated in TEST).
concurrent.futures.wait() returns for a cancelled future only after its executor acknowledges the cancellation, which S3AioExecutor did only when the event loop ran the function. A rollback on the event loop thread therefore blocked, and a future cancelled on loop shutdown was never acknowledged. Wait only for the parts that could not be cancelled, and let run() or settle(), whichever comes first, claim the future, with settle() acknowledging a cancellation. The tests release the running part when the pending one is cancelled instead of relying on a sleep. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| future.cancel() | ||
| # A part that is still uploading when the upload is aborted may | ||
| # be stored after the abort, so wait for the running parts first. | ||
| wait(futures) |
There was a problem hiding this comment.
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 01a100d4-75b5-7911-9d04-668a5b270311. Static review only: no tests, builds, edits, or network access. Snapshot: detached worktree at 0191c1f683ff6ee0283d4cbdf431ecb8f876b87b, diff git diff aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe..0191c1f683ff6ee0283d4cbdf431ecb8f876b87b; 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 multipart upload and copy, completion failures, abort ordering, executor cancellation; transactions, autocommit=False, rollback, event-loop/worker deadlocks; constructor validation, garbage collection, small/large append, read mode, block size boundaries; new tests, callers, docs, installed fsspec lifecycle.
- P1 (introduced) — waiting on cancelled futures can deadlock cleanup (
s3.py:1443,s3.py:2532,s3_executor.py:131).Future.cancel()leaves a futureCANCELLED, andwait()returns for it only afterset_running_or_notify_cancel(), whichS3AioExecutorcalls only whenrun()executes. (a) AnAioS3Filerollback on the event loop thread blocks inwait(), since the loop cannot dispatchrun(); the same when the rollback occupies the only default-executor worker. (b) On loop shutdown,settle()only callsfuture.cancel(), so the future is never acknowledged and a laterwait()hangs;test_loop_shutdownignoredwait()'snot_done. - P2 (introduced) — the new concurrency tests depend on a 100 ms window (
test_s3_executor.py:45,test_s3.py:720,test_s3.py:2443, andtest_loop_shutdown): a descheduled controlling thread lets the pending function start before it is cancelled. - P2 (pre-existing) —
AioS3FileSystem._copy_object_with_multipart_upload()(s3_async.py:321) usesgather()and never aborts on a part or completion failure.
Author verification: 1 and 2 confirmed (concurrent.futures.wait counts only CANCELLED_AND_NOTIFIED/FINISHED as done). 3 is pre-existing and is #973's second item; deferred to it.
There was a problem hiding this comment.
Repair 35b07a9 for findings 1 and 2, with both self-review perspectives on the repair
Range: git range-diff aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe..0191c1f683ff6ee0283d4cbdf431ecb8f876b87b aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe..35b07a9837396e1753fcda5477e36998534126ef (same base; one new commit).
Repair:
_finish_multipart_upload()andS3File.discard()wait only for the futures thatcancel()could not cancel (running or done):wait([f for f in futures if not f.cancel()]).S3AioExecutor.submit(): a one-shotthreading.Lockletsrun()orsettle(), whichever comes first, claim the future.settle()on a cancelled task cancels the future and callsset_running_or_notify_cancel()so waiters are notified, or sets the task's error.- Tests: running parts are released by the pending part's cancellation (done callback) instead of a 100 ms window;
test_loop_shutdownassertswait()returns nonot_done; newtest_finish_multipart_upload_does_not_wait_for_cancelled_partsandTestS3File::test_discard_on_event_loop_thread(rollback on the loop thread withS3AioExecutorparts, guarded by a 5 s join). The three fail on 0191c1f and pass on 35b07a9.
Behavior (round-1 perspective): the claim covers each race — user cancel() before run() (run() claims, set_running_or_notify_cancel() returns False and notifies); task cancelled while run() is queued (settle() claims; run() later returns at once); task cancelled while run() runs (run() already claimed; settle() returns; the thread resolves the future). ThreadPoolExecutor callers keep their behavior except that cancelled parts are no longer waited for; the copy path's with executor: still drains the queue on exit. A failed or finished future returns False from cancel() and is done for wait().
Claims (round-2 perspective): docstrings of _finish_multipart_upload(), discard() and S3AioExecutor still hold; the PR body now describes the narrowed wait and the claim, and lists the new tests and the 35b07a9 results.
Validation at 35b07a9: just lint pass; offline targeted tests 3 consecutive runs pass; live tests/pyathena/filesystem/ 300 passed.
Out of scope (pre-existing): S3File.commit() on the event loop thread still blocks in future.result() for S3AioExecutor parts that the loop has not started, as on the base. Finding 3 is deferred to #973.
There was a problem hiding this comment.
Independent follow-up (relayed) on 35b07a9 — FINDINGS (1), and repair 5060d5a
Reviewer: OpenAI Codex CLI 0.160.0, model gpt-6-astra, effort high, sandbox read-only; session 01a100dd-7ae4-7052-a396-01e41329eede. Static review of git diff 0191c1f6..35b07a98 (same base) at a detached snapshot of 35b07a9, which stayed unchanged.
Reviewer result: finding 1 resolved by source inspection — successful cancellations are excluded from both waits; worker/user-cancel and worker/callback interleavings each settle the future exactly once; the new cancelled-part, event-loop rollback and shutdown tests would expose the original blocking. No production regression and no behavior beyond the intended scope. One P2 remained (test_s3.py:725, test_s3.py:2484): the ordering tests released the running part on the pending part's cancellation plus a 50 ms sleep, so they could pass with the wait() removed if the aborting thread was descheduled.
Author verification: confirmed (a false pass, not a false failure of correct code).
Repair (tests only): both tests patch pyathena.filesystem.s3.wait with a wrapper that releases the running part and delegates to the real wait(); the sleep and cancellation callback are gone, and the tests assert the waited futures ([failed, running], [running]). Mutation check: with only the wait() call removed from both sites, both tests fail; restored, 3 consecutive offline runs pass; just lint passes. No production code changed, so the live run at 35b07a9 (300 passed) still covers the source.
There was a problem hiding this comment.
Independent follow-up (relayed) on 5060d5a — CLEAN
Reviewer: OpenAI Codex CLI 0.160.0, model gpt-6-astra, effort high, sandbox read-only; session 01a100e0-9075-7753-aa32-c4c70700c34e. Static review of git diff 35b07a98..5060d5a4 (test-only) at a detached snapshot of 5060d5a, which stayed unchanged.
Reviewer result: pyathena.filesystem.s3.wait is the right patch target and the wrappers delegate to the real wait() without recursion; correct code cancels the pending part before releasing and waiting for the running one; removing wait() fails the call assertion even if scheduling satisfies the event order, and omitting the running future fails the argument assertion; a production regression does not hang beyond the workers' 5 s fallback timeout. Limit noted by the reviewer: the 5 s timeout is a fallback release, so "released only by the wait" holds within that budget.
The running part was released by the cancellation of the pending one and a short sleep, so the tests could pass without the wait. A wrapper around wait() now releases it and records the waited futures. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
WHAT
Fix two error paths of the S3 multipart write and copy code, and the write block size validation that
S3File.__init__shares with them.S3FileSystem._finish_multipart_upload()cancels the parts that have not started, waits withconcurrent.futures.wait()for the parts that could not be cancelled (running or done), and then sendsAbortMultipartUpload.S3File.discard()(a transaction rollback) does the same. Cancelled parts are not waited for:wait()returns for them only after their executor acknowledges the cancellation, which anS3AioExecutorblocked by the caller's event loop thread could never do. This coversS3Filewrites and appends,cp_file()'s multipart copy, andAioS3File, which inherits both.S3AioExecutorfutures: the future returned byrun_coroutine_threadsafe()could be cancelled while its function kept running in its thread, sowait()returned at once andAioS3Filestill aborted while parts were in flight.submit()now returns a future that the function's thread starts and resolves, so, as withThreadPoolExecutor, it cannot be cancelled once running. The function's thread or, if the event loop task ends before the function starts (for example, on loop shutdown), the task's done callback claims the future exactly once; the callback cancels it and notifies its waiters, or sets the task's error, instead of leaving it pending.S3File.__init__validates first: the path, version and write block size are checked, and the existing object of an append is looked up (exists,info, andcatfor a small object), beforeAbstractBufferedFile.__init__. A failedopen()therefore sends no request for aValueErrorand leaves no half-initialized file whose garbage collection would close it.MULTIPART_UPLOAD_MAX_PART_SIZE(5 GiB) now raisesValueErrorlike one below 5 MiB. This applies toopen()'sblock_sizeand the filesystem'sdefault_block_size.S3Fileand both_copy_object_with_multipart_upload()now say "between 5 MiB (5242880 bytes) and 5 GiB (5368709120 bytes), inclusive" and include the given value.docs/filesystem.mdstates the write block size range.Release-note items:
open()for writing with ablock_size(ordefault_block_size) above 5 GiB now raisesValueError. Before, a write that filled a block submitted a part larger than the 5 GiB maximum part size of the multipart upload limits (documented; not reproduced). A write smaller than the block size and at most 5 GB used to succeed with a singlePutObjectand now fails atopen()as well.pipe_file()of data that fits in onePutObjectdoes not open a file and is not affected.ValueErrormessages changed.S3AioExecutorfutures can no longer be cancelled once their function has started.WHY
Closes #976, closes #926, closes #952.
Future.cancel()does not stop a running part. The AbortMultipartUpload documentation states that parts still uploading during the abort might succeed and that the upload may have to be aborted again to free their storage. Nothing aborted again, so such parts outlived the abort.S3File.__init__checked the write block size after the base class initializer. The failed file hadclosed == False, so its garbage collection calledclose(), which flushed and raisedAttributeError: 'S3File' object has no attribute 'append_block'. Defining the attributes earlier is not enough: as the issue notes, a failed file with them reachescommit()and creates an empty object. An append whoseexists()/info()lookup failed (for example,PermissionError) had the same half-initialized state, withmultipart_uploadmissing.Not changed here:
S3File.close()and the write path. Whichever merges second rebases on the other.exceptblock in_finish_multipart_upload()andS3File.commit().TEST
Tested commit: 5060d5a unless noted (5060d5a changes only tests after 35b07a9).
just format,just lint,just docs lint: pass.--noconftest):TestS3FileSystem::test_finish_multipart_upload_waits_for_running_partsandTestS3File::test_discard_waits_for_running_parts: a running part finishes before the abort, and a queued part is cancelled. The running part is released by a wrapper around the productionwait(), which records the waited futures ([failed, running],[running]); with only thewait()call removed, both tests fail.test_finish_multipart_upload_does_not_wait_for_cancelled_parts(a never-started future) andTestS3File::test_discard_on_event_loop_thread(a rollback on the event loop thread withS3AioExecutorparts): no blocking, the parts are cancelled, and the abort is sent.TestS3AioExecutor(newtests/pyathena/filesystem/test_s3_executor.py): results and errors, a running future that cannot be cancelled and is waited for, a pending future cancelled and acknowledged on event loop shutdown (wait()returns nonot_done), and the task's error set on the future when the default executor is shut down.TestS3FileSystem::test_open_invalid_for_writing(wb/ab/xb× block size 5 MiB − 1, 5 GiB + 1, missing key):ValueError, no request, and no unraisable exception aftergc.collect().test_open_append_lookup_failure: an append whoseexists()fails leaves no unraisable exception.test_open_block_size_limits_for_writing: 5 MiB and 5 GiB are accepted.test_copy_object_with_multipart_upload_invalid_block_size(sync and async): the new message.test_finish_multipart_upload_does_not_wait_for_cancelled_parts, andtest_discard_on_event_loop_thread, which pass on the base as well (the base does not wait). The last two fail on 0191c1f, before the repair, together withtest_loop_shutdown.tests/pyathena/filesystem/directory offline has the same failures as master (only the AWS integration tests, which need credentials), plus the new tests failing on master.['create_multipart_upload', 'part 2 stored', 'part 3 stored', 'abort_multipart_upload']), and the failedopen()calls print no "Exception ignored" line and send no request.uv run --env-file .env pytest -n 4 -p no:cacheprovider tests/pyathena/filesystem/at 35b07a9: 300 passed (at b2c0fe7: 298 passed).🤖 Generated with Claude Code