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
52 changes: 45 additions & 7 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -300,13 +300,15 @@ directories below the bucket level) and is always a no-op.
## Typed S3 operations

`S3FileSystem.core` is an `S3Core`, the typed operations that the filesystem sends
its listing, lookup, read, write, delete, multipart upload and copy requests with. It can also be
built on a boto3 S3 client. Each operation sends one request (one per page for the
iterators and `list_object_annotations()`); `plan_multipart_copy()` and
`copy_object_annotation()`, described below, send several. The requests are sent with
the retry policy. An operation raises `FileNotFoundError` for a missing bucket or
multipart upload, or for a missing object or version that it reads, and caches
nothing. Requests sent through `fs.core` do not invalidate the filesystem's cache: call
its listing, lookup, read, write, delete, multipart upload, copy, tagging, ACL and
metadata replacement requests with. It can also be built on a boto3 S3 client. Each
operation sends one request (one per page for the iterators and
`list_object_annotations()`); `plan_multipart_copy()` and `copy_object_annotation()`,
described below, send several, and `generate_presigned_url()` signs a URL locally
without a request. The requests are sent with the retry policy. An operation raises
`FileNotFoundError` for a missing bucket or multipart upload, or for a missing object or
version that it reads, and caches nothing. Requests sent through `fs.core` do not
invalidate the filesystem's cache: call
`fs.invalidate_cache()` after a change, or make it through the filesystem.

```python
Expand Down Expand Up @@ -357,6 +359,42 @@ for batch in S3DeleteBatch.from_paths(paths):
print(error) # path (code: message)
```

`get_object_tagging()` returns the tags of an object as a dictionary, and
`put_object_tagging()` replaces all of its tags with the given ones. Both act on the
version ID of the path, if any, including `null`. `put_object_acl()` and
`put_bucket_acl()` apply a canned ACL, which must be in `S3Core.OBJECT_ACLS` or
`S3Core.BUCKET_ACLS`, respectively; another value raises `ValueError`, and an object
ACL also applies to the version ID of the path.

`replace_object_metadata()` replaces the user-defined metadata of an object, given the
`head_object()` result of the same path, by copying the object onto itself. The copy
rewrites the object, or creates a new version in a versioned bucket, so the path must
not have a version ID. The copy retains the content headers (`CacheControl`,
`ContentDisposition`, `ContentEncoding`, `ContentLanguage`, `ContentType`), `Expires`,
`WebsiteRedirectLocation`, `StorageClass`, and the `ServerSideEncryption`,
`SSEKMSKeyId` and `BucketKeyEnabled` of the object. Additional CopyObject parameters
take precedence over the retained fields, and any encryption parameter replaces all of
the retained encryption fields.

```python
path = S3Path.parse("s3://YOUR_S3_BUCKET/path/to/object")
core.put_object_tagging(path, {**core.get_object_tagging(path), "tag2": "value2"})

head = core.head_object(path)
core.replace_object_metadata(path, head, {**head, "attr1": "value1"})
```

`generate_presigned_url()` signs a URL locally, for `get_object` unless another client
method is given, and sends no request. Its parameters take precedence over the
`Bucket`, `Key` and `VersionId` of the path. The `request_kwargs` of the core, such as
`RequestPayer`, are not included in the signed parameters; pass them as parameters
of `generate_presigned_url()` to sign them.

```python
url = core.generate_presigned_url(path, expires_in=600)
upload_url = core.generate_presigned_url(path, "put_object", ContentType="text/csv")
```

### Multipart writer

