From 5d9482b19cdb0135e499936c442160305c6773d9 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 15:08:25 +0900 Subject: [PATCH 1/4] Accept version_id in S3FileSystem.open() and cat_file() open() passed version_id=None to S3File and forwarded the remaining keywords, so open(path, version_id=...) raised TypeError; cat_file() passed the version parsed from the path to _get_object() along with the keywords and raised the same way. AioS3FileSystem.open() had the same collision, and get_file() reached it through open(). open() now hands version_id to S3File, which already rejects a version that differs from the one in the path. cat_file() resolves it like info(): the version in the path takes precedence, and a range is measured against the size of that version. S3File looked up the object after the base class initializer, which takes the read size from info() of the path alone, so a version given only as an argument was read with the size of the latest version. The object is now looked up first and its size passed to the initializer, which also drops the second info() call on open. Closes #936 Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 58 +++++++++++-------- pyathena/filesystem/s3_async.py | 1 - tests/pyathena/filesystem/test_s3.py | 66 ++++++++++++++++++++++ tests/pyathena/filesystem/test_s3_async.py | 20 ++++++- 4 files changed, 118 insertions(+), 27 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index ba7abb984..01e6c34b1 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -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, @@ -2183,16 +2187,6 @@ def __init__( 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,16 +2203,13 @@ 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.") - 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. @@ -2226,7 +2217,26 @@ def __init__( 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 +2249,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]] = [] diff --git a/pyathena/filesystem/s3_async.py b/pyathena/filesystem/s3_async.py index cb0592c25..2d6dd3e5a 100644 --- a/pyathena/filesystem/s3_async.py +++ b/pyathena/filesystem/s3_async.py @@ -366,7 +366,6 @@ def _open( self._sync_fs, path, mode, - version_id=None, max_workers=max_workers, executor=S3AioExecutor(loop=self._loop), block_size=block_size, diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 2d0fa6cc0..16c0b0778 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -438,6 +438,53 @@ def test_open_max_workers(self): with fs.open("s3://bucket/key", "wb", max_workers=2) as f: assert f.max_workers == 2 + def test_open_version_id(self): + fs = self._make_fs() + fs.default_cache_type = "bytes" + fs.info = mock.MagicMock(return_value=self._file_object("key")) + fs.info.return_value.size = 4 + + with fs.open("s3://bucket/key", "rb", version_id="v1") as f: + assert f.version_id == "v1" + # The size is that of the requested version, not the latest one. + assert f.size == 4 + fs.info.assert_called_once_with("bucket/key", version_id="v1") + # The argument must match the version in the path. + with pytest.raises(ValueError, match="do not match"): + fs.open("s3://bucket/key?versionId=v2", "rb", version_id="v1") + + @pytest.mark.parametrize( + ("path", "expected"), + [ + ("s3://bucket/key", "v1"), + # The version in the path takes precedence, as in info(). + ("s3://bucket/key?versionId=v2", "v2"), + ], + ) + def test_cat_file_version_id(self, path, expected): + fs = self._make_fs() + fs.info = mock.MagicMock(return_value=self._file_object("key")) + fs.info.return_value.size = 10 + + fs._call.return_value = {"Body": io.BytesIO(b"data")} + assert fs.cat_file(path, version_id="v1") == b"data" + fs._call.assert_called_once_with( + fs._client.get_object, Bucket="bucket", Key="key", VersionId=expected + ) + + # A range is resolved against the size of the same version. + fs._call.reset_mock() + fs._call.return_value = {"Body": io.BytesIO(b"ta")} + assert fs.cat_file(path, start=2, end=4, version_id="v1") == b"ta" + fs.info.assert_called_once_with(path, version_id=expected) + fs._call.assert_called_once_with( + fs._client.get_object, + Bucket="bucket", + Key="key", + Range="bytes=2-3", + VersionId=expected, + ) + def test_finish_multipart_upload(self): fs = self._make_fs() fs._complete_multipart_upload = mock.MagicMock() @@ -1676,6 +1723,25 @@ def test_object_version_info(self, fs): # An unversioned bucket reports the "null" version. assert version.version_id + def test_read_version_id(self, fs): + path = ( + f"s3://{ENV.s3_staging_bucket}/{ENV.s3_staging_key}{ENV.schema}/" + f"filesystem/test_read_version_id/{uuid.uuid4()}" + ) + data = b"0123456789" + fs.pipe(path, data) + # An unversioned bucket reports the "null" version, which can be + # read explicitly. + version_id = fs.object_version_info(path)[0].version_id + + assert fs.cat_file(path, version_id=version_id) == data + assert fs.cat_file(path, start=2, end=5, version_id=version_id) == data[2:5] + with fs.open(path, "rb", version_id=version_id) as f: + assert f.read() == data + # The version reaches S3, which rejects an unknown one. + with pytest.raises(OSError, match="Invalid version id"): + fs.cat_file(path, version_id="invalid") + @pytest.mark.parametrize("fs", [{"version_aware": True}], indirect=True) def test_version_aware_read(self, fs): # On an unversioned bucket, the version-aware mode is a no-op for diff --git a/tests/pyathena/filesystem/test_s3_async.py b/tests/pyathena/filesystem/test_s3_async.py index fc0c354f5..a8c167ac3 100644 --- a/tests/pyathena/filesystem/test_s3_async.py +++ b/tests/pyathena/filesystem/test_s3_async.py @@ -15,7 +15,7 @@ from pyathena.filesystem.s3 import S3File, S3FileSystem from pyathena.filesystem.s3_async import AioS3File, AioS3FileSystem -from pyathena.filesystem.s3_object import S3ObjectType, S3StorageClass +from pyathena.filesystem.s3_object import S3Object, S3ObjectType, S3StorageClass from tests import ENV from tests.pyathena.conftest import connect @@ -901,6 +901,24 @@ def test_open_max_workers(self): assert isinstance(f, AioS3File) assert f.max_workers == 2 + def test_open_version_id(self): + fs = AioS3FileSystem(connection=mock.MagicMock(), skip_instance_cache=True) + fs._sync_fs.info = mock.MagicMock( + return_value=S3Object( + init={"Key": "key"}, + type=S3ObjectType.S3_OBJECT_TYPE_FILE, + bucket="bucket", + key="key", + ) + ) + fs._sync_fs.info.return_value.size = 4 + + with fs.open("s3://bucket/key", "rb", version_id="v1") as f: + assert isinstance(f, AioS3File) + assert f.version_id == "v1" + assert f.size == 4 + fs._sync_fs.info.assert_called_once_with("bucket/key", version_id="v1") + @pytest.mark.parametrize( ("objects", "target"), [ From 40f8c4ca3a059e741e73345eb05281b151d23ce3 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 15:13:55 +0900 Subject: [PATCH 2/4] Look up a version without the cache of the latest version info() looked up a version given as an argument in the directory cache as if it were the object path: a cached entry or a parent listing returned the latest version, and a HeadObject result was cached under the object path, so later lookups of the latest version returned it. A version in the path was looked up again from its own cache entry, whose name has no version, and returned as a directory; with a cached parent listing it raised FileNotFoundError. open() and cat_file() now look versions up this way, so a version skips the cached entries of the latest version and is cached under its version-qualified path, as with the ?versionId= suffix. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 13 ++++++---- tests/pyathena/filesystem/test_s3.py | 36 ++++++++++++++++++++++++++++ 2 files changed, 45 insertions(+), 4 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 01e6c34b1..afb0cd0ee 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -585,7 +585,9 @@ def info(self, path: str, **kwargs) -> S3Object: exists, with a ListObjectsV2 request (``Delimiter="/"``, ``MaxKeys=1``) that checks whether it is a key prefix; a bucket path is looked up with HeadBucket. With ``version_aware``, a cached file - entry without a version ID is looked up again. + entry without a version ID is looked up again. A version is looked up + with HeadObject and cached under its version-qualified path, because + listings and the entry of the object path describe the latest version. Args: path: S3 path (e.g., "s3://bucket" or "s3://bucket/key"). @@ -617,7 +619,7 @@ def info(self, path: str, **kwargs) -> S3Object: key=None, version_id=None, ) - if not refresh: + if not refresh and not version_id: caches: list[S3Object] | S3Object | None = self._ls_from_cache(path) if caches is not None: if isinstance(caches, list): @@ -630,7 +632,6 @@ def info(self, path: str, **kwargs) -> S3Object: if cache: if ( self.version_aware - and not version_id and cache.get("type") == S3ObjectType.S3_OBJECT_TYPE_FILE and not cache.get("version_id") ): @@ -645,7 +646,11 @@ def info(self, path: str, **kwargs) -> S3Object: bucket, key.rstrip("/") if key else None, version_id ) if key: - object_info = self._head_object(path, refresh=refresh, version_id=version_id) + # Cache a version under the same path as the ?versionId= suffix. + head_path = ( + f"{path}?versionId={version_id}" if version_id and not path_version_id else path + ) + object_info = self._head_object(head_path, refresh=refresh) if object_info: return object_info else: diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 16c0b0778..1ab7d2979 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -438,6 +438,42 @@ def test_open_max_workers(self): with fs.open("s3://bucket/key", "wb", max_workers=2) as f: assert f.max_workers == 2 + @pytest.mark.parametrize( + "dircache", + [ + {}, + # The entry of the object path and the listing of its parent + # describe the latest version. + {"bucket/dir/key": _file_object("dir/key")}, + {"bucket/dir": [_file_object("dir/key")]}, + ], + ) + @pytest.mark.parametrize( + ("path", "kwargs"), + [ + ("s3://bucket/dir/key", {"version_id": "v1"}), + ("s3://bucket/dir/key?versionId=v1", {}), + ], + ) + def test_info_version_id(self, dircache, path, kwargs): + fs = self._make_fs() + fs.dircache.update(dircache) + fs._call.return_value = {"ContentLength": 4, "ETag": '"etag"'} + + for _ in range(2): + info = fs.info(path, **kwargs) + assert (info.type, info.size, info.version_id) == ( + S3ObjectType.S3_OBJECT_TYPE_FILE, + 4, + "v1", + ) + # The version is looked up once and cached under its + # version-qualified path, apart from the latest version. + fs._call.assert_called_once_with( + fs._client.head_object, Bucket="bucket", Key="dir/key", VersionId="v1" + ) + assert fs.dircache == {**dircache, "bucket/dir/key?versionId=v1": info} + def test_open_version_id(self): fs = self._make_fs() fs.default_cache_type = "bytes" From dffe1e52461252f341203d41996593d5a12103db Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 15:26:47 +0900 Subject: [PATCH 3/4] Leave version caching in info() to #957 Revert "Look up a version without the cache of the latest version" (40f8c4ca). #957 changes how info() caches explicit versions for #932, so the overlapping change is dropped here. This reverts commit 40f8c4ca3a059e741e73345eb05281b151d23ce3. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 13 ++++------ tests/pyathena/filesystem/test_s3.py | 36 ---------------------------- 2 files changed, 4 insertions(+), 45 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index afb0cd0ee..01e6c34b1 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -585,9 +585,7 @@ def info(self, path: str, **kwargs) -> S3Object: exists, with a ListObjectsV2 request (``Delimiter="/"``, ``MaxKeys=1``) that checks whether it is a key prefix; a bucket path is looked up with HeadBucket. With ``version_aware``, a cached file - entry without a version ID is looked up again. A version is looked up - with HeadObject and cached under its version-qualified path, because - listings and the entry of the object path describe the latest version. + entry without a version ID is looked up again. Args: path: S3 path (e.g., "s3://bucket" or "s3://bucket/key"). @@ -619,7 +617,7 @@ def info(self, path: str, **kwargs) -> S3Object: key=None, version_id=None, ) - if not refresh and not version_id: + if not refresh: caches: list[S3Object] | S3Object | None = self._ls_from_cache(path) if caches is not None: if isinstance(caches, list): @@ -632,6 +630,7 @@ def info(self, path: str, **kwargs) -> S3Object: if cache: if ( self.version_aware + and not version_id and cache.get("type") == S3ObjectType.S3_OBJECT_TYPE_FILE and not cache.get("version_id") ): @@ -646,11 +645,7 @@ def info(self, path: str, **kwargs) -> S3Object: bucket, key.rstrip("/") if key else None, version_id ) if key: - # Cache a version under the same path as the ?versionId= suffix. - head_path = ( - f"{path}?versionId={version_id}" if version_id and not path_version_id else path - ) - object_info = self._head_object(head_path, refresh=refresh) + object_info = self._head_object(path, refresh=refresh, version_id=version_id) if object_info: return object_info else: diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 1ab7d2979..16c0b0778 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -438,42 +438,6 @@ def test_open_max_workers(self): with fs.open("s3://bucket/key", "wb", max_workers=2) as f: assert f.max_workers == 2 - @pytest.mark.parametrize( - "dircache", - [ - {}, - # The entry of the object path and the listing of its parent - # describe the latest version. - {"bucket/dir/key": _file_object("dir/key")}, - {"bucket/dir": [_file_object("dir/key")]}, - ], - ) - @pytest.mark.parametrize( - ("path", "kwargs"), - [ - ("s3://bucket/dir/key", {"version_id": "v1"}), - ("s3://bucket/dir/key?versionId=v1", {}), - ], - ) - def test_info_version_id(self, dircache, path, kwargs): - fs = self._make_fs() - fs.dircache.update(dircache) - fs._call.return_value = {"ContentLength": 4, "ETag": '"etag"'} - - for _ in range(2): - info = fs.info(path, **kwargs) - assert (info.type, info.size, info.version_id) == ( - S3ObjectType.S3_OBJECT_TYPE_FILE, - 4, - "v1", - ) - # The version is looked up once and cached under its - # version-qualified path, apart from the latest version. - fs._call.assert_called_once_with( - fs._client.head_object, Bucket="bucket", Key="dir/key", VersionId="v1" - ) - assert fs.dircache == {**dircache, "bucket/dir/key?versionId=v1": info} - def test_open_version_id(self): fs = self._make_fs() fs.default_cache_type = "bytes" From bbf3759dee3bf1f4f9dea4fe7a26084972c06b07 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 15:39:22 +0900 Subject: [PATCH 4/4] Keep the version in the path of S3File and reject it for writing A version_id argument was kept only in S3File.version_id, so fsspec reopened an unpickled file at the path without the version, and an append read the existing content with cat(self.path) and copied it with UploadPartCopy from the latest version while sizing it from the given version. The version is now carried in the path, as with the ?versionId= suffix. Writing to or appending from a version has no meaning for S3, and pipe_file(), touch() and copy already reject a versioned target, so S3File now raises ValueError for a version in either form outside read mode. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 12 ++++++++++-- tests/pyathena/filesystem/test_s3.py | 22 +++++++++++++++++++++- tests/pyathena/filesystem/test_s3_async.py | 2 +- 3 files changed, 32 insertions(+), 4 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 01e6c34b1..827a623e5 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -2162,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 @@ -2181,7 +2182,8 @@ 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) @@ -2203,6 +2205,12 @@ def __init__( self.version_id = path_version_id else: self.version_id = version_id + if self.version_id and "r" not in mode: + raise ValueError("Cannot write to the file with the version specified.") + 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}" self._details: S3Object | dict[str, Any] = {} if "r" in mode: diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 16c0b0778..704ffc3e2 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -448,11 +448,31 @@ def test_open_version_id(self): assert f.version_id == "v1" # The size is that of the requested version, not the latest one. assert f.size == 4 - fs.info.assert_called_once_with("bucket/key", version_id="v1") + fs.info.assert_called_once_with("bucket/key?versionId=v1", version_id="v1") + # The version is carried in the path, which fsspec reopens an + # unpickled file with. + assert f.path == "bucket/key?versionId=v1" + assert f.__reduce__()[1][1] == "bucket/key?versionId=v1" # The argument must match the version in the path. with pytest.raises(ValueError, match="do not match"): fs.open("s3://bucket/key?versionId=v2", "rb", version_id="v1") + @pytest.mark.parametrize("mode", ["wb", "ab", "xb"]) + @pytest.mark.parametrize( + ("path", "kwargs"), + [ + ("s3://bucket/key", {"version_id": "v1"}), + ("s3://bucket/key?versionId=v1", {}), + ], + ) + def test_open_version_id_for_writing(self, mode, path, kwargs): + fs = self._make_fs() + fs.default_cache_type = "bytes" + fs._call.side_effect = AssertionError("No request is expected.") + + with pytest.raises(ValueError, match="version specified"): + fs.open(path, mode, **kwargs) + @pytest.mark.parametrize( ("path", "expected"), [ diff --git a/tests/pyathena/filesystem/test_s3_async.py b/tests/pyathena/filesystem/test_s3_async.py index a8c167ac3..ad30d9e8f 100644 --- a/tests/pyathena/filesystem/test_s3_async.py +++ b/tests/pyathena/filesystem/test_s3_async.py @@ -917,7 +917,7 @@ def test_open_version_id(self): assert isinstance(f, AioS3File) assert f.version_id == "v1" assert f.size == 4 - fs._sync_fs.info.assert_called_once_with("bucket/key", version_id="v1") + fs._sync_fs.info.assert_called_once_with("bucket/key?versionId=v1", version_id="v1") @pytest.mark.parametrize( ("objects", "target"),