Skip to content

Compress pipe_file() values up front, and route and key them consistently - #1038

Merged
laughingman7743 merged 11 commits into
masterfrom
fix/1037-pipe-file-compression
Oct 3, 2026
Merged

laughingman7743 merged 11 commits into
masterfrom
fix/1037-pipe-file-compression

Conversation

@laughingman7743

@laughingman7743 laughingman7743 commented Oct 3, 2026 •

Copy link
Copy Markdown
Member

WHAT

S3FileSystem.pipe_file() and AioS3FileSystem._pipe_file_in_transaction() compress the value up front when fsspec's compression argument is given, and upload the compressed bytes through the usual paths.

  • pipe_file() normalizes the path as open() does (self._strip_protocol(path), which also drops a trailing slash) before anything else, so the single-request path writes the same key as the buffered path. A path with a trailing slash used to be written to the key with the slash by PutObject and without it through open(), depending on the size of the value; compression could switch a call between the two.

  • The callers resolve the codec as open() does: fsspec.core.get_compression() on that normalized path, including "infer" from its extension. When "infer" finds no codec, the value is uploaded as it is, as open() does.

  • A new public CompressedBuffer in pyathena/filesystem/s3.py (not added to the API reference, like the other support classes): CompressedBuffer.compress(value, compression) compresses a value with a codec name of fsspec.compression.compr and returns bytes; an unknown name (also "infer") raises ValueError. The buffer is a BytesIO whose close() keeps the data, because some codecs close the file they write to (the zstandard stream writer, fsspec's zstd codec when neither compression.zstd nor backports.zstd is available).

  • compression is no longer passed to open() or to S3. The routing between the single-request and buffered paths, and the multipart size check, use the compressed size.

  • pipe_file() chooses between the single-request and buffered paths by the size of the value in bytes (memoryview(value).nbytes, already used for the multipart size check) instead of len(value), which counts the items of a memoryview. test_pipe_file_non_contiguous_memoryview pinned the item count (4 items, 8 bytes, block_size=6 → PutObject); it now uses block_size=8 for the same PutObject case.

  • docs/filesystem.md: the trailing-slash paragraph now says that pipe, pipe_file, and put_file also write the object dir for dir/ (fsspec's put() resolves directory destinations itself and is not described); and it states that pipe/pipe_file compress the data before uploading it with the compression argument, and that the sizes of its write section apply to the compressed data.

Behavior changes (release note):

  • pipe_file(..., compression=...) for a value up to the block size outside a transaction used to fail with ParamValidationError: Unknown parameter in input: "compression"; it now uploads the compressed value.
  • A failed write of a compressed value on the buffered path used to raise AttributeError: 'GzipFile' object has no attribute '_close_without_commit', hiding the original error, and to upload a 10-byte object holding only the gzip header, in and outside a transaction (checked offline on master fe21250a); it now raises the original error and leaves the existing object unchanged, as Write pipe_file() data without committing a failed write #1003 does for uncompressed values.
  • An unsupported codec raises ValueError before any request.
  • A typed memoryview (e.g. memoryview(data).cast("I")) larger than the block size in bytes, but not in items, used to go to PutObject, which fails above 5 GiB; it is now uploaded through the buffered (multipart) path.
  • pipe_file() called directly with a value up to the block size, outside a transaction, and a path with a trailing slash (s3://bucket/dir/) used to write the key dir/; it now writes dir, as open(), put_file(), larger values, and fsspec's pipe() (which strips the path before calling pipe_file()) already did. touch() still writes the key as given, e.g. to create a folder marker.

WHY

Closes #1037.

The routing by bytes is a pre-existing defect that the independent review of this PR found (reported in its review thread); the maintainer asked to fix it here.

Found in the review of PR #1035 (#1034): moving the write-and-close helper to S3File breaks compressed buffered writes while open() returns a compression wrapper. With this change open() always returns the S3File, and #1035 will be rebased onto it.

TEST

Tested commit: 73cdb62 (code); 67a3db4 and c6297d3 change only docs/filesystem.md (just docs lint passed)

  • just lint: passed. just docs lint and just docs build: passed on 8462270 (docs unchanged since).
  • uv run --env-file .env pytest -n 8 tests/pyathena/filesystem/: 565 passed (live S3).
  • New offline tests: test_pipe_file_compression (gzip, "infer", and "infer" with a trailing slash, up to and above the block size, in and outside a transaction), test_pipe_file_compression_multipart (compressed data above the block size, in and outside a transaction), test_pipe_file_compression_non_contiguous_memoryview, test_pipe_file_compression_inferred_none, test_pipe_file_unsupported_compression, test_pipe_file_compression_failed_write (in and outside a transaction), aio test_pipe_file_compression (gzip and "infer" with a trailing slash, in and outside a transaction), TestCompressedBuffer (gzip/bz2/xz round trips, a non-contiguous memoryview, unsupported names incl. "infer", and a codec that closes its file, which fails without the close() override), test_pipe_file_trailing_slash (the key written for s3://bucket/dir/key/, up to and above the block size, in and outside a transaction; the single-request case fails without the normalization, as do the trailing-slash cases of test_pipe_file_compression, which now also assert the key), and test_pipe_file_memoryview_routed_by_bytes (a cast("I") memoryview of block size + 4 bytes goes to the multipart path; fails with len(value) routing).
  • With the source changes reverted (tests kept), the 7 cases for the defects fail: the single-request path (2), "infer" without a codec, the unsupported codec, the failed write (2), and aio outside a transaction. The 7 buffered-path cases pass on master too, which already compressed those.
  • Every codec registered with zstandard and lz4 installed in a temporary environment (zip, bz2, gzip, lzma, xz, lz4, zstd) round-trips through CompressedBuffer.compress() and the codec's reader; zstd failed with a plain BytesIO.
  • The routing by bytes is checked offline only: a typed memoryview over 5 GiB was not uploaded to S3.
  • Live check of the trailing slash (ad hoc, objects deleted): pipe_file() of 1 byte and of 6 MiB, and put_file(), to paths ending in / wrote the keys small, big, and put without the slash; fsspec's put() to dir/ wrote dir/local.txt, and pipe() to piped/ wrote piped.
  • Live check against S3 (ad hoc, not committed): pipe_file() with compression="infer" on a .gz key (7 bytes, single request) and compression="gzip" (6 MiB of random bytes, multipart), each in and outside a transaction; every object read back with cat_file() decompressed to the value. The objects were deleted.

🤖 Generated with Claude Code

pipe_file() passed fsspec's compression argument to PutObject on the
single-request path, which botocore rejects, and to open() on the
buffered path, which returns a compression wrapper. A failed write on the
buffered path then raised AttributeError from _close_without_commit() on
the wrapper, hiding the original error, and still uploaded an empty
compressed object.

Compress the value in memory with the codec that open() uses, including
"infer" from the path, and upload the compressed bytes through the usual
paths. The same applies to the transaction path of AioS3FileSystem.

Closes #1037

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
laughingman7743 and others added 2 commits October 3, 2026 23:52
The zstandard stream writer, fsspec's zstd codec without compression.zstd
or backports.zstd, closes the file that it writes to, so reading the
compressed bytes from the BytesIO afterwards raised ValueError. Write to a
BytesIO whose close() keeps the data.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Comment thread pyathena/filesystem/s3.py Outdated
_logger = logging.getLogger(__name__)


class _CompressedBuffer(BytesIO):

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

Base fe21250a3e228e3d38e61898334c5baacc1285c2 (merge-base with master). Initial pass on e25850c869662e5e938b486cbfdf923d3784894f; repair ebb855d7cc6759eec2981f00f5d2601eb4bd27c7, rechecked on 846227079ea492e3edddb55ea7fd4e0911d1bee8.

Covered: _compress(), S3FileSystem.pipe_file(), AioS3FileSystem._pipe_file_in_transaction() (and _pipe_file() outside a transaction, which delegates to the sync pipe_file()), the new tests.

Finding (repaired in ebb855d7): _compress() wrote to a plain BytesIO. fsspec's zstd codec falls back to the zstandard stream writer when neither compression.zstd nor backports.zstd is available (fsspec 2026.9.0 compression.py:158-177). That writer closes the file it writes to, so buffer.getvalue() raised ValueError: I/O operation on closed file. Reproduced with uv run --with zstandard. _CompressedBuffer keeps the data on close(). test_compress_codec_closing_its_file fails with a plain BytesIO. With zstandard and lz4 installed, every registered codec (zip, bz2, gzip, lzma, xz, lz4, zstd) round-trips through _compress() and the codec's reader.

Checked, no finding:

  • compression is popped before the size check, so the routing and _check_multipart_upload_size() use the compressed size, and compression reaches neither open() nor PutObject.
  • compression=None given explicitly is popped and ignored, as open() ignored it.
  • "infer" without a codec returns the value unchanged, so it takes the uncompressed paths as before.
  • A non-contiguous memoryview is converted before compressing (codecs need a contiguous buffer).
  • Memory: the value is already in memory; the compressed copy is one more buffer, the same as the wrapper's buffering on the old path.
  • put_file() never passed compression to open() and is unchanged.

Validation: just lint passed; tests/pyathena/filesystem/ 545 passed against live S3 on ebb855d7.

Comment thread docs/filesystem.md
through the buffered file path. Inside an
[fsspec transaction](https://filesystem-spec.readthedocs.io/en/latest/features.html#transactions),
writes are deferred until the transaction commits and are discarded on rollback.
With the `compression` argument (a codec of fsspec, or `"infer"` from the extension of

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 (text only, repaired)

Base fe21250a3e228e3d38e61898334c5baacc1285c2, head 846227079ea492e3edddb55ea7fd4e0911d1bee8.

Claims checked:

  • Single-request path on master raises ParamValidationError: Unknown parameter in input: "compression": the PutObject request on master holds compression (_make_api_call captured), and botocore validate_parameters against the PutObject input shape rejects it. Not sent to live S3.
  • Failed compressed write on master: raises AttributeError: 'GzipFile' object has no attribute '_close_without_commit' with the original error as context, and uploads an object, both outside a transaction (incompressible-size path) and inside one, checked offline on fe21250a.
  • "7 cases fail with the source reverted": rerun with only pyathena/ reverted.
  • Live round trip (ad hoc script, objects deleted): compression="infer" on a .gz key (single request) and compression="gzip" on 6 MiB of random bytes (multipart), in and outside a transaction, decompress to the values.
  • Docstrings of pipe_file() and _pipe_file_in_transaction() match the code; this docs sentence matches the routing on the compressed size.

Finding (repaired): the issue, the PR body, and the relayed review on #1035 said a failed write uploaded "an empty gzip object (10 bytes)". The object holds only the 10-byte gzip header, which is not a valid gzip file; an empty gzip file is 20 bytes. All three now say "a 10-byte object holding only the gzip header".

Callers: fsspec pipe() passes compression per file through pipe_file(); AioS3FileSystem.pipe_file() outside a transaction delegates to the sync method. Operationally: no extra S3 requests; small compressed values now take one PutObject instead of failing.

fsspec's open() infers the codec from the path without the protocol,
which also drops a trailing slash. pipe_file() inferred it from the path
as given, so "key.gz/" was uploaded uncompressed. Strip the protocol
first, also in the transaction path of AioS3FileSystem.

Cover compressed multipart uploads and non-contiguous memoryviews too.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Comment thread pyathena/filesystem/s3.py Outdated
Raises:
ValueError: If the codec is not supported.
"""
compression = get_compression(path, compression)

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): FINDINGS — regression

Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only, model reported gpt-6-astra), session 01a10244-528d-7e50-8ab1-84e0bb7d40a2. Static review only. Base fe21250a3e228e3d38e61898334c5baacc1285c2, head 846227079ea492e3edddb55ea7fd4e0911d1bee8, detached snapshot (unchanged afterwards). The prompt held the diff, the intent, and conventions, without the PR number, PR text, or prior findings.

Covered: sync and async pipe/pipe_file; single-request and multipart uploads; transaction commit/rollback; mode="create"; block sizes and failure cleanup; bytes, bytearray, contiguous/non-contiguous and non-byte-format memoryviews; compression None, "infer", unsupported names; all codec registration branches incl. writers that close their buffer; changed tests, docstrings, docs, conventions.

P2 — Compression inference differs from open() for trailing-slash paths. s3.py:84 infers compression from the unnormalized path. fsspec open() first strips trailing slashes, then infers compression. Consequently, inside a transaction, fs.pipe_file("s3://bucket/key.gz/", b"data", compression="infer") now stores uncompressed bytes at key.gz; previously it stored gzip data. Reading that object with inferred compression fails. Both sync and async transaction paths are affected. Infer from the same normalized path that open() uses, and add a regression case for it.

The added tests assert observable uploaded bytes, keyword removal, and absence of uploads after write failures. [...] However, every successful compression case fits within one compressed block, so successful compressed multipart uploads and compressed memoryview variants remain uncovered.

Author verification: confirmed. fsspec 2026.9.0 AbstractFileSystem.open() calls self._strip_protocol(path) before get_compression(path, compression), and _strip_protocol("s3://bucket/key.gz/") is bucket/key.gz; on 84622707 the transaction write sent b"data" uncompressed to key.gz.

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.

Repaired in 135e64872fad15cc949a5f659f4a71f6fec49cc4: pipe_file() and AioS3FileSystem._pipe_file_in_transaction() pass self._strip_protocol(path) to _compress(), and its docstring says so. test_pipe_file_compression gains a "s3://bucket/key.gz/" + "infer" case (4 parametrizations, which fail without the repair). The coverage gap is closed by test_pipe_file_compression_multipart (compressed data above the block size, in and outside a transaction) and test_pipe_file_compression_non_contiguous_memoryview.

Self-review of the repair:

  • Round 1 (behavior): _strip_protocol() is fsspec's default for both filesystems (neither overrides it) and is what open() calls; a ?versionId= suffix stays in the path, as in open(), so inference matches there too. Only "infer" uses the path.
  • Round 2 (claims): PR body updated (inference path, tests, tested commit, 552 passed).
  • Validation: just lint passed; tests/pyathena/filesystem/ 552 passed against live S3 on 135e6487.

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): FINDINGS (test coverage only)

Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only, model reported gpt-6-astra), session 01a1024b-b4ed-7260-9acb-50ada0979845. Static review only; scope git range-diff fe21250a..84622707 fe21250a..135e6487 and git diff 84622707..135e6487; snapshot at 135e64872fad15cc949a5f659f4a71f6fec49cc4, unchanged afterwards. Reported scope deviation: the reviewer read local review guidance and searched the agent memory index outside the given roots (no matches); it changed nothing.

Earlier point (1): resolved across all three routes. Both compression callers now use the same _strip_protocol() normalization as open(); nontransaction aio delegates to sync. This covers bare paths, s3, s3a, and trailing slashes. Query suffixes remain in the inference input, matching open(); version-qualified writes are rejected.
Earlier point (2): addressed for sync. [...] These assert observable behavior.

P3 — Remaining test-coverage gap. test_s3_async.py:281 still tests only small bytes with explicit gzip. [...] Reverting just s3_async.py:203 would leave these compression tests passing while an aio transaction writing s3://bucket/key.gz/ with compression="infer" again stores uncompressed data at key.gz. Add corresponding aio cases.

Repaired in 6d90268d601bd27f631049e1d72eb0e23325f451 (tests only): aio test_pipe_file_compression gains the "s3://bucket/key.gz/" + "infer" case, in and outside a transaction. With only the aio call reverted to the unstripped path, the transaction case fails. Self-review: round 1, test-only change, no production code touched; round 2, PR body updated (aio test cases, tested commit). just lint passed; tests/pyathena/filesystem/ 554 passed against live S3 on 6d90268d.

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 of the test-only repair (relayed): CLEAN

Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only), session 01a10253-7f50-71c0-a134-43a7e597c557. Static review only; scope git diff 135e6487..6d90268d; snapshot at 6d90268d601bd27f631049e1d72eb0e23325f451, unchanged afterwards.

CLEAN. The trailing-slash inference case runs both inside and outside a transaction. Using the unstripped path in the aio helper would leave b"data" uncompressed, causing the transaction case's gzip.decompress assertion to fail. Assertions verify exactly one upload, no forwarded compression keyword, and an uploaded gzip payload that decompresses to b"data".

Comment thread pyathena/filesystem/s3.py Outdated
version, if the compression is not supported, or if the data
takes more than ``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
"""
compression = kwargs.pop("compression", None)

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): pre-existing issue (reviewer and scope: see the comment on pyathena/filesystem/s3.py:84)

P2 — Typed memoryviews can bypass multipart routing. s3.py:1828 compares len(value)—an item count—with byte thresholds. Outside transactions, a 6 GiB memoryview cast to "I", with block_size=5 GiB, passes the size guard but takes the single PutObject path, exceeding its size limit. This remains reachable with compression omitted, None, or "infer" resolving to no codec. Use the byte count for routing too.

Author: pre-existing and unchanged by this PR: the routing by len(value) predates it (kept in #1003), and a compressed value is bytes, so it is not affected. Deferred out of #1037's scope; to be reported to the maintainer.

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.

Fixed here at the maintainer's request, in 407ec3edcd953f420b05c877eccda1abd7b7444d: pipe_file() computes size = memoryview(value).nbytes once and uses it for both _check_multipart_upload_size() and the routing. The PR was returned to Draft for this change.

  • test_pipe_file_memoryview_routed_by_bytes: a cast("I") memoryview of block size + 4 bytes goes to the multipart path, and the uploaded parts join to the data. It fails with the len(value) routing.
  • test_pipe_file_non_contiguous_memoryview pinned the item count (4 items of 2 bytes, block_size=6 → PutObject). It now uses block_size=8, keeping its purpose: a non-contiguous memoryview within the block size is converted to bytes for PutObject.

Self-review of the change:

  • Round 1 (behavior): for bytes and bytearray, nbytes == len(), so their routing is unchanged; only multi-byte memoryviews change. A non-contiguous view's nbytes is the size of its tobytes(). On the buffered path, BytesIO.write() takes any buffer and returns the byte count, as the new test shows. The aio transaction path is always buffered and has no routing; aio outside a transaction delegates to this method.
  • Round 2 (claims): PR title, WHAT, release notes, and TEST updated; issue pipe_file() rejects compression for small values and uploads an empty object when a compressed write fails #1037 notes the addition. The 5 GiB single-PUT limit is the documented S3 PutObject maximum; the over-5-GiB case was not uploaded (offline test only, stated in the PR).
  • Validation: just lint passed; tests/pyathena/filesystem/ 555 passed against live S3 on 407ec3ed.

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 of 407ec3ed (relayed): CLEAN

Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only), session 01a10268-a4e0-7930-9efa-75d21c2765b4. Static review only; scope git diff 6d90268d..407ec3ed with the surrounding pipe_file(), _write_and_close(), _upload_chunk(), aio _pipe_file()/_pipe_file_in_transaction(), and the changed tests; snapshot at 407ec3edcd953f420b05c877eccda1abd7b7444d, unchanged afterwards.

nbytes handles bytes, bytearray, and memoryviews independently of format, shape, or contiguity. Compression precedes sizing. No remaining raw-value item-count routing was found in the local upload paths. [...] The adjusted test preserves non-contiguous PutObject conversion coverage with an eight-byte boundary. The new test would fail under the old condition and checks observable behavior: no PutObject, exact uploaded payload, and multipart completion. The changed comments are accurate.

CLEAN — the reported pre-existing routing defect is resolved. No regressions or additional pre-existing defects identified in scope.

S3File.write is inherited from fsspec, whose implementation is absent from this tree, so that dependency implementation could not be inspected.

Author: test_pipe_file_memoryview_routed_by_bytes runs the real fsspec write() (only the S3 requests are mocked), and the joined parts equal the data.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@laughingman7743
laughingman7743 marked this pull request as ready for review October 3, 2026 15:16
@laughingman7743
laughingman7743 marked this pull request as draft October 3, 2026 15:32
pipe_file() chose between PutObject and the buffered path by len(value),
which counts the items of a memoryview. A typed memoryview larger than
the block size in bytes, but not in items, went to PutObject, which
accepts at most 5 GiB. Compare the size in bytes, as the multipart size
check already does.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@laughingman7743 laughingman7743 changed the title Compress pipe_file() values up front for the compression argument Compress pipe_file() values up front and route them by their size in bytes Oct 3, 2026
@laughingman7743
laughingman7743 marked this pull request as ready for review October 3, 2026 15:40
@laughingman7743
laughingman7743 marked this pull request as draft October 3, 2026 15:53
_compress() took the path only to resolve "infer" and returned the value
unchanged when no codec was inferred. Resolve the codec in the callers
with get_compression(), as open() does, and compress with the public
CompressedBuffer.compress(), which takes a codec name and always returns
bytes. s3_async no longer imports a private name from s3.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Comment thread pyathena/filesystem/s3.py
pass

@classmethod
def compress(cls, value: bytes | bytearray | memoryview, compression: str) -> bytes:

@laughingman7743 laughingman7743 Oct 3, 2026 •

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 (rounds 1 and 2) of edc338db — public CompressedBuffer API: CLEAN

Scope: git diff 407ec3ed..edc338db (base fe21250a3e228e3d38e61898334c5baacc1285c2). The maintainer asked to redesign the private _compress() / _CompressedBuffer as a public API: CompressedBuffer.compress(value, compression) -> bytes takes a codec name only, and the callers resolve the codec as open() does. The PR was returned to Draft for this change.

Round 1 (behavior):

  • pipe_file() and aio _pipe_file_in_transaction() call get_compression(self._strip_protocol(path), kwargs.pop("compression", None)). None → None, and "infer" without a codec → None, so nothing is compressed, as before. An unknown name raises ValueError("Compression type … not supported") from get_compression() before any request, as before.
  • compress() validates the name itself for direct callers ("infer" is not a codec name and raises), converts a non-contiguous memoryview, writes through the codec into cls(), and always returns bytes.
  • s3_async imports CompressedBuffer instead of the private _compress.
  • While verifying the closing-codec test, a git checkout -- pyathena/filesystem/s3.py reverted the uncommitted s3.py part of this change; it was re-applied by the same scripted edit, the diff was re-read, and the suites were rerun before the commit.

Round 2 (claims): "public, not in the API reference like the other support classes": docs/api/filesystem.rst lists none of S3ClientError, S3Metadata, S3ObjectVersion. Docstrings match the code; PR body updated (WHAT, tests, tested commit).

Validation on edc338db: just lint passed; tests/pyathena/filesystem/ 561 passed against live S3; TestCompressedBuffer 7 cases, of which test_compress_codec_closing_its_file fails without the close() override; with zstandard and lz4 installed, all 7 registered codecs round-trip through CompressedBuffer.compress().

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Comment thread pyathena/filesystem/s3.py Outdated
"""
# The codec is taken as open() takes it, also from the extension of
# the path without the protocol for "infer".
compression = get_compression(self._strip_protocol(path), kwargs.pop("compression", None))

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 of the full PR at edc338db (relayed): FINDINGS — finding 1 of 2

Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only, model reported gpt-6-astra), session 01a1027c-e4b5-74e3-aed0-7a2ba6dafcc4. Static review only. Scope git diff fe21250a..edc338db (full PR, as the public CompressedBuffer API expanded the contract); snapshot at edc338db70f70ce38d2267953be9d898119c1c33, unchanged afterwards. The prompt held the diff, intent, and conventions, without the PR number, PR text, or prior findings, and limited reading to the snapshot and the fsspec source.

Covered: sync/aio pipe/pipe_file, transactions, create mode, failures, block sizes, multipart limits; bytes/bytearray/memoryview variants; None/"infer"/unknown codecs, fsspec codec adapters, sink-closing finalization; CompressedBuffer's public contract, conventions, docs, tests.

P2 — Compression-based routing can change the destination key. s3.py:1825. Outside a transaction, call pipe_file("s3://bucket/key.gz/", b"a" * (5 * 2**20 + 1), compression="gzip"). Previously, the uncompressed size selected open(), which strips the trailing slash and targets key.gz. Now compression reduces the size below the threshold, selecting PutObject with the original path and targeting key.gz/. [...] Both routes should use the same destination normalization. The new trailing-slash tests check decompression but omit destination-key assertions.

Pre-existing versus introduced: the underlying mismatch between raw-path parsing at s3.py:1833 and normalization through open() already existed for uncompressed writes. Finding 1 identifies the newly introduced destination change when compression switches an existing call between those routes.

Author verification: the pre-existing mismatch is confirmed on master fe21250a without compression: pipe_file("s3://bucket/key/", …) writes the key key/ for 10 bytes and key for block size + 1 bytes. A fix means choosing one destination for a trailing-slash path on both routes, which changes uncompressed writes too; that is a design decision for the maintainer, so this PR is kept in Draft until it is decided.

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.

Maintainer's decision: write a path with a trailing slash to the key without it on both paths, as open() does (s3fs was considered: its pipe_file() keeps the slash on both paths while its open() drops it, so it is inconsistent itself).

Repaired in 73cdb62537e951d8aa58a9baaf75cfe04599a9e1: pipe_file() starts with path = self._strip_protocol(path), and the routing, the codec of "infer", the single-request key, the mode="create" exists() check, and the cache invalidation all use that path. Docstring and the trailing-slash paragraph of docs/filesystem.md updated; release note added to the PR body.

  • test_pipe_file_trailing_slash asserts the key written for s3://bucket/dir/key/ (up to and above the block size, in and outside a transaction), and the trailing-slash cases of test_pipe_file_compression now assert the key (key.gz). Without the normalization, the 3 single-request cases outside a transaction fail, including the scenario above (compressible block size + 1 bytes).
  • Live check (ad hoc, objects deleted): pipe_file() of 1 byte and of 6 MiB, and put_file(), to paths ending in / wrote small, big, and put.

Self-review of the repair:

  • Round 1 (behavior): "s3://bucket" still raises "Cannot write to a bucket" and a ?versionId= path still raises (the query stays in the normalized path; test_pipe_file_invalid_path_raises passes). s3a:// is stripped too. Aio outside a transaction delegates to this method; aio in a transaction goes through open(), which already normalizes. Error messages now show the path without the protocol, like the errors raised from open().
  • Round 2 (claims): "touch() still writes the key as given" is from code reading (touch() parses the raw path); the docs sentence on pipe/put is backed by the live check.
  • Validation: just lint and just docs lint passed; tests/pyathena/filesystem/ 565 passed against live S3 on 73cdb625.

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 of 73cdb625 (relayed): FINDINGS (docs only)

Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only), session 01a10290-bde7-7370-9a5f-25f9a2291403. Static review only; scope git diff edc338db..73cdb625 and git range-diff fe21250a..edc338db fe21250a..73cdb625; snapshot at 73cdb62537e951d8aa58a9baaf75cfe04599a9e1, unchanged afterwards.

Both original review points are resolved: Normalization makes the upload destination independent of size and compression. Aio delegates to this implementation outside transactions and already normalizes through open() inside transactions. No implementation regression found in the traced paths. The close() docstring accurately describes the retained no-op behavior. The changed tests assert observable upload keys and decompressed payloads. [...]

P3 (low) — put() does not write the documented destination. docs/filesystem.md:112 newly states that put targeting dir/ writes object dir. For an existing local file, fs.put("local.txt", "s3://bucket/dir/") instead writes dir/local.txt. Inherited fsspec put() detects the directory destination before normalization and appends the source basename through other_paths; async _put() behaves likewise. Replace put with put_file here and distinguish put's directory-target behavior.

Repaired in 67a3db4c921a4843cb8f37a137be9deb630b4107 (docs only): "Opening it for writing, pipe, pipe_file, and put_file write the object dir, and put writes the file into the directory dir." Live check (ad hoc, objects deleted): put() to dir/ wrote dir/local.txt; pipe() to piped/ wrote piped. Also verified while fixing: fsspec pipe() strips the path before calling pipe_file(), so the release note now limits the trailing-slash change to direct pipe_file() calls. Self-review: docs/PR text only, no behavior change; just docs lint passed.

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 of the docs repair 67a3db4c (relayed): FINDINGS (docs only)

Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only), session 01a10296-b8f5-7321-a823-9216e8e2a9e3. Static review only; scope git diff 73cdb625..67a3db4c, checked against open/_open, pipe_file, put_file, fsspec pipe/put/open, and s3_async.py; snapshot at 67a3db4c921a4843cb8f37a137be9deb630b4107, unchanged afterwards.

docs/filesystem.md:112: "put writes the file into the directory dir" needs qualification to a string destination. With a readable local file.txt, fs.put(["file.txt"], ["s3://bucket/dir/"]) writes object dir [...] Both fsspec implementations bypass directory expansion for paired lists (spec.py:1119, asyn.py:688), then PyAthena's put_file strips the slash through open. [...]

The remaining paragraph statements match the inspected sync and async paths.

Repaired in c6297d35b3552af208fb931ac8a60201178f47d2 by removing the put statement: the paragraph now says only "Opening it for writing, pipe, pipe_file, and put_file write the object dir", the statements that this review found accurate. fsspec's put() resolves directory destinations itself and is not described. No further independent pass was run for this deletion, as it leaves only the reviewed statements. just docs lint passed.

Comment thread pyathena/filesystem/s3.py
"""

@override
def close(self) -> None:

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 of the full PR at edc338db (relayed): finding 2 of 2 (reviewer and scope: see the comment on s3.py:1816)

P3 — The public close() override lacks its changed-behavior docstring. CompressedBuffer.close() leaves the buffer open: closed remains false and subsequent writes remain possible. Method-level help can inherit the incompatible BytesIO.close() description. docs/contributing.md:85 requires a docstring when an override changes behavior. Document the no-op semantics directly on this method.

Repaired in d8ec2207678815b2dee85386d379edd523fd804d: """Do nothing, so that the data can still be read after a codec closes the buffer.""". Self-review: docstring-only change; round 1 no behavior change, round 2 the sentence matches the code. just lint passed; TestCompressedBuffer 7 passed.

The single PutObject request wrote a path with a trailing slash to the
key with the slash, while the buffered path, through open(), wrote it to
the key without it, so the key depended on the size of the value. With
the compression argument, compressing could now switch a call between
the two. Normalize the path as open() does before choosing, so both
paths write the same key.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@laughingman7743 laughingman7743 changed the title Compress pipe_file() values up front and route them by their size in bytes Compress pipe_file() values up front, and route and key them consistently Oct 3, 2026
laughingman7743 and others added 2 commits October 4, 2026 01:26
fsspec's put() treats a destination ending in a slash as a directory and
writes the file under it, unlike put_file().

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
fsspec's put() resolves a destination ending in a slash itself, as a
directory for a single source but not for paired lists, so the paragraph
no longer describes it.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@laughingman7743
laughingman7743 marked this pull request as ready for review October 3, 2026 16:34
@laughingman7743
laughingman7743 merged commit abc99b0 into master Oct 3, 2026
14 checks passed
@laughingman7743
laughingman7743 deleted the fix/1037-pipe-file-compression branch October 3, 2026 16:57
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

pipe_file() rejects compression for small values and uploads an empty object when a compressed write fails

1 participant