Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
178 changes: 122 additions & 56 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import os.path
import re
import time
from collections import Counter

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.

Rebase onto a73c6db (conflict resolution): fe21250a..a7306a10 → a73c6db4..ace8248c44e166e8dd22a07431863a0e37edc497.

from collections.abc import Callable, Iterator, Mapping
from concurrent.futures import Future, as_completed, wait
from copy import deepcopy
Expand Down Expand Up @@ -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)

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 one (behavior and implementation; base aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe, initial head 2371d9686d0be4391a25c4b73a39d97621393e8e, repairs through a54580b2eda18fc6c16d3f0f8cd7df665489ef04)

Simplification (applied): instead of a new prefix filter inside _extract_parent_directories(), the unlimited branch passes the base key rebased below the last slash of the prefix. This gives the same set (the altitude pass compared 3,000 random key/prefix cases) and leaves the helper unchanged from master. Also applied: _find_levels() handles maxdepth < 1 itself, so the caller's ternary goes; the hand-written listing table in test_find_maxdepth_counts_levels_like_fsspec is replaced by the shared _serve_keys().

# 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]

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 one (behavior and implementation; base aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe, initial head 2371d9686d0be4391a25c4b73a39d97621393e8e, repairs through a54580b2eda18fc6c16d3f0f8cd7df665489ef04)

Deferred, with reasons:

  • A key that is both an object dir and a prefix dir/: the root is added as a directory, while info() returns the file. Subdirectories in both branches were already treated this way, and checking would cost a HeadObject on every find().
  • A dir//x key gives duplicate dir entries. This is the pre-existing double-slash normalization in _ls_dirs(); prefix boundaries are tracked in clear_multipart_uploads() and object_version_info() match sibling keys by string prefix #981.
  • The root entry of a bucket path is a synthetic directory rather than the HeadBucket info, and an empty bucket gets no root. Matching fsspec exactly would cost a HeadBucket request. The docstring (s3.py:828) now states the condition: objects exist below the path.
  • A prefix with a leading slash or empty segments is not normalized. That is malformed input.
  • Sibling directories are listed one at a time. This is pre-existing design.
  • The list is copied twice with withdirs in the unlimited branch. The cost is negligible next to the listing requests.

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): OpenAI Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max, --sandbox read-only, session 01a100e3-fa7e-7202-a731-1fa30509bda3. This was a static review of a detached snapshot at d016e99ab96f0a014aa5a0e7457c99fdaf812a86, diffed against the merge-base aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe. The prompt carried no PR number, description, commit messages, or prior findings, and asked for no edits, builds, tests, network, or GitHub access. Afterwards, the snapshot and the PR worktree were both unchanged.

Reviewer output, verbatim:

Surfaces covered:

  • Sync/async find, depth and prefix handling, directory inclusion, fallback, and detailed results.
  • fsspec glob, expand_path, rm, copy, and du; bucket, version-aware, and trailing-slash paths.
  • Listing caches, refresh/invalidation, request costs, and regression-test assertions.

FINDINGS

  1. P2 — Introduced regression: recursive bucket-glob copying fails.
    pyathena/filesystem/s3.py:779
    Given an existing src bucket containing only a.txt, fs.copy("s3://src/**", "s3://dst/out/", recursive=True) now expands to include src itself. fsspec retains directories for unlimited recursive copying, so cp_file receives the bucket and raises ValueError("Cannot copy buckets.") at line 1222. The default error handler only ignores FileNotFoundError. At the base revision, this glob expanded only to the object and copying succeeded. The async copy path has the same rejection. Preserve the required root inclusion while making copying tolerate these directory entries.

  2. P2 — Pre-existing omission, still violates the intended behavior: empty bucket roots disappear.
    pyathena/filesystem/s3.py:780
    For an existing empty bucket, fs.find("s3://bucket", withdirs=True) returns [], as does the maxdepth=1 variant, instead of ["bucket"]. Because key is None, neither branch reaches info(path). Consequently, glob("s3://bucket/**") also omits the existing directory. This behavior predates the diff but remains an uncovered part of the requested fallback/root-inclusion fix.

The added assertions check observable paths and request counts and would catch the three original defect classes. They do not cover either scenario above. Integration additions reuse existing fixtures and add modest request costs.

