-
Notifications
You must be signed in to change notification settings - Fork 116
Align S3FileSystem.find() with fsspec for the path itself and prefix levels #993
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
688b087
579c120
535bd17
3c882ce
089e86e
0d6d680
60d7164
8e4be2a
e37e046
bf24179
68cec03
ace8248
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 |
|---|---|---|
|
|
@@ -10,6 +10,7 @@ | |
| import os.path | ||
| import re | ||
| import time | ||
| from collections import Counter | ||
| from collections.abc import Callable, Iterator, Mapping | ||
| from concurrent.futures import Future, as_completed, wait | ||
| from copy import deepcopy | ||
|
|
@@ -899,54 +900,66 @@ 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) | ||
| refresh = kwargs.pop("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="/", refresh=refresh) | ||
|
|
||
| for item in current_items: | ||
| if item.type == S3ObjectType.S3_OBJECT_TYPE_FILE: | ||
| # Add files | ||
| result.append(item) | ||
| elif item.type == S3ObjectType.S3_OBJECT_TYPE_DIRECTORY: | ||
| # Add directory if withdirs is True | ||
| if withdirs: | ||
| result.append(item) | ||
|
|
||
| # Recursively explore subdirectory if depth allows | ||
| if maxdepth > 1: | ||
| sub_path = f"s3://{bucket}/{item.key}" | ||
| sub_results = self._find( | ||
| sub_path, maxdepth=maxdepth - 1, withdirs=withdirs, **kwargs | ||
| ) | ||
| result.extend(sub_results) | ||
|
|
||
| return result | ||
|
|
||
| # For unlimited depth, use the original approach (get all files at once) | ||
| files = self._ls_dirs(path, prefix=prefix, delimiter="", refresh=refresh) | ||
| if not files and key: | ||
| # The entries listed with the prefix lie as many levels further | ||
| # below the path as the prefix has slashes. | ||
| files = self._find_levels( | ||
| path, maxdepth - prefix.count("/"), prefix=prefix, refresh=refresh | ||
| ) | ||
| else: | ||
| files = self._ls_dirs(path, prefix=prefix, delimiter="", refresh=refresh) | ||
| # S3 doesn't return directory entries without a delimiter, so the | ||
| # directories are derived from the listed keys, below the last | ||
| # slash of the prefix, as with maxdepth. | ||
| if withdirs: | ||
| base_key = "/".join(k for k in (key, prefix.rpartition("/")[0]) if k) | ||
|
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 (behavior and implementation; base Simplification (applied): instead of a new |
||
| # Build a new list; files may be the cached listing. | ||
| files = files + self._extract_parent_directories(files, bucket, base_key) | ||
|
|
||
| if files: | ||
| # Something is listed below the path, so the path is a directory, | ||
| # which fsspec includes with the directories. A bucket is not | ||
| # included, since cp_file() cannot copy it in a recursive copy. | ||
| if withdirs and key: | ||
| files = [self._directory_object(bucket, key), *files] | ||
|
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 (behavior and implementation; base Deferred, with reasons:
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): OpenAI Codex CLI 0.160.0, model Reviewer output, verbatim:
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 of the independent-review findings (head
Repair self-review: behavior: the aio path delegates to the same
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 review (relayed): OpenAI Codex CLI 0.160.0, model Reviewer output, verbatim:
Author verification. These are offline probes with the mocked client, run on master, d016e99, and 39e41ad.
|
||
| elif key: | ||
| # As in fsspec, the path itself is returned if it is an object, | ||
| # or with withdirs if it is a 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 two (claims, callers, and operational effects; base Correction (d016e99): the comment said the fallback returns the path only if it is an object. AWS operator / adversarial: retries are unchanged ( Evidence: at a54580b, Result: FINDINGS (wording only), repaired. No deferrals beyond those recorded in round one. |
||
| try: | ||
| files = [self.info(path, refresh=refresh)] | ||
| except FileNotFoundError: | ||
| files = [] | ||
|
|
||
| # If withdirs is True, we need to derive directories from file paths | ||
| if withdirs: | ||
| # 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: | ||
| if not withdirs: | ||
| files = [f for f in files if f.type != S3ObjectType.S3_OBJECT_TYPE_DIRECTORY] | ||
|
|
||
| return files | ||
|
|
||
| def _find_levels( | ||
| self, path: str, maxdepth: int, prefix: str = "", refresh: bool = False | ||
| ) -> list[S3Object]: | ||
| """List the entries below a path level by level with ``Delimiter="/"``. | ||
|
|
||
| Args: | ||
| path: S3 path to search under. | ||
| maxdepth: Number of levels to list. | ||
| prefix: Key prefix, relative to the path, to filter the first | ||
| level by. | ||
| refresh: If True, bypass the cache and list from S3. | ||
|
|
||
| Returns: | ||
| The objects and directories found. | ||
| """ | ||
| if maxdepth < 1: | ||
| return [] | ||
| result: list[S3Object] = [] | ||
| for item in self._ls_dirs(path, prefix=prefix, delimiter="/", refresh=refresh): | ||
| result.append(item) | ||
| if item.type == S3ObjectType.S3_OBJECT_TYPE_DIRECTORY: | ||
| result.extend(self._find_levels(item.name, maxdepth - 1, refresh=refresh)) | ||
|
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 (behavior and implementation; base Scope: the full diff, covering FINDINGS (repaired):
|
||
| return result | ||
|
|
||
| def find( | ||
| self, | ||
| path: str, | ||
|
|
@@ -959,7 +972,10 @@ def find( | |
|
|
||
| Recursively searches for files under the specified path, with optional | ||
| depth limiting and directory inclusion. Uses efficient S3 list operations | ||
| with delimiter handling for performance. | ||
| with delimiter handling for performance. As in fsspec, the result | ||
| includes the path itself if withdirs is True and objects exist below | ||
| it, unless it is a bucket, or if it is an object and nothing is listed | ||
| below it. | ||
|
|
||
| Args: | ||
| path: S3 path to search under (e.g., "s3://bucket/prefix"). | ||
|
|
@@ -970,8 +986,9 @@ def find( | |
| detail: If True, return dict of {path: S3Object}; if False, return list of paths. | ||
| **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. | ||
| by. Each slash in the prefix counts as one level of maxdepth. | ||
| With withdirs, the directories above the prefix, such as | ||
| ``sub`` for ``sub/deep/``, are not included. | ||
|
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 (claims, callers, and operational effects; base Claims checked:
Corrections (d016e99): this docstring said that with |
||
| refresh: If True, bypass the cache and list from S3. | ||
|
|
||
| Returns: | ||
|
|
@@ -1458,8 +1475,9 @@ def mv(self, path1, path2, recursive=False, maxdepth=None, **kwargs) -> None: | |
| ``AbstractFileSystem.mv()`` instead removes ``path1`` by expanding it | ||
| again, which also deletes copies placed where ``path1`` matches them | ||
| and files that ``maxdepth`` kept from being copied. A file whose | ||
| destination is the file itself is left in place, and directories, | ||
| which S3 does not store as objects, are not copied. | ||
| destination is the file itself, or the ``null`` version of a file | ||
| moved to the file, is left in place, and directories, which S3 does | ||
| not store as objects, are not copied. | ||
|
|
||
| Args: | ||
| path1: Source S3 path, glob pattern, or list of paths. | ||
|
|
@@ -1471,8 +1489,9 @@ def mv(self, path1, path2, recursive=False, maxdepth=None, **kwargs) -> None: | |
|
|
||
| Raises: | ||
| ValueError: If two sources have the same destination, or a | ||
| destination is another source, which is checked before | ||
| anything is copied. | ||
| destination is another source, including one left in place, | ||
| which is checked before anything is copied. A directory with | ||
| no object at its key, which is not copied, does not conflict. | ||
| """ | ||
| if path1 == path2: | ||
| return | ||
|
|
@@ -1497,11 +1516,14 @@ def _move_paths( | |
|
|
||
| Returns: | ||
| The source and destination paths, except the sources whose | ||
| destination is the source itself. | ||
| destination is the source itself or, for a ``null`` version, the | ||
| key of the source. | ||
|
|
||
| Raises: | ||
| ValueError: If two sources have the same destination, or a | ||
| destination is another source. | ||
| destination is another source, including one left in place, | ||
| except for a directory with no object at its key, which is not | ||
| copied. | ||
| """ | ||
| if isinstance(path1, list) and isinstance(path2, list): | ||
| paths1, paths2 = path1, path2 | ||
|
|
@@ -1524,19 +1546,63 @@ def _move_paths( | |
| ) | ||
| ) | ||
| paths2 = other_paths(paths1, path2, exists=exists, flatten=not source_is_str) | ||
| # The paths are copied as given, and compared without the protocol. | ||
| pairs = [ | ||
| (p1, p2) | ||
| # The paths are copied as given, and compared by what they name. | ||
| named = [ | ||
| (p1, p2, self._move_target(p1), self._move_target(p2)) | ||
| for p1, p2 in zip(paths1, paths2, strict=False) | ||
| if self._strip_protocol(p1) != self._strip_protocol(p2) | ||
| ] | ||
| destinations = [self._strip_protocol(p2) for _, p2 in pairs] | ||
| if len(set(destinations)) != len(destinations): | ||
| raise ValueError("Cannot move several paths to the same destination.") | ||
| if {self._strip_protocol(p1) for p1, _ in pairs}.intersection(destinations): | ||
| raise ValueError("Cannot move a path onto another path that is moved.") | ||
| pairs = [(p1, p2) for p1, p2, source, dest in named if source != dest] | ||
| moved = [(p1, source, dest) for p1, _, source, dest in named if source != dest] | ||
| # The sources left in place count too; a copy onto one overwrites it. | ||
| sources = {source for _, _, source, _ in named} | ||
| counts = Counter(dest for _, _, dest in moved) | ||
| # A source with another source below it may be a directory. | ||
| directories: set[str] = set() | ||
| for source in sources: | ||
| parent = source.rpartition("/")[0] | ||
| while parent and parent not in directories: | ||
| directories.add(parent) | ||
| parent = parent.rpartition("/")[0] | ||
| # A directory without an object at its key is not copied, so it | ||
| # writes no destination and is left out of the checks. A path with a | ||
| # version always names an object. | ||
| writers = [ | ||
| (source, dest) | ||
| for p1, source, dest in moved | ||
| if not ( | ||
| (counts[dest] > 1 or dest in sources) | ||
| and source in directories | ||
| and not self.parse_path(p1)[2] | ||
|
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. mv() data-safety fixes folded in at the maintainer's request (the pre-existing findings 2 and 3 from the earlier follow-up, plus Self-review.
Independent follow-up on 1444eda (relayed): Codex
Disposition.
Independent follow-up on a7306a1 (relayed): same setup, session
Validation at a7306a1: |
||
| and self._head_object(source) is None | ||
| ) | ||
| ] | ||
| counts = Counter(dest for _, dest in writers) | ||
| for _, dest in writers: | ||
| if counts[dest] > 1: | ||
| raise ValueError("Cannot move several paths to the same destination.") | ||
| if dest in sources: | ||
| raise ValueError("Cannot move a path onto another path that is moved.") | ||
| return pairs | ||
|
|
||
| def _move_target(self, path: str) -> str: | ||
| """Return what a path of a move names, for comparing the paths. | ||
|
|
||
| A write to a key replaces its ``null`` version, which the objects of a | ||
| bucket without versioning have, so that version names the key itself. | ||
|
|
||
| Args: | ||
| path: S3 path, possibly with a version ID. | ||
|
|
||
| Returns: | ||
| The path in ``bucket/key`` form, with the version ID unless it is | ||
| ``null``. | ||
| """ | ||
| bucket, key, version_id = self.parse_path(path) | ||
| target = f"{bucket}/{key}" if key else bucket | ||
| if version_id and version_id != "null": | ||
| return f"{target}?versionId={version_id}" | ||
| return target | ||
|
|
||
| def cp_file( | ||
| self, path1: str, path2: str, recursive=False, maxdepth=None, on_error=None, **kwargs | ||
| ): | ||
|
|
||
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.
Rebase onto a73c6db (conflict resolution):
fe21250a..a7306a10→a73c6db4..ace8248c44e166e8dd22a07431863a0e37edc497.s3.py. Send the lookup parameters of a file with its lookups #1024 addedimport timenext to this PR'sfrom collections import Counter, and both are kept.git range-diffshows all 12 commits unchanged except for that context line in e37e046._head_object()an optionallookup_kwargsand caches versions under one?versionId=spelling. The call in_move_paths()keeps its meaning._move_paths()sends HeadObject without lookup parameters, as_copy_file()'sinfo(path1)andexpand_path()on master do, so this adds no new divergence.find()ormv().just lintpassed.uv run --env-file .env pytest -n 4 tests/pyathena/filesystem/against real S3: 581 passed.