Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
70 changes: 43 additions & 27 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Copy link
Copy Markdown
Member Author

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.md versioning section.

Claims checked:

  • get_file() accepts version_id through open() (s3.py get_file → self.open(rpath, "rb", **kwargs)): offline run at head sends VersionId=v1 on HeadObject and GetObject.
  • AioS3FileSystem.cat_file() delegates to the sync cat_file(): offline run at head sends VersionId=v1.
  • fsspec base initializer reads the size from info(self.path) (fsspec 2026.9.0 AbstractBufferedFile.__init__/details), so the reorder in S3File.__init__ is needed; open now calls info() once (asserted in test_open_version_id).
  • info() suffix-form directory/FileNotFoundError and kwarg-form stale/poisoned cache claims: reproduced offline on the pre-repair code; test_info_version_id fails there (6/6).
  • docs/filesystem.md:139-140 ("Explicit versions can always be read with the ?versionId= suffix or the version_id argument") is now true for open()/cat_file()/get_file(); no other docs mention the argument.
  • Test counts in the PR body (12 -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 via rm(path?versionId=) already invalidates that key. Read-mode open drops one info() 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.

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:
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand All @@ -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:
Expand All @@ -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.")

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed) — Codex CLI 0.160.0, model gpt-6-astra, reasoning effort high, sandbox read-only, session 01a10067-f3d3-75c1-8e58-0382a9615c1c. Static review only (no tests/builds/network). Snapshot: detached worktree at 40f8c4c, diff e0e85da09435257a18c8b8832b3e10ba7840b697..40f8c4ca; prompt omitted PR number/body, commit messages and prior findings; snapshot unchanged afterwards.

Reviewer result (verbatim):

Covered the full diff, sync/async open, cat_file and get_file, fsspec initialization and serialization, read/write/append modes, version pinning, cache invalidation, listings, docstrings, and added tests. Static review only; no tests, builds, network access, or modifications.

FINDINGS

Introduced regression

  1. P2 — Version-qualified metadata survives object overwrite. s3.py:651
    In an unversioned bucket, call info(path, version_id="null"), then overwrite that object through pipe_file(path, new_data). The new cache key is path?versionId=null, but invalidation removes only the plain path and ancestors. Subsequent versioned info returns the old size/ETag; buffered reads can fail their stale IfMatch, and ranged cat_file can calculate incorrect offsets. Previously, the keyword lookup cached under the plain path and was invalidated. Invalidate affected version entries, particularly the mutable "null" version.

Pre-existing implementation defects exposed by the newly accepted keyword

  1. P2 — Append selects versioned metadata but reads latest-version content. s3.py:2250
    Given historical v1 = b"old" and latest v2 = b"new", open(path, "ab", version_id="v1") followed by write(b"!") writes b"new!". Metadata uses v1, but cat(self.path) has neither the suffix nor the argument. Multipart append likewise copies the unqualified path at s3.py:2294. Reject version arguments for append or propagate the selected version to every source read/copy. The underlying omission predates this diff; filesystem open previously rejected this keyword.

  2. P2 — Pickling loses the selected version. s3.py:2228
    version_id is absent from both the stored path and the kwargs passed to fsspec. Its inherited __reduce__ reconstructs the handle using those fields. With equally sized historical and current objects, pickle/unpickle of open(path, "rb", version_id="v1") therefore silently reads the current object. Preserve version identity in reconstruction. This omission already affected directly constructed S3File instances; the diff exposes it through filesystem open.

The added tests do not cover these scenarios. The integration test reads a single "null" version, and the mocked open tests verify metadata selection without reading distinct historical/latest payloads.