Static review only; no files changed, tests/builds run, or network accessed.

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 of the independent-review findings (head 39e41ad54…; the maintainer chose the approach)

  1. Recursive bucket-glob copy (introduced regression): verified and repaired in 39e41ad. fsspec 2026.9.0's copy() filters directories only when it is not recursive or has a maxdepth (spec.py:1212). cp_file() raises ValueError("Cannot copy buckets.") for a path without a key (s3.py:1222), and that error is not ignored the way FileNotFoundError is. _find() now adds the root only for key paths (if withdirs and key), so a bucket path is left out, as on master. The find() docstring and the PR body list this as a difference from fsspec. test_find_withdirs_omits_bucket checks find() with and without maxdepth, and expand_path("s3://bucket/**", recursive=True), which copy() uses; it fails before the repair. A non-bucket root goes through cp_file() → CopyObject NoSuchKey → FileNotFoundError, which recursive copy() ignores, the same as the subdirectory entries already did (tracked in Recursive get(), copy() and mv() mishandle directory entries #974).
  2. Empty bucket root (pre-existing): not changed. The bucket root is now excluded in every case, so an empty bucket follows the same rule, and the PR body records it as a difference from fsspec.

Repair self-review: behavior: the aio path delegates to the same _find(), and rm() skips keyless entries, so it is unaffected. Claims: the commit message, docstring, and PR body were checked against spec.py:1212 and s3.py:1222. Validation at 39e41ad: just lint passed; 30 offline find/glob tests passed; the live targeted run (-k "find or glob or rm or expand or copy or cp") gave 66 passed. The full filesystem suite result is recorded in the PR body.

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 review (relayed): OpenAI Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max, --sandbox read-only, session 01a100f1-b978-7873-8621-c2e783f89bfd. This was a static review of d016e99a..39e41ad505b058e17ab6ab3dd9ce8f3d00d5f4d4 (same base aa0fc914) on a detached snapshot, which remained clean, as did the PR worktree.

Reviewer output, verbatim:

Reviewed d016e99a..39e41ad5 at the requested HEAD. Covered surfaces:

  • _find/find: bucket versus key paths, bounded/unbounded depth, directory inclusion, and list/detail results.
  • fsspec glob → expand_path → copy/cp_file/get, including destination mapping and exception handling.
  • rm: sync and async deletion already filter keyless entries.
  • du: sync explicitly retains the zero-size bucket entry; async follows _find and now omits it. Byte totals remain unchanged.
  • Aio delegation through s3_async.py:335, inherited async callers, and generated sync wrappers.
  • The new test covers bucket omission and wildcard expansion, but does not check transfer destinations.

FINDINGS

P2 — Removing the bucket changes recursive transfer destinations. s3.py:779

With only dir/file in the bucket, expansion of s3://bucket/** changes from:

bucket, bucket/dir, bucket/dir/file

to:

bucket/dir, bucket/dir/file

fsspec’s other_paths at utils.py:414 consequently strips bucket/dir instead of bucket. Thus await aio._get("s3://bucket/**", "/tmp/out", recursive=True) now downloads to /tmp/out/file, whereas it previously downloaded to /tmp/out/dir/file. Sync get has the same regression when its required local directories already exist. Copy uses the same mapping. Bucket-relative destination mapping needs preserving while excluding bucket copy operations.

Remaining failures, distinguished from that regression:

  • Original failure only partly resolved: with the same tree, copy("s3://bucket/**", "s3://destination/", recursive=True) maps bucket/dir to the bare destination bucket. s3.py:1223 still raises ValueError("Cannot copy buckets."). This operation already failed before the repair; the source bucket is removed, but a destination bucket still reaches cp_file.
  • Pre-existing literal-bucket case: copy("s3://bucket", dst, recursive=True) still includes the bucket because spec.py:1285 adds the existing literal input back.

Static source review only; no edits, builds, tests, GitHub, or network access. Worktree remained clean.

Author verification. These are offline probes with the mocked client, run on master, d016e99, and 39e41ad.

  • Bucket-glob destination mapping: the result on 39e41ad is identical to master. For s3://bucket/** with only dir/file, master and 39e41ad both expand to [bucket/dir, bucket/dir/file], so the destination is /tmp/out/file; d016e99 gave /tmp/out/dir/file. This is a remaining difference from fsspec, not a regression against master. For a non-bucket root the PR improves the mapping. With only d/sub/nested, master mapped get("s3://bucket/d/**", out) to out/nested and dropped sub; 39e41ad maps it to out/sub/nested, as fsspec does.
  • Copy to a bare destination bucket and literal bucket copy: both fail the same way on master. They are pre-existing.
  • While tracing, the author found a further regression against master. With a flat directory (d/a, d/b), mv("s3://bucket/d/**", "s3://bucket/out/", recursive=True) succeeds on master but raises FileNotFoundError('d') on 39e41ad, without copying anything. The root d now reaches cp_file(), CopyObject fails with NoSuchKey, and mv() uses on_error="raise". Recursive copy() still succeeds but sends one more failing CopyObject (2→3). A literal mv("s3://bucket/d", recursive=True) already fails this way on master, which is Recursive get(), copy() and mv() mishandle directory entries #974. Open PR Follow fsspec's get_file() contract for directories and file objects #990 makes cp_file()/_cp_file() skip directory entries, which resolves this; how to proceed is with the maintainer.

