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
1 change: 1 addition & 0 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,7 @@ operations raise natural errors instead of botocore's `ClientError`:
| `404` / `NoSuchKey` / `NoSuchBucket` | `FileNotFoundError` |
| `403` / `AccessDenied` | `PermissionError` |
| `BucketAlreadyExists` / `BucketAlreadyOwnedByYou` | `FileExistsError` |
| `PreconditionFailed` of an `If-None-Match` condition | `FileExistsError` |
| `RequestTimeout` | `TimeoutError` |
| Others | `OSError` with the matching `errno` |

Expand Down
52 changes: 43 additions & 9 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -1568,7 +1568,8 @@ def pipe_file(
path: S3 path (s3://bucket/key) to write to.
value: The bytes to write.
mode: "overwrite" (default) or "create". With "create", raise
FileExistsError when the object already exists.
FileExistsError when the object already exists, including
one created during the write, which is not replaced.
**kwargs: Additional parameters passed to the PutObject API
(e.g., ContentType, StorageClass) on the single-request
path. The ``block_size``, ``max_workers``, and
Expand All @@ -1577,7 +1578,8 @@ def pipe_file(

Raises:
FileExistsError: If the mode is "create" and the path already
exists.
exists, or an object is created at it before the write is
committed.
ValueError: If the path does not contain a key or specifies a
version, or if the data takes more than
``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
Expand All @@ -1589,15 +1591,20 @@ def pipe_file(
# Defer to the buffered open() path, which keeps the
# deferred-commit semantics of fsspec transactions and uploads
# large data as a parallel multipart upload.
super().pipe_file(path, value, mode=mode, **kwargs)
with self.open(path, "xb" if mode == "create" else "wb", **kwargs) as f:
f.write(value)
return
bucket, key, version_id = self.parse_path(path)
if version_id:
raise ValueError("Cannot write to the file with the version specified.")
if not key:
raise ValueError("Cannot write to a bucket.")
if mode == "create" and self.exists(path):
raise FileExistsError(path)
if mode == "create":
# Checked up front, as open() does in "xb" mode, and with
# IfNoneMatch for an object created since.
if self.exists(path):
raise FileExistsError(path)
kwargs["IfNoneMatch"] = "*"
if not isinstance(value, bytes):
# Accept bytes-like values (bytearray, memoryview) as the
# buffered path does.
Expand Down Expand Up @@ -1753,7 +1760,14 @@ def cat_file(
return b""
raise

def put_file(self, lpath: str, rpath: str, callback=_DEFAULT_CALLBACK, **kwargs):
def put_file(
self,
lpath: str,
rpath: str,
callback=_DEFAULT_CALLBACK,
mode: str = "overwrite",
**kwargs,
):
"""Upload a local file to S3.

Uploads a file from the local filesystem to an S3 location. Supports
Expand All @@ -1764,11 +1778,18 @@ def put_file(self, lpath: str, rpath: str, callback=_DEFAULT_CALLBACK, **kwargs)
lpath: Local file path to upload.
rpath: S3 destination path (s3://bucket/key).
callback: Progress callback for tracking upload progress.
mode: "overwrite" (default) or "create". With "create", the file
is written as with ``open()`` in ``xb`` mode: raise
FileExistsError when the object already exists, including
one created during the upload, which is not replaced.
**kwargs: Additional S3 parameters (e.g., ContentType, StorageClass).
The ``block_size``, ``max_workers``, and ``s3_additional_kwargs``
parameters of ``open()`` are also accepted.

Raises:
FileExistsError: If the mode is "create" and the path already
exists, or an object is created at it before the upload is
committed.
ValueError: If the file takes more than
``MULTIPART_UPLOAD_MAX_PARTS`` blocks.

Expand Down Expand Up @@ -1801,7 +1822,7 @@ def put_file(self, lpath: str, rpath: str, callback=_DEFAULT_CALLBACK, **kwargs)
with (
self.open(
rpath,
"wb",
"xb" if mode == "create" else "wb",
block_size=block_size,
max_workers=max_workers,
s3_additional_kwargs=s3_additional_kwargs,
Expand Down Expand Up @@ -2589,12 +2610,15 @@ def __init__(
existing object smaller than ``MULTIPART_UPLOAD_MIN_PART_SIZE`` is
read into the write buffer; a larger one is copied with
``UploadPartCopy`` as the first parts of a multipart upload, whatever
the block size.
the block size. In exclusive-create mode, the object must not exist
when the file is opened, and the upload is committed with
``IfNoneMatch="*"`` so that it does not replace an object created in
the meantime.

Args:
fs: The filesystem that the file belongs to.
path: S3 path (s3://bucket/key) of the file.
mode: The file mode, such as ``rb``, ``wb`` or ``ab``.
mode: The file mode: ``rb``, ``wb``, ``ab``, or ``xb``.
version_id: The version ID to read. Must match the version ID in
the path if both are given. A version cannot be given, in
either form, for writing or appending.
Expand All @@ -2618,6 +2642,8 @@ def __init__(
which take precedence over ``s3_additional_kwargs``.

Raises:
FileExistsError: If an object exists at the path in
exclusive-create mode.
FileNotFoundError: If no object exists at the path when reading,
including when the path is a prefix.
ValueError: If the path has no key, the version IDs do not match,
Expand Down Expand Up @@ -2691,6 +2717,12 @@ def __init__(
# Too small to be a part of a multipart upload: rewritten
# from the buffer.
append_data = fs.cat(path)
elif "x" in mode:

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, resolved without code change)

Base 28826be7d37983cb54050e849b32a4fe2331ea7d, head bf5743b6d3fca17a1d37ce716cad6312a858044c.

Covered: S3File.__init__ x-mode check and IfNoneMatch routing through _get_request_kwargs to PutObject / CompleteMultipartUpload / touch() (not CreateMultipartUpload / UploadPart; checked against the botocore input shapes); commit() failure paths (412 on PutObject; 412 on completion → _finish_multipart_upload abort, multipart_upload reset); pipe_file() single-request and buffered paths; put_file() sync, async, and async in-transaction; S3ClientError translation (If-None-Match vs If-Match vs no condition); fsspec Transaction.complete() on a failed commit (re-queues and discards the rest, clears _intrans); fsspec put()/_put() forwarding mode through **kwargs; executor on a FileExistsError at open (created lazily, as on the read path's FileNotFoundError); botocore floor 1.41.2 has IfNoneMatch on both operations.

Finding: the up-front exists() propagates PermissionError for a key that HeadObject denies (pyathena/filesystem/s3.py:917-923 catches only FileNotFoundError), so a PutObject-only principal cannot use xb/create. Pre-existing for pipe_file(mode=\"create\"). Maintainer chose to keep it (permission errors are not hidden); recorded in the PR body.

Tests: the regression tests fail on the original code (15 cases), and the live test exercises the real 412 path. The async test_put_file_mode avoids the shared internal S3FileSystem instance (pre-existing test pollution, noted in the PR body).

# Checked up front so that no data is uploaded for an existing
# object, and on commit with IfNoneMatch for one created since.
if fs.exists(path):
raise FileExistsError(path)
self.s3_additional_kwargs.update({"IfNoneMatch": "*"})

self._executor: S3Executor = executor or S3ThreadPoolExecutor(max_workers=max_workers)
super().__init__(
Expand Down Expand Up @@ -2880,6 +2912,8 @@ def commit(self) -> None:
completion fails. Invalidates the cache of the path afterwards.

Raises:
FileExistsError: If an object was created at the path after the
file was opened in exclusive-create mode.
RuntimeError: If parts were submitted but no multipart upload is
initialized.
"""
Expand Down
39 changes: 29 additions & 10 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -180,8 +180,10 @@ def _pipe_file_in_transaction(
Args:
path: S3 path (s3://bucket/key) to write to.
value: The bytes to write.
mode: "overwrite" or "create". With "create", raise
FileExistsError when the object already exists.
mode: "overwrite" or "create". With "create", the file is
opened in ``xb`` mode: raise FileExistsError when the object
already exists, including one created before the
transaction is committed, which is not replaced.
**kwargs: Additional parameters passed to ``open()``.

Raises:
Expand All @@ -193,19 +195,30 @@ def _pipe_file_in_transaction(
block_size = kwargs.get("block_size") or self._sync_fs.default_block_size
# The size in bytes; the length of a memoryview counts its items.
self._sync_fs._check_multipart_upload_size(path, memoryview(value).nbytes, block_size)
if mode == "create" and self._sync_fs.exists(path):
raise FileExistsError(path)
with self.open(path, "wb", **kwargs) as f:
with self.open(path, "xb" if mode == "create" else "wb", **kwargs) as f:
f.write(value)

async def _put_file(self, lpath: str, rpath: str, callback=_DEFAULT_CALLBACK, **kwargs) -> None:
async def _put_file(
self,
lpath: str,
rpath: str,
callback=_DEFAULT_CALLBACK,
mode: str = "overwrite",
**kwargs,
) -> None:
if self._intrans:
# See _pipe_file.
await asyncio.to_thread(self._put_file_in_transaction, lpath, rpath, callback, **kwargs)
await asyncio.to_thread(
self._put_file_in_transaction, lpath, rpath, callback, mode, **kwargs
)
return
await asyncio.to_thread(self._sync_fs.put_file, lpath, rpath, callback=callback, **kwargs)
await asyncio.to_thread(
self._sync_fs.put_file, lpath, rpath, callback=callback, mode=mode, **kwargs
)

def _put_file_in_transaction(self, lpath: str, rpath: str, callback, **kwargs) -> None:
def _put_file_in_transaction(
self, lpath: str, rpath: str, callback, mode: str, **kwargs
) -> None:
"""Upload a local file as a file of this filesystem's transaction.

Mirrors :meth:`S3FileSystem.put_file`, but writes through ``open()``
Expand All @@ -215,11 +228,17 @@ def _put_file_in_transaction(self, lpath: str, rpath: str, callback, **kwargs) -
lpath: Local file path to upload.
rpath: S3 destination path (s3://bucket/key).
callback: Progress callback for tracking upload progress.
mode: "overwrite" or "create". With "create", the file is
opened in ``xb`` mode: raise FileExistsError when the object
already exists, including one created before the
transaction is committed, which is not replaced.
**kwargs: Additional S3 parameters (e.g., ContentType, StorageClass).
The ``block_size``, ``max_workers``, and ``s3_additional_kwargs``
parameters of ``open()`` are also accepted.

Raises:
FileExistsError: If the mode is "create" and the path already
exists.
ValueError: If the file takes more than
``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
"""
Expand All @@ -244,7 +263,7 @@ def _put_file_in_transaction(self, lpath: str, rpath: str, callback, **kwargs) -
with (
self.open(
rpath,
"wb",
"xb" if mode == "create" else "wb",
block_size=block_size,
max_workers=max_workers,
s3_additional_kwargs=s3_additional_kwargs,
Expand Down
8 changes: 7 additions & 1 deletion pyathena/filesystem/s3_errors.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,9 @@ class S3ClientError:
as properties, along with :attr:`os_error`, the equivalent standard
Python exception. The error is mapped by its S3 error code first, then
by its HTTP status code; if neither is recognized, a generic ``OSError``
with the original error message is used.
with the original error message is used. A ``PreconditionFailed`` error
of an ``If-None-Match`` condition, from a write that must not replace an
existing object, is mapped to ``FileExistsError``.

Example:
>>> try:
Expand Down Expand Up @@ -105,11 +107,15 @@ def __init__(self, error: botocore.exceptions.ClientError) -> None:
error_info = error.response.get("Error", {})
self._code: str = str(error_info.get("Code", ""))
self._message: str = str(error_info.get("Message", error))
self._condition: str = str(error_info.get("Condition", ""))
status_code = error.response.get("ResponseMetadata", {}).get("HTTPStatusCode")
self._http_status_code: int | None = int(status_code) if status_code is not None else None
self._os_error: OSError = self._translate()

def _translate(self) -> OSError:
if self._code == "PreconditionFailed" and self._condition == "If-None-Match":

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 2 (claims, callers, operations) — FINDINGS (2, repaired in 08bcff5 and the PR body)

Base 28826be7d37983cb54050e849b32a4fe2331ea7d, head bf5743b6d3fca17a1d37ce716cad6312a858044c (repair head 08bcff53).

Claims checked:

  • 412 shapes for PutObject/CompleteMultipartUpload/If-Match and abort-after-failed-completion: measured on S3 before implementation; live test_exclusive_create reproduces the commit-time 412.
  • mode already dropped on master since Route S3 request parameters to the operations that accept them #1000: test_put_file_mode[overwrite] (sync and async) passes with the source changes reverted.
  • 15 failing tests without the fix, 453 passing live: match the recorded runs on bf5743b's tree.
  • Write pipe_file() data without committing a failed write #1003 overlap: its diff changes super().pipe_file(...) and the async with self.open(path, \"wb\"); the described resolution matches.
  • fsspec 2026.9.0 (uv.lock): _Cached.__call__ pops skip_instance_cache (spec.py:108), confirming the test-isolation note; Transaction.complete() re-raises a failed commit and discards the rest.
  • Callers: put_file's new mode sits where fsspec's base has it (after callback); no positional argument beyond callback was possible before. fsspec put()/_put() pass mode by keyword.
  • Retries: RetryConfig retries throttling codes only, so a 412 is not retried.

Findings and repairs:

  1. docs/filesystem.md error-translation table omitted the new FileExistsError mapping — row added (08bcff5), just docs lint passes.
  2. The PR body did not state the operational cost or edge cases of the up-front check — added: HeadObject + ListObjectsV2 per open of a missing key (info() at pyathena/filesystem/s3.py:704-725), a same-name prefix counts as existing, and S3's documented 409 ConditionalRequestConflict stays OSError(EBUSY) (unmeasured).

# A conditional write (IfNoneMatch="*") found an existing object.
return FileExistsError(self._message)
exception = self._ERROR_CODE_TO_EXCEPTION.get(self._code)
if exception:
return exception(self._message)
Expand Down
Loading
Loading