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
23 changes: 23 additions & 0 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -348,6 +348,29 @@ instead. `copy_object_annotation()` copies one annotation onto the destination a
the upload completes, with GetObjectAnnotation and PutObjectAnnotation. The
filesystems' `cp_file()`, `copy()` and `mv()` run these plans.

The multipart primitives take the `S3MultipartUpload` returned by creation,
keeping its bucket, key, upload ID and checksum configuration together.
The core retains no upload state.
`upload_part()` passes the upload's checksum algorithm to the SDK;
completion selects the matching part checksum and sends the upload's
`ChecksumType` when present, including `FULL_OBJECT`.
Without a creation algorithm, completion sends the ETag and part number without
checksums that the SDK may add to part uploads.
`abort_multipart_upload()` also accepts uploads returned by
`list_multipart_uploads()`.

```python
upload = core.create_multipart_upload(
S3Path("YOUR_S3_BUCKET", "path/to/object"), ChecksumAlgorithm="SHA256"
)
try:
part = core.upload_part(upload, 1, b"data")
result = core.complete_multipart_upload(upload, [part])
except Exception:
core.abort_multipart_upload(upload)
raise
```

## Async filesystem

`AioS3FileSystem` provides the same functionality on top of fsspec's
Expand Down
75 changes: 36 additions & 39 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -1784,18 +1784,14 @@ def _copy_object_with_multipart_upload(
wait([creation])
if creation.exception() is None:
self._abort_multipart_upload(
plan.destination.bucket,
cast(str, plan.destination.key),
cast(str, creation.result().upload_id),
creation.result(),
plan.abort_params,
)
raise
upload_id = cast(str, multipart_upload.upload_id)
futures = [
executor.submit(
self.core.upload_part_copy,
path=plan.destination,
upload_id=upload_id,
upload=multipart_upload,
part_number=i + 1,
source=plan.source,
range_=range_,
Expand All @@ -1804,9 +1800,7 @@ def _copy_object_with_multipart_upload(
for i, range_ in enumerate(plan.ranges)
]
completed = self._finish_multipart_upload(
bucket=plan.destination.bucket,
key=cast(str, plan.destination.key),
upload_id=upload_id,
upload=multipart_upload,
futures=futures,
# Filtered again for the completion and the abort, which
# leaves the plan's parameters of each unchanged.
Expand Down Expand Up @@ -1930,9 +1924,7 @@ def pipe_file(

def _finish_multipart_upload(
self,
bucket: str,
key: str,
upload_id: str,
upload: S3MultipartUpload,
futures: list[Future[S3MultipartUploadPart]],
request_kwargs: Mapping[str, Any] | None = None,
abort: bool = True,
Expand All @@ -1946,9 +1938,7 @@ def _finish_multipart_upload(
false. The original error is then re-raised.

Args:
bucket: S3 bucket name.
key: Object key being uploaded.
upload_id: Unique identifier for the multipart upload.
upload: The multipart upload returned by creation.
futures: Futures of the part uploads, in part-number order.
request_kwargs: Parameters of the upload, such as
``RequestPayer`` or the SSE-C parameters; the completion and
Expand All @@ -1964,8 +1954,7 @@ def _finish_multipart_upload(
# The futures are in part-number order.
parts = [future.result() for future in futures]
return self.core.complete_multipart_upload(
S3Path(bucket, key),
upload_id,
upload,
parts,
**self.core.operation_params("complete_multipart_upload", request_kwargs),
)
Expand All @@ -1976,33 +1965,31 @@ def _finish_multipart_upload(
# be stored after the abort, so wait for the parts that could not
# be cancelled first.
wait([future for future in futures if not future.cancel()])
self._abort_multipart_upload(bucket, key, upload_id, request_kwargs)
self._abort_multipart_upload(upload, request_kwargs)
raise

def _abort_multipart_upload(
self, bucket: str, key: str, upload_id: str, request_kwargs: Mapping[str, Any]
self, upload: S3MultipartUpload, request_kwargs: Mapping[str, Any]
) -> None:
"""Abort a failed multipart upload, logging an error of the abort.

The error of the abort is logged instead of raised, so that the
caller can re-raise the error that made the upload fail.

Args:
bucket: S3 bucket name.
key: Object key being uploaded.
upload_id: Unique identifier for the multipart upload.
upload: The multipart upload returned by creation.
request_kwargs: Parameters of the upload; the abort receives
those that it accepts.
"""
try:
self.core.abort_multipart_upload(
S3Path(bucket, key),
upload_id,
upload,
**self.core.operation_params("abort_multipart_upload", request_kwargs),
)
except Exception:
_logger.exception(
f"Failed to abort multipart upload {upload_id} to s3://{bucket}/{key}."
f"Failed to abort multipart upload {upload.upload_id} "
f"to s3://{upload.bucket}/{upload.key}."
)

def cat_file(
Expand Down Expand Up @@ -2623,6 +2610,9 @@ def object_version_info(
def clear_multipart_uploads(self, path: str) -> None:
"""Abort any incomplete multipart uploads in the bucket.

Uploads completed or aborted after listing are already cleared.
All abort results are checked before another abort error is raised.

Args:
path: S3 bucket or key path (e.g., "bucket", "s3://bucket" or
"s3://bucket/prefix"). If the path contains a key, only the
Expand All @@ -2632,17 +2622,30 @@ def clear_multipart_uploads(self, path: str) -> None:
uploads = self.list_multipart_uploads(path)
if not uploads:
return
error: Exception | None = None
with self._create_executor(max_workers=self.max_workers) as executor:
futures = [
executor.submit(
self.core.abort_multipart_upload,
S3Path(cast(str, upload.bucket), cast(str, upload.key)),
cast(str, upload.upload_id),
upload,
)
for upload in uploads
]
for future in as_completed(futures):
future.result()
try:
future.result()
except Exception as e:
cause = e.__cause__
if (
isinstance(e, FileNotFoundError)
and isinstance(cause, botocore.exceptions.ClientError)
and S3ClientError(cause).code == "NoSuchUpload"

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 (compatibility, operations, and claims): CLEAN.

Base: e8f2853
Head: fa673be

Audited the complete six-file inventory and PR/docstring claims. The API serialization now explicitly omits None values; existing checksum properties, copy metadata, request precedence, input order, and method signatures are preserved. Verified all ten checksum fields in the minimum supported botocore 1.43.31. Core error chaining permits distinguishing NoSuchUpload from NoSuchBucket; other abort failures propagate after result draining without changing retries or adding S3 requests.

Traced synchronous finishing, asynchronous copy, aio synchronous wrappers, and direct abort/discard callers. Live tests use unique object paths and validate round-trip data plus a controlled list/abort race. The PR's requirement claim was narrowed to the reported SHA256/CRC32 cases; the eight other algorithms have offline coverage only. New API properties appear in the rendered documentation; strict Sphinx diagnostics have the same count as the unchanged parent. Broader filesystem validation and Ready-triggered CI remain pending, so this is not a full-runtime verdict.

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 follow-up (full compatibility, operations, and claims scope): CLEAN.

Base: 6d268b1
Head: 534e9a4

Audited the complete nine-file diff and the revised PR, docstrings, and related filesystem documentation as a separate pass. Claims checked: matching creation-algorithm checksums only; no-argument core completion preserves ETag/part-number requests even when parts contain SDK CRC32; upload and CopyPartResult parsing retains all ten fields; serialization omits None; only chained NoSuchUpload is treated as cleared; missing bucket, permission, unclassified errors, and direct core abort errors remain observable after ordinary abort-result draining.

The new core option is keyword-only and optional; existing positional calls and uppercase AWS parameters retain their behavior and precedence. The private sync finish option follows its existing positional arguments. Traced async copy and synchronous wrappers/open paths to the same creation response. The supported minimum botocore model (1.43.31) accepts all ten fields; the locked runtime is 1.43.102, whose checksum-default implementation and measured default failures justify excluding incidental checksums. This changes no retry policy, pagination, AWS request count, or cleanup prefix scope.

Evidence: 177 rebased offline cases and just lint passed; just docs lint passed; the normal rebased Sphinx build succeeded with 195 inherited warnings. Unique-path live tests verify byte contents and a controlled listed-upload race for both filesystem classes. Live algorithm claims are limited to default/SHA256/CRC32; eight other algorithms are covered only by offline request tests. The earlier eight default runtime failures are resolved; corrected copy Stubber fixtures passed separately. Final rebased AWS coverage is running, and CI is still Draft/offline. No unrun live or CI check is claimed as 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.

Final validation update for the reviewed revision.

Base: 6d268b1
Head: 534e9a4

No source changes after the full review. Lint and docs lint passed; normal rendered docs succeeded with 195 inherited warnings. All 177 focused offline cases passed. The affected 272-case filesystem selection completed as 268 passes before deliberate AWS serialization and the remaining four CRC32 append/list-abort race cases passed after resumption. No failures remain in that selection.

Published head matches this reviewed commit, the worktree is clean, all Draft/offline CI checks passed, and GitHub reports MERGEABLE. Both self-review rounds and the independent full review are complete; Ready-triggered AWS CI will be checked separately.

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.

Ready-triggered CI completed successfully for reviewed head 534e9a4: https://github.com/pyathena-dev/PyAthena/actions/runs/37182904092

Python 3.14 PyAthena suite: 2304 passed, 1 skipped, 7 warnings. All applicable PR checks passed (lint, license headers, offline checks, docs lint/build, and AWS tests). SQLAlchemy/Spark are excluded by this filesystem-only diff's workflow path conditions. The PR is Ready and MERGEABLE; the reviewed worktree remains clean. Recent master changes affect cursor conversion/result paths and add an unrelated query constant in tests/pyathena/util.py; they do not alter the reviewed filesystem, multipart API, or test helper contracts. No source change follows the completed reviews.

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: test-convention repair, CLEAN.

Base: 6d268b1
Previously reviewed head: 534e9a4
Current head: 745664a

Audited the four-file repair and updated description separately from implementation review. Verified the stated class/file conventions directly against neighboring existing tests: TestS3Core, TestS3MultipartUploadPart, TestS3FileSystem and TestAioS3FileSystem. No added module-level test functions, fs_class branching, separate regression class, or replacement fs fixture remain. The live tests use the existing class-scoped fs fixtures and distinct operation methods.

Checked fixture scope and caching, sync versus async entry points, the preserved 2-filesystem x 4-operation x 3-algorithm matrix plus two races, unique object paths, cleanup, and offline error/request assertions. Dummy credentials and single-worker Stubber ordering remain confined to offline cases. Existing live fixture settings replace the special new fixture settings and have now been exercised against AWS.

Evidence at this head: format/lint passed, 157 model/core/error and 20 filesystem offline cases passed, and all 26 revised live cases passed. Production code, documented API behavior, supported dependencies, and request/retry semantics are unchanged by this repair. Corrected PR validation commands for moved tests; the prior 2304-case CI result is not claimed for this new head. Current CI and the independent follow-up remain separate requirements.

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, final bounded follow-up: CLEAN.

Base: 6d268b1
Old head: 745664a
Head: a9e7dc8

Separately audited claims, compatibility, and evidence for the callback/path follow-up. The synchronous mock targets fs.list_multipart_uploads; the asynchronous mock targets the delegated sync_fs method. Both asserts execute after successful clear, and outside-context listing still uses the restored method. Unique async prefixes differ only in naming and retain account/bucket/schema scope and UUID isolation. The earlier class/fixture conventions and 26-case live matrix are preserved.

The two strengthened live races and current format/lint passed. No claim is made that earlier 177/26 counts or earlier full CI were rerun on this head. Updated validation will name their tested revisions and final CI will be collected on this published head. No API, AWS retry, or production contract changes are introduced.

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.

Final convention-repair validation: complete.

Base: 6d268b1
Head: a9e7dc8
AWS CI: https://github.com/pyathena-dev/PyAthena/actions/runs/37185094919
Python 3.14: 2318 passed, 1 skipped, 13 warnings.

All applicable checks passed, including final AWS tests, lint, license headers, offline tooling, and docs lint/build. SQLAlchemy/Spark are excluded by the current diff's workflow path conditions. PR is Ready and MERGEABLE, at the independently reviewed head; the worktree is clean.

Existing classes/fixtures, separate sync/aio files and entry points, distinct operation methods, async naming, and explicit race-callback assertions are now in place. Both self-review perspectives and the requested Claude Max / claude-opus-5-5 independent follow-ups are recorded above as CLEAN. Earlier local 177/26 results retain their explicit tested-revision boundaries; the current full CI establishes validation on this final head.

):
continue
if error is None:
error = e
if error is not None:
raise error

def created(self, path: str) -> datetime:
"""Return the creation time of the path.
Expand Down Expand Up @@ -3287,8 +3290,7 @@ def _initiate_upload(self) -> None:
self.multipart_upload_parts.append(
self._executor.submit(
self.fs.core.upload_part_copy,
path=path,
upload_id=cast(str, self.multipart_upload.upload_id),
upload=self.multipart_upload,
part_number=i + 1,
# The existing object is copied into the upload.
source=path,
Expand All @@ -3300,8 +3302,7 @@ def _initiate_upload(self) -> None:
self.multipart_upload_parts.append(
self._executor.submit(
self.fs.core.upload_part_copy,
path=path,
upload_id=cast(str, self.multipart_upload.upload_id),
upload=self.multipart_upload,
part_number=1,
source=path,
**self._get_request_kwargs("upload_part_copy"),
Expand Down Expand Up @@ -3365,8 +3366,7 @@ def _upload_chunk(self, final: bool = False) -> bool:
self.multipart_upload_parts.append(
self._executor.submit(
self.fs.core.upload_part,
path=S3Path(self.bucket, self.key),
upload_id=cast(str, self.multipart_upload.upload_id),
upload=self.multipart_upload,
part_number=part_number,
body=upload,
**self._get_request_kwargs("upload_part"),
Expand Down Expand Up @@ -3421,9 +3421,7 @@ def commit(self) -> None:
upload_id = cast(str, self.multipart_upload.upload_id)
try:
self.fs._finish_multipart_upload(
bucket=self.bucket,
key=self.key,
upload_id=upload_id,
upload=self.multipart_upload,
futures=self.multipart_upload_parts,
request_kwargs=self.s3_additional_kwargs,
abort=False,
Expand Down Expand Up @@ -3456,8 +3454,7 @@ def discard(self) -> None:
# be cancelled first.
wait([f for f in self.multipart_upload_parts if not f.cancel()])
self.fs.core.abort_multipart_upload(
S3Path(self.bucket, self.key),
cast(str, self.multipart_upload.upload_id),
self.multipart_upload,
**self._get_request_kwargs("abort_multipart_upload"),
)

Expand Down
13 changes: 4 additions & 9 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -585,7 +585,7 @@ async def _copy_object_with_multipart_upload(
# See S3FileSystem._copy_object_with_multipart_upload.
await asyncio.to_thread(self.core.copy_object, plan.source, plan.destination, **kwargs)
return
upload_id: str
multipart_upload: S3MultipartUpload

semaphore = asyncio.Semaphore(max_workers)
failed = False
Expand All @@ -599,8 +599,7 @@ async def _upload_part(i: int, range_: tuple[int, int]) -> S3MultipartUploadPart
try:
return await asyncio.to_thread(
self.core.upload_part_copy,
path=plan.destination,
upload_id=upload_id,
upload=multipart_upload,
part_number=i + 1,
source=plan.source,
range_=range_,
Expand Down Expand Up @@ -630,9 +629,7 @@ async def _abort() -> None:
return
await asyncio.to_thread(
self._sync_fs._abort_multipart_upload,
plan.destination.bucket,
cast(str, plan.destination.key),
cast(str, creation.result().upload_id),
creation.result(),

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 — claim audit, compatibility, AWS effects, and evidence: CLEAN.

Frozen full diff: 200762088e7b0bac45054f22e32a16aa0ce95dbb..f9e2a831c907408dcba44846c8496d3dfb993e42 (actual PR base: master). All nine changed files were inventoried, with their direct callers and regression tests.

Audited the final PR description, changed docs/docstrings/comments, and the maintainer contract separately from round one. The destination, upload ID, and checksum configuration come from one upload object; no per-upload state is cached in S3Core. Old primitive signatures and the separate completion checksum_algorithm option are removed intentionally for the unreleased 4.0.0 API. The decision and release-note item are recorded on #1063.

Checked AWS primary API documentation for UploadPart algorithm consistency, list-entry checksum fields, and completion ChecksumType. Botocore 1.43.31 exposes all required fields. ListMultipartUploads supplies Key/UploadId and optional checksum algorithm/type; the adapter adds Bucket, and abort sends identity only. FULL_OBJECT is retained and forwarded, rather than inferred from incidental part checksums. SDK algorithm support is not overstated by the ten-field serialization coverage.

Reviewed request filtering, SSE-C/requester-pays/conditional writes, source version pinning, retry/error translation, futures/tasks and cancellation ordering, prefix-scoped test cleanup, and cache invalidation boundaries. The #1074 planning/scheduling/cancellation contract is preserved. No new requests, retries, or upload state were added to the primitives.

Evidence boundary: 702 offline cases and 34 targeted live filesystem cases passed on the unchanged implementation; current direct-core/live checksum-value checks remain pending. Docs lint passed and a rendered build succeeded with 205 warnings, without claiming a warning-free build. Current CI and remaining live results are required before delivery. No actionable finding.

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.

Local validation completed at f9e2a83: the direct-core/listing selection passed all 12 cases in 75.05s. Nine use live AWS (six default/SHA256-COMPOSITE/CRC32-FULL_OBJECT upload/copy cases and three listing/cleanup cases); three are offline sibling-prefix cases. Stored FULL_OBJECT CRC32 matches the checksum calculated over the complete data, for both uploaded and copied first parts. Listing preserves the creation algorithm/type, and clear_multipart_uploads aborts those entries without leaking checksum fields into AbortMultipartUpload. The separate 34 live filesystem cases passed on the unchanged implementation at 82fbe04. Format/lint passed on this head, and current offline CI checks, including rendered docs, all passed. The independent review and Ready-triggered full AWS CI remain pending.

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 repair claim/operational audit: CLEAN. Frozen base/head: 2007620..49a438f. Bounded repair comparison: f9e2a83..49a438f; both objects exist, range-diff confirms the six previous patches are unchanged and the repair affects only s3_core.py, test_s3_core.py, and test_s3.py. Verified the AWS CreateMultipartUpload response Bucket excludes the access-point alias/ARN, so retaining the request endpoint is necessary for subsequent authorization and cleanup. The creation docstring states this identity contract; checksum configuration remains from the service response. No upload state, extra requests, retries or changes to cancellation were added. The stronger assertion now observes the actual creation call and mapping propagation rather than only a hard-coded mock result. Full-object CRC32 and checksum-bearing listing measurements remain valid for the unchanged checksum/cleanup logic; current full CI remains required. No live access-point claim is made. Independent follow-up is pending.

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.

Final current-head CI — PASS. Frozen review base/head: 2007620..49a438f.

Ready-triggered Test run https://github.com/pyathena-dev/PyAthena/actions/runs/37191012018 completed successfully. The Python 3.14 PyAthena job, including AWS integration tests, reported 2,424 passed, 1 skipped, and 13 warnings in 464.03s. Lint, offline checks, license headers, and documentation lint/build have all completed successfully on this head. SQLAlchemy compliance and Spark suites are intentionally excluded by the filesystem-only path filters in test.yaml; the older Draft run's skipped test job is superseded by this successful Ready run.

Both self-review rounds and repair follow-ups remain CLEAN. The requested independent Claude Max / claude-opus-5-5 review and bounded repair follow-up are recorded separately; the follow-up is CLEAN (static only). Author runtime validation is provided by this current-head CI and the earlier targeted AWS runs. Access-point alias/ARN coverage remains offline only.

The published PR is Ready, OPEN, and MERGEABLE at the reviewed head; the dedicated worktree is clean.

plan.abort_params,
)

Expand All @@ -648,7 +645,6 @@ async def _abort() -> None:
# shield keeps a cancellation from cancelling the creation, whose
# thread would keep running, so that _abort() can wait for it.
multipart_upload = await asyncio.shield(creation)
upload_id = cast(str, multipart_upload.upload_id)
tasks = [asyncio.ensure_future(_upload_part(i, r)) for i, r in enumerate(plan.ranges)]
# Unlike gather, wait does not cancel the parts when this task is
# cancelled; their threads would keep copying, so they are waited
Expand All @@ -662,8 +658,7 @@ async def _abort() -> None:
completion = asyncio.ensure_future(
asyncio.to_thread(
self.core.complete_multipart_upload,
plan.destination,
upload_id,
multipart_upload,
cast(list[S3MultipartUploadPart], parts),
**plan.complete_params,
)
Expand Down
Loading
Loading