elif key:
# As in fsspec, the path itself is returned if it is an object,
# or with withdirs if it is a directory.

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 two (claims, callers, and operational effects; base aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe, head a54580b2a→d016e99ab96f0a014aa5a0e7457c99fdaf812a86)

Correction (d016e99): the comment said the fallback returns the path only if it is an object. info() also returns a directory, which is kept with withdirs, for example when a prefix filters out everything below an existing directory. The comment now covers both.

AWS operator / adversarial: retries are unchanged (_call → retry_api_call). The fallback goes through info(), which reads the existing listing cache first, so a cached parent listing can answer it with no request, and a stale cache entry is served the same way it was on master's unlimited branch. The subdirectory listings of maxdepth are now cached under stripped paths, so invalidate_cache() drops them (round one's repair).

Evidence: at a54580b, pytest -n 4 tests/pyathena/filesystem/ against real S3 gave 299 passed, and the targeted run gave 65 passed. d016e99 changes only a docstring and a comment (just lint passed). The other suites are left to CI.

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

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 one (behavior and implementation; base aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe, initial head 2371d9686d0be4391a25c4b73a39d97621393e8e, repairs through a54580b2eda18fc6c16d3f0f8cd7df665489ef04)

Scope: the full diff, covering _find(), _find_levels(), _extract_parent_directories(), the find() docstring, and the offline and live tests. The callers through fsspec are glob(), expand_path() (used by rm(), copy() and get()), du(), and AioS3FileSystem._find()/_glob(), which delegate to the sync _find(). Round one included a /code-review pass and a /simplify pass (reuse, simplification, efficiency, altitude).

FINDINGS (repaired):

  1. _find_levels() recursed with f"s3://{bucket}/{item.key}". _ls_dirs() caches under the raw path, so the subdirectory listings went under ("s3://bucket/dir/sub", "/"), a key that invalidate_cache() never drops. The old recursion through _find() stripped the protocol. Repro: after find(maxdepth=2) and invalidate_cache("s3://bucket/dir/sub/nested"), the entry stayed. Repaired in a54580b: the recursion uses item.name. test_find_maxdepth_listings_follow_invalidation fails on 533d401 and passes now.
  2. Without withdirs, _find_levels() dropped directories. A path with only subdirectories below it therefore looked empty and went to the info() fallback, 1→3 requests. Repaired in 533d401: _find_levels() returns all entries and _find() filters them at the end. The prefix="s" case of test_find_directory_without_extra_requests covers it.

return result

def find(
self,
path: str,
Expand All @@ -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").
Expand All @@ -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.

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 two (claims, callers, and operational effects; base aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe, head a54580b2a→d016e99ab96f0a014aa5a0e7457c99fdaf812a86)

Claims checked:

  • Every row of the PR's behavior table, and every request count, was re-run on master and on a54580b with the mocked client. All match. The request counts are fake-provider counts, not AWS measurements, and the PR says so.
  • "expand_path() checks exists(path) only when find() did not return the path": confirmed in fsspec 2026.9.0, spec.py expand_path (if p not in out and (recursive is False or self.exists(p))). copy() (spec.py:1209) and get() (spec.py:1040) go through expand_path(), and S3FileSystem.rm() calls it directly. The expanded set is unchanged (the root was added through exists() before); only the requests drop from 3 to 1.
  • "The fallback raises PermissionError where HeadObject is denied, as the unlimited branch already did": info() → _head_object() catches only FileNotFoundError, and _call() maps 403 to PermissionError. That is the same code path the unlimited branch had on master.
  • Existing callers: the signatures of _find(), find(), and _extract_parent_directories() match master. AioS3FileSystem._find() still passes kwargs through. The result order now starts with the root, but nothing depends on the order, and fsspec sorts. du(withdirs=True) adds a size-0 root, so the totals are unchanged. fsspec is unpinned in pyproject.toml. Older fsspec _glob() does not pass prefix, which only avoids the prefix rules.
  • Docs: docs/filesystem.md:65 shows plain fs.find(path), which this PR does not change. No other docs mention withdirs, prefix, or find(maxdepth).

Corrections (d016e99): this docstring said that with withdirs only the directories starting with the prefix are included, which contradicts the included root. It now says that the directories above the prefix are not included. The PR body was updated the same way, and it now lists the two remaining differences from fsspec (empty bucket, object-and-prefix key).

refresh: If True, bypass the cache and list from S3.

Returns:
Expand Down Expand Up @@ -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.
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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]

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.

