-
Notifications
You must be signed in to change notification settings - Fork 116
Keep the existing object when appending within a larger block size #929
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
e26b398
33907e1
2f0a1c3
09f34c3
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 |
|---|---|---|
|
|
@@ -2111,8 +2111,9 @@ def __init__( | |
| In read mode, the object is looked up with ``info()`` and the reads | ||
| are made conditional on its ETag (``IfMatch``). In append mode, an | ||
| existing object smaller than ``MULTIPART_UPLOAD_MIN_PART_SIZE`` is | ||
| read into the write buffer; a larger one is copied as the first parts | ||
| once a multipart upload starts. | ||
| read into the write buffer; a larger one is copied with | ||
| ``UploadPartCopy`` as the first parts of a multipart upload, whatever | ||
| the block size. | ||
|
|
||
| Args: | ||
| fs: The filesystem that the file belongs to. | ||
|
|
@@ -2188,11 +2189,15 @@ def __init__( | |
| self.s3_additional_kwargs.update({"IfMatch": etag}) | ||
| self._details = info | ||
| elif "a" in mode and self.fs.exists(path): | ||
| self.append_block = True | ||
| info = self.fs.info(self.path, version_id=self.version_id) | ||
| loc = info.get("size", 0) | ||
| if loc < self.fs.MULTIPART_UPLOAD_MIN_PART_SIZE: | ||
| # Too small to be a part of a multipart upload: rewrite it | ||
| # from the buffer. | ||
| self.write(self.fs.cat(self.path)) | ||
| else: | ||
| # Copied with UploadPartCopy as the leading part(s). | ||
| self.append_block = True | ||
| self.loc = loc | ||
| self.s3_additional_kwargs.update(info.to_api_repr()) | ||
| self._details = info | ||
|
|
@@ -2208,8 +2213,10 @@ def close(self) -> None: | |
| self._executor.shutdown() | ||
|
|
||
| def _initiate_upload(self) -> None: | ||
| if self.tell() < self.blocksize: | ||
| if not self.append_block and self.tell() < self.blocksize: | ||
|
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): FINDINGS Reviewer: OpenAI Codex CLI 0.160.0, model Covered: existing object absent / 0 / < 5 MiB / == 5 MiB / > 5 MiB / > 5 GiB; empty, short, threshold-crossing, and successive appends; block-size boundaries; part ordering; autocommit, deferred commit, discard, metadata forwarding; Regression (P2): "Existing object 6 MiB, append 1 byte, block size 16 MiB, Pre-existing, not made worse (P2 each):
Test notes: "both added integration cases and three of the four unit cases would fail before the change … The unit reconstruction does not validate S3 part limits or completion contents. Append transactions and discard remain uncovered."
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. Repair: 33907e1 ( Regression: verified and fixed. Botocore rejects the extra parameters ( New tests:
Self-review of the repair:
Pre-existing items 1–4: deferred. None of them is made worse by this PR, and each needs its own reproduction; items 1, 2, and 4 involve objects or parts ≥ 5 GiB. I will verify them and show them to the maintainer before filing issues.
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 follow-up review (relayed): FINDINGS Reviewer: OpenAI Codex CLI 0.160.0, model Covered: append metadata, multipart init, close, and discard for 5 MiB and 16 MiB blocks; all three abort call sites and
Repair: 2f0a1c3.
Self-review of the repair. Behavior: this restores the pre-PR forwarding of
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 follow-up review (relayed): CLEAN Reviewer: OpenAI Codex CLI 0.160.0, model Covered: Result: "no actionable regression … The changed call supplies Pre-existing, unchanged by this PR (deferred, to be verified separately):
|
||
| # Files smaller than block size in size cannot be multipart uploaded. | ||
| # An append to an object copied with UploadPartCopy always uses | ||
| # a multipart upload, whatever the block size. | ||
| return | ||
|
|
||
| self.multipart_upload = self.fs._create_multipart_upload( | ||
|
|
@@ -2259,7 +2266,7 @@ def _upload_chunk(self, final: bool = False) -> bool: | |
| # can still read the bytes; resetting it there would upload an empty | ||
| # object for small files. Mid-stream chunks (final=False) return True so | ||
| # fsspec clears the already-uploaded buffer between parts. | ||
| if self.tell() < self.blocksize: | ||
| if not self.append_block and self.tell() < self.blocksize: | ||
|
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 1 (implementation behavior): CLEAN Scope: Covered:
Out of scope (pre-existing, unchanged): the ranged copy for existing objects > 5 GiB ( |
||
| # Files smaller than block size in size cannot be multipart uploaded. | ||
| if self.autocommit and final: | ||
| self.commit() | ||
|
|
@@ -2360,12 +2367,19 @@ def discard(self) -> None: | |
| if self.multipart_upload: | ||
| for f in self.multipart_upload_parts: | ||
| f.cancel() | ||
| # s3_additional_kwargs also holds object parameters (e.g., the | ||
| # existing object's metadata in append mode) that | ||
| # AbortMultipartUpload rejects. | ||
| self.fs._call( | ||
| "abort_multipart_upload", | ||
| Bucket=self.bucket, | ||
| Key=self.key, | ||
| UploadId=self.multipart_upload.upload_id, | ||
| **self.s3_additional_kwargs, | ||
| **{ | ||
| k: v | ||
| for k, v in self.s3_additional_kwargs.items() | ||
| if k in ("RequestPayer", "ExpectedBucketOwner") | ||
| }, | ||
| ) | ||
|
|
||
| self.multipart_upload = None | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -676,6 +676,65 @@ def test_append(self, fs, base, exp): | |
| assert len(actual) == len(data + extra) | ||
| assert actual == data + extra | ||
|
|
||
| @pytest.mark.parametrize( | ||
| ("size", "extra_size", "block_size"), | ||
| [ | ||
| # GH-921: an existing object of at least 5 MiB, appended within a | ||
| # larger block size, is copied with UploadPartCopy. | ||
| (6 * 2**20, 5, 16 * 2**20), | ||
| # An existing object smaller than 5 MiB is rewritten from the | ||
| # buffer, not copied as well, when the append crosses the block size. | ||
| (2**10, 5 * 2**20, None), | ||
| ], | ||
| ) | ||
| def test_append_with_block_size(self, fs, size, extra_size, block_size): | ||
| data = b"a" * size | ||
| extra = b"b" * extra_size | ||
| path = ( | ||
| f"s3://{ENV.s3_staging_bucket}/{ENV.s3_staging_key}{ENV.schema}/" | ||
| f"filesystem/test_append_with_block_size/{uuid.uuid4()}" | ||
| ) | ||
| fs.pipe_file(path, data) | ||
| with fs.open(path, "ab", block_size=block_size) as f: | ||
| f.write(extra) | ||
| # Check the size and the bytes at the ends and around the boundary | ||
|
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 of 09f34c3 (test transfer reduction), rounds 1 + 2: CLEAN Scope: Behavior and test quality:
Claims and operations: the PR body's TEST section now states the per-run transfer: about 23 MiB uploaded (free), a few bytes + 1 KiB read, server-side copies, and the 1-day lifecycle expiry of objects and incomplete uploads (
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 follow-up review (relayed): CLEAN Reviewer: OpenAI Codex CLI 0.160.0, model Covered: both changed tests and the append/commit/discard paths; detection of the original defects (replacement, duplication, rollback Result: "All three named regressions remain detectable." Stated limit: "The append samples do not prove full interior equality: same-size corruption confined to unsampled bytes could escape." Author note on the limit: accepted as the cost trade-off. Full-content equality is still asserted offline by |
||
| # instead of reading the whole object back, to keep the transfer small. | ||
| assert fs.info(path, refresh=True).size == size + extra_size | ||
| assert fs.cat_file(path, start=0, end=1) == b"a" | ||
| assert fs.cat_file(path, start=size - 1, end=size + 1) == b"ab" | ||
| assert fs.cat_file(path, start=-1) == b"b" | ||
|
|
||
| @pytest.mark.parametrize("block_size", [None, 16 * 2**20]) | ||
| def test_append_transaction_rollback(self, fs, block_size): | ||
| # Raising inside the transaction aborts the multipart upload that | ||
| # copies the existing object and leaves the object unchanged. | ||
| data = b"a" * (6 * 2**20) | ||
| path = ( | ||
| f"s3://{ENV.s3_staging_bucket}/{ENV.s3_staging_key}{ENV.schema}/" | ||
| f"filesystem/test_append_transaction_rollback/{uuid.uuid4()}" | ||
| ) | ||
| fs.pipe_file(path, data) | ||
| before = fs.info(path, refresh=True) | ||
|
|
||
| def append_then_fail(): | ||
| with fs.transaction: | ||
| f = fs.open(path, "ab", block_size=block_size) | ||
| f.write(b"b" * 5) | ||
| f.close() | ||
| raise RuntimeError("rollback") | ||
|
|
||
| with pytest.raises(RuntimeError): | ||
| append_then_fail() | ||
| # A committed append (a multipart upload, or the appended bytes alone) | ||
| # would change the ETag and the size, so the object is not read back. | ||
| after = fs.info(path, refresh=True) | ||
| assert (after.etag, after.last_modified, after.size) == ( | ||
| before.etag, | ||
| before.last_modified, | ||
| before.size, | ||
| ) | ||
| assert not fs.list_multipart_uploads(path) | ||
|
|
||
| def test_ls_buckets(self, fs): | ||
| fs.invalidate_cache() | ||
| actual = fs.ls("s3://") | ||
|
|
@@ -1511,6 +1570,7 @@ def _make_write_file(data: bytes, autocommit: bool): | |
| file.s3_additional_kwargs = {} | ||
| file.autocommit = autocommit | ||
| file.blocksize = S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE | ||
| file.append_block = False | ||
| file.multipart_upload = None | ||
| file.multipart_upload_parts = [] | ||
| file.buffer = io.BytesIO(data) | ||
|
|
@@ -1534,6 +1594,102 @@ def _make_multipart_write_file(data: bytes, autocommit: bool): | |
| ) | ||
| return file | ||
|
|
||
| @staticmethod | ||
| def _make_append_fs(existing: bytes): | ||
| # A mocked filesystem holding an existing object, with a minimum part | ||
| # size of 4 bytes so that the append paths can be exercised with tiny | ||
| # data and no AWS access. | ||
| fs = mock.MagicMock(spec=S3FileSystem) | ||
| fs.MULTIPART_UPLOAD_MIN_PART_SIZE = 4 | ||
| fs.MULTIPART_UPLOAD_MAX_PART_SIZE = 64 | ||
| fs.exists.return_value = True | ||
| fs.info.return_value = S3Object( | ||
| init={"ContentLength": len(existing)}, | ||
| type=S3ObjectType.S3_OBJECT_TYPE_FILE, | ||
| bucket="bucket", | ||
| key="key.txt", | ||
| ) | ||
| fs.cat.return_value = existing | ||
| fs._create_multipart_upload.return_value = SimpleNamespace(upload_id="uploadid") | ||
|
|
||
| def part(**kw): | ||
| return SimpleNamespace(etag=f'"e{kw["part_number"]}"', part_number=kw["part_number"]) | ||
|
|
||
| fs._upload_part.side_effect = part | ||
| fs._upload_part_copy.side_effect = part | ||
| return fs | ||
|
|
||
| @staticmethod | ||
| def _uploaded_object(fs, existing: bytes) -> bytes: | ||
| # Rebuild the object S3 would store from the mocked upload calls. | ||
| # A part copy without a range copies the whole existing object. | ||
| if fs._put_object.called: | ||
| fs._create_multipart_upload.assert_not_called() | ||
| return fs._put_object.call_args.kwargs["body"] | ||
| fs._finish_multipart_upload.assert_called_once() | ||
| parts = [(c.kwargs["part_number"], existing) for c in fs._upload_part_copy.call_args_list] | ||
| parts += [ | ||
| (c.kwargs["part_number"], c.kwargs["body"]) for c in fs._upload_part.call_args_list | ||
| ] | ||
| part_numbers = sorted(n for n, _ in parts) | ||
| assert part_numbers == list(range(1, len(parts) + 1)) | ||
| return b"".join(body for _, body in sorted(parts)) | ||
|
|
||
| @pytest.mark.parametrize( | ||
| ("existing", "appended", "multipart", "part_copy"), | ||
| [ | ||
| # Smaller than the minimum part size: read into the buffer. | ||
| (b"aa", b"bb", False, False), | ||
| # GH-921: an existing object of at least the minimum part size is | ||
| # copied with UploadPartCopy even when the block size is larger | ||
| # than the whole object. | ||
| (b"a" * 6, b"bb", True, True), | ||
| (b"a" * 6, b"", True, True), | ||
| # An existing object read into the buffer is not copied again | ||
| # when the append crosses the block size. | ||
| (b"aa", b"b" * 16, True, False), | ||
| ], | ||
| ) | ||
| def test_append(self, existing, appended, multipart, part_copy): | ||
| fs = self._make_append_fs(existing) | ||
|
|
||
| with S3File(fs, "s3://bucket/key.txt", mode="ab", block_size=16) as f: | ||
| f.write(appended) | ||
|
|
||
| assert self._uploaded_object(fs, existing) == existing + appended | ||
| assert fs._create_multipart_upload.called is multipart | ||
| assert fs._upload_part_copy.called is part_copy | ||
| fs.touch.assert_not_called() | ||
|
|
||
| def test_append_discard(self): | ||
| # Rolling back an append aborts its multipart upload without the | ||
| # existing object's metadata, which AbortMultipartUpload rejects, | ||
| # but with the request parameters it accepts. | ||
| fs = self._make_append_fs(b"a" * 6) | ||
| f = S3File( | ||
| fs, | ||
| "s3://bucket/key.txt", | ||
| mode="ab", | ||
| block_size=16, | ||
| autocommit=False, | ||
| s3_additional_kwargs={"RequestPayer": "requester", "ExpectedBucketOwner": "123"}, | ||
| ) | ||
| f.write(b"bb") | ||
| f.close() | ||
|
|
||
| f.discard() | ||
|
|
||
| fs._call.assert_called_once_with( | ||
| "abort_multipart_upload", | ||
| Bucket="bucket", | ||
| Key="key.txt", | ||
| UploadId="uploadid", | ||
| RequestPayer="requester", | ||
| ExpectedBucketOwner="123", | ||
| ) | ||
| fs._finish_multipart_upload.assert_not_called() | ||
| fs._put_object.assert_not_called() | ||
|
|
||
| @pytest.mark.parametrize( | ||
| ("objects", "target"), | ||
| [ | ||
|
|
||
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 2 (claims, callers, operations): FINDINGS → repaired (PR body only)
Scope:
git diff 9b74767cf65e0ef513338fb2bb2da678cdfbcd2a..e26b3980ee627d23c71fbc6a80b34b5a08be64b7, plus the PR body and the commit message.Claims checked:
append_blockwas introduced by 4fbced8 (Support for writing with s3 file system #539); the first containing tag is v3.8.0. ✔MULTIPART_UPLOAD_MIN_PART_SIZEis copied withUploadPartCopyas the first parts, whatever the block size. This matchess3.py:2200/:2216/:2269. ✔ (It replaces Document every public API in pyathena/ and check docstrings with ruff #919's "once a multipart upload starts", which described the old behavior.)EntityTooSmall": measured. Offline, the fake givesb''andb'aaaa…'; on real S3,CompleteMultipartUploadreturnedEntityTooSmall. ✔block_size=None→default_block_size(s3.py:1925). ✔AioS3FileSystemgets the same fix":AioS3File(s3_async.py:547) overrides no methods. ✔s3.py(3 of 4 cases fail as described) and the append tests on e26b398 (114 passed). The body now states which commit each result comes from.Caller and operator view:
append_blockis an unprefixed attribute, but no caller in the repository reads it. It now means "copied server-side", the same meaning s3fs gives it.x-amz-copy-source-if-match, so a concurrent overwrite betweeninfo()and the copy is not detected. Recorded, not changed here.docs/searched).