Skip to content

Wait for running parts before aborting multipart uploads, and validate S3File before init - #992

Merged
laughingman7743 merged 4 commits into
masterfrom
fix/976-multipart-abort-in-flight
Oct 3, 2026
Merged

laughingman7743 merged 4 commits into
masterfrom
fix/976-multipart-abort-in-flight

Conversation

@laughingman7743

@laughingman7743 laughingman7743 commented Oct 3, 2026 •

Copy link
Copy Markdown
Member

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.

  • Abort after the running parts: when a part or the completion fails, S3FileSystem._finish_multipart_upload() cancels the parts that have not started, waits with concurrent.futures.wait() for the parts that could not be cancelled (running or done), and then sends AbortMultipartUpload. 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 an S3AioExecutor blocked by the caller's event loop thread could never do. This covers S3File writes and appends, cp_file()'s multipart copy, and AioS3File, which inherits both.
  • S3AioExecutor futures: the future returned by run_coroutine_threadsafe() could be cancelled while its function kept running in its thread, so wait() returned at once and AioS3File still aborted while parts were in flight. submit() now returns a future that the function's thread starts and resolves, so, as with ThreadPoolExecutor, 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, and cat for a small object), before AbstractBufferedFile.__init__. A failed open() therefore sends no request for a ValueError and leaves no half-initialized file whose garbage collection would close it.
  • Write block size upper bound (S3File accepts a write block size larger than the maximum part size #952): a write block size above MULTIPART_UPLOAD_MAX_PART_SIZE (5 GiB) now raises ValueError like one below 5 MiB. This applies to open()'s block_size and the filesystem's default_block_size.
  • Block size messages (S3 filesystem block size error messages do not match their checks #926): S3File and 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: docs/filesystem.md states the write block size range.

Release-note items:

  • An open() for writing with a block_size (or default_block_size) above 5 GiB now raises ValueError. 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 single PutObject and now fails at open() as well. pipe_file() of data that fits in one PutObject does not open a file and is not affected.
  • The block size ValueError messages changed.
  • A failed multipart upload or copy, and a rolled-back transaction, now wait for the parts that are already uploading before aborting, so the error is raised after those requests end. Each part request is bounded by the botocore timeouts and the filesystem's retry configuration, as before.
  • S3AioExecutor futures 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 had closed == False, so its garbage collection called close(), which flushed and raised AttributeError: 'S3File' object has no attribute 'append_block'. Defining the attributes earlier is not enough: as the issue notes, a failed file with them reaches commit() and creates an empty object. An append whose exists()/info() lookup failed (for example, PermissionError) had the same half-initialized state, with multipart_upload missing.
  • Since the issue was filed (at e0e85da), PR Accept version_id in S3FileSystem.open() and cat_file() #958 (bbf3759) moved the path and version checks before the base class initializer, so on the base of this PR a missing key no longer leaves a half-initialized file; the block size check and the append lookups still did. The tests keep the missing-key case for the no-request contract, and it passes on the base as well.

Not changed here:

TEST

Tested commit: 5060d5a unless noted (5060d5a changes only tests after 35b07a9).

  • just format, just lint, just docs lint: pass.
  • New offline tests: 3 consecutive runs pass.
  • Offline unit tests (dummy region and credentials, --noconftest):
    • TestS3FileSystem::test_finish_multipart_upload_waits_for_running_parts and TestS3File::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 production wait(), which records the waited futures ([failed, running], [running]); with only the wait() call removed, both tests fail.
    • test_finish_multipart_upload_does_not_wait_for_cancelled_parts (a never-started future) and TestS3File::test_discard_on_event_loop_thread (a rollback on the event loop thread with S3AioExecutor parts): no blocking, the parts are cancelled, and the abort is sent.
    • TestS3AioExecutor (new tests/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 no not_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 after gc.collect(). test_open_append_lookup_failure: an append whose exists() 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.
    • With the source changes reverted, all of these fail except the missing-key cases (see WHY), test_finish_multipart_upload_does_not_wait_for_cancelled_parts, and test_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 with test_loop_shutdown.
    • The whole 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.
  • Issue reproduction: the abort is now sent after parts 2 and 3 are stored (['create_multipart_upload', 'part 2 stored', 'part 3 stored', 'abort_multipart_upload']), and the failed open() calls print no "Exception ignored" line and send no request.
  • AWS integration tests (local, against the CI account): uv run --env-file .env pytest -n 4 -p no:cacheprovider tests/pyathena/filesystem/ at 35b07a9: 300 passed (at b2c0fe7: 298 passed).
  • Not covered: the in-flight abort against real S3 (it needs a part to fail while others are uploading); the tests use mocked requests with real thread pools.

🤖 Generated with Claude Code

laughingman7743 and others added 2 commits October 3, 2026 16:59
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():

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) — 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() after cancel(): 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, so wait()/result() no longer depend on the loop running; a task cancelled or failed before the function starts settles the future (covered by test_loop_shutdown and test_executor_shut_down). settle cannot see a RUNNING future with a task exception, because run catches BaseException.
  • S3File.__init__: the append lookups use path instead of self.path and no version_id; a version is rejected for writing earlier, so the requests are the same as before. self.loc is still overwritten with the existing size after write().
  • Live: uv run --env-file .env pytest -n 4 tests/pyathena/filesystem/ at b2c0fe7: 298 passed.

Comment thread pyathena/filesystem/s3.py
# 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 (

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: 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() with default_block_size > 5 GiB now fail at open() 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.py pipe_file routes on min(block_size, MULTIPART_UPLOAD_MAX_PART_SIZE)). Added this to the release-note items. Read-mode block sizes (e.g. S3FSResultSet's default_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 in docs/ 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>
Comment thread pyathena/filesystem/s3.py Outdated
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)

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 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.

  1. P1 (introduced) — waiting on cancelled futures can deadlock cleanup (s3.py:1443, s3.py:2532, s3_executor.py:131). Future.cancel() leaves a future CANCELLED, and wait() returns for it only after set_running_or_notify_cancel(), which S3AioExecutor calls only when run() executes. (a) An AioS3File rollback on the event loop thread blocks in wait(), since the loop cannot dispatch run(); the same when the rollback occupies the only default-executor worker. (b) On loop shutdown, settle() only calls future.cancel(), so the future is never acknowledged and a later wait() hangs; test_loop_shutdown ignored wait()'s not_done.
  2. 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, and test_loop_shutdown): a descheduled controlling thread lets the pending function start before it is cancelled.
  3. P2 (pre-existing) — AioS3FileSystem._copy_object_with_multipart_upload() (s3_async.py:321) uses gather() 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.

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 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() and S3File.discard() wait only for the futures that cancel() could not cancel (running or done): wait([f for f in futures if not f.cancel()]).
  • S3AioExecutor.submit(): a one-shot threading.Lock lets run() or settle(), whichever comes first, claim the future. settle() on a cancelled task cancels the future and calls set_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_shutdown asserts wait() returns no not_done; new test_finish_multipart_upload_does_not_wait_for_cancelled_parts and TestS3File::test_discard_on_event_loop_thread (rollback on the loop thread with S3AioExecutor parts, 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.

@laughingman7743 laughingman7743 Oct 3, 2026 •

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 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.

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 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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

1 participant