Skip to content
Merged
76 changes: 59 additions & 17 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -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={
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)

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): 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/find cache 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:

  1. 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.
  2. s3.py:449 (_ls_dirs): an empty refresh=True listing neither replaces nor evicts the old entry, so the next plain ls() returns the deleted objects again. Confirmed. This also qualifies the new docstring wording "A complete listing of the path is cached".
  3. 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, so exists("b/d/key") can stay True. Confirmed.
  4. s3.py:595 (info): info(p, version_id="v1") followed by info(p, version_id="v2") returns v1's cached metadata, because neither the cache lookup nor _head_object compares the version. Confirmed.
  5. s3.py:696 (_find): maxdepth is off by one compared with fsspec 2026.9.0 (walk requires maxdepth >= 1, and maxdepth=1 means only the direct entries). Here maxdepth=0 lists the direct entries and maxdepth=1 descends one more level. test_find_maxdepth encodes the current behavior. Confirmed.

Disposition: pending the maintainer's decision on folding versus separate issues; follow-ups will be recorded in this thread.

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.

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 (maxdepth off by one compared with fsspec).

Repair (7d84691..81e2310):

  • _find() builds a new list (files + self._extract_parent_directories(...)) instead of files.extend(...) on the list that _ls_dirs() may have returned from the cache.
  • _ls_dirs(): for a cacheable listing (no prefix/next_token), a non-empty result is stored and an empty result pops (path, delimiter). The docstring now says this.
  • Tests: test_ls_dirs_empty_refresh_evicts_cached_listing and test_find_withdirs_does_not_modify_cached_listing (offline). Both fail on 7d84691 and pass on 81e2310.

