-
Notifications
You must be signed in to change notification settings - Fork 116
Fix stale S3FileSystem listing and object caches #928
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
f7e8690
4d21894
ab9f48a
38a2e28
4ad7674
c982550
a95e84f
51a4362
b16c58f
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 |
|---|---|---|
|
|
@@ -315,6 +315,11 @@ def _head_bucket(self, bucket, refresh: bool = False) -> S3Object | None: | |
| Bucket=bucket, | ||
| ) | ||
| except FileNotFoundError: | ||
| self.dircache.pop(bucket, None) | ||
| # Evict the cached bucket listing only if it still lists the bucket. | ||
| buckets = self.dircache.get("") | ||
| if buckets and any(b.name == bucket for b in buckets): | ||
| self.dircache.pop("", None) | ||
| return None | ||
| file = S3Object( | ||
| init={ | ||
|
|
@@ -352,6 +357,7 @@ def _head_object( | |
| **request, | ||
| ) | ||
| except FileNotFoundError: | ||
| self.dircache.pop(path, None) | ||
| return None | ||
| if self.version_aware and not version_id: | ||
| # Pin the version of the object so that subsequent reads see | ||
|
|
@@ -404,13 +410,34 @@ def _ls_dirs( | |
| max_keys: int | None = None, | ||
| refresh: bool = False, | ||
| ) -> list[S3Object]: | ||
| """List the objects and common prefixes under a path. | ||
|
|
||
| A complete, non-empty listing of the path is cached under | ||
| ``(path, delimiter)``, and an empty one evicts it. | ||
| ``invalidate_cache`` drops it when the path or a path under it is | ||
| invalidated. | ||
|
|
||
| Args: | ||
| path: The bucket or directory path to list. | ||
| prefix: Key prefix to filter by, relative to the path. A prefixed | ||
| listing is neither read from nor written to the cache. | ||
| delimiter: Delimiter to group keys by; ``""`` lists recursively. | ||
| next_token: Continuation token to start listing from. A listing | ||
| that starts from a token is neither read from nor written to | ||
| the cache. | ||
| max_keys: Maximum number of keys per ListObjectsV2 request. | ||
| refresh: If True, bypass the cache and list from S3. | ||
|
|
||
| Returns: | ||
| The listed directories and files. | ||
| """ | ||
| bucket, key, version_id = self.parse_path(path) | ||
| use_cache = not prefix and not next_token | ||
| if key: | ||
| prefix = f"{key}/{prefix if prefix else ''}" | ||
|
|
||
| # Create a cache key that includes the delimiter | ||
| cache_key = (path, delimiter) | ||
| if cache_key in self.dircache and not refresh: | ||
| if use_cache and cache_key in self.dircache and not refresh: | ||
| return cast(list[S3Object], self.dircache[cache_key]) | ||
|
|
||
| files: list[S3Object] = [] | ||
|
|
@@ -444,8 +471,11 @@ def _ls_dirs( | |
| next_token = response.get("NextContinuationToken") | ||
| if not next_token: | ||
| break | ||
| if files: | ||
| self.dircache[cache_key] = files | ||
| if use_cache: | ||
| if files: | ||
| self.dircache[cache_key] = files | ||
| else: | ||
| self.dircache.pop(cache_key, None) | ||
| return files | ||
|
|
||
| def ls( | ||
|
|
@@ -696,13 +726,15 @@ def _find( | |
| raise ValueError("Cannot traverse all files in S3.") | ||
| bucket, key, _ = self.parse_path(path) | ||
| prefix = kwargs.pop("prefix", "") | ||
| # Keep refresh in kwargs so that the recursive calls also refresh. | ||
| refresh = kwargs.get("refresh", False) | ||
|
|
||
| # When maxdepth is specified, use a recursive approach with delimiter | ||
| if maxdepth is not None: | ||
| result: list[S3Object] = [] | ||
|
|
||
| # List files and directories at current level | ||
| current_items = self._ls_dirs(path, prefix=prefix, delimiter="/") | ||
| current_items = self._ls_dirs(path, prefix=prefix, delimiter="/", refresh=refresh) | ||
|
|
||
| for item in current_items: | ||
| if item.type == S3ObjectType.S3_OBJECT_TYPE_FILE: | ||
|
|
@@ -724,16 +756,17 @@ def _find( | |
| return result | ||
|
|
||
| # For unlimited depth, use the original approach (get all files at once) | ||
| files = self._ls_dirs(path, prefix=prefix, delimiter="") | ||
| files = self._ls_dirs(path, prefix=prefix, delimiter="", refresh=refresh) | ||
| if not files and key: | ||
| try: | ||
| files = [self.info(path)] | ||
| files = [self.info(path, refresh=refresh)] | ||
| except FileNotFoundError: | ||
| files = [] | ||
|
|
||
| # If withdirs is True, we need to derive directories from file paths | ||
| if withdirs: | ||
| files.extend(self._extract_parent_directories(files, bucket, key)) | ||
| # Build a new list; files may be the cached listing. | ||
| files = files + self._extract_parent_directories(files, bucket, key) | ||
|
|
||
| # Filter directories if withdirs is False (default) | ||
| if withdirs is False or withdirs is None: | ||
|
|
@@ -760,7 +793,11 @@ def find( | |
| maxdepth: Maximum depth to recurse (None for unlimited). | ||
| withdirs: Whether to include directories in results (None = default behavior). | ||
| detail: If True, return dict of {path: S3Object}; if False, return list of paths. | ||
| **kwargs: Additional arguments. | ||
| **kwargs: Additional arguments including: | ||
| prefix: Key prefix, relative to the path, to filter the listed keys | ||
| by. Without maxdepth, if nothing is listed and the path itself is | ||
| an object, that object is returned regardless of the prefix. | ||
| refresh: If True, bypass the cache and list from S3. | ||
|
|
||
| Returns: | ||
| Dictionary mapping paths to S3Objects (if detail=True) or | ||
|
|
@@ -784,7 +821,8 @@ def exists(self, path: str, **kwargs) -> bool: | |
|
|
||
| Args: | ||
| path: S3 path to check (e.g., "s3://bucket" or "s3://bucket/key"). | ||
| **kwargs: Additional arguments (unused). | ||
| **kwargs: Additional arguments including: | ||
| refresh: If True, bypass the cache and query S3. | ||
|
|
||
| Returns: | ||
| True if the path exists, False otherwise. | ||
|
|
@@ -794,29 +832,30 @@ def exists(self, path: str, **kwargs) -> bool: | |
| >>> fs.exists("s3://my-bucket/file.txt") | ||
| >>> fs.exists("s3://my-bucket/") | ||
| """ | ||
| refresh = kwargs.pop("refresh", False) | ||
| path = self._strip_protocol(path) | ||
| if path in ["", "/"]: | ||
| # The root always exists. | ||
| return True | ||
| bucket, key, _ = self.parse_path(path) | ||
| if key: | ||
| try: | ||
| if self._ls_from_cache(path): | ||
| if not refresh and self._ls_from_cache(path): | ||
| return True | ||
| info = self.info(path) | ||
| info = self.info(path, refresh=refresh) | ||
| return bool(info) | ||
| except FileNotFoundError: | ||
| return False | ||
| elif self.dircache.get(bucket, False): | ||
| return True | ||
| else: | ||
| 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) | ||
| return bool(file) | ||
| file = self._head_bucket(bucket, refresh=refresh) | ||
| return bool(file) | ||
|
|
||
| def rm_file(self, path: str, **kwargs) -> None: | ||
| """Delete an S3 object with DeleteObject. | ||
|
|
@@ -1882,6 +1921,9 @@ def invalidate_cache(self, path: str | None = None) -> None: | |
| path = self._strip_protocol(path) | ||
| while path: | ||
| self.dircache.pop(path, None) | ||
| # _ls_dirs caches listings under (path, delimiter). | ||
| for delimiter in ("/", ""): | ||
|
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 (behavior and implementation): FINDINGS (1, repaired) Scope: base 775874c, head 6e60732 (full diff: Covered:
Finding (repaired in 2fba19f): Deferred (pre-existing, out of scope): |
||
| self.dircache.pop((path, delimiter), None) | ||
| path = self._parent(path) | ||
|
|
||
| def _ls_from_cache(self, path: str) -> list[S3Object] | S3Object | None: | ||
|
|
||
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 review (relayed): CLEAN for this diff; 5 pre-existing issues reported.
Reviewer: Codex CLI 0.160.0 (
codex exec -s read-only, model reported as gpt-6-astra), a different model from the author (Claude). Static review: source inspection only, with no tests, builds, edits, network, GitHub or PR-discussion access. The prompt contained no PR number, description, commit messages or prior findings.Scope: base 775874c, head 7d84691 (detached snapshot plus literal diff). The snapshot and the PR worktree were unchanged afterwards.
Covered (reviewer):
_ls_dirs/ls/findcache keys, pagination, prefix, maxdepth, withdirs, refresh and path normalization;info/exists/_ls_from_cache, version ids and invalidation; sync and async writes (delete, copy, upload, touch, bucket ops, metadata/ACL/tags, buffered commit); new tests, docstrings and comments.Result: "no actionable regression introduced by this diff". All new tests would fail against the original implementation; the integration test verifies observable listing results.
Pre-existing issues reported, each verified by the author against base 775874c:
s3.py:715(_find):find(d, withdirs=True)extends the list object returned from the cache in place, so repeating the call duplicates the derived directories. Confirmed by reading the code.s3.py:449(_ls_dirs): an emptyrefresh=Truelisting neither replaces nor evicts the old entry, so the next plainls()returns the deleted objects again. Confirmed. This also qualifies the new docstring wording "A complete listing of the path is cached".s3.py:1790(invalidate_cache):rm_file("b/d/key?versionId=v1")invalidates the version-qualified path and its ancestors but not the cached"b/d/key"HeadObject entry, soexists("b/d/key")can stay True. Confirmed.s3.py:595(info):info(p, version_id="v1")followed byinfo(p, version_id="v2")returns v1's cached metadata, because neither the cache lookup nor_head_objectcompares the version. Confirmed.s3.py:696(_find):maxdepthis off by one compared with fsspec 2026.9.0 (walkrequiresmaxdepth >= 1, andmaxdepth=1means only the direct entries). Heremaxdepth=0lists the direct entries andmaxdepth=1descends one more level.test_find_maxdepthencodes the current behavior. Confirmed.Disposition: pending the maintainer's decision on folding versus separate issues; follow-ups will be recorded in this thread.
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.
Repair 81e2310 (findings 1 and 2) and its self-review
Maintainer decision: fold findings 1 and 2 (listing cache) into this PR. Findings 3 to 5 are filed as #931 (stale HeadObject entry after deleting a version), #932 (
info(version_id=...)ignores the version in the cache) and #933 (maxdepthoff by one compared with fsspec).Repair (7d84691..81e2310):
_find()builds a new list (files + self._extract_parent_directories(...)) instead offiles.extend(...)on the list that_ls_dirs()may have returned from the cache._ls_dirs(): for a cacheable listing (noprefix/next_token), a non-empty result is stored and an empty result pops(path, delimiter). The docstring now says this.test_ls_dirs_empty_refresh_evicts_cached_listingandtest_find_withdirs_does_not_modify_cached_listing(offline). Both fail on 7d84691 and pass on 81e2310.Self-review of the repair:
find()calls never evict the unprefixed listing.ls(file_path)listsfile_path/(empty), so it pops a key that never holds an entry, then falls back to_head_objectas before. Other cached-list consumers (lsreturnslist(files),find(maxdepth=...)buildsresult,findbuilds new lists or dicts) do not mutate the cache._ls_buckets()returns its cached list, butls()copies it.just lintpasses; offline unit tests 6 passed; livetest_s3.py+test_s3_async.py(-k "reflect_changes or invalidate or ls or partial_listing or empty_refresh or find or rm or touch or exists or info or glob") 55 passed on the 81e2310 tree.An independent follow-up on this range is next.
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 (relayed) on 7d84691..81e2310: CLEAN; 2 more pre-existing issues
Reviewer: Codex CLI 0.160.0 (
codex exec -s read-only, model reported as gpt-6-astra, session 01a0ffd0-7964-7343-a710-deb9db9c117b). Static review of the follow-up diff with the full diff as context; the snapshot at 81e2310 and the PR worktree were unchanged afterwards.Covered:
_ls_dirs()cache hits, pagination, refresh, both delimiters and the prefix/token bypass; thels()file fallback, both_find()branches and async delegation; the two new tests (they fail on the old code); the changed docstring. Result: both targeted defects are fixed, no introduced regression, and no remaining caller mutates the cached list.Pre-existing issues reported (both verified by the author):
s3.py:337_head_object: the not-found branch leaves the cached entry, so afterls(p, refresh=True)returns[]for an externally deleted file, the nextls(p)returns it again._head_buckethas the same pattern.s3.py:710_find:refreshis not forwarded, sofind(d, refresh=True)returns the cached listing.Repair 2 (maintainer: fix both in this PR), now on head c982550 after rebasing onto master 9b74767 (#919). The rebase had no conflicts; upstream only adds docstrings to
pyathena/filesystem(549 insertions, 0 deletions)._head_object()/_head_bucket()pop the path/bucket entry when HeadObject/HeadBucket raisesFileNotFoundError._find()readsrefreshfrom kwargs (kept there so the recursivemaxdepthcalls inherit it) and forwards it to both_ls_dirs()calls and to theinfo()fallback.find()docstring documentsprefixandrefresh. The tests gain a_file_object()helper for the new unit tests.test_find_refresh_bypasses_cached_listings(recursive andmaxdepth=1, including a stale subdirectory listing) andtest_refresh_evicts_cached_object_and_bucket_not_found. Both fail on 81e2310 and pass now.Self-review of repair 2:
version_idwhose HEAD is not found, the unversioned entry is evicted, which costs only a later HEAD.PermissionErrorand other errors do not evict.refreshis now honored byglob(), which callsfind()with the caller's kwargs. The async_findforwards kwargs to the sync_find, so it inherits the fix.find()docstring match the code. The PR title and body were updated (WHAT/WHY/TEST, tested commit c982550).info()docstring says a cached listing of the path or its parent decides directory or not-found. With tuple-keyed listings, that applies only to the bucket list under"". This belongs to the docs/docstring audit Check the user guides and docstrings against the implementation #927, outside this diff.just lint(with the pydocstyle rules) passes; offline unit tests 8 passed; livetest_s3.py+test_s3_async.pyfiltered run 57 passed on c982550.An independent follow-up on this range is next.
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 2 (relayed) on 38a2e28..c982550 (rebased series): FINDINGS (1, P3); 3 pre-existing issues
Reviewer: Codex CLI 0.160.0 (
codex exec -s read-only, model reported as gpt-6-astra, session 01a0ffdc-1cea-7c52-8017-c4c6273e0bd7). Static review with the range-diff (775874c..81e2310 vs 9b74767..c982550), the follow-up diff and the full diff. The snapshot at c982550 and the PR worktree were unchanged afterwards.Covered:
_head_object,_head_bucket, cache lookup,info,exists,ls, bothfindpaths, inherited sync/asyncglob, and the async wrappers. Results: A evicts the direct entries onFileNotFoundErrorand keeps them on other errors. B reaches both delimiters, every recursive level and theinfo()fallback. The rebase contains only the expected docstring merge. The new tests fail on the old code and assert observable results.find()prefixdoc claimed unconditional filtering, butfind("s3://bucket/file", prefix="unmatched")returns the object through theinfo()fallback.ls("s3://"), a deleted bucket comes back from the cached bucket listingdircache[""]even afterinfo(bucket, refresh=True); (2)exists(path, refresh=True)ignoresrefresh; (3) a keywordversion_idshares the path cache key. (3) is S3FileSystem.info(version_id=...) returns cached metadata of another version #932.Repair 3 (a95e84f): the docstring now says that if nothing is listed and the path itself is an object, it is returned regardless of the prefix.
_head_bucket()also popsdircache[""]on not found, completing A, asmkdir()/rmdir()already do for buckets.test_refresh_evicts_cached_object_and_bucket_not_foundnow seeds the bucket listing too, and fails without the change.Repair 4 (51a4362, maintainer: fix (2) in this PR):
exists()popsrefresh. Withrefresh=Trueit skips the_ls_from_cache/dircache checks and passesrefreshtoinfo()(keys) and_head_bucket()(buckets). Withoutrefresh, the logic is unchanged. Docstring updated. Newtest_exists_refresh_bypasses_cachefails without the change.Self-review of repairs 3 and 4:
mkdir, etc.) callexists()withoutrefresh, so their path is unchanged. Withrefresh=Trueon a key prefix,info()falls through to theMaxKeys=1listing and returns True. Popping""on a not-found bucket costs one later ListBuckets. The async_existsforwards kwargs to the syncexists.exists()/find()docstrings and the updated PR body (tested commit 51a4362) match the code.just lintpasses; offline unit tests 9 passed (plus the offline mkdir/rmdir unit tests); livetest_s3.py+test_s3_async.pyfiltered run (addsmkdir or bucket) 63 passed on 51a4362.An independent follow-up on a95e84f and 51a4362 is next.
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 3 (relayed) on c982550..51a4362: FINDINGS (2, both introduced by repairs 3/4)
Reviewer: Codex CLI 0.160.0 (
codex exec -s read-only, model reported as gpt-6-astra, session 01a0ffe9-3cbc-79b3-9403-2847a5d46510). Static review with the follow-up diff and full diff 9b74767..51a4362. The snapshot at 51a4362 and the PR worktree were unchanged afterwards.Covered:
find()/_find(),info(), the object/bucket HEAD paths, cache lookup and eviction;exists()with and without refresh; internal callers (mkdir,makedirs,rmdir,touch,pipe_file, append); async_exists(). Refresh forwarding works, the stale deleted-bucket lookup is repaired, and the tests assert observable behavior and fail on the previous code.s3.py:320: poppingdircache[""]on any missing bucket also drops a listing that never contained it. Afterls("s3://"),exists("s3://missing")makes laterexists()of listed buckets issue HeadBucket, which raisesPermissionErrorfor callers withs3:ListAllMyBucketsbut nos3:ListBucket. The same happens throughmkdir("s3://missing/prefix", create_parents=False).s3.py:796: the new prefix fallback note does not hold withmaxdepth; the depth-limited branch never falls back toinfo(path).Repair 5 (b16c58f):
_head_bucket()evicts the bucket listing only if it lists the missing bucket (b.name == bucket; bucketS3Objectnames are the bucket name). Theprefixnote now starts with "Without maxdepth". Newtest_missing_bucket_keeps_bucket_listing_without_itfails on 51a4362.test_refresh_evicts_cached_object_and_bucket_not_found(listing contains the bucket) still passes.Self-review: only
_ls_buckets()writesdircache[""], always asS3Objectbuckets, sob.nameis safe. Behavior without a missing bucket is unchanged. The commit message, docstring and PR body (tested commit b16c58f) match.just lintpasses; offline unit tests 10 passed; live filteredtest_s3.py+test_s3_async.pyrun 64 passed on b16c58f.An independent follow-up on b16c58f is next.
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 4 (relayed) on 51a4362..b16c58f: CLEAN
Reviewer: Codex CLI 0.160.0 (
codex exec -s read-only, model reported as gpt-6-astra, session 01a0fff1-e2f7-7183-bb91-e5968b561c7d). Static review (no edits, builds, tests, network or GitHub) of the follow-up diff, with the full diff 9b74767..b16c58f as context. The snapshot at b16c58f and the PR worktree were unchanged afterwards.Covered:
_head_bucket(),info(),exists(),mkdir(),rmdir()and the async wrappers. Refresh reaches HeadBucket. A missing bucket loses its direct cache entry and any bucket listing that contains it, so later lookups cannot return it from the cache._ls_buckets()is the only writer ofdircache[""]and storeslist[S3Object]; absent or empty entries are safe. Successful bucket creation and deletion still evict the listing. The new test fails on the previous code atassert fs.exists("s3://bucket"). The changed docstring is accurate: theinfo(path)fallback only runs whenmaxdepth is None.Result: no actionable findings, and no further pre-existing issues reported.
This completes the independent review for head b16c58f. Next: confirm the offline checks and mergeability, then mark the PR Ready for AWS CI.