`S3MultipartWriter` provides synchronous multipart requests and part planning without fsspec.
Expand Down
146 changes: 15 additions & 131 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -161,34 +161,6 @@ class S3FileSystem(AbstractFileSystem):
"""

DEFAULT_BLOCK_SIZE: int = 5 * 2**20 # 5MiB
# https://docs.aws.amazon.com/AmazonS3/latest/userguide/acl-overview.html#canned-acl
OBJECT_ACLS: frozenset[str] = frozenset(
{
"private",
"public-read",
"public-read-write",
"authenticated-read",
"aws-exec-read",
"bucket-owner-read",
"bucket-owner-full-control",
}
)
BUCKET_ACLS: frozenset[str] = frozenset(
{"private", "public-read", "public-read-write", "authenticated-read"}
)
# https://docs.aws.amazon.com/AmazonS3/latest/API/API_CopyObject.html
# The CopyObject parameters that set the encryption of the copy.
_SSE_COPY_PARAMS: frozenset[str] = frozenset(
{
"ServerSideEncryption",
"SSEKMSKeyId",
"SSEKMSEncryptionContext",
"BucketKeyEnabled",
"SSECustomerAlgorithm",
"SSECustomerKey",
"SSECustomerKeyMD5",
}
)
PATTERN_PATH: Pattern[str] = S3Path.PATTERN

protocol = ("s3", "s3a")
Expand Down Expand Up @@ -1282,8 +1254,8 @@ def mkdir(self, path: str, create_parents: bool = True, **kwargs) -> None:
"Set allow_bucket_creation=True on the filesystem to enable it."
)
acl = kwargs.pop("acl", "")
if acl and acl not in self.BUCKET_ACLS:
raise ValueError(f"ACL not in {self.BUCKET_ACLS}.")
if acl and acl not in self.core.BUCKET_ACLS:
raise ValueError(f"ACL not in {self.core.BUCKET_ACLS}.")
request: dict[str, Any] = {"Bucket": s3_path.bucket}
if acl:
request.update({"ACL": acl})
Expand Down Expand Up @@ -2197,23 +2169,9 @@ def sign(self, path: str, expiration: int = 3600, **kwargs):
... client_method="put_object"
... )
"""
s3_path = S3Path.parse(path)
client_method = kwargs.pop("client_method", "get_object")
params = {"Bucket": s3_path.bucket, "Key": s3_path.key}
if s3_path.version_id:
params.update({"VersionId": s3_path.version_id})
if kwargs:
params.update(kwargs)
request = {
"ClientMethod": client_method,
"Params": params,
"ExpiresIn": expiration,
}

