-
Notifications
You must be signed in to change notification settings - Fork 116
Check the part limit in AioS3FileSystem transaction writes #999
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
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 |
|---|---|---|
|
|
@@ -187,7 +187,12 @@ def _pipe_file_in_transaction( | |
| Raises: | ||
| FileExistsError: If the mode is "create" and the path already | ||
| exists. | ||
| ValueError: If the data takes more than | ||
| ``MULTIPART_UPLOAD_MAX_PARTS`` blocks. | ||
| """ | ||
| block_size = kwargs.get("block_size") or self._sync_fs.default_block_size | ||
| # The size in bytes; the length of a memoryview counts its items. | ||
| self._sync_fs._check_multipart_upload_size(path, memoryview(value).nbytes, block_size) | ||
|
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) — reviewer: Codex CLI 0.160.0, model Covered (reviewer): equivalence with the sync Coverage limits noted by the reviewer (not defects): the forwarding test checks argument routing only (its 8-byte block works only with the mocked |
||
| if mode == "create" and self._sync_fs.exists(path): | ||
| raise FileExistsError(path) | ||
| with self.open(path, "wb", **kwargs) as f: | ||
|
|
@@ -211,21 +216,29 @@ def _put_file_in_transaction(self, lpath: str, rpath: str, callback, **kwargs) - | |
| rpath: S3 destination path (s3://bucket/key). | ||
| callback: Progress callback for tracking upload progress. | ||
| **kwargs: Additional S3 parameters (e.g., ContentType, StorageClass). | ||
| The ``block_size`` parameter of ``open()`` is also accepted. | ||
|
|
||
| Raises: | ||
| ValueError: If the file takes more than | ||
| ``MULTIPART_UPLOAD_MAX_PARTS`` blocks. | ||
| """ | ||
| if os.path.isdir(lpath): | ||
| return | ||
| _, key, _ = self.parse_path(rpath) | ||
| if not key: | ||
| return | ||
|
|
||
| callback.set_size(os.path.getsize(lpath)) | ||
| size = os.path.getsize(lpath) | ||
| block_size = kwargs.pop("block_size", None) or self._sync_fs.default_block_size | ||
|
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, operations) — base Claims checked: #988's branch head does not contain #968 ( |
||
| self._sync_fs._check_multipart_upload_size(rpath, size, block_size) | ||
| callback.set_size(size) | ||
| if "ContentType" not in kwargs: | ||
| content_type, _ = mimetypes.guess_type(lpath) | ||
| if content_type is not None: | ||
| kwargs["ContentType"] = content_type | ||
|
|
||
| with ( | ||
| self.open(rpath, "wb", s3_additional_kwargs=kwargs) as remote, | ||
| self.open(rpath, "wb", block_size=block_size, s3_additional_kwargs=kwargs) as remote, | ||
| open(lpath, "rb") as local, | ||
| ): | ||
| while data := local.read(remote.blocksize): | ||
|
|
||
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) — base
92c9e3e67180eb3d52ef7b85579a24540a30fd64, headbab08235f0050aebdc4cf2d85a8fcba08b6dc59b. Result: CLEAN.Covered:
_pipe_file_in_transaction()and_put_file_in_transaction()againstS3FileSystem.pipe_file()/put_file()(same check, same order: before themode="create"existence check, beforecallback.set_size(), beforeopen()); theblock_sizethatpipestill passes toopen()throughkwargs; the error propagating throughasyncio.to_threadand the generated sync wrappers; no file opened, so nothing joinstransaction.files; non-transaction paths unchanged (they delegate to the sync methods). Tests: 3 new cases fail on the base and pass here; offlinetests/pyathena/filesystem/262 passed vs 259 on the base, same 110 credential-less failures. Out of scope: the non-contiguous memoryview failure in the samewith self.open(...)pattern (#997);max_workersrouting (#969).