-
Notifications
You must be signed in to change notification settings - Fork 116
Clamp byte ranges to the object in cat_file() and S3File reads #987
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
a906f0d
793565c
721a0e2
ee4f5ee
fd97af3
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 | ||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
|
@@ -1458,6 +1458,13 @@ def cat_file( | |||||||||||||
| ) -> bytes: | ||||||||||||||
| """Read the contents of an S3 object with GetObject. | ||||||||||||||
|
|
||||||||||||||
| ``start`` and ``end`` select bytes like a slice of the object: an | ||||||||||||||
| empty range, or one that starts at or past the end of the object, | ||||||||||||||
| returns ``b""``, and an end past the object reads up to its end. | ||||||||||||||
| Non-negative offsets are sent to S3 as they are; a negative offset is | ||||||||||||||
| resolved against the size from :meth:`info`, which also checks that | ||||||||||||||
| the object exists for an empty range. | ||||||||||||||
|
|
||||||||||||||
| Args: | ||||||||||||||
| path: S3 path (s3://bucket/key) of the object. | ||||||||||||||
| start: Byte offset to start reading at. A negative value counts | ||||||||||||||
|
|
@@ -1470,38 +1477,51 @@ def cat_file( | |||||||||||||
|
|
||||||||||||||
| Returns: | ||||||||||||||
| The bytes read from the object. | ||||||||||||||
|
|
||||||||||||||
| Raises: | ||||||||||||||
| FileNotFoundError: If the key does not exist. | ||||||||||||||
| """ | ||||||||||||||
| 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 and start < 0) or ( | ||||||||||||||
| end is not None and (end < 0 or (start or 0) >= end) | ||||||||||||||
|
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 narrow follow-up review (relayed): Codex CLI 0.160.0, model Reviewer output, verbatim: CLEAN — no actionable regressions found in Surfaces covered:
Static comparison of the changed tests predicts seven distinguishing cases:
The other parameterized cases retain their previous behavior. Missing-version handling was traced through source; the new missing-object tests mock Pre-existing limitation: lookup freshness still follows Static review only; no tests, builds, writes, or network access. HEAD remained |
||||||||||||||
| ): | ||||||||||||||
| # A negative offset needs the size of the object, and an empty | ||||||||||||||
| # range sends no GetObject request that would report a missing | ||||||||||||||
| # object. | ||||||||||||||
| info = self.info(path, version_id=version_id) | ||||||||||||||
| if info.get("type") == S3ObjectType.S3_OBJECT_TYPE_DIRECTORY: | ||||||||||||||
|
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 one (implementation behavior). Base Covered: Finding: Out of scope, pre-existing:
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. Repaired in 2452e4f:
|
||||||||||||||
| # There is no object to read, as GetObject reports for the | ||||||||||||||
| # other ranges. | ||||||||||||||
| raise FileNotFoundError(path) | ||||||||||||||
| start, end, _ = slice(start, end).indices(info.get("size", 0)) | ||||||||||||||
|
|
||||||||||||||
| ranges: tuple[int, int | None] | None = None | ||||||||||||||
| if start is not None or end is not None: | ||||||||||||||
| size = self.info(path, version_id=version_id).get("size", 0) | ||||||||||||||
| if start is None: | ||||||||||||||
| range_start = 0 | ||||||||||||||
| elif start < 0: | ||||||||||||||
| range_start = size + start | ||||||||||||||
| else: | ||||||||||||||
| range_start = start | ||||||||||||||
|
|
||||||||||||||
| if end is None: | ||||||||||||||
| range_end = size | ||||||||||||||
| elif end < 0: | ||||||||||||||
| range_end = size + end | ||||||||||||||
| else: | ||||||||||||||
| range_end = end | ||||||||||||||
|
|
||||||||||||||
| ranges = (range_start, range_end) | ||||||||||||||
| else: | ||||||||||||||
| ranges = None | ||||||||||||||
|
|
||||||||||||||
| return self._get_object( | ||||||||||||||
| bucket=bucket, | ||||||||||||||
| key=cast(str, key), | ||||||||||||||
| ranges=ranges, | ||||||||||||||
| version_id=version_id, | ||||||||||||||
| **kwargs, | ||||||||||||||
| )[1] | ||||||||||||||
| start = start or 0 | ||||||||||||||
| if end is not None and start >= end: | ||||||||||||||
| # S3 would return the whole object for an empty range. | ||||||||||||||
| return b"" | ||||||||||||||
| ranges = (start, end) | ||||||||||||||
| try: | ||||||||||||||
| return self._get_object( | ||||||||||||||
| bucket=bucket, | ||||||||||||||
| key=cast(str, key), | ||||||||||||||
| ranges=ranges, | ||||||||||||||
| version_id=version_id, | ||||||||||||||
| **kwargs, | ||||||||||||||
| )[1] | ||||||||||||||
| except OSError as e: | ||||||||||||||
| if ( | ||||||||||||||
| ranges | ||||||||||||||
| and isinstance(e.__cause__, botocore.exceptions.ClientError) | ||||||||||||||
| and S3ClientError(e.__cause__).code == "InvalidRange" | ||||||||||||||
|
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 one, after the independent-review repair (implementation behavior). This is a full pass, because the repair expanded the contract. Base Covered:
Finding: the condition Out of scope: recursive |
||||||||||||||
| ): | ||||||||||||||
| # The range starts at or past the end of the object. | ||||||||||||||
| return b"" | ||||||||||||||
| raise | ||||||||||||||
|
|
||||||||||||||
| def put_file(self, lpath: str, rpath: str, callback=_DEFAULT_CALLBACK, **kwargs): | ||||||||||||||
| """Upload a local file to S3. | ||||||||||||||
|
|
@@ -1567,7 +1587,9 @@ def get_file(self, rpath: str, lpath: str, callback=_DEFAULT_CALLBACK, outfile=N | |||||||||||||
| if os.path.isdir(lpath): | ||||||||||||||
| return | ||||||||||||||
|
|
||||||||||||||
| with open(lpath, "wb") as local, self.open(rpath, "rb", **kwargs) as remote: | ||||||||||||||
| # The remote file is opened first so that no local file is created | ||||||||||||||
| # when open() finds no object at the path. | ||||||||||||||
| with self.open(rpath, "rb", **kwargs) as remote, open(lpath, "wb") as local: | ||||||||||||||
| callback.set_size(remote.size) | ||||||||||||||
| while data := remote.read(remote.blocksize): | ||||||||||||||
| local.write(data) | ||||||||||||||
|
|
@@ -2079,12 +2101,33 @@ def _get_object( | |||||||||||||
| self, | ||||||||||||||
| bucket: str, | ||||||||||||||
| key: str, | ||||||||||||||
| ranges: tuple[int, int] | None = None, | ||||||||||||||
| ranges: tuple[int, int | None] | None = None, | ||||||||||||||
| version_id: str | None = None, | ||||||||||||||
| **kwargs, | ||||||||||||||
| ) -> tuple[int, bytes]: | ||||||||||||||
| """Read an object or a byte range of it with GetObject. | ||||||||||||||
|
|
||||||||||||||
| Args: | ||||||||||||||
| bucket: The bucket name. | ||||||||||||||
| key: The object key. | ||||||||||||||
| ranges: The ``(start, end)`` byte range to read, with an exclusive | ||||||||||||||
| end or ``None`` to read to the end of the object, or ``None`` | ||||||||||||||
| to read the whole object. | ||||||||||||||
| version_id: The version ID to read, or ``None`` for the latest. | ||||||||||||||
| **kwargs: Additional parameters passed to the GetObject API. | ||||||||||||||
|
|
||||||||||||||
| Returns: | ||||||||||||||
| Tuple of the start of the range (0 for the whole object) and the | ||||||||||||||
| bytes read. | ||||||||||||||
|
|
||||||||||||||
| Raises: | ||||||||||||||
| ValueError: If the range is empty. S3 ignores a range whose last | ||||||||||||||
| byte precedes its first byte and returns the whole object. | ||||||||||||||
| """ | ||||||||||||||
| request = {"Bucket": bucket, "Key": key} | ||||||||||||||
| if ranges: | ||||||||||||||
| if ranges[1] is not None and ranges[0] >= ranges[1]: | ||||||||||||||
| raise ValueError(f"Invalid empty range: {ranges}.") | ||||||||||||||
| range_ = S3File._format_ranges(ranges) | ||||||||||||||
| request.update({"Range": range_}) | ||||||||||||||
| else: | ||||||||||||||
|
|
@@ -2270,6 +2313,8 @@ def __init__( | |||||||||||||
| **kwargs: Accepted for compatibility; not used. | ||||||||||||||
|
|
||||||||||||||
| Raises: | ||||||||||||||
| FileNotFoundError: If no object exists at the path when reading, | ||||||||||||||
| including when the path is a prefix. | ||||||||||||||
| ValueError: If the path has no key, the version IDs do not match, | ||||||||||||||
| a version is given for writing, or the block size is not | ||||||||||||||
| between ``MULTIPART_UPLOAD_MIN_PART_SIZE`` and | ||||||||||||||
|
|
@@ -2321,6 +2366,9 @@ def __init__( | |||||||||||||
| # 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 info.get("type") == S3ObjectType.S3_OBJECT_TYPE_DIRECTORY: | ||||||||||||||
| # A prefix has no object to read. | ||||||||||||||
|
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, after the independent-review repair (claims, callers, operations). This is a full pass. Base Claims checked:
Corrections to the PR body:
Callers and operations:
Imprecision not rewritten: commit 20d0913 says that the synthesized size "made reads of keys ending in '/' empty". That describes 0b0e150; at 2452e4f, the same read raised
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 after rebase (rounds one and two, narrow). Old
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 rebase follow-up review (relayed): Codex CLI 0.160.0, model Reviewer output, verbatim: CLEAN — no actionable findings in the rebase or the inspected upstream interactions. Surfaces covered:
Regressions: None found. Static review only; no edits, builds, tests, GitHub, or network access. HEAD remains |
||||||||||||||
| raise FileNotFoundError(path) | ||||||||||||||
| 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 | ||||||||||||||
|
|
@@ -2602,6 +2650,23 @@ def setxattr(self, copy_kwargs: dict[str, Any] | None = None, **kwargs) -> None: | |||||||||||||
| self.fs.setxattr(self.path, copy_kwargs=copy_kwargs, **kwargs) | ||||||||||||||
|
|
||||||||||||||
| def _fetch_range(self, start: int, end: int) -> bytes: | ||||||||||||||
| """Read a byte range of the object for the fsspec cache. | ||||||||||||||
|
|
||||||||||||||
| The range is clamped to the size of the object, since fsspec caches | ||||||||||||||
| may request a range that is empty or reaches past the end of the | ||||||||||||||
| object. S3 would answer the former with the whole object and a range | ||||||||||||||
| starting past the end with an ``InvalidRange`` error. | ||||||||||||||
|
|
||||||||||||||
| Args: | ||||||||||||||
| start: The offset of the first byte to read. | ||||||||||||||
| end: The offset to stop reading at (exclusive). | ||||||||||||||
|
|
||||||||||||||
| Returns: | ||||||||||||||
| The bytes read, empty if the clamped range is empty. | ||||||||||||||
| """ | ||||||||||||||
| end = min(end, self.size) | ||||||||||||||
| if start >= end: | ||||||||||||||
| return b"" | ||||||||||||||
| ranges = self._get_ranges( | ||||||||||||||
| start, end, max_workers=self.max_workers, worker_block_size=self.blocksize | ||||||||||||||
| ) | ||||||||||||||
|
|
@@ -2629,8 +2694,18 @@ def _fetch_range(self, start: int, end: int) -> bytes: | |||||||||||||
| return object_ | ||||||||||||||
|
|
||||||||||||||
| @staticmethod | ||||||||||||||
| def _format_ranges(ranges: tuple[int, int]): | ||||||||||||||
| return f"bytes={ranges[0]}-{ranges[1] - 1}" | ||||||||||||||
| def _format_ranges(ranges: tuple[int, int | None]) -> str: | ||||||||||||||
| """Format a byte range as the value of an HTTP ``Range`` header. | ||||||||||||||
|
|
||||||||||||||
| Args: | ||||||||||||||
| ranges: The ``(start, end)`` byte range, with an exclusive end or | ||||||||||||||
| ``None`` for the end of the object. | ||||||||||||||
|
|
||||||||||||||
| Returns: | ||||||||||||||
| The range, such as ``bytes=0-99`` or ``bytes=100-``. | ||||||||||||||
| """ | ||||||||||||||
| start, end = ranges | ||||||||||||||
| return f"bytes={start}-" if end is None else f"bytes={start}-{end - 1}" | ||||||||||||||
|
|
||||||||||||||
| @staticmethod | ||||||||||||||
| def _get_ranges( | ||||||||||||||
|
|
||||||||||||||
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.
Independent follow-up review (relayed): Codex CLI 0.160.0, model
gpt-6-astra, reasoning effortmax,codex exec -s read-only, session01a100df-9260-77a2-907b-ca289cf98428. Full review of the expanded scope on a detached snapshot of head02f8ebdde0c337ac59d6bc71f6d479eea1da0c6cagainst baseaa0fc9146f4e683d6cdf96b894f963f9dd8f7abe. The prompt contained the literal diff and the repository conventions, without PR text, commit messages or prior findings. This is a static review. Afterwards, the snapshot and the PR worktree were unchanged at that head. Result: FINDINGS.Reviewer output, verbatim:
Covered the supplied diff, synchronous and async
cat_file/cat_ranges/get_file,_get_objectcallers, version handling, prefixes, buckets, empty objects, error translation/retries, docstrings, and tests. Traced all ten built-in fsspec caches:none,bytes,readahead,adaptive,first,mmap,blockcache,background,all, andparts.Used installed fsspec 2026.9.0 source from the main checkout’s
.venv. The s3fs comparison uses API knowledge; s3fs was not installed.FINDINGS
P2 — Regression: empty ranges suppress missing-object errors.
pyathena/filesystem/s3.py:1491
With an uncached, nonexistent key,
fs.cat_file("s3://bucket/missing", start=0, end=0)now returnsb""without checking existence. Previously,info()raisedFileNotFoundError. A prefix without an object or a nonexistent version similarly returns success. This also propagates throughcat_rangesand the async wrappers, contradicting the newRaisesdocumentation and fsspec’s open-before-read behavior.Validate the target before this early return. The new tests miss this combination: missing-key coverage uses a nonempty range, and directory coverage uses negative offsets.
P3 — New comment overstates protection against a pre-existing download failure.
pyathena/filesystem/s3.py:1577
Cache metadata for a nonempty object, delete it externally, then call
get_file()with the default cache.open()succeeds using cached metadata; the destination is created or truncated; the first read raisesFileNotFoundError, leaving an empty local file. Thus the new comment’s unconditional “no local file is left behind” claim is false.The underlying behavior predates this change. Narrow the comment to failures detected during
open(). The directory test does not exercise a failure deferred until the first read.The fake deliberately rejects malformed ranges instead of returning S3’s whole-object response. Static tracing shows this still detects the targeted EOF defects; it does not falsely pass them. However, it does not model missing keys/versions or conditional-request failures.
No files changed; no builds, tests, GitHub access, or network access performed.
Author verification:
02f8ebdd, a missing key, a prefix, and a missing version with(0, 0),(5, 3)or(None, 0)returnedb""without any lookup. On master,info()raisedFileNotFoundError.open()succeeds from the cache and the first read fails.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.
Repaired in 97ce192:
endis notNoneand(start or 0) >= end) now looks up the object withinfo(), as negative offsets do. A missing object or version raisesFileNotFoundErrorfrominfo(), and a prefix raises it through the directory check. Neither sends a GetObject request.get_file()comment now says that no local file is created "when open() finds no object at the path". The failure deferred to the first read stays as before; it is pre-existing and not changed here.Validation on 97ce192:
just lint: passed.uv run --env-file .env pytest -n 4 tests/pyathena/filesystem/: 306 passed.test_cat_file_empty_range_missing((0, 0),(5, 3),(None, 0)) andtest_cat_file_range_directory[5-5]fail on02f8ebdd.test_cat_file_rangenow asserts thatinfo()is called exactly for negative offsets and empty ranges.FileNotFoundErrorfor all three empty ranges, and an existing object's(5, 5)returnsb"".Self-review of this repair:
info()call only for empty ranges. Master calledinfo()for every ranged read, so this restores master's existence check. Non-empty non-negative ranges still skipinfo().cat_file()docstring now says thatinfo()also checks existence for an empty range. The PR body is updated. The existence check may use cached metadata, as master'sinfo()did.