_logger.debug(f"Generate signed url: {s3_path.uri}")
return self._call(
self._client.generate_presigned_url,
**request,
return self.core.generate_presigned_url(
S3Path.parse(path), client_method, expiration, **kwargs
)

def metadata(self, path: str, **kwargs) -> S3Metadata:
Expand Down Expand Up @@ -2295,43 +2253,7 @@ def setxattr(self, path: str, copy_kwargs: dict[str, Any] | None = None, **kw_ar
metadata.pop(k, None)
else:
metadata[k] = v

# With the REPLACE directive, S3 does not copy what the request
# omits: the system-defined metadata is dropped, and the copy is
# written as STANDARD with the default encryption of the bucket.
kept: dict[str, Any] = {
"CacheControl": head.cache_control,
"ContentDisposition": head.content_disposition,
"ContentEncoding": head.content_encoding,
"ContentLanguage": head.content_language,
"ContentType": head.content_type,
"Expires": head.expires,
"WebsiteRedirectLocation": head.website_redirect_location,
"StorageClass": head.storage_class,
}
copy_kwargs = copy_kwargs if copy_kwargs else {}
if not self._SSE_COPY_PARAMS.intersection(copy_kwargs):
kept.update(
{
"ServerSideEncryption": head.server_side_encryption,
"SSEKMSKeyId": head.sse_kms_key_id,
"BucketKeyEnabled": head.bucket_key_enabled,
}
)

_logger.debug(f"Set object metadata: {s3_path.uri}")
self._call(
self._client.copy_object,
CopySource={"Bucket": s3_path.bucket, "Key": s3_path.key},
Bucket=s3_path.bucket,
Key=s3_path.key,
Metadata=metadata,
MetadataDirective="REPLACE",
**{
**{k: v for k, v in kept.items() if v is not None},
**copy_kwargs,
},
)
self.core.replace_object_metadata(s3_path, head, metadata, **(copy_kwargs or {}))
self.invalidate_cache(path)

def get_tags(self, path: str) -> dict[str, str]:
Expand All @@ -2346,16 +2268,7 @@ def get_tags(self, path: str) -> dict[str, str]:
s3_path = S3Path.parse(path)
if not s3_path.key:
raise ValueError("Cannot get tags of a bucket.")
request: dict[str, Any] = {"Bucket": s3_path.bucket, "Key": s3_path.key}
if s3_path.version_id:
request.update({"VersionId": s3_path.version_id})

_logger.debug(f"Get object tagging: {s3_path.uri}")
response = self._call(
self._client.get_object_tagging,
**request,
)
return {v["Key"]: v["Value"] for v in response["TagSet"]}
return self.core.get_object_tagging(s3_path)

def put_tags(self, path: str, tags: dict[str, str], mode: str = "o") -> None:
"""Set the tags for the given existing key.
Expand All @@ -2377,26 +2290,12 @@ def put_tags(self, path: str, tags: dict[str, str], mode: str = "o") -> None:
if not s3_path.key:
raise ValueError("Cannot put tags of a bucket.")
if mode == "m":
existing_tags = self.get_tags(path)
existing_tags.update(tags)
new_tags = [{"Key": k, "Value": v} for k, v in existing_tags.items()]
new_tags = {**self.core.get_object_tagging(s3_path), **tags}
elif mode == "o":
new_tags = [{"Key": k, "Value": v} for k, v in tags.items()]
new_tags = tags
else:
raise ValueError(f"Mode must be {{'o', 'm'}}, not {mode}.")
request: dict[str, Any] = {
"Bucket": s3_path.bucket,
"Key": s3_path.key,
"Tagging": {"TagSet": new_tags},
}
if s3_path.version_id:
request.update({"VersionId": s3_path.version_id})

_logger.debug(f"Put object tagging: {s3_path.uri}")
self._call(
self._client.put_object_tagging,
**request,
)
self.core.put_object_tagging(s3_path, new_tags)

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): base 855d4a7, head c3a4e4f.

Claims checked:

  • "The requests they send are unchanged": traced each operation. test_put_tags_merge and test_put_tags_invalid_mode also pass on master 855d4a7, unmodified, and the existing Stubber tests for setxattr, the mock tests for chmod, and the url() request test pass unchanged.
  • Tag order and the merge semantics of mode m: dict order is preserved, and the merge yields the same TagSet as existing.update(tags).
  • The core class docstring and docs "each operation sends one request" with the named exceptions: generate_presigned_url sends none (it is local and listed as an exception). Each new operation otherwise sends one request.
  • Release note for the removed attributes: git grep finds no other users in the repository or docs.
  • The S3Metadata STANDARD note: matches test_setxattr_omits_unset_system_metadata on master.

Existing callers:

  • The public signatures and return shapes of get_tags/put_tags/chmod/setxattr/sign are unchanged.
  • A subclass or test that overrides only S3FileSystem._call no longer intercepts these requests, which now go through core.call. This is the same trade-off as the earlier Separate an S3 core from the fsspec adapter in pyathena.filesystem #1053 steps, and _call stays a delegate.

AWS operator: the same retry layer (S3Core.call) and the same request count.

Documentation reader: docs/filesystem.md lists the new request kinds. The setxattr note is still accurate.

Evidence limits: static plus offline Stubber/mock evidence only. The live AWS tests (test_sign, test_metadata_getxattr_setxattr, test_get_and_put_tags, test_chmod_validation, and the aio equivalents) have not run on this head; they run when the PR is marked Ready. just docs build builds committed refs and was not used as evidence.

Result: CLEAN (no corrections needed beyond round one's repairs).

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.

Docs follow-up (09e4bce and 3b445a1, docs/filesystem.md only)

At the maintainer's request, the "Typed S3 operations" section now states the contracts of the new operations:

  • put_object_tagging() replaces all tags. Tagging and object ACLs act on the version ID of the path, including null.
  • ACLs are canned only and are checked against S3Core.OBJECT_ACLS and S3Core.BUCKET_ACLS; another value raises ValueError.
  • replace_object_metadata() takes the head_object() result of the same path. It rewrites the object, or creates a new version in a versioned bucket. The section lists the retained fields, says that parameters take precedence over them, and says that any encryption parameter replaces the retained encryption.
  • generate_presigned_url() signs locally, and its parameters take precedence over Bucket/Key/VersionId. request_kwargs such as RequestPayer are not signed unless they are passed as parameters.

The section also has two small examples.

Self-check:

  • Documentation-reader and claim perspectives: every statement was checked against s3_core.py.
  • The examples ran offline against Stubber with this code. They send the documented requests, and an explicit RequestPayer appears in the signed URL.
  • just docs lint: 0 errors.

Independent follow-up:

  • On 962ed09e..09e4bceb: Codex (codex exec -s read-only) CLEAN, static review. claude-fable-5-1 (profile max, effort high, plan mode) CLEAN, static review, with one optional note: "pass them to the call" could be read as S3Core.call().
  • That wording was fixed in 3b445a1 ("pass them as parameters of generate_presigned_url()"). On 09e4bceb..3b445a18, Codex and Fable were both CLEAN.

The live AWS run by the coordinator at 962ed09 passed: pytest -n 8 tests/pyathena/filesystem/ tests/pyathena/s3fs/ tests/pyathena/aio/s3fs/ → 1156 passed, 1 skipped. The later commits change only the docs.

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 master 5a0a5cb2e4894e1bb49edfc9bd587859dc97bf1d (after #1084 and #1091): 3b445a1841cd872069caa3c71bc6a484bd672212 → 5c233f8e04c8ba08328ee87aa2454d7fcfd88d20, force-pushed with an explicit lease.

  • Range-diff (855d4a7b..3b445a18 vs 5a0a5cb2..5c233f8e): the five commits are unchanged except for the conflict in the intro paragraph of docs/filesystem.md.
    • The resolution keeps Move object GET and PUT requests into S3Core #1091's "read, write" and this PR's "tagging, ACL and metadata replacement".
    • It also keeps this PR's generate_presigned_url() sentence.
    • The source and test patches applied without conflicts.
  • Upstream contracts checked:
  • Validation at 5c233f8e:
    • just lint and just docs lint pass.
    • Offline pytest --noconftest -n 8 tests/pyathena/filesystem/ gives 927 passed, 164 failed and 9 skipped, against master 5a0a5cb2 with 899 passed, 164 failed and 9 skipped. The failing tests are identical on both and need AWS.
    • The full AWS suite runs in CI when the PR is marked Ready.
  • Out of scope for this PR: Fix null-version moves in versioning-enabled S3 buckets #1084 added a get_bucket_versioning request and a call to the private S3Core._is_directory_bucket() in the adapter's _move_pairs (s3.py:1464). The final Finish the public S3Core API extraction from filesystem adapters #1086 SDK-call audit will cover them.


def chmod(self, path: str, acl: str, recursive: bool = False, **kwargs) -> None:
"""Set the Access Control on a bucket/key.
Expand All @@ -2414,10 +2313,10 @@ def chmod(self, path: str, acl: str, recursive: bool = False, **kwargs) -> None:
s3_path = S3Path.parse(path)
# Validate before any ACL is applied so that a recursive call cannot
# partially apply object ACLs and then fail on the bucket ACL.
if not s3_path.key and acl not in self.BUCKET_ACLS:
raise ValueError(f"ACL not in {self.BUCKET_ACLS}.")
if s3_path.key and acl not in self.OBJECT_ACLS:
raise ValueError(f"ACL not in {self.OBJECT_ACLS}.")
if not s3_path.key and acl not in self.core.BUCKET_ACLS:
raise ValueError(f"ACL not in {self.core.BUCKET_ACLS}.")
if s3_path.key and acl not in self.core.OBJECT_ACLS:
raise ValueError(f"ACL not in {self.core.OBJECT_ACLS}.")
if recursive:
with self._create_executor(max_workers=self.max_workers) as executor:
futures = [
Expand All @@ -2431,24 +2330,9 @@ def chmod(self, path: str, acl: str, recursive: bool = False, **kwargs) -> None:
# below it have ACLs.
return
if s3_path.key:
request: dict[str, Any] = {"Bucket": s3_path.bucket, "Key": s3_path.key, "ACL": acl}
if s3_path.version_id:
request.update({"VersionId": s3_path.version_id})

_logger.debug(f"Put object acl: {s3_path.uri}")
self._call(
self._client.put_object_acl,
**request,
**kwargs,
)
self.core.put_object_acl(s3_path, acl, **kwargs)
else:
_logger.debug(f"Put bucket acl: {s3_path.uri}")
self._call(
self._client.put_bucket_acl,
Bucket=s3_path.bucket,
ACL=acl,
**kwargs,
)
self.core.put_bucket_acl(s3_path.bucket, acl, **kwargs)

def list_multipart_uploads(self, path: str) -> list[S3MultipartUpload]:
"""List in-progress (incomplete) multipart uploads in a bucket.
Expand Down
Loading
Loading