Self-review of the repair:

  • Behavior: the empty-result eviction runs only for cacheable listings, so prefixed find() calls never evict the unprefixed listing. ls(file_path) lists file_path/ (empty), so it pops a key that never holds an entry, then falls back to _head_object as before. Other cached-list consumers (ls returns list(files), find(maxdepth=...) builds result, find builds new lists or dicts) do not mutate the cache. _ls_buckets() returns its cached list, but ls() copies it.
  • Claims: the commit message and the docstring ("A complete, non-empty listing ... is cached ..., and an empty one evicts it") match the code. The PR body WHAT/WHY/TEST were updated; TEST now states commit 81e2310.
  • Validation: just lint passes; offline unit tests 6 passed; live test_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.

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 (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; the ls() 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):

  • A, s3.py:337 _head_object: the not-found branch leaves the cached entry, so after ls(p, refresh=True) returns [] for an externally deleted file, the next ls(p) returns it again. _head_bucket has the same pattern.
  • B, s3.py:710 _find: refresh is not forwarded, so find(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 raises FileNotFoundError.
  • _find() reads refresh from kwargs (kept there so the recursive maxdepth calls inherit it) and forwards it to both _ls_dirs() calls and to the info() fallback.
  • find() docstring documents prefix and refresh. The tests gain a _file_object() helper for the new unit tests.
  • Tests: test_find_refresh_bypasses_cached_listings (recursive and maxdepth=1, including a stale subdirectory listing) and test_refresh_evicts_cached_object_and_bucket_not_found. Both fail on 81e2310 and pass now.

Self-review of repair 2:

  • Behavior: the pop runs only after a real HEAD returned not found, so uncached misses are no-ops. With an explicit version_id whose HEAD is not found, the unversioned entry is evicted, which costs only a later HEAD. PermissionError and other errors do not evict. refresh is now honored by glob(), which calls find() with the caller's kwargs. The async _find forwards kwargs to the sync _find, so it inherits the fix.
  • Claims: the commit messages and the find() docstring match the code. The PR title and body were updated (WHAT/WHY/TEST, tested commit c982550).
  • Deferred (pre-existing, from Document every public API in pyathena/ and check docstrings with ruff #919): the new 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.
  • Validation: just lint (with the pydocstyle rules) passes; offline unit tests 8 passed; live test_s3.py + test_s3_async.py filtered run 57 passed on c982550.

An independent follow-up on this range is next.

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 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, both find paths, inherited sync/async glob, and the async wrappers. Results: A evicts the direct entries on FileNotFoundError and keeps them on other errors. B reaches both delimiters, every recursive level and the info() fallback. The rebase contains only the expected docstring merge. The new tests fail on the old code and assert observable results.

  • Finding (P3, introduced): the new find() prefix doc claimed unconditional filtering, but find("s3://bucket/file", prefix="unmatched") returns the object through the info() fallback.
  • Pre-existing: (1) after ls("s3://"), a deleted bucket comes back from the cached bucket listing dircache[""] even after info(bucket, refresh=True); (2) exists(path, refresh=True) ignores refresh; (3) a keyword version_id shares 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 pops dircache[""] on not found, completing A, as mkdir()/rmdir() already do for buckets. test_refresh_evicts_cached_object_and_bucket_not_found now seeds the bucket listing too, and fails without the change.
Repair 4 (51a4362, maintainer: fix (2) in this PR): exists() pops refresh. With refresh=True it skips the _ls_from_cache/dircache checks and passes refresh to info() (keys) and _head_bucket() (buckets). Without refresh, the logic is unchanged. Docstring updated. New test_exists_refresh_bypasses_cache fails without the change.

Self-review of repairs 3 and 4:

  • Behavior: internal callers (mkdir, etc.) call exists() without refresh, so their path is unchanged. With refresh=True on a key prefix, info() falls through to the MaxKeys=1 listing and returns True. Popping "" on a not-found bucket costs one later ListBuckets. The async _exists forwards kwargs to the sync exists.
  • Claims: the commit messages, the exists()/find() docstrings and the updated PR body (tested commit 51a4362) match the code.
  • Validation: just lint passes; offline unit tests 9 passed (plus the offline mkdir/rmdir unit tests); live test_s3.py + test_s3_async.py filtered run (adds mkdir or bucket) 63 passed on 51a4362.

An independent follow-up on a95e84f and 51a4362 is next.

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 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.

  1. P2, s3.py:320: popping dircache[""] on any missing bucket also drops a listing that never contained it. After ls("s3://"), exists("s3://missing") makes later exists() of listed buckets issue HeadBucket, which raises PermissionError for callers with s3:ListAllMyBuckets but no s3:ListBucket. The same happens through mkdir("s3://missing/prefix", create_parents=False).
  2. P3, s3.py:796: the new prefix fallback note does not hold with maxdepth; the depth-limited branch never falls back to info(path).

Repair 5 (b16c58f): _head_bucket() evicts the bucket listing only if it lists the missing bucket (b.name == bucket; bucket S3Object names are the bucket name). The prefix note now starts with "Without maxdepth". New test_missing_bucket_keeps_bucket_listing_without_it fails on 51a4362. test_refresh_evicts_cached_object_and_bucket_not_found (listing contains the bucket) still passes.
Self-review: only _ls_buckets() writes dircache[""], always as S3Object buckets, so b.name is safe. Behavior without a missing bucket is unchanged. The commit message, docstring and PR body (tested commit b16c58f) match. just lint passes; offline unit tests 10 passed; live filtered test_s3.py + test_s3_async.py run 64 passed on b16c58f.

An independent follow-up on b16c58f is next.

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 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 of dircache[""] and stores list[S3Object]; absent or empty entries are safe. Successful bucket creation and deletion still evict the listing. The new test fails on the previous code at assert fs.exists("s3://bucket"). The changed docstring is accurate: the info(path) fallback only runs when maxdepth 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.

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] = []
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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:
Expand All @@ -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:
Expand All @@ -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
Expand All @@ -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.
Expand All @@ -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.
Expand Down Expand Up @@ -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 ("/", ""):

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 (behavior and implementation): FINDINGS (1, repaired)

Scope: base 775874c, head 6e60732 (full diff: pyathena/filesystem/s3.py, tests/pyathena/filesystem/test_s3.py).

Covered:

  • Cache key normalization: ls() passes _strip_protocol(path).rstrip("/"), and _find() passes _strip_protocol(path) (fsspec strips the trailing slash), so invalidate_cache() builds the same path string for both "/" (ls, maxdepth find) and "" (recursive find) entries. Recursive _find() subcalls use s3://bucket/{item.key} with no trailing slash.
  • Write paths: rm/_rm, touch, cp_file, pipe_file, put_file, setxattr, S3File.commit, bucket mkdir/rmdir, and the async _rm/_cp_file all call invalidate_cache(). put_tags/chmod do not change listing fields.
  • Async: S3FileSystemAsync shares the sync dircache and delegates _ls, _find (including the prefix kwarg) and invalidate_cache to the sync instance.
  • use_cache is computed before prefix is merged with key; prefixed and next_token listings neither read nor write the cache. The unprefixed recursive subcalls of find(maxdepth=...) still cache. No caller passes next_token/max_keys today.
  • Tests: the 3 new unit tests and the integration test fail without the s3.py change (verified locally).

Finding (repaired in 2fba19f): test_invalidate_cache_drops_listings_of_path_and_parents asserted that a descendant listing ("bucket/a/b/c.txt/d", "") is kept. fsspec documents invalidate_cache(path) as "listings at or under given path". This PR keeps the existing path-and-parents walk (pre-v3.15.0 and s3fs behavior), so the test should not pin descendant retention. The assertion was removed. The offline unit tests and just lint were rerun and pass.

Deferred (pre-existing, out of scope): invalidate_cache(path) does not drop listings under path. Doing so would require scanning every dircache key on each write.

self.dircache.pop((path, delimiter), None)
path = self._parent(path)

def _ls_from_cache(self, path: str) -> list[S3Object] | S3Object | None:
Expand Down
164 changes: 164 additions & 0 deletions tests/pyathena/filesystem/test_s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,16 @@ def _make_fs():
fs.version_aware = False
return fs

@staticmethod
def _file_object(key):
# Build a listed file entry in the bucket named "bucket".
return S3Object(
init={"Key": key},
type=S3ObjectType.S3_OBJECT_TYPE_FILE,
bucket="bucket",
key=key,
)

def test_get_client_compatible_with_s3fs(self):
# Only constructs a boto3 client; no AWS access.
fs = S3FileSystem(
Expand Down Expand Up @@ -192,6 +202,137 @@ def test_ls_from_cache_with_cached_object(self):
# every cache value is a listing and raises TypeError here).
assert fs._ls_from_cache("bucket/key/child") is None

def test_invalidate_cache_drops_listings_of_path_and_parents(self):
fs = self._make_fs()
invalidated = [
"bucket/a/b/c.txt",
("bucket/a/b", "/"),
("bucket/a/b", ""),
("bucket/a", "/"),
("bucket/a", ""),
("bucket", "/"),
("bucket", ""),
]
kept = ["", ("bucket/a/x", "/")]
for cache_key in invalidated + kept:
fs.dircache[cache_key] = []

fs.invalidate_cache("s3://bucket/a/b/c.txt")
assert list(fs.dircache) == kept

@pytest.mark.parametrize(
("prefix", "next_token"),
[
("test_", None),
("", "token"),
],
)
def test_ls_dirs_partial_listing_bypasses_cache(self, prefix, next_token):
fs = self._make_fs()
cached = self._file_object("dir/cached")
fs.dircache[("bucket/dir", "")] = [cached]
fs._call.return_value = {"Contents": [{"Key": "dir/test_1"}]}

files = fs._ls_dirs("bucket/dir", prefix=prefix, delimiter="", next_token=next_token)
assert [f.name for f in files] == ["bucket/dir/test_1"]
assert fs.dircache[("bucket/dir", "")] == [cached]

# A complete listing of the path is still served from the cache.
fs._call.reset_mock()
assert fs._ls_dirs("bucket/dir", delimiter="") == [cached]
fs._call.assert_not_called()

def test_ls_dirs_empty_refresh_evicts_cached_listing(self):
fs = self._make_fs()
fs.dircache[("bucket/dir", "/")] = [self._file_object("dir/deleted")]
fs._call.return_value = {}

assert fs._ls_dirs("bucket/dir", refresh=True) == []
# The next listing must not return the deleted object from the cache.
fs._call.reset_mock()
assert fs._ls_dirs("bucket/dir") == []
fs._call.assert_called_once()

def test_find_withdirs_does_not_modify_cached_listing(self):
fs = self._make_fs()
fs.dircache[("bucket/dir", "")] = [self._file_object("dir/sub/file")]

expected = ["bucket/dir/sub", "bucket/dir/sub/file"]
assert sorted(fs.find("s3://bucket/dir", withdirs=True)) == expected
assert sorted(fs.find("s3://bucket/dir", withdirs=True)) == expected
assert fs.find("s3://bucket/dir") == ["bucket/dir/sub/file"]
fs._call.assert_not_called()

def test_find_refresh_bypasses_cached_listings(self):
fs = self._make_fs()
fs.dircache[("bucket/dir", "")] = [self._file_object("dir/old")]
fs.dircache[("bucket/dir", "/")] = [self._file_object("dir/old")]
fs.dircache[("bucket/dir/sub", "/")] = [self._file_object("dir/sub/old")]
responses = {
("dir/", ""): {"Contents": [{"Key": "dir/sub/new"}]},
("dir/", "/"): {"CommonPrefixes": [{"Prefix": "dir/sub/"}]},
("dir/sub/", "/"): {"Contents": [{"Key": "dir/sub/new"}]},
}
fs._call.side_effect = lambda method, **kwargs: responses[
(kwargs["Prefix"], kwargs["Delimiter"])
]

assert fs.find("s3://bucket/dir", refresh=True) == ["bucket/dir/sub/new"]
# The subdirectory listings of maxdepth are refreshed as well.
assert fs.find("s3://bucket/dir", maxdepth=1, refresh=True) == ["bucket/dir/sub/new"]

def test_refresh_evicts_cached_object_and_bucket_not_found(self):
fs = self._make_fs()
fs.dircache["bucket/key"] = self._file_object("key")
fs.dircache["bucket"] = fs._directory_object("bucket", None)
fs.dircache[""] = [fs._directory_object("bucket", None)]

def call(method, **kwargs):
if method in (fs._client.head_object, fs._client.head_bucket):
raise FileNotFoundError
return {}

fs._call.side_effect = call

assert fs.ls("s3://bucket/key", refresh=True) == []
# The next lookups must not return the deleted object and bucket from the cache.
assert fs.ls("s3://bucket/key") == []
assert not fs.exists("s3://bucket/key")
with pytest.raises(FileNotFoundError):
fs.info("s3://bucket", refresh=True)
assert not fs.exists("s3://bucket")

def test_exists_refresh_bypasses_cache(self):
fs = self._make_fs()
fs.dircache["bucket/key"] = self._file_object("key")
fs.dircache["bucket"] = fs._directory_object("bucket", None)
fs.dircache[""] = [fs._directory_object("bucket", None)]

def call(method, **kwargs):
if method in (fs._client.head_object, fs._client.head_bucket):
raise FileNotFoundError
return {}

fs._call.side_effect = call

assert fs.exists("s3://bucket/key")
assert fs.exists("s3://bucket")
fs._call.assert_not_called()

assert not fs.exists("s3://bucket/key", refresh=True)
assert not fs.exists("s3://bucket", refresh=True)

def test_missing_bucket_keeps_bucket_listing_without_it(self):
fs = self._make_fs()
fs.dircache[""] = [fs._directory_object("bucket", None)]
fs._call.side_effect = FileNotFoundError

assert not fs.exists("s3://missing")
# Other buckets are still answered from the cached bucket listing.
fs._call.reset_mock()
assert fs.exists("s3://bucket")
fs._call.assert_not_called()

def test_mkdir_creates_bucket(self):
fs = self._make_fs()
fs.allow_bucket_creation = True
Expand Down Expand Up @@ -719,6 +860,29 @@ def test_ls_dirs(self, fs):
assert test_1_detail[0].name == fs._strip_protocol(f"{dir_}/prefix/test_1")
assert test_1_detail[0].size == 1

def test_ls_and_find_reflect_changes_through_the_filesystem(self, fs):
dir_ = (
f"s3://{ENV.s3_staging_bucket}/{ENV.s3_staging_key}{ENV.schema}/"
f"filesystem/test_ls_and_find_reflect_changes/{uuid.uuid4()}"
)
path = fs._strip_protocol(dir_)
fs.touch(f"{dir_}/a.txt")
fs.touch(f"{dir_}/b.txt")
assert sorted(fs.ls(dir_)) == [f"{path}/a.txt", f"{path}/b.txt"]
assert sorted(fs.find(dir_)) == [f"{path}/a.txt", f"{path}/b.txt"]

fs.rm(f"{dir_}/a.txt")
assert fs.ls(dir_) == [f"{path}/b.txt"]
assert fs.find(dir_) == [f"{path}/b.txt"]

fs.touch(f"{dir_}/c.txt")
assert sorted(fs.ls(dir_)) == [f"{path}/b.txt", f"{path}/c.txt"]
assert sorted(fs.find(dir_)) == [f"{path}/b.txt", f"{path}/c.txt"]
# A prefixed find must not be served from the unprefixed listing.
assert fs.find(dir_, prefix="c") == [f"{path}/c.txt"]

fs.rm(dir_, recursive=True)

def test_info_bucket(self, fs):
dir_ = f"s3://{ENV.s3_staging_bucket}"
bucket, key, version_id = fs.parse_path(dir_)
Expand Down
Loading