-
Notifications
You must be signed in to change notification settings - Fork 116
Accept version_id in S3FileSystem.open() and cat_file() #958
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
5d9482b
40f8c4c
dffe1e5
bbf3759
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 |
|---|---|---|
|
|
@@ -1394,14 +1394,19 @@ def cat_file( | |
| from the end of the object. | ||
| end: Byte offset to stop reading at (exclusive). A negative value | ||
| counts from the end of the object. | ||
| **kwargs: Additional parameters passed to the GetObject API. | ||
| **kwargs: Additional parameters passed to the GetObject API, | ||
| except ``version_id``: the version ID to read when the path | ||
| has none. | ||
|
|
||
| Returns: | ||
| The bytes read from the object. | ||
| """ | ||
| bucket, key, version_id = self.parse_path(path) | ||
| bucket, key, path_version_id = self.parse_path(path) | ||
| version_id = kwargs.pop("version_id", None) | ||
| if path_version_id: | ||
| version_id = path_version_id | ||
| if start is not None or end is not None: | ||
| size = self.info(path).get("size", 0) | ||
| size = self.info(path, version_id=version_id).get("size", 0) | ||
| if start is None: | ||
| range_start = 0 | ||
| elif start < 0: | ||
|
|
@@ -1972,7 +1977,6 @@ def _open( | |
| self, | ||
| path, | ||
| mode, | ||
| version_id=None, | ||
| max_workers=max_workers, | ||
| executor=self._create_executor(max_workers=max_workers), | ||
| block_size=block_size, | ||
|
|
@@ -2158,7 +2162,8 @@ def __init__( | |
| path: S3 path (s3://bucket/key) of the file. | ||
| mode: The file mode, such as ``rb``, ``wb`` or ``ab``. | ||
| version_id: The version ID to read. Must match the version ID in | ||
| the path if both are given. | ||
| the path if both are given. A version cannot be given, in | ||
| either form, for writing or appending. | ||
| max_workers: The number of parallel workers for range reads and | ||
| part copies. | ||
| executor: The executor for parallel operations. If None, a new | ||
|
|
@@ -2177,22 +2182,13 @@ def __init__( | |
|
|
||
| Raises: | ||
| ValueError: If the path has no key, the version IDs do not match, | ||
| or the block size is too small for writing. | ||
| a version is given for writing, or the block size is too small | ||
| for writing. | ||
| """ | ||
| self.max_workers = max_workers | ||
| self._executor: S3Executor = executor or S3ThreadPoolExecutor(max_workers=max_workers) | ||
| self.s3_additional_kwargs = s3_additional_kwargs if s3_additional_kwargs else {} | ||
|
|
||
| super().__init__( | ||
| fs=fs, | ||
| path=path, | ||
| mode=mode, | ||
| block_size=block_size, | ||
| autocommit=autocommit, | ||
| cache_type=cache_type, | ||
| cache_options=cache_options, | ||
| size=size, | ||
| ) | ||
| bucket, key, path_version_id = S3FileSystem.parse_path(path) | ||
| self.bucket = bucket | ||
| if not key: | ||
|
|
@@ -2209,24 +2205,46 @@ def __init__( | |
| self.version_id = path_version_id | ||
| else: | ||
| self.version_id = version_id | ||
| if "r" not in mode and block_size < self.fs.MULTIPART_UPLOAD_MIN_PART_SIZE: | ||
| # When writing occurs, the block size should not be smaller | ||
| # than the minimum size of a part in a multipart upload. | ||
| raise ValueError(f"Block size must be >= {self.fs.MULTIPART_UPLOAD_MIN_PART_SIZE}MB.") | ||
| if self.version_id and "r" not in mode: | ||
| raise ValueError("Cannot write to the file with the version specified.") | ||
|
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) — Codex CLI 0.160.0, model Reviewer result (verbatim): Covered the full diff, sync/async FINDINGS Introduced regression
Pre-existing implementation defects exposed by the newly accepted keyword
The added tests do not cover these scenarios. The integration test reads a single Author verification and disposition:
|
||
| if self.version_id and not path_version_id: | ||
| # 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}" | ||
|
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 (repair Covered: the revert (dffe1e5 restores
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 2 (repair, claims) — CLEAN Claims checked:
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) — Codex CLI 0.160.0, model Reviewer result (verbatim): Surfaces covered: full and repair diffs; sync/aio argument forwarding; version precedence and mismatches; read/write/append/ CLEAN. No introduced defects found by static inspection. All three repairs are verified:
Test coverage remains limited: the pickle assertion checks the serialized path without reopening; range tests use positive bounds that do not exercise version-dependent size; the integration test reads the sole No files changed, tests/builds run, or network access performed. Author disposition: no action. The coverage limits match those stated in the PR body (the CI bucket is unversioned); the cat_file unit test asserts that the size lookup carries the version, which is the version-dependent part of a range read. |
||
|
|
||
| self.append_block = False | ||
| self._details: S3Object | dict[str, Any] | ||
| self._details: S3Object | dict[str, Any] = {} | ||
| if "r" in mode: | ||
| info = self.fs.info(self.path, version_id=self.version_id) | ||
| if self.fs.version_aware and not self.version_id: | ||
| # Looked up before the base class initializer, which would | ||
| # otherwise take the size from the latest version of the object. | ||
| info = fs.info(path, version_id=self.version_id) | ||
| if fs.version_aware and not self.version_id: | ||
| # Pin the version observed at open time so that reads are | ||
| # consistent even if the object is overwritten. info() heads | ||
| # the object when the cached entry carries no version. | ||
| self.version_id = info.get("version_id") | ||
| if etag := info.get("etag"): | ||
| self.s3_additional_kwargs.update({"IfMatch": etag}) | ||
| self._details = info | ||
| elif "a" in mode and self.fs.exists(path): | ||
| if size is None: | ||
| size = info.get("size") | ||
|
|
||
| super().__init__( | ||
| fs=fs, | ||
| path=path, | ||
| mode=mode, | ||
| block_size=block_size, | ||
| autocommit=autocommit, | ||
| cache_type=cache_type, | ||
| cache_options=cache_options, | ||
| size=size, | ||
| ) | ||
| if "r" not in mode and block_size < self.fs.MULTIPART_UPLOAD_MIN_PART_SIZE: | ||
| # When writing occurs, the block size should not be smaller | ||
| # than the minimum size of a part in a multipart upload. | ||
| raise ValueError(f"Block size must be >= {self.fs.MULTIPART_UPLOAD_MIN_PART_SIZE}MB.") | ||
|
|
||
| self.append_block = False | ||
| if "a" in mode and self.fs.exists(path): | ||
| 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: | ||
|
|
@@ -2239,8 +2257,6 @@ def __init__( | |
| self.loc = loc | ||
| self.s3_additional_kwargs.update(info.to_api_repr()) | ||
| self._details = info | ||
| else: | ||
| self._details = {} | ||
|
|
||
| self.multipart_upload: S3MultipartUpload | None = None | ||
| self.multipart_upload_parts: list[Future[S3MultipartUploadPart]] = [] | ||
|
|
||
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) — CLEAN
Scope: full
git diff e0e85da09435257a18c8b8832b3e10ba7840b697..40f8c4ca3a059e741e73345eb05281b151d23ce3, PR body, both commit messages, changed docstrings/comments,docs/filesystem.mdversioning section.Claims checked:
get_file()acceptsversion_idthroughopen()(s3.pyget_file→self.open(rpath, "rb", **kwargs)): offline run at head sendsVersionId=v1on HeadObject and GetObject.AioS3FileSystem.cat_file()delegates to the synccat_file(): offline run at head sendsVersionId=v1.info(self.path)(fsspec 2026.9.0AbstractBufferedFile.__init__/details), so the reorder inS3File.__init__is needed; open now callsinfo()once (asserted intest_open_version_id).info()suffix-form directory/FileNotFoundErrorand kwarg-form stale/poisoned cache claims: reproduced offline on the pre-repair code;test_info_version_idfails there (6/6).docs/filesystem.md:139-140("Explicit versions can always be read with the?versionId=suffix or theversion_idargument") is now true foropen()/cat_file()/get_file(); no other docs mention the argument.-k version_id, 253 filesystem against AWS) re-run at this head.Callers/operations: unversioned
info()lookups are unchanged; versioned lookups now issue one HeadObject per version (cached under the version-qualified key) instead of reading latest-version cache entries, which were wrong. Version entries are immutable, and deleting a version viarm(path?versionId=)already invalidates that key. Read-mode open drops oneinfo()call.Limits: the CI bucket is unversioned, so reading an older version with a different size is covered by unit tests only; the integration test reads the
"null"version and checks that S3 rejects an unknown one.