mv() data-safety fixes folded in at the maintainer's request (the pre-existing findings 2 and 3 from the earlier follow-up, plus mv([b?versionId=null], [b]), which deleted b). The commits are 1444eda and a7306a1, on master fe21250.

Self-review.

  • Round one (behavior): _move_target() compares bucket/key, keeping the version ID except for null. The sources left in place count for the "destination is another source" check. parse_path() accepts all four versionId spellings. An invalid path now raises at the check instead of later in cp_file(). The aio _mv() shares _move_paths(). No requests are added.
  • Round two (claims and operations): in a bucket with versioning enabled, the null version is an older version that a write does not replace. PyAthena does not call GetBucketVersioning, so there a null version moved onto its key is left in place, and [a, b?versionId=null] → [b, out] raises. The PR body and the release note state this. The "deleted b" claim follows from the requests: CopyObject b(null) → b, then DeleteObjects b with VersionId=null, which in a bucket without versioning is the object itself.

Independent follow-up on 1444eda (relayed): Codex gpt-6-astra, effort max, read-only, session 01a10268-75a8-7290-823f-c0f3375816c2.

Covered 1444edad: path normalization, all four version-ID spellings, invalid paths, sources left in place, duplicate destinations, directory exemptions, copy/delete ordering, async parity, request costs, and new tests.

The stated examples are repaired. Normalization adds no requests; directory exemptions retain conditional HeadObject calls, with no bucket-versioning lookup. Async uses the same preflight checks. The new tests verify rejection/no-request behavior and version-specific request parameters.

FINDINGS

P1 — Regression: directory exemption probes the wrong version. s3.py:1475 supplies normalized identities to the HeadObject check at line 1494, removing versionId=null.

Concrete scenario: versioned bucket src contains a readable null version of d, hidden by a current delete marker, plus objects d/x and a. Bucket dst is unversioned:

fs.mv(
    ["src/d?versionId=null", "src/d/x", "src/a"],
    ["dst/out", "dst/x", "dst/out"],
)

d/x makes normalized src/d a directory candidate. Unqualified HeadObject returns missing, so d → out is excluded from duplicate-destination checks. However, the returned pairs retain the original version-qualified source: its copy succeeds, then a → out overwrites it. Deletion permanently removes the original null version, losing its contents.

The parent rejected this duplicate destination. Async admits the same conflicting copies. This exceeds the accepted conservative null-version trade-off. Preserve the original source version for the existence probe, or exclude explicit versions from directory exemptions.

The new version test misses this because _serve_keys models keys without versions or delete markers.

Pre-existing limitation — scalar version-qualified sources undergo glob expansion. At s3.py:1452, mv("bucket/b?versionId=null", "bucket/out") treats ? as a wildcard and raises FileNotFoundError when only b exists. Both-list inputs bypass this. This behavior is unchanged by the commit.

Static review only; no tests, builds, edits, or network access. HEAD remained unchanged and the worktree clean.

Disposition.

Independent follow-up on a7306a1 (relayed): same setup, session 01a1026f-9a5d-7d22-8c50-8c1fa01747ea.

Covered a7306a1 against its parent:

  • Repair and collisions: s3.py:1495 checks the original path, so ?versionId=null cannot receive the directory exemption. A hidden null version with descendant keys now remains subject to both duplicate-destination and destination-is-source checks.
  • Data safety: Traced normalization, protocol aliases, null/non-null versions, sources left in place, ordinary/multipart copies, and deletion. No remaining collision bypass found within the reviewed scope.
  • Request costs: The added parsing makes no requests. Versioned candidates skip the exemption’s HeadObject; unversioned handling is unchanged.
  • Async parity: s3_async.py:363 awaits the same validation before scheduling any copies.
  • New test: test_s3.py:1606 exercises the null-version exemption hole through duplicate destinations and asserts no copy/delete requests. By inspection, it would fail against the parent. Its mock does not model stored versions/delete markers, so coverage is of preflight rejection.

CLEAN — the prior finding is resolved. No regression introduced by this commit or additional pre-existing defect identified within this scope.

Static review only; no files changed, builds/tests run, or GitHub/network access.

Validation at a7306a1: just lint passed; 48 offline find/glob/mv tests passed; the live tests/pyathena/filesystem/ run gave 564 passed. The snapshots and the PR worktree stayed clean.

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
):
Expand Down
5 changes: 3 additions & 2 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -354,8 +354,9 @@ async def _mv(self, path1, path2, recursive=False, maxdepth=None, **kwargs) -> N
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
Expand Down
Loading
Loading