Author verification and disposition:

  1. Confirmed (the null version is mutable in unversioned/suspended buckets). Version caching in info() belongs to Cache each explicit version of an object separately in info() #957 (S3FileSystem.info(version_id=...) returns cached metadata of another version #932), which does not cache null; per the maintainer, the overlapping info() change was reverted here (dffe1e5) and Accept version_id in S3FileSystem.open() and cat_file() #958 depends on Cache each explicit version of an object separately in info() #957.
  2. Confirmed: the argument form left self.path unversioned, so cat(self.path) / copy_source=self.path used the latest version. Per the maintainer, S3File now rejects a version (either form) for wb/ab/xb (bbf3759, this line). Release-noted as breaking.
  3. Confirmed: fsspec __reduce__ reopens with self.path and self.kwargs. S3File now carries the argument version in the path (bbf3759); test_open_version_id asserts f.path and __reduce__() carry it. The version pinned by version_aware is still not carried (pre-existing, unchanged).

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}"

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 1 (repair 40f8c4ca..bbf3759dee3bf1f4f9dea4fe7a26084972c06b07) — CLEAN

Covered: the revert (dffe1e5 restores info() to master; git diff e0e85da09435257a18c8b8832b3e10ba7840b697..bbf3759dee3bf1f4f9dea4fe7a26084972c06b07 no longer touches info()), and S3File.__init__ (bbf3759): rejection runs before any request (asserted with _call raising); self.version_id at that point comes only from the argument or path (version_aware pinning is later and read-only); the folded path feeds info(), fsspec self.path, cat/sign/metadata/getxattr, all of which already accept the suffix form; put_file()/get_file() route through open().
Tests: test_open_version_id_for_writing (6 cases) fails fast without the change (DID NOT RAISE / unexpected request); test_open_version_id and the aio test fail without it (path lacks the version).
Validation: just lint; -k version_id 12 passed at bbf3759; tests/pyathena/filesystem/ 253 passed against AWS before the last test-only tweak (_call.side_effect).

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 2 (repair, claims) — CLEAN

Claims checked: pipe_file(), touch() and copy reject versioned targets (s3.py:1318, touch, _copy: "Cannot copy to a versioned file."); suffix-form append started from that version (cat(self.path) and CopySource="bucket/key?versionId=..." passed through unchanged by _upload_part_copy); fsspec 2026.9.0 AbstractBufferedFile.__reduce__ → reopen(fs, self.path, ..., self.kwargs).
PR body corrected: info() bullets and release note removed; breaking release note for versioned writes added; dependency on #957 stated (master's info() returns a directory on the second lookup of a ?versionId= path and caches a kwarg version under the plain path).
Correction to round 1 above: the "cached parent listing" case does not occur with current code, because _ls_dirs caches listings under (path, delimiter) keys (pointed out by the #957 session).

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up review (relayed) — Codex CLI 0.160.0, model gpt-6-astra, effort high, sandbox read-only, session 01a1007f-1822-7323-9838-18948f7c5926. Static review only. Snapshot: detached worktree at bbf3759 (unchanged afterwards); full diff e0e85da0..bbf3759d plus repair diff 40f8c4ca..bbf3759d.

Reviewer result (verbatim):

Surfaces covered: full and repair diffs; sync/aio argument forwarding; version precedence and mismatches; read/write/append/xb modes; text wrapping; version_aware pinning; size initialization; pickle reopening; put_file/get_file; transaction commit/rollback; changed docstrings, comments, and tests.

CLEAN. No introduced defects found by static inspection.

All three repairs are verified:

  1. info() caching changes are absent from the full diff; existing caching behavior was excluded.
  2. Explicit versions, supplied by argument or suffix, are rejected for wb, ab, and xb before metadata or write requests.
  3. Argument versions are preserved in S3File.path, which fsspec’s pickle reopening mechanism uses.

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 "null" version rather than distinct historical content. These limit regression evidence but do not establish an implementation defect.

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:
Expand All @@ -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]] = []
Expand Down
1 change: 0 additions & 1 deletion pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
86 changes: 86 additions & 0 deletions tests/pyathena/filesystem/test_s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -438,6 +438,73 @@ 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?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"),
[
("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()
Expand Down Expand Up @@ -1676,6 +1743,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
Expand Down
20 changes: 19 additions & 1 deletion tests/pyathena/filesystem/test_s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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?versionId=v1", version_id="v1")

@pytest.mark.parametrize(
("objects", "target"),
[
Expand Down
Loading