-
Notifications
You must be signed in to change notification settings - Fork 116
Answer info() and exists() from cached parent listings and fix bucket lookups #1006
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
2efc5a1
1b26dd1
fed2d2a
eebe98d
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 |
|---|---|---|
|
|
@@ -359,7 +359,9 @@ def _head_object( | |
|
|
||
| The result is cached under the path, or under the version-qualified | ||
| path for an explicit version. An explicitly requested ``"null"`` | ||
| version is not cached. A missing object evicts its entry. | ||
| version is not cached. A missing object evicts its entry and, unless | ||
| a version was requested, the cached listing of its parent that still | ||
| lists it. | ||
|
|
||
| Args: | ||
| path: The object path, optionally with a versionId query. | ||
|
|
@@ -394,6 +396,13 @@ def _head_object( | |
| ) | ||
| except FileNotFoundError: | ||
| self._evict_cache(path) | ||
| if not version_id: | ||
| # Evict the cached listing of the parent only if it still | ||
| # lists the path. | ||
| parent_key = (self._parent(path), "/") | ||
| files = self.dircache.get(parent_key) | ||
| if files and any(f.name == path for f in files): | ||
| self._evict_cache(parent_key) | ||
| return None | ||
| if self.version_aware and not version_id: | ||
| # Pin the version of the object so that subsequent reads see | ||
|
|
@@ -628,12 +637,18 @@ def info(self, path: str, **kwargs) -> S3Object: | |
| """Return information about an S3 path. | ||
|
|
||
| Uses the directory cache first: a cached entry for the path is | ||
| returned, a cached listing of the path itself makes it a directory, | ||
| and a cached listing of its parent without it means it does not exist. | ||
| returned, the entry of the path in a cached listing of its parent is | ||
| returned, preferring an object to a key prefix of the same name as | ||
| HeadObject does, and a cached listing of its parent without it means | ||
| it does not exist. | ||
| The cached bucket listing holds only the buckets that the caller owns, | ||
| so a bucket missing from it is looked up with HeadBucket. | ||
| Otherwise, a key path is looked up with HeadObject and, if no object | ||
| 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 | ||
| is looked up with HeadBucket. If these requests find a listed object | ||
| missing, or find a key prefix, the cached listing of the parent is | ||
| removed. With ``version_aware``, a cached file | ||
| entry without a version ID is looked up again. With an explicit | ||
| version, the cached entries of the path are skipped, and the | ||
| HeadObject result is cached under the version-qualified path apart | ||
|
|
@@ -647,16 +662,16 @@ def info(self, path: str, **kwargs) -> S3Object: | |
| version_id: The version ID to look up when the path has none. | ||
|
|
||
| Returns: | ||
| S3Object describing the bucket, directory, or file. | ||
| S3Object describing the bucket, directory, or file. The root path | ||
| (``""``, ``"/"`` or ``"s3://"``) is a directory. | ||
|
|
||
| Raises: | ||
| FileNotFoundError: If the path does not exist. | ||
| """ | ||
| refresh = kwargs.pop("refresh", False) | ||
| path = self._strip_protocol(path) | ||
| bucket, key, path_version_id = self.parse_path(path) | ||
| version_id = path_version_id if path_version_id else kwargs.pop("version_id", None) | ||
| if path in ["/", ""]: | ||
| # parse_path rejects the root path. | ||
| return S3Object( | ||
| init={ | ||
| "ContentLength": 0, | ||
|
|
@@ -666,17 +681,24 @@ def info(self, path: str, **kwargs) -> S3Object: | |
| "LastModified": None, | ||
| }, | ||
| type=S3ObjectType.S3_OBJECT_TYPE_DIRECTORY, | ||
| bucket=bucket, | ||
| bucket="", | ||
| key=None, | ||
| version_id=None, | ||
| ) | ||
| bucket, key, path_version_id = self.parse_path(path) | ||
| version_id = path_version_id if path_version_id else kwargs.pop("version_id", None) | ||
| # Cached entries describe the current version of a path, so an | ||
| # explicit version uses only the HeadObject cache of that version. | ||
| 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): | ||
| cache = next((c for c in caches if c.name == path), None) | ||
| matches = [c for c in caches if c.name == path] | ||
| # A key can be both an object and a key prefix. | ||
| cache = next( | ||
| (c for c in matches if c.type == S3ObjectType.S3_OBJECT_TYPE_FILE), | ||
| next(iter(matches), None), | ||
| ) | ||
| elif caches.name == path: | ||
| cache = caches | ||
| else: | ||
|
|
@@ -720,6 +742,9 @@ def info(self, path: str, **kwargs) -> S3Object: | |
| or response.get("Contents", []) | ||
| or response.get("CommonPrefixes", []) | ||
| ): | ||
| # Nothing caches the key prefix, and the cached listing of the | ||
| # parent may predate it. | ||
| self._evict_cache((self._parent(path), "/")) | ||
| return self._directory_object(bucket, key.rstrip("/") if key else None, version_id) | ||
| raise FileNotFoundError(path) | ||
|
|
||
|
|
@@ -898,7 +923,8 @@ def exists(self, path: str, **kwargs) -> bool: | |
| refresh: If True, bypass the cache and query S3. | ||
|
|
||
| Returns: | ||
| True if the path exists, False otherwise. | ||
| True if the path exists, False otherwise. A bucket that HeadBucket | ||
| denies access to (403) exists. | ||
|
|
||
| Example: | ||
| >>> fs = S3FileSystem() | ||
|
|
@@ -919,15 +945,14 @@ def exists(self, path: str, **kwargs) -> bool: | |
| return bool(info) | ||
| except FileNotFoundError: | ||
| return False | ||
| if not refresh: | ||
| if self.dircache.get(bucket, False): | ||
| return True | ||
| try: | ||
| if self._ls_from_cache(bucket): | ||
| return True | ||
| except FileNotFoundError: | ||
| pass | ||
| file = self._head_bucket(bucket, refresh=refresh) | ||
| if not refresh and self._ls_from_cache(bucket): | ||
| return True | ||
| try: | ||
| file = self._head_bucket(bucket, refresh=refresh) | ||
| except PermissionError: | ||
| # HeadBucket answers 403 for a bucket that exists but that the | ||
| # caller may not access. | ||
| return True | ||
| return bool(file) | ||
|
|
||
| def rm_file(self, path: str, **kwargs) -> None: | ||
|
|
@@ -1231,8 +1256,8 @@ def mkdir(self, path: str, create_parents: bool = True, **kwargs) -> None: | |
| ) | ||
| except botocore.exceptions.ParamValidationError as e: | ||
| raise ValueError(f"Bucket create failed {bucket!r}: {e}") from e | ||
| # invalidate_cache walks parent paths and never pops the root | ||
| # entry itself, so evict the cached bucket listing directly. | ||
| # invalidate_cache of the bucket keeps the cached bucket | ||
| # listing, so evict it directly. | ||
| self._evict_cache("") | ||
| self.invalidate_cache(bucket) | ||
| else: | ||
|
|
@@ -1304,8 +1329,8 @@ def rmdir(self, path: str) -> None: | |
| Bucket=bucket, | ||
| ) | ||
| self.invalidate_cache(bucket) | ||
| # invalidate_cache walks parent paths and never pops the root | ||
| # entry itself, so evict the cached bucket listing directly. | ||
| # invalidate_cache of the bucket keeps the cached bucket listing, | ||
| # so evict it directly. | ||
| self._evict_cache("") | ||
|
|
||
| def touch(self, path: str, truncate: bool = True, **kwargs) -> dict[str, Any]: | ||
|
|
@@ -2222,6 +2247,8 @@ def modified(self, path: str) -> datetime: | |
| def invalidate_cache(self, path: str | None = None) -> None: | ||
| """Remove the cached entries of the path and its parent paths. | ||
|
|
||
| The cached bucket listing is removed only by the root path (``""``, | ||
| ``"/"`` or ``"s3://"``), not by the paths of buckets or keys. | ||
| A version-qualified path invalidates the version under every query | ||
| spelling that ``parse_path`` accepts, and also the object path without | ||
| the version, because deleting or copying a version can change the | ||
|
|
@@ -2234,6 +2261,8 @@ def invalidate_cache(self, path: str | None = None) -> None: | |
| self.dircache.clear() | ||
| else: | ||
| path = self._strip_protocol(path) | ||
| if not path: | ||
| self._evict_cache("") | ||
| while path: | ||
| # parse_path does not accept "?" in keys, so it starts the | ||
| # versionId query. | ||
|
|
@@ -2272,17 +2301,38 @@ def _evict_cache(self, key: str | tuple[str, str]) -> None: | |
| def _ls_from_cache(self, path: str) -> list[S3Object] | S3Object | None: | ||
| """Check the dircache for a cached entry of the path. | ||
|
|
||
| fsspec's implementation assumes every dircache value is a listing, | ||
| but S3FileSystem also caches a single S3Object under the object's own | ||
| path (HeadObject/HeadBucket results). Guard the parent lookup so that | ||
| looking up a child path of a cached object does not fail, and fall | ||
| through to the S3 API instead. | ||
| fsspec's implementation looks up listings under the path itself, but | ||
| S3FileSystem caches a single S3Object under the path of an object or | ||
| a bucket (HeadObject/HeadBucket results), the bucket listing under | ||
| ``""``, and the other listings under ``(path, delimiter)`` (see | ||
| ``_ls_dirs``). | ||
|
|
||
| Args: | ||
| path: The path without the protocol. | ||
|
|
||
| Returns: | ||
| The cached entry of the path, the entries of a cached parent | ||
| listing named as the path, or None if no cached entry describes | ||
| the path. A listing of the path itself is not used, because it | ||
| cannot tell whether an object of the same name exists. A | ||
| version-qualified path uses only its own entry, because listings | ||
| describe the current versions. | ||
|
|
||
| Raises: | ||
| FileNotFoundError: If a cached listing of the parent directory of | ||
| a key path does not contain the path. | ||
| """ | ||
| cache = self.dircache.get(path.rstrip("/")) | ||
| if cache is not None: | ||
| return cast("list[S3Object] | S3Object", cache) | ||
| parent_cache = self.dircache.get(self._parent(path)) | ||
| if isinstance(parent_cache, list): | ||
| _, key, version_id = self.parse_path(path) | ||
| if version_id: | ||
| return None | ||
| if key: | ||
| parent_cache = self.dircache.get((self._parent(path), "/")) | ||
|
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) — base Covered: Write paths: Result: FINDINGS (1).
No other findings: the version-qualified guard, the version_aware re-head, the same-name preference (object over prefix) and the bucket fallback each have a test. The |
||
| else: | ||
| parent_cache = self.dircache.get("") | ||
| if parent_cache is not None: | ||
| files = [ | ||
| f | ||
| for f in parent_cache | ||
|
|
@@ -2294,7 +2344,10 @@ def _ls_from_cache(self, path: str) -> list[S3Object] | S3Object | None: | |
| ] | ||
| if files: | ||
| return files | ||
| raise FileNotFoundError(path) | ||
| if key: | ||
| raise FileNotFoundError(path) | ||
| # The bucket listing holds only the buckets that the caller owns, | ||
| # so a bucket missing from it is looked up with HeadBucket. | ||
| return None | ||
|
|
||
| def _open( | ||
|
|
@@ -2613,9 +2666,15 @@ def __init__( | |
| self._details = info | ||
| if size is None: | ||
| size = info.get("size") | ||
| elif "a" in mode and fs.exists(path): | ||
| append_info = fs.info(path) | ||
| if append_info.get("size", 0) < fs.MULTIPART_UPLOAD_MIN_PART_SIZE: | ||
| elif "a" in mode: | ||
| # The rewritten object keeps the metadata of the existing one, | ||
| # which a cached listing entry lacks, so look up the object. | ||
| with contextlib.suppress(FileNotFoundError): | ||
| append_info = fs.info(path, refresh=True) | ||
|
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 gpt-6-astra, session Covered: Result: FINDINGS
The reviewer noted that the own-listing, version_aware and explicit-version tests also pass on the merge-base; they are compatibility guards. Verification by the author: all three reproduced (3 also on master). Repairs for 1 and 2 are in 1b26dd1 (this line and the
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 1b26dd1 — self-review of the repair, both perspectives (range: Round 1 (behavior):
Round 2 (claims/operations):
Validation at 1b26dd1: Result: CLEAN. An independent follow-up on the repair is requested next.
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 (relayed): Codex CLI 0.160.0, model gpt-6-astra, session Result: FINDINGS. Both original regressions are fixed and both new tests fail on 2efc5a1. One P2 remains, introduced relative to the merge-base: the parent-listing eviction applied only to Repair
Self-review of the repair:
Tests:
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. Second independent follow-up (relayed): Codex CLI 0.160.0, model gpt-6-astra, session Covered: Result: CLEAN.
Independent review complete for head fed2d2a.
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 Self-review of the repair:
Tests:
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 (relayed): Codex CLI 0.160.0, model gpt-6-astra, session Covered: root normalization, Result: CLEAN.
Independent review complete for head eebe98d. |
||
| if ( | ||
| append_info is not None | ||
| and append_info.get("size", 0) < fs.MULTIPART_UPLOAD_MIN_PART_SIZE | ||
| ): | ||
| # Too small to be a part of a multipart upload: rewritten | ||
| # from the buffer. | ||
| append_data = fs.cat(path) | ||
|
|
||
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) — base
9b2f033708b81dce7323dfda61e8b2ec95df23da, head2efc5a10e90a8332d998ba6159a0a673b16c96a7.Claims checked:
git tag --contains 681e7496→ first tag is v3.15.0. ✔isdir/isfile/sizecallinfo().ducallsfind()andisdir(), and reachesinfo()only through_find's object fallback. Corrected the PR body to name onlyisdir()/isfile()/size()._infouses fsspec_ls_from_cache, which raisesFileNotFoundErrorfrom a parent listing without the path. ✔ s3fs also answers "directory" from a listing of the path itself, which this PR deliberately does not do.ls/pipe/info→PermissionErrorafterexists()returned True. ✔docs/filesystem.md:101-104("info/isfile/opentreatdir/asdir: the objectdirif it exists, otherwise the directory"): still holds with a cached parent listing, because the object entry is preferred (this line). It would not have held if the path's own listing were used. No other docs mention the listing cache orexists()of buckets. The 403 →PermissionErrortable stays true for the other operations.Adversarial callers:
//has a parent with a trailing slash, which never matches a stripped listing key, so it falls through to HeadObject (no falseFileNotFoundError).dir/is listed as a CommonPrefix of the parent, so it does not shadow the objectdir.Evidence: the live
tests/pyathena/filesystemrun (422 passed) was local at this head. Draft CI ran only check/lint/offline (all pass); AWS jobs run on Ready. Theexists()403 path is offline-only.Result: FINDINGS (1, PR-text only, repaired). No code change.