Use case
Step 3 of #1053. Step 2 (#1059, PR #1060) added the public, synchronous S3Core with typed listing and lookup operations. The decisions recorded on #1053 still apply: the core is public API; it is synchronous and the adapters schedule its requests (no aiobotocore); the cache stays in fsspec's DirCache in the adapter.
Copying, moving and deleting are still implemented inside S3FileSystem, where S3 requests, S3 rules, scheduling and fsspec semantics are mixed in private helpers, and AioS3FileSystem reaches into those helpers through self._sync_fs. Line numbers are from pyathena/filesystem/s3.py on master 8bbf8cb, unless another file is named.
- Delete.
_delete_objects_requests() (:1217) applies S3's rules: group the keys by bucket, send at most DELETE_OBJECTS_MAX_KEYS keys per request, and add a VersionId per entry.
_delete_objects_request() (:1256) sends one request and invalidates the cache in finally. _delete_objects_path() (:1275) turns an entry back into a path.
_raise_delete_objects_errors() (:1288) merges the request exceptions with the per-object Errors.
_delete_object() (:1164) sends DeleteObject and ignores its **kwargs.
_delete_objects() (:1194) schedules the batches on an executor; aio _delete_objects() schedules them with asyncio.gather.
- Multipart primitives.
_create_multipart_upload(), _upload_part_copy(), _upload_part() and _complete_multipart_upload() (:3404-3494) are used by S3File writes and appends (:3785, :3805, :3818, :3884), by the multipart copy, and by aio.
_abort_multipart_upload() (:2402) logs and swallows errors; that is a policy of its callers, not a primitive. S3File.discard() (:3975-3984) sends AbortMultipartUpload through fs._call and lets errors propagate.
_finish_multipart_upload() (:2349) is the synchronous wait/complete/abort orchestration, used by the multipart copy (:1943), pipe_file() (:2385) and S3File.commit() (:3941). aio has its own form (s3_async.py :687-731).
- Copy.
_copy_object() (:1814) sends CopyObject.
_copy_object_with_multipart_upload() (:1841) has a second implementation in aio (s3_async.py _copy_object_with_multipart_upload()).
- Both use the following helpers:
_get_multipart_copy_kwargs() (:1990): HeadObject and GetObjectTagging of the source, plus the metadata and tag directives that make a multipart copy behave like CopyObject.
_get_copy_source_kwargs() (:1970), _is_directory_bucket() (:1956) and _COPY_METADATA_PARAMS (:207).
_get_copy_ranges() (:2209), which the S3File append also uses (:3797).
_copies_annotations(), _list_object_annotations() and _copy_object_annotation() (:2099-2207).
- The block size validation, which is duplicated in aio.
- The S3 limits
MULTIPART_UPLOAD_* and DELETE_OBJECTS_MAX_KEYS (:166-174) are repeated on AioS3FileSystem (s3_async.py :81).
- Pairing and expansion (fsspec semantics).
_copy_paths() (:1590) is used by copy()/get() and by aio _copy()/_get(), which reach it through _sync_fs.
_move_paths() (:1528) and _move_target() (:1703) pair and check moves.
_expand_delete_paths() (:1127) expands the paths of rm().
expand_path() (:979) and aio _expand_path() (s3_async.py :783-820) are near copies of each other, and both use _split_version_paths() (:1017).
Proposed change
Move the S3 requests and S3 rules of copy and delete into the core, each on the class that owns its rule. Keep fsspec's pairing, expansion, cache invalidation and scheduling in the adapters. No behavior change.
Core (s3_core.py)
New entries are frozen dataclasses with from_response(), as in step 2.
-
Delete.
S3DeleteBatch(bucket, objects: tuple[S3Path, ...], quiet=True) owns the batching rule: MAX_KEYS = 1000, and S3DeleteBatch.from_paths(paths, quiet=True) -> list[S3DeleteBatch] groups the paths by bucket.
S3Core.delete_objects(batch: S3DeleteBatch, **params) -> S3DeleteResult sends one batch, which is valid by construction.
S3DeleteResult(bucket, deleted, errors) and S3DeleteError(path, code, message) (str() gives "path (code: message)"). With Quiet=True, the default, S3 returns no Deleted entries. A 200 response with Errors is a result, not an exception.
S3Core.delete_object(path: S3Path, **params). rm_file() keeps sending no parameters.
-
Multipart primitives. All five move to the core, one request each:
create_multipart_upload(path, **params);
upload_part(path, upload_id, part_number, body, **params);
upload_part_copy(path, upload_id, part_number, source, range_=None, **params);
complete_multipart_upload(path, upload_id, parts, **params);
abort_multipart_upload(path, upload_id, **params), which raises. Logging and swallowing an abort failure stays in the callers.
They return the existing S3MultipartUpload, S3MultipartUploadPart and S3CompleteMultipartUpload. The limits become S3Core class attributes, and the filesystems keep their aliases. S3Core.part_ranges(size, block_size) is the range split for a multipart copy and for an append.
-
Copy.
S3Core.copy_object(source: S3Path, destination: S3Path, **params).
S3Core.plan_multipart_copy(source, destination, block_size=None, **params) -> S3MultipartCopyPlan makes the reads only (HeadObject, GetObjectTagging, ListObjectAnnotations), validates block_size, and returns a frozen plan with these fields:
source, with its version pinned unless it is null;
destination, size and ranges;
create_params, part_params, complete_params and abort_params, each already filtered for its operation;
annotations, empty when the directive or the bucket excludes them;
fits_single_request, for a source whose head size fits one CopyObject.
S3Core.list_object_annotations(path, **params) and S3Core.copy_object_annotation(name, source, destination, version_id, etag, **params).
- The copy-source parameter mapping, the directory-bucket rule and the copied metadata fields move into the planner.
-
Move target. S3Path.target -> S3Path is the object that a write to the path replaces: a null version names its key. It replaces _move_target().
Adapters
-
Scheduling stays in the adapters. One sync and one aio orchestration of the multipart copy each run a plan from the core. The sync one schedules parts with the executor, through _finish_multipart_upload(). The aio one uses asyncio.to_thread with gather/wait, its cancellation handling, the failed flag, the semaphore and the cleanup tasks, all unchanged. Only the callables they schedule change.
-
Deletes are scheduled one batch at a time. Each scheduled unit is try: core.delete_objects(batch) finally: invalidate the batch's paths, so that a failed or cancelled batch still invalidates the cache. The cross-batch error policy of rm() stays one adapter method: re-raise the first exception, with a note listing every S3DeleteError, or raise OSError if there was no exception. aio calls it through _sync_fs.
-
Pairing gets a named owner. It is one component, S3PathPairing, built once by S3FileSystem and reached from aio through _sync_fs, with three methods:
copy_pairs(path1, path2, recursive, maxdepth, isdir=None), which replaces _copy_paths();
move_pairs(...), which replaces _move_paths();
delete_paths(path, recursive, maxdepth), which replaces _expand_delete_paths().
Expansion stays in the fsspec overrides expand_path()/_expand_path(), because the aio one awaits _exists. _split_version_paths() is either inlined into them or kept with the sync adapter as the expansion rule.
-
Cache invalidation, skipping directories in a recursive copy(), and the rm()/mv()/copy()/get() signatures stay in the adapters.
Sub-steps (one PR each)
- Delete:
S3DeleteBatch, S3DeleteResult, S3DeleteError, and delete_object/delete_objects on the core. The sync and aio rm() schedule the batches.
- Multipart primitives: all five move to the core, together with the limits and
part_ranges(). S3File, pipe_file() and the multipart copy switch to fs.core.* mechanically.
- Copy:
copy_object, plan_multipart_copy/S3MultipartCopyPlan and the annotation operations. The sync and aio multipart copy are reduced to scheduling.
- Pairing:
S3PathPairing and S3Path.target.
Step 4 of #1053 then only covers the multipart writer of S3File (and _check_multipart_upload_size(), which aio reaches through _sync_fs at s3_async.py :219 and :276).
Open questions:
- The names:
S3DeleteBatch, S3DeleteResult, S3DeleteError, S3MultipartCopyPlan, plan_multipart_copy, part_ranges, S3PathPairing, S3Path.target.
- Whether
S3PathPairing is public (part of the adapter's API) or an internal component of the filesystem package.
Out of scope: the multipart writer of S3File, the cache, and any behavior change.
Validation plan (if implementing)
Stubber tests for each new core operation and planner:
- the batches at 1,000 keys and across buckets;
- the per-object errors of a 200 response;
- the version pinning and the single-request fit of the copy plan;
- the filtered parameters of each operation;
- error translation.
- The existing offline and S3 integration tests in
tests/pyathena/filesystem/ pass with the same request sequences in sync and aio, including the cancellation tests of the aio multipart copy and the cache invalidation after a failed delete batch.
Use case
Step 3 of #1053. Step 2 (#1059, PR #1060) added the public, synchronous
S3Corewith typed listing and lookup operations. The decisions recorded on #1053 still apply: the core is public API; it is synchronous and the adapters schedule its requests (no aiobotocore); the cache stays in fsspec'sDirCachein the adapter.Copying, moving and deleting are still implemented inside
S3FileSystem, where S3 requests, S3 rules, scheduling and fsspec semantics are mixed in private helpers, andAioS3FileSystemreaches into those helpers throughself._sync_fs. Line numbers are frompyathena/filesystem/s3.pyon master 8bbf8cb, unless another file is named._delete_objects_requests()(:1217) applies S3's rules: group the keys by bucket, send at mostDELETE_OBJECTS_MAX_KEYSkeys per request, and add aVersionIdper entry._delete_objects_request()(:1256) sends one request and invalidates the cache infinally._delete_objects_path()(:1275) turns an entry back into a path._raise_delete_objects_errors()(:1288) merges the request exceptions with the per-objectErrors._delete_object()(:1164) sends DeleteObject and ignores its**kwargs._delete_objects()(:1194) schedules the batches on an executor; aio_delete_objects()schedules them withasyncio.gather._create_multipart_upload(),_upload_part_copy(),_upload_part()and_complete_multipart_upload()(:3404-3494) are used byS3Filewrites and appends (:3785, :3805, :3818, :3884), by the multipart copy, and by aio._abort_multipart_upload()(:2402) logs and swallows errors; that is a policy of its callers, not a primitive.S3File.discard()(:3975-3984) sends AbortMultipartUpload throughfs._calland lets errors propagate._finish_multipart_upload()(:2349) is the synchronous wait/complete/abort orchestration, used by the multipart copy (:1943),pipe_file()(:2385) andS3File.commit()(:3941). aio has its own form (s3_async.py:687-731)._copy_object()(:1814) sends CopyObject._copy_object_with_multipart_upload()(:1841) has a second implementation in aio (s3_async.py_copy_object_with_multipart_upload())._get_multipart_copy_kwargs()(:1990): HeadObject and GetObjectTagging of the source, plus the metadata and tag directives that make a multipart copy behave like CopyObject._get_copy_source_kwargs()(:1970),_is_directory_bucket()(:1956) and_COPY_METADATA_PARAMS(:207)._get_copy_ranges()(:2209), which theS3Fileappend also uses (:3797)._copies_annotations(),_list_object_annotations()and_copy_object_annotation()(:2099-2207).MULTIPART_UPLOAD_*andDELETE_OBJECTS_MAX_KEYS(:166-174) are repeated onAioS3FileSystem(s3_async.py:81)._copy_paths()(:1590) is used bycopy()/get()and by aio_copy()/_get(), which reach it through_sync_fs._move_paths()(:1528) and_move_target()(:1703) pair and check moves._expand_delete_paths()(:1127) expands the paths ofrm().expand_path()(:979) and aio_expand_path()(s3_async.py:783-820) are near copies of each other, and both use_split_version_paths()(:1017).Proposed change
Move the S3 requests and S3 rules of copy and delete into the core, each on the class that owns its rule. Keep fsspec's pairing, expansion, cache invalidation and scheduling in the adapters. No behavior change.
Core (
s3_core.py)New entries are frozen dataclasses with
from_response(), as in step 2.Delete.
S3DeleteBatch(bucket, objects: tuple[S3Path, ...], quiet=True)owns the batching rule:MAX_KEYS = 1000, andS3DeleteBatch.from_paths(paths, quiet=True) -> list[S3DeleteBatch]groups the paths by bucket.S3Core.delete_objects(batch: S3DeleteBatch, **params) -> S3DeleteResultsends one batch, which is valid by construction.S3DeleteResult(bucket, deleted, errors)andS3DeleteError(path, code, message)(str()gives"path (code: message)"). WithQuiet=True, the default, S3 returns noDeletedentries. A 200 response withErrorsis a result, not an exception.S3Core.delete_object(path: S3Path, **params).rm_file()keeps sending no parameters.Multipart primitives. All five move to the core, one request each:
create_multipart_upload(path, **params);upload_part(path, upload_id, part_number, body, **params);upload_part_copy(path, upload_id, part_number, source, range_=None, **params);complete_multipart_upload(path, upload_id, parts, **params);abort_multipart_upload(path, upload_id, **params), which raises. Logging and swallowing an abort failure stays in the callers.They return the existing
S3MultipartUpload,S3MultipartUploadPartandS3CompleteMultipartUpload. The limits becomeS3Coreclass attributes, and the filesystems keep their aliases.S3Core.part_ranges(size, block_size)is the range split for a multipart copy and for an append.Copy.
S3Core.copy_object(source: S3Path, destination: S3Path, **params).S3Core.plan_multipart_copy(source, destination, block_size=None, **params) -> S3MultipartCopyPlanmakes the reads only (HeadObject, GetObjectTagging, ListObjectAnnotations), validatesblock_size, and returns a frozen plan with these fields:source, with its version pinned unless it isnull;destination,sizeandranges;create_params,part_params,complete_paramsandabort_params, each already filtered for its operation;annotations, empty when the directive or the bucket excludes them;fits_single_request, for a source whose head size fits one CopyObject.S3Core.list_object_annotations(path, **params)andS3Core.copy_object_annotation(name, source, destination, version_id, etag, **params).Move target.
S3Path.target -> S3Pathis the object that a write to the path replaces: anullversion names its key. It replaces_move_target().Adapters
Scheduling stays in the adapters. One sync and one aio orchestration of the multipart copy each run a plan from the core. The sync one schedules parts with the executor, through
_finish_multipart_upload(). The aio one usesasyncio.to_threadwithgather/wait, its cancellation handling, thefailedflag, the semaphore and the cleanup tasks, all unchanged. Only the callables they schedule change.Deletes are scheduled one batch at a time. Each scheduled unit is
try: core.delete_objects(batch) finally: invalidate the batch's paths, so that a failed or cancelled batch still invalidates the cache. The cross-batch error policy ofrm()stays one adapter method: re-raise the first exception, with a note listing everyS3DeleteError, or raiseOSErrorif there was no exception. aio calls it through_sync_fs.Pairing gets a named owner. It is one component,
S3PathPairing, built once byS3FileSystemand reached from aio through_sync_fs, with three methods:copy_pairs(path1, path2, recursive, maxdepth, isdir=None), which replaces_copy_paths();move_pairs(...), which replaces_move_paths();delete_paths(path, recursive, maxdepth), which replaces_expand_delete_paths().Expansion stays in the fsspec overrides
expand_path()/_expand_path(), because the aio one awaits_exists._split_version_paths()is either inlined into them or kept with the sync adapter as the expansion rule.Cache invalidation, skipping directories in a recursive
copy(), and therm()/mv()/copy()/get()signatures stay in the adapters.Sub-steps (one PR each)
S3DeleteBatch,S3DeleteResult,S3DeleteError, anddelete_object/delete_objectson the core. The sync and aiorm()schedule the batches.part_ranges().S3File,pipe_file()and the multipart copy switch tofs.core.*mechanically.copy_object,plan_multipart_copy/S3MultipartCopyPlanand the annotation operations. The sync and aio multipart copy are reduced to scheduling.S3PathPairingandS3Path.target.Step 4 of #1053 then only covers the multipart writer of
S3File(and_check_multipart_upload_size(), which aio reaches through_sync_fsats3_async.py:219 and :276).Open questions:
S3DeleteBatch,S3DeleteResult,S3DeleteError,S3MultipartCopyPlan,plan_multipart_copy,part_ranges,S3PathPairing,S3Path.target.S3PathPairingis public (part of the adapter's API) or an internal component of the filesystem package.Out of scope: the multipart writer of
S3File, the cache, and any behavior change.Validation plan (if implementing)
Stubbertests for each new core operation and planner:tests/pyathena/filesystem/pass with the same request sequences in sync and aio, including the cancellation tests of the aio multipart copy and the cache invalidation after a failed delete batch.