zeroshade commented on code in PR #1620:
URL: https://github.com/apache/iceberg-go/pull/1620#discussion_r4158511963


##########
encryption/go_envelope_manager.go:
##########
@@ -0,0 +1,889 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package encryption
+
+import (
+       "context"
+       "crypto/aes"
+       "crypto/cipher"
+       "crypto/rand"
+       "encoding/binary"
+       "encoding/json"
+       "errors"
+       "fmt"
+       "io"
+       "io/fs"
+       "math"
+       "sync"
+
+       icebergio "github.com/apache/iceberg-go/io"
+)
+
+// Defaults for [GoEnvelopeEncryptionManager].
+const (
+       // GoEnvelopeDefaultDEKLength is the default length, in bytes, of the
+       // per-file data encryption key (DEK) generated for AES-256-GCM.
+       GoEnvelopeDefaultDEKLength = 32
+
+       // GoEnvelopeDefaultBlockSize is the default plaintext block size, in
+       // bytes, used to split a file into independently authenticated AES-GCM
+       // blocks. Blocks allow random access (Seek/ReadAt) without buffering or
+       // decrypting the whole file.
+       //
+       // This is not Java's Ciphers.PLAIN_BLOCK_SIZE (1 MiB, a compile-time
+       // constant there); see the [GoEnvelopeEncryptionManager] doc comment.
+       GoEnvelopeDefaultBlockSize = 64 * 1024

Review Comment:
   Make the block size a fixed 1 MiB, matching Java's 
`Ciphers.PLAIN_BLOCK_SIZE` and iceberg-rust's `PLAIN_BLOCK_SIZE`. Java's 
`AesGcmInputStream.validateHeader` rejects any other header value, so a 64 KiB 
stream can't be read there regardless of key metadata.
   
   My call on the knob:
   - Drop `WithBlockSize`, `GoEnvelopeMaxBlockSize` and `ErrInvalidBlockSize`.
   - The writer always emits 1 MiB; the reader requires the header to say 1 MiB 
and returns `ErrInvalidStreamHeader` otherwise.
   
   That also removes the attacker-sized-allocation surface we went back and 
forth on. The tests built on `WithBlockSize(16)` (:69, :81, :91, ...) can use 
~2.5 MiB payloads, as iceberg-rust's suite does.
   
   @tanmayrauth this supersedes the block-size validation from your 07-31 
thread: the knob goes away instead.
   



##########
encryption/go_envelope_manager.go:
##########
@@ -0,0 +1,889 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package encryption
+
+import (
+       "context"
+       "crypto/aes"
+       "crypto/cipher"
+       "crypto/rand"
+       "encoding/binary"
+       "encoding/json"
+       "errors"
+       "fmt"
+       "io"
+       "io/fs"
+       "math"
+       "sync"
+
+       icebergio "github.com/apache/iceberg-go/io"
+)
+
+// Defaults for [GoEnvelopeEncryptionManager].
+const (
+       // GoEnvelopeDefaultDEKLength is the default length, in bytes, of the
+       // per-file data encryption key (DEK) generated for AES-256-GCM.
+       GoEnvelopeDefaultDEKLength = 32
+
+       // GoEnvelopeDefaultBlockSize is the default plaintext block size, in
+       // bytes, used to split a file into independently authenticated AES-GCM
+       // blocks. Blocks allow random access (Seek/ReadAt) without buffering or
+       // decrypting the whole file.
+       //
+       // This is not Java's Ciphers.PLAIN_BLOCK_SIZE (1 MiB, a compile-time
+       // constant there); see the [GoEnvelopeEncryptionManager] doc comment.
+       GoEnvelopeDefaultBlockSize = 64 * 1024
+
+       // GoEnvelopeMaxBlockSize is the largest plaintext block size accepted
+       // from the AES GCM Stream header on read, and the largest that
+       // [WithBlockSize] will configure for writing. The header's block-length
+       // field is unauthenticated, untrusted data (the Iceberg AES GCM Stream
+       // spec's "File length" note applies equally here); without a ceiling, a
+       // crafted header claiming e.g. 1<<40 bytes would force a huge
+       // allocation (see readBlock) before any authentication check can run.
+       GoEnvelopeMaxBlockSize = 128 * 1024 * 1024
+)
+
+// Sentinel errors returned by [GoEnvelopeEncryptionManager].
+var (
+       // ErrKeyIDRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile] when keyID is
+       // empty. GoEnvelopeEncryptionManager always encrypts, so it requires a
+       // KEK to wrap the generated DEK; use [PlaintextEncryptionManager] for
+       // unencrypted tables instead of passing an empty keyID here.
+       ErrKeyIDRequired = errors.New("encryption: GoEnvelopeEncryptionManager 
requires a non-empty keyID")
+
+       // ErrKeyMetadataRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when keyMetadata
+       // is empty. GoEnvelopeEncryptionManager always decrypts, so it requires
+       // the per-file key metadata produced by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile].
+       ErrKeyMetadataRequired = errors.New("encryption: 
GoEnvelopeEncryptionManager requires non-empty key metadata")
+
+       // ErrUnsupportedKeyMetadataVersion is returned when key metadata was
+       // produced by a newer, incompatible encoding version.
+       ErrUnsupportedKeyMetadataVersion = errors.New("encryption: unsupported 
key metadata version")
+
+       // ErrInvalidBlockSize is returned when a configured block size is not
+       // positive or exceeds [GoEnvelopeMaxBlockSize].
+       ErrInvalidBlockSize = errors.New("encryption: block size must be 
positive and at most GoEnvelopeMaxBlockSize")
+
+       // ErrInvalidStreamHeader is returned when the AES GCM Stream header
+       // (the "AGS1" magic and little-endian block-length fields at the start
+       // of an encrypted file, per the Iceberg AES GCM Stream spec) is
+       // missing, truncated, or specifies a block length outside the
+       // supported range.
+       ErrInvalidStreamHeader = errors.New("encryption: invalid AES GCM Stream 
header")
+
+       // ErrInvalidKeyMetadata is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when decoded key
+       // metadata fails basic sanity checks (e.g. a negative plaintext length
+       // or a missing AAD prefix). Key metadata is untrusted input on a crypto
+       // read path, so it is validated rather than trusted blindly.
+       ErrInvalidKeyMetadata = errors.New("encryption: invalid key metadata")
+
+       // ErrOutputFileClosed is returned by [goEnvelopeOutputFile.Write] when
+       // called after Close, or after a previous flush has poisoned the 
writer.
+       // It wraps [fs.ErrClosed] so callers can test with errors.Is(err, 
fs.ErrClosed).
+       ErrOutputFileClosed = fmt.Errorf("encryption: write to closed 
GoEnvelopeEncryptionManager output file: %w", fs.ErrClosed)
+
+       // ErrBlockTruncated is returned by [goEnvelopeInputFile.ReadAt] when 
the
+       // underlying storage returns fewer ciphertext bytes for a block than 
its
+       // recorded length requires, indicating the file was truncated at rest.
+       // This is distinct from [ErrCiphertextTooShort], which 
[KeyManagementClient]
+       // implementations use for a too-short wrapped key or KMS-encrypted
+       // payload; keeping them separate lets a caller tell a malformed KMS 
blob
+       // apart from a short block read.
+       ErrBlockTruncated = errors.New("encryption: block truncated: read fewer 
ciphertext bytes than expected")
+)
+
+// Constants describing the Iceberg AES GCM Stream ("AGS1") wire format used
+// for the ciphertext produced by [GoEnvelopeEncryptionManager]. See
+// https://iceberg.apache.org/gcm-stream-spec/ for the full specification.
+const (
+       // gcmStreamMagic identifies an AES GCM Stream version 1 file.
+       gcmStreamMagic = "AGS1"
+
+       // gcmStreamHeaderLength is the length, in bytes, of the magic plus the
+       // little-endian block-length field written at the start of every file.
+       gcmStreamHeaderLength = len(gcmStreamMagic) + 4
+
+       // gcmStreamNonceLength is the length, in bytes, of the random AES-GCM
+       // nonce stored at the start of every cipher block.
+       gcmStreamNonceLength = 12
+
+       // gcmStreamTagLength is the length, in bytes, of the AES-GCM
+       // authentication tag appended to every cipher block's ciphertext.
+       gcmStreamTagLength = 16
+
+       // gcmStreamBlockOverhead is the number of ciphertext bytes added to
+       // each block beyond its plaintext length (nonce + tag).
+       gcmStreamBlockOverhead = gcmStreamNonceLength + gcmStreamTagLength
+
+       // gcmStreamAADPrefixLength is the length, in bytes, of the random
+       // per-file AAD prefix generated for new output files.
+       gcmStreamAADPrefixLength = 16
+)
+
+// goEnvelopeKeyMetadataVersion is the current encoding version written by
+// [GoEnvelopeEncryptionManager]. It is bumped whenever the on-disk layout of
+// goEnvelopeKeyMetadata changes incompatibly.
+const goEnvelopeKeyMetadataVersion = 1
+
+// goEnvelopeKeyMetadata is the JSON-encoded structure stored as the opaque
+// [EncryptionKeyMetadata] for files produced by [GoEnvelopeEncryptionManager].
+//
+// This encoding is Go-specific, not Java's Avro-encoded StandardKeyMetadata;
+// see the [GoEnvelopeEncryptionManager] doc comment. Per the table spec,
+// DataFile/ManifestFile key_metadata is explicitly "implementation-specific";
+// what must be Iceberg AES GCM Stream compliant - and is - is the wire
+// format of the encrypted byte stream itself (magic, block framing, nonce
+// placement, and AAD; see the constants above).
+type goEnvelopeKeyMetadata struct {

Review Comment:
   Replace the JSON blob with Java's `StandardKeyMetadata` encoding: a version 
byte `0x01` followed by an Avro binary record `{encryption_key: bytes, 
aad_prefix: bytes, file_length: long}`.
   
   - In Java's schema `aad_prefix` and `file_length` are optional `["null", T]` 
unions, so a present value is written as union branch 1 followed by the value. 
`github.com/twmb/avro` is already a dependency, or you can hand-encode the 
three fields.
   - `encryption_key` is the raw DEK. No KMS-wrapped key and no key-id go in 
the per-file metadata; see the comment on :280.
   - `file_length` is the encrypted length. `Close` knows it, so always write 
it. On read, derive the plaintext length from it the way Java's 
`AesGcmInputStream.calculatePlaintextLength` does, instead of storing 
`plaintext-length`. Java treats the field as optional; when it's absent, the 
trusted length has to come from the caller (the manifest entry). Don't fall 
back to `Stat().Size()`, per the AGS1 "File length" rule.
   
   Those bytes then match pyiceberg's and iceberg-rust's `StandardKeyMetadata`. 
Bytes produced by Java's `StandardKeyMetadata.buffer()` make the regression 
fixture.
   



##########
encryption/go_envelope_manager.go:
##########
@@ -0,0 +1,889 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package encryption
+
+import (
+       "context"
+       "crypto/aes"
+       "crypto/cipher"
+       "crypto/rand"
+       "encoding/binary"
+       "encoding/json"
+       "errors"
+       "fmt"
+       "io"
+       "io/fs"
+       "math"
+       "sync"
+
+       icebergio "github.com/apache/iceberg-go/io"
+)
+
+// Defaults for [GoEnvelopeEncryptionManager].
+const (
+       // GoEnvelopeDefaultDEKLength is the default length, in bytes, of the
+       // per-file data encryption key (DEK) generated for AES-256-GCM.
+       GoEnvelopeDefaultDEKLength = 32
+
+       // GoEnvelopeDefaultBlockSize is the default plaintext block size, in
+       // bytes, used to split a file into independently authenticated AES-GCM
+       // blocks. Blocks allow random access (Seek/ReadAt) without buffering or
+       // decrypting the whole file.
+       //
+       // This is not Java's Ciphers.PLAIN_BLOCK_SIZE (1 MiB, a compile-time
+       // constant there); see the [GoEnvelopeEncryptionManager] doc comment.
+       GoEnvelopeDefaultBlockSize = 64 * 1024
+
+       // GoEnvelopeMaxBlockSize is the largest plaintext block size accepted
+       // from the AES GCM Stream header on read, and the largest that
+       // [WithBlockSize] will configure for writing. The header's block-length
+       // field is unauthenticated, untrusted data (the Iceberg AES GCM Stream
+       // spec's "File length" note applies equally here); without a ceiling, a
+       // crafted header claiming e.g. 1<<40 bytes would force a huge
+       // allocation (see readBlock) before any authentication check can run.
+       GoEnvelopeMaxBlockSize = 128 * 1024 * 1024
+)
+
+// Sentinel errors returned by [GoEnvelopeEncryptionManager].
+var (
+       // ErrKeyIDRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile] when keyID is
+       // empty. GoEnvelopeEncryptionManager always encrypts, so it requires a
+       // KEK to wrap the generated DEK; use [PlaintextEncryptionManager] for
+       // unencrypted tables instead of passing an empty keyID here.
+       ErrKeyIDRequired = errors.New("encryption: GoEnvelopeEncryptionManager 
requires a non-empty keyID")
+
+       // ErrKeyMetadataRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when keyMetadata
+       // is empty. GoEnvelopeEncryptionManager always decrypts, so it requires
+       // the per-file key metadata produced by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile].
+       ErrKeyMetadataRequired = errors.New("encryption: 
GoEnvelopeEncryptionManager requires non-empty key metadata")
+
+       // ErrUnsupportedKeyMetadataVersion is returned when key metadata was
+       // produced by a newer, incompatible encoding version.
+       ErrUnsupportedKeyMetadataVersion = errors.New("encryption: unsupported 
key metadata version")
+
+       // ErrInvalidBlockSize is returned when a configured block size is not
+       // positive or exceeds [GoEnvelopeMaxBlockSize].
+       ErrInvalidBlockSize = errors.New("encryption: block size must be 
positive and at most GoEnvelopeMaxBlockSize")
+
+       // ErrInvalidStreamHeader is returned when the AES GCM Stream header
+       // (the "AGS1" magic and little-endian block-length fields at the start
+       // of an encrypted file, per the Iceberg AES GCM Stream spec) is
+       // missing, truncated, or specifies a block length outside the
+       // supported range.
+       ErrInvalidStreamHeader = errors.New("encryption: invalid AES GCM Stream 
header")
+
+       // ErrInvalidKeyMetadata is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when decoded key
+       // metadata fails basic sanity checks (e.g. a negative plaintext length
+       // or a missing AAD prefix). Key metadata is untrusted input on a crypto
+       // read path, so it is validated rather than trusted blindly.
+       ErrInvalidKeyMetadata = errors.New("encryption: invalid key metadata")
+
+       // ErrOutputFileClosed is returned by [goEnvelopeOutputFile.Write] when
+       // called after Close, or after a previous flush has poisoned the 
writer.
+       // It wraps [fs.ErrClosed] so callers can test with errors.Is(err, 
fs.ErrClosed).
+       ErrOutputFileClosed = fmt.Errorf("encryption: write to closed 
GoEnvelopeEncryptionManager output file: %w", fs.ErrClosed)

Review Comment:
   The doc says this error is returned "after a previous flush has poisoned the 
writer", but `Write` returns the sticky `f.err` before it ever checks `closed` 
(:513-515). After a failed flush, `errors.Is(err, ErrOutputFileClosed)` and 
`errors.Is(err, fs.ErrClosed)` are both false. Drop that clause.
   



##########
encryption/go_envelope_manager.go:
##########
@@ -0,0 +1,889 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package encryption
+
+import (
+       "context"
+       "crypto/aes"
+       "crypto/cipher"
+       "crypto/rand"
+       "encoding/binary"
+       "encoding/json"
+       "errors"
+       "fmt"
+       "io"
+       "io/fs"
+       "math"
+       "sync"
+
+       icebergio "github.com/apache/iceberg-go/io"
+)
+
+// Defaults for [GoEnvelopeEncryptionManager].
+const (
+       // GoEnvelopeDefaultDEKLength is the default length, in bytes, of the
+       // per-file data encryption key (DEK) generated for AES-256-GCM.
+       GoEnvelopeDefaultDEKLength = 32
+
+       // GoEnvelopeDefaultBlockSize is the default plaintext block size, in
+       // bytes, used to split a file into independently authenticated AES-GCM
+       // blocks. Blocks allow random access (Seek/ReadAt) without buffering or
+       // decrypting the whole file.
+       //
+       // This is not Java's Ciphers.PLAIN_BLOCK_SIZE (1 MiB, a compile-time
+       // constant there); see the [GoEnvelopeEncryptionManager] doc comment.
+       GoEnvelopeDefaultBlockSize = 64 * 1024
+
+       // GoEnvelopeMaxBlockSize is the largest plaintext block size accepted
+       // from the AES GCM Stream header on read, and the largest that
+       // [WithBlockSize] will configure for writing. The header's block-length
+       // field is unauthenticated, untrusted data (the Iceberg AES GCM Stream
+       // spec's "File length" note applies equally here); without a ceiling, a
+       // crafted header claiming e.g. 1<<40 bytes would force a huge
+       // allocation (see readBlock) before any authentication check can run.
+       GoEnvelopeMaxBlockSize = 128 * 1024 * 1024
+)
+
+// Sentinel errors returned by [GoEnvelopeEncryptionManager].
+var (
+       // ErrKeyIDRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile] when keyID is
+       // empty. GoEnvelopeEncryptionManager always encrypts, so it requires a
+       // KEK to wrap the generated DEK; use [PlaintextEncryptionManager] for
+       // unencrypted tables instead of passing an empty keyID here.
+       ErrKeyIDRequired = errors.New("encryption: GoEnvelopeEncryptionManager 
requires a non-empty keyID")
+
+       // ErrKeyMetadataRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when keyMetadata
+       // is empty. GoEnvelopeEncryptionManager always decrypts, so it requires
+       // the per-file key metadata produced by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile].
+       ErrKeyMetadataRequired = errors.New("encryption: 
GoEnvelopeEncryptionManager requires non-empty key metadata")
+
+       // ErrUnsupportedKeyMetadataVersion is returned when key metadata was
+       // produced by a newer, incompatible encoding version.
+       ErrUnsupportedKeyMetadataVersion = errors.New("encryption: unsupported 
key metadata version")
+
+       // ErrInvalidBlockSize is returned when a configured block size is not
+       // positive or exceeds [GoEnvelopeMaxBlockSize].
+       ErrInvalidBlockSize = errors.New("encryption: block size must be 
positive and at most GoEnvelopeMaxBlockSize")
+
+       // ErrInvalidStreamHeader is returned when the AES GCM Stream header
+       // (the "AGS1" magic and little-endian block-length fields at the start
+       // of an encrypted file, per the Iceberg AES GCM Stream spec) is
+       // missing, truncated, or specifies a block length outside the
+       // supported range.
+       ErrInvalidStreamHeader = errors.New("encryption: invalid AES GCM Stream 
header")
+
+       // ErrInvalidKeyMetadata is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when decoded key
+       // metadata fails basic sanity checks (e.g. a negative plaintext length
+       // or a missing AAD prefix). Key metadata is untrusted input on a crypto
+       // read path, so it is validated rather than trusted blindly.
+       ErrInvalidKeyMetadata = errors.New("encryption: invalid key metadata")
+
+       // ErrOutputFileClosed is returned by [goEnvelopeOutputFile.Write] when
+       // called after Close, or after a previous flush has poisoned the 
writer.
+       // It wraps [fs.ErrClosed] so callers can test with errors.Is(err, 
fs.ErrClosed).
+       ErrOutputFileClosed = fmt.Errorf("encryption: write to closed 
GoEnvelopeEncryptionManager output file: %w", fs.ErrClosed)
+
+       // ErrBlockTruncated is returned by [goEnvelopeInputFile.ReadAt] when 
the
+       // underlying storage returns fewer ciphertext bytes for a block than 
its
+       // recorded length requires, indicating the file was truncated at rest.
+       // This is distinct from [ErrCiphertextTooShort], which 
[KeyManagementClient]
+       // implementations use for a too-short wrapped key or KMS-encrypted
+       // payload; keeping them separate lets a caller tell a malformed KMS 
blob
+       // apart from a short block read.
+       ErrBlockTruncated = errors.New("encryption: block truncated: read fewer 
ciphertext bytes than expected")
+)
+
+// Constants describing the Iceberg AES GCM Stream ("AGS1") wire format used
+// for the ciphertext produced by [GoEnvelopeEncryptionManager]. See
+// https://iceberg.apache.org/gcm-stream-spec/ for the full specification.
+const (
+       // gcmStreamMagic identifies an AES GCM Stream version 1 file.
+       gcmStreamMagic = "AGS1"
+
+       // gcmStreamHeaderLength is the length, in bytes, of the magic plus the
+       // little-endian block-length field written at the start of every file.
+       gcmStreamHeaderLength = len(gcmStreamMagic) + 4
+
+       // gcmStreamNonceLength is the length, in bytes, of the random AES-GCM
+       // nonce stored at the start of every cipher block.
+       gcmStreamNonceLength = 12
+
+       // gcmStreamTagLength is the length, in bytes, of the AES-GCM
+       // authentication tag appended to every cipher block's ciphertext.
+       gcmStreamTagLength = 16
+
+       // gcmStreamBlockOverhead is the number of ciphertext bytes added to
+       // each block beyond its plaintext length (nonce + tag).
+       gcmStreamBlockOverhead = gcmStreamNonceLength + gcmStreamTagLength
+
+       // gcmStreamAADPrefixLength is the length, in bytes, of the random
+       // per-file AAD prefix generated for new output files.
+       gcmStreamAADPrefixLength = 16
+)
+
+// goEnvelopeKeyMetadataVersion is the current encoding version written by
+// [GoEnvelopeEncryptionManager]. It is bumped whenever the on-disk layout of
+// goEnvelopeKeyMetadata changes incompatibly.
+const goEnvelopeKeyMetadataVersion = 1
+
+// goEnvelopeKeyMetadata is the JSON-encoded structure stored as the opaque
+// [EncryptionKeyMetadata] for files produced by [GoEnvelopeEncryptionManager].
+//
+// This encoding is Go-specific, not Java's Avro-encoded StandardKeyMetadata;
+// see the [GoEnvelopeEncryptionManager] doc comment. Per the table spec,
+// DataFile/ManifestFile key_metadata is explicitly "implementation-specific";
+// what must be Iceberg AES GCM Stream compliant - and is - is the wire
+// format of the encrypted byte stream itself (magic, block framing, nonce
+// placement, and AAD; see the constants above).
+type goEnvelopeKeyMetadata struct {
+       Version    int    `json:"v"`
+       KeyID      string `json:"key-id"`
+       WrappedKey []byte `json:"wrapped-key"`
+
+       // BlockSize is the trusted plaintext block size the file was written
+       // with. It is validated bounded before any file I/O, and is what's
+       // actually used to size reads; the stream header's block-length field
+       // is untrusted (unauthenticated, attacker-influenced storage bytes) and
+       // is only compared against this value, never used directly to size an
+       // allocation.
+       BlockSize int64 `json:"block-size"`
+
+       // AADPrefix is combined with each block's little-endian index to form
+       // the AES GCM Stream additional authenticated data, binding every
+       // ciphertext block to this file and to its position so that blocks
+       // cannot be silently reordered, replayed from another file, or spliced
+       // in from elsewhere in the same file. It is not secret.
+       AADPrefix []byte `json:"aad-prefix"`
+
+       // PlaintextLength is the trusted total plaintext size used to compute
+       // the block count on read. Per the AES GCM Stream spec's "File length"
+       // note, a reader must use a length from a trusted source rather than
+       // the underlying storage's reported size, since storage size alone
+       // cannot distinguish a genuinely short file from one truncated by an
+       // attacker who does not also control this metadata.
+       PlaintextLength int64 `json:"plaintext-length"`
+}
+
+// GoEnvelopeEncryptionManager is a generic, format-agnostic 
[EncryptionManager]
+// that provides envelope encryption for arbitrary files (e.g. manifests,
+// manifest lists, Puffin statistics) using a [KeyManagementClient] to wrap
+// and unwrap a fresh AES-256-GCM data encryption key (DEK) per file.
+//
+// Each file is split into fixed-size plaintext blocks and written using the
+// Iceberg AES GCM Stream ("AGS1") format: a magic/block-length header
+// followed by independently authenticated blocks, each carrying its own
+// random nonce and authenticated with an AAD that binds it to the file and
+// to its position. This bounds memory usage and supports random access
+// (Seek/ReadAt) on the decrypted file without buffering or decrypting more
+// than the requested blocks.
+//
+// GoEnvelopeEncryptionManager always encrypts and always decrypts: it fails
+// closed, returning [ErrKeyIDRequired] or [ErrKeyMetadataRequired] rather
+// than silently falling back to plaintext. Use [PlaintextEncryptionManager]
+// for tables or files that are not encrypted.
+//
+// Not interoperable with Java's StandardEncryptionManager: the AGS1 byte
+// stream framing (magic, block layout, nonce placement, AAD) matches the
+// spec, but the two implementations are not cross-readable. Java's reader
+// hard-requires a 1 MiB header block size (a compile-time constant there),
+// while [GoEnvelopeDefaultBlockSize] is 64 KiB, and Java decodes key
+// metadata as a 1-byte version plus Avro-encoded StandardKeyMetadata, while
+// this type encodes key metadata as JSON (goEnvelopeKeyMetadata). Files
+// written by one are not readable by the other.
+type GoEnvelopeEncryptionManager struct {
+       kms       KeyManagementClient
+       dekLength int
+       blockSize int
+}
+
+var _ EncryptionManager = (*GoEnvelopeEncryptionManager)(nil)
+
+// GoEnvelopeManagerOption configures a [GoEnvelopeEncryptionManager] created
+// by [NewGoEnvelopeEncryptionManager].
+type GoEnvelopeManagerOption func(*GoEnvelopeEncryptionManager)
+
+// WithDEKLength overrides the default data encryption key length (in bytes).
+// Valid AES key lengths are 16, 24, or 32 bytes.
+func WithDEKLength(length int) GoEnvelopeManagerOption {
+       return func(m *GoEnvelopeEncryptionManager) { m.dekLength = length }
+}

Review Comment:
   `WithDEKLength` takes a generic name in package `encryption` but only 
configures this type. Once this is the Java-compatible manager, and the KEK 
layer adds more options, name it for what it maps to, the 
`encryption.data-key-length` table property. Either share one option type 
across the package's managers or use a type-specific name. Renaming after a 
release breaks callers.
   
   Two related points:
   - Java's default for that property is 16 bytes 
(`ENCRYPTION_DEK_LENGTH_DEFAULT`), and iceberg-rust also defaults to 128-bit. 
Match it unless there's a reason not to.
   - Once the format matches, the `StandardEncryptionManager` name is accurate 
again, so renaming back is fine.
   



##########
encryption/go_envelope_manager.go:
##########
@@ -0,0 +1,889 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package encryption
+
+import (
+       "context"
+       "crypto/aes"
+       "crypto/cipher"
+       "crypto/rand"
+       "encoding/binary"
+       "encoding/json"
+       "errors"
+       "fmt"
+       "io"
+       "io/fs"
+       "math"
+       "sync"
+
+       icebergio "github.com/apache/iceberg-go/io"
+)
+
+// Defaults for [GoEnvelopeEncryptionManager].
+const (
+       // GoEnvelopeDefaultDEKLength is the default length, in bytes, of the
+       // per-file data encryption key (DEK) generated for AES-256-GCM.
+       GoEnvelopeDefaultDEKLength = 32
+
+       // GoEnvelopeDefaultBlockSize is the default plaintext block size, in
+       // bytes, used to split a file into independently authenticated AES-GCM
+       // blocks. Blocks allow random access (Seek/ReadAt) without buffering or
+       // decrypting the whole file.
+       //
+       // This is not Java's Ciphers.PLAIN_BLOCK_SIZE (1 MiB, a compile-time
+       // constant there); see the [GoEnvelopeEncryptionManager] doc comment.
+       GoEnvelopeDefaultBlockSize = 64 * 1024
+
+       // GoEnvelopeMaxBlockSize is the largest plaintext block size accepted
+       // from the AES GCM Stream header on read, and the largest that
+       // [WithBlockSize] will configure for writing. The header's block-length
+       // field is unauthenticated, untrusted data (the Iceberg AES GCM Stream
+       // spec's "File length" note applies equally here); without a ceiling, a
+       // crafted header claiming e.g. 1<<40 bytes would force a huge
+       // allocation (see readBlock) before any authentication check can run.
+       GoEnvelopeMaxBlockSize = 128 * 1024 * 1024
+)
+
+// Sentinel errors returned by [GoEnvelopeEncryptionManager].
+var (
+       // ErrKeyIDRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile] when keyID is
+       // empty. GoEnvelopeEncryptionManager always encrypts, so it requires a
+       // KEK to wrap the generated DEK; use [PlaintextEncryptionManager] for
+       // unencrypted tables instead of passing an empty keyID here.
+       ErrKeyIDRequired = errors.New("encryption: GoEnvelopeEncryptionManager 
requires a non-empty keyID")
+
+       // ErrKeyMetadataRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when keyMetadata
+       // is empty. GoEnvelopeEncryptionManager always decrypts, so it requires
+       // the per-file key metadata produced by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile].
+       ErrKeyMetadataRequired = errors.New("encryption: 
GoEnvelopeEncryptionManager requires non-empty key metadata")
+
+       // ErrUnsupportedKeyMetadataVersion is returned when key metadata was
+       // produced by a newer, incompatible encoding version.
+       ErrUnsupportedKeyMetadataVersion = errors.New("encryption: unsupported 
key metadata version")
+
+       // ErrInvalidBlockSize is returned when a configured block size is not
+       // positive or exceeds [GoEnvelopeMaxBlockSize].
+       ErrInvalidBlockSize = errors.New("encryption: block size must be 
positive and at most GoEnvelopeMaxBlockSize")
+
+       // ErrInvalidStreamHeader is returned when the AES GCM Stream header
+       // (the "AGS1" magic and little-endian block-length fields at the start
+       // of an encrypted file, per the Iceberg AES GCM Stream spec) is
+       // missing, truncated, or specifies a block length outside the
+       // supported range.
+       ErrInvalidStreamHeader = errors.New("encryption: invalid AES GCM Stream 
header")
+
+       // ErrInvalidKeyMetadata is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when decoded key
+       // metadata fails basic sanity checks (e.g. a negative plaintext length
+       // or a missing AAD prefix). Key metadata is untrusted input on a crypto
+       // read path, so it is validated rather than trusted blindly.
+       ErrInvalidKeyMetadata = errors.New("encryption: invalid key metadata")
+
+       // ErrOutputFileClosed is returned by [goEnvelopeOutputFile.Write] when
+       // called after Close, or after a previous flush has poisoned the 
writer.
+       // It wraps [fs.ErrClosed] so callers can test with errors.Is(err, 
fs.ErrClosed).
+       ErrOutputFileClosed = fmt.Errorf("encryption: write to closed 
GoEnvelopeEncryptionManager output file: %w", fs.ErrClosed)
+
+       // ErrBlockTruncated is returned by [goEnvelopeInputFile.ReadAt] when 
the
+       // underlying storage returns fewer ciphertext bytes for a block than 
its
+       // recorded length requires, indicating the file was truncated at rest.
+       // This is distinct from [ErrCiphertextTooShort], which 
[KeyManagementClient]
+       // implementations use for a too-short wrapped key or KMS-encrypted
+       // payload; keeping them separate lets a caller tell a malformed KMS 
blob
+       // apart from a short block read.
+       ErrBlockTruncated = errors.New("encryption: block truncated: read fewer 
ciphertext bytes than expected")
+)
+
+// Constants describing the Iceberg AES GCM Stream ("AGS1") wire format used
+// for the ciphertext produced by [GoEnvelopeEncryptionManager]. See
+// https://iceberg.apache.org/gcm-stream-spec/ for the full specification.
+const (
+       // gcmStreamMagic identifies an AES GCM Stream version 1 file.
+       gcmStreamMagic = "AGS1"
+
+       // gcmStreamHeaderLength is the length, in bytes, of the magic plus the
+       // little-endian block-length field written at the start of every file.
+       gcmStreamHeaderLength = len(gcmStreamMagic) + 4
+
+       // gcmStreamNonceLength is the length, in bytes, of the random AES-GCM
+       // nonce stored at the start of every cipher block.
+       gcmStreamNonceLength = 12
+
+       // gcmStreamTagLength is the length, in bytes, of the AES-GCM
+       // authentication tag appended to every cipher block's ciphertext.
+       gcmStreamTagLength = 16
+
+       // gcmStreamBlockOverhead is the number of ciphertext bytes added to
+       // each block beyond its plaintext length (nonce + tag).
+       gcmStreamBlockOverhead = gcmStreamNonceLength + gcmStreamTagLength
+
+       // gcmStreamAADPrefixLength is the length, in bytes, of the random
+       // per-file AAD prefix generated for new output files.
+       gcmStreamAADPrefixLength = 16
+)
+
+// goEnvelopeKeyMetadataVersion is the current encoding version written by
+// [GoEnvelopeEncryptionManager]. It is bumped whenever the on-disk layout of
+// goEnvelopeKeyMetadata changes incompatibly.
+const goEnvelopeKeyMetadataVersion = 1
+
+// goEnvelopeKeyMetadata is the JSON-encoded structure stored as the opaque
+// [EncryptionKeyMetadata] for files produced by [GoEnvelopeEncryptionManager].
+//
+// This encoding is Go-specific, not Java's Avro-encoded StandardKeyMetadata;
+// see the [GoEnvelopeEncryptionManager] doc comment. Per the table spec,
+// DataFile/ManifestFile key_metadata is explicitly "implementation-specific";
+// what must be Iceberg AES GCM Stream compliant - and is - is the wire
+// format of the encrypted byte stream itself (magic, block framing, nonce
+// placement, and AAD; see the constants above).
+type goEnvelopeKeyMetadata struct {
+       Version    int    `json:"v"`
+       KeyID      string `json:"key-id"`
+       WrappedKey []byte `json:"wrapped-key"`
+
+       // BlockSize is the trusted plaintext block size the file was written
+       // with. It is validated bounded before any file I/O, and is what's
+       // actually used to size reads; the stream header's block-length field
+       // is untrusted (unauthenticated, attacker-influenced storage bytes) and
+       // is only compared against this value, never used directly to size an
+       // allocation.
+       BlockSize int64 `json:"block-size"`
+
+       // AADPrefix is combined with each block's little-endian index to form
+       // the AES GCM Stream additional authenticated data, binding every
+       // ciphertext block to this file and to its position so that blocks
+       // cannot be silently reordered, replayed from another file, or spliced
+       // in from elsewhere in the same file. It is not secret.
+       AADPrefix []byte `json:"aad-prefix"`
+
+       // PlaintextLength is the trusted total plaintext size used to compute
+       // the block count on read. Per the AES GCM Stream spec's "File length"
+       // note, a reader must use a length from a trusted source rather than
+       // the underlying storage's reported size, since storage size alone
+       // cannot distinguish a genuinely short file from one truncated by an
+       // attacker who does not also control this metadata.
+       PlaintextLength int64 `json:"plaintext-length"`
+}
+
+// GoEnvelopeEncryptionManager is a generic, format-agnostic 
[EncryptionManager]
+// that provides envelope encryption for arbitrary files (e.g. manifests,
+// manifest lists, Puffin statistics) using a [KeyManagementClient] to wrap
+// and unwrap a fresh AES-256-GCM data encryption key (DEK) per file.
+//
+// Each file is split into fixed-size plaintext blocks and written using the
+// Iceberg AES GCM Stream ("AGS1") format: a magic/block-length header
+// followed by independently authenticated blocks, each carrying its own
+// random nonce and authenticated with an AAD that binds it to the file and
+// to its position. This bounds memory usage and supports random access
+// (Seek/ReadAt) on the decrypted file without buffering or decrypting more
+// than the requested blocks.
+//
+// GoEnvelopeEncryptionManager always encrypts and always decrypts: it fails
+// closed, returning [ErrKeyIDRequired] or [ErrKeyMetadataRequired] rather
+// than silently falling back to plaintext. Use [PlaintextEncryptionManager]
+// for tables or files that are not encrypted.
+//
+// Not interoperable with Java's StandardEncryptionManager: the AGS1 byte
+// stream framing (magic, block layout, nonce placement, AAD) matches the
+// spec, but the two implementations are not cross-readable. Java's reader
+// hard-requires a 1 MiB header block size (a compile-time constant there),
+// while [GoEnvelopeDefaultBlockSize] is 64 KiB, and Java decodes key
+// metadata as a 1-byte version plus Avro-encoded StandardKeyMetadata, while
+// this type encodes key metadata as JSON (goEnvelopeKeyMetadata). Files
+// written by one are not readable by the other.
+type GoEnvelopeEncryptionManager struct {
+       kms       KeyManagementClient
+       dekLength int
+       blockSize int
+}
+
+var _ EncryptionManager = (*GoEnvelopeEncryptionManager)(nil)
+
+// GoEnvelopeManagerOption configures a [GoEnvelopeEncryptionManager] created
+// by [NewGoEnvelopeEncryptionManager].
+type GoEnvelopeManagerOption func(*GoEnvelopeEncryptionManager)
+
+// WithDEKLength overrides the default data encryption key length (in bytes).
+// Valid AES key lengths are 16, 24, or 32 bytes.
+func WithDEKLength(length int) GoEnvelopeManagerOption {
+       return func(m *GoEnvelopeEncryptionManager) { m.dekLength = length }
+}
+
+// WithBlockSize overrides the default plaintext block size (in bytes) used
+// to split files for independent block-level authentication. size must be
+// positive and at most [GoEnvelopeMaxBlockSize].
+func WithBlockSize(size int) GoEnvelopeManagerOption {
+       return func(m *GoEnvelopeEncryptionManager) { m.blockSize = size }
+}
+
+// NewGoEnvelopeEncryptionManager creates a [GoEnvelopeEncryptionManager]
+// backed by kms. kms must not be nil; NewGoEnvelopeEncryptionManager panics
+// if it is.
+func NewGoEnvelopeEncryptionManager(kms KeyManagementClient, opts 
...GoEnvelopeManagerOption) *GoEnvelopeEncryptionManager {
+       if kms == nil {
+               panic("encryption: NewGoEnvelopeEncryptionManager: kms must not 
be nil")
+       }
+
+       m := &GoEnvelopeEncryptionManager{
+               kms:       kms,
+               dekLength: GoEnvelopeDefaultDEKLength,
+               blockSize: GoEnvelopeDefaultBlockSize,
+       }
+       for _, opt := range opts {
+               opt(m)
+       }
+
+       return m
+}
+
+// NewEncryptedOutputFile creates a new AES-GCM block-encrypted output file.
+// keyID identifies the KEK used to wrap the freshly generated per-file DEK,
+// and must be non-empty; otherwise [ErrKeyIDRequired] is returned.
+func (m *GoEnvelopeEncryptionManager) NewEncryptedOutputFile(ctx 
context.Context, writer icebergio.FileWriter, keyID string) 
(EncryptedOutputFile, error) {
+       if keyID == "" {
+               return nil, ErrKeyIDRequired
+       }
+       if m.blockSize <= 0 || m.blockSize > GoEnvelopeMaxBlockSize {
+               return nil, fmt.Errorf("%w: got %d", ErrInvalidBlockSize, 
m.blockSize)
+       }
+       switch m.dekLength {
+       case 16, 24, 32:
+       default:
+               return nil, fmt.Errorf("%w: DEK length must be 16, 24, or 32 
bytes; got %d", ErrInvalidKeyLength, m.dekLength)
+       }
+
+       // The (key, nonce) uniqueness this block format relies on requires a
+       // freshly generated DEK for every file: never cache or reuse
+       // plainDEK/wrappedDEK across calls to NewEncryptedOutputFile.
+       var (
+               plainDEK, wrappedDEK []byte
+               err                  error
+       )
+       if m.kms.SupportsKeyGeneration() {
+               plainDEK, wrappedDEK, err = m.kms.GenerateKey(ctx, keyID, 
m.dekLength)
+               if err != nil {
+                       return nil, fmt.Errorf("encryption: failed to generate 
DEK: %w", err)
+               }
+       } else {
+               plainDEK = make([]byte, m.dekLength)
+               if _, err = io.ReadFull(rand.Reader, plainDEK); err != nil {
+                       return nil, fmt.Errorf("encryption: failed to generate 
DEK: %w", err)
+               }
+               if wrappedDEK, err = m.kms.WrapKey(ctx, keyID, plainDEK); err 
!= nil {
+                       return nil, fmt.Errorf("encryption: failed to wrap DEK: 
%w", err)
+               }
+       }
+
+       aead, err := newGoEnvelopeAEAD(plainDEK)
+       if err != nil {
+               return nil, err
+       }
+
+       aadPrefix := make([]byte, gcmStreamAADPrefixLength)
+       if _, err := io.ReadFull(rand.Reader, aadPrefix); err != nil {
+               return nil, fmt.Errorf("encryption: failed to generate AAD 
prefix: %w", err)
+       }
+
+       // Write the AES GCM Stream header (magic + little-endian block length)
+       // up front, before any ciphertext blocks, per the format spec.
+       header := make([]byte, gcmStreamHeaderLength)
+       copy(header, gcmStreamMagic)
+       binary.LittleEndian.PutUint32(header[len(gcmStreamMagic):], 
uint32(m.blockSize)) //nolint:gosec // bounded by ErrInvalidBlockSize above
+       if _, err := writer.Write(header); err != nil {
+               return nil, fmt.Errorf("encryption: failed to write stream 
header: %w", err)
+       }
+
+       return &goEnvelopeOutputFile{
+               FileWriter: writer,
+               aead:       aead,
+               aadPrefix:  aadPrefix,
+               blockSize:  m.blockSize,
+               keyID:      keyID,
+               wrappedKey: wrappedDEK,
+       }, nil
+}
+
+// NewDecryptedInputFile wraps file for transparent block-level AES-GCM
+// decryption. keyMetadata must be the non-empty blob produced by
+// [GoEnvelopeEncryptionManager.NewEncryptedOutputFile]; otherwise
+// [ErrKeyMetadataRequired] is returned.
+func (m *GoEnvelopeEncryptionManager) NewDecryptedInputFile(ctx 
context.Context, file icebergio.File, keyMetadata EncryptionKeyMetadata) 
(EncryptedInputFile, error) {
+       if len(keyMetadata) == 0 {
+               return nil, ErrKeyMetadataRequired
+       }
+
+       var meta goEnvelopeKeyMetadata
+       if err := json.Unmarshal(keyMetadata, &meta); err != nil {
+               return nil, fmt.Errorf("encryption: failed to decode key 
metadata: %w", err)
+       }
+       if meta.Version != goEnvelopeKeyMetadataVersion {
+               return nil, fmt.Errorf("%w: %d", 
ErrUnsupportedKeyMetadataVersion, meta.Version)
+       }
+       if meta.BlockSize <= 0 || meta.BlockSize > GoEnvelopeMaxBlockSize {
+               return nil, fmt.Errorf("%w: block-size must be positive and at 
most %d, got %d", ErrInvalidKeyMetadata, GoEnvelopeMaxBlockSize, meta.BlockSize)
+       }
+       if meta.PlaintextLength < 0 {
+               return nil, fmt.Errorf("%w: plaintext-length must be 
non-negative, got %d", ErrInvalidKeyMetadata, meta.PlaintextLength)
+       }
+       // AADPrefix is bounded to exactly the length this implementation's own
+       // writer emits (gcmStreamAADPrefixLength): key metadata is Go-specific
+       // (not cross-implementation compatible; see the type doc comment), so
+       // there is no other length a valid file could carry. Without this
+       // ceiling, an oversized prefix would be re-allocated on every
+       // gcmStreamBlockAAD call in readBlock.
+       if len(meta.AADPrefix) != gcmStreamAADPrefixLength {
+               return nil, fmt.Errorf("%w: aad-prefix must be exactly %d 
bytes, got %d", ErrInvalidKeyMetadata, gcmStreamAADPrefixLength, 
len(meta.AADPrefix))
+       }
+
+       plainDEK, err := m.kms.UnwrapKey(ctx, meta.KeyID, meta.WrappedKey)
+       if err != nil {
+               return nil, fmt.Errorf("encryption: failed to unwrap DEK: %w", 
err)
+       }
+
+       aead, err := newGoEnvelopeAEAD(plainDEK)
+       if err != nil {
+               return nil, err
+       }
+
+       // headerBlockSize is untrusted (unauthenticated bytes read from
+       // storage): it is only ever compared against the trusted, already
+       // bounded meta.BlockSize below, never used to size a read or 
allocation.
+       headerBlockSize, err := readGCMStreamHeader(file)
+       if err != nil {
+               return nil, err
+       }
+       if headerBlockSize != meta.BlockSize {
+               return nil, fmt.Errorf("%w: stream header block length %d does 
not match key metadata block-size %d", ErrInvalidStreamHeader, headerBlockSize, 
meta.BlockSize)
+       }
+
+       return &goEnvelopeInputFile{
+               underlying:      file,
+               aead:            aead,
+               aadPrefix:       meta.AADPrefix,
+               blockSize:       meta.BlockSize,
+               plaintextLength: meta.PlaintextLength,
+               keyMetadata:     keyMetadata,
+       }, nil
+}
+
+func newGoEnvelopeAEAD(key []byte) (cipher.AEAD, error) {
+       block, err := aes.NewCipher(key)
+       if err != nil {
+               return nil, fmt.Errorf("%w: %w", ErrInvalidKeyLength, err)
+       }
+       gcm, err := cipher.NewGCM(block)
+       if err != nil {
+               return nil, fmt.Errorf("encryption: failed to create GCM: %w", 
err)
+       }
+
+       return gcm, nil
+}
+
+// readGCMStreamHeader reads and validates the AES GCM Stream magic and
+// block-length header at the start of file, returning the plaintext block
+// length. The header is untrusted, unauthenticated data, so the returned
+// length is bounded to [GoEnvelopeMaxBlockSize] before any allocation sized
+// by it takes place.
+func readGCMStreamHeader(file icebergio.File) (int64, error) {
+       header := make([]byte, gcmStreamHeaderLength)
+       n, err := file.ReadAt(header, 0)
+       if err != nil && !errors.Is(err, io.EOF) {
+               return 0, fmt.Errorf("encryption: failed to read stream header: 
%w", err)
+       }
+       if n != gcmStreamHeaderLength {
+               return 0, fmt.Errorf("%w: expected %d header bytes, got %d", 
ErrInvalidStreamHeader, gcmStreamHeaderLength, n)
+       }
+       if string(header[:len(gcmStreamMagic)]) != gcmStreamMagic {
+               return 0, fmt.Errorf("%w: missing %q magic", 
ErrInvalidStreamHeader, gcmStreamMagic)
+       }
+
+       blockSize := 
int64(binary.LittleEndian.Uint32(header[len(gcmStreamMagic):]))
+       if blockSize <= 0 || blockSize > GoEnvelopeMaxBlockSize {
+               return 0, fmt.Errorf("%w: block length %d out of supported 
range (0, %d]", ErrInvalidStreamHeader, blockSize, GoEnvelopeMaxBlockSize)
+       }
+
+       return blockSize, nil
+}
+
+// gcmStreamBlockAAD derives the AES-GCM additional authenticated data for
+// blockIndex: the per-file AAD prefix followed by the 4-byte little-endian
+// block index, per the Iceberg AES GCM Stream spec. This binds every
+// ciphertext block to this file and to its position, which matters because
+// blocks carry independent random nonces (rather than a nonce derived from
+// the block index): without the index in the AAD, an attacker able to
+// tamper with ciphertext at rest could silently reorder or splice blocks.
+func gcmStreamBlockAAD(prefix []byte, blockIndex uint32) []byte {
+       aad := make([]byte, len(prefix)+4)
+       copy(aad, prefix)
+       binary.LittleEndian.PutUint32(aad[len(prefix):], blockIndex)

Review Comment:
   Nothing in the suite checks this against the spec. If I change this line to 
`binary.BigEndian` (the spec requires 4-byte little-endian), every test still 
passes, because the writer and reader share this helper.
   
   Add a test that decodes the output with only `crypto/aes` and 
`crypto/cipher`:
   - parse the `AGS1` header
   - split each block into `nonce‖ciphertext‖tag`
   - build the AAD by hand as `aadPrefix‖LE32(i)` and call `Open`
   
   A decoder like that fails at block 1 under the BigEndian change. If you can 
produce one, also check in a Java-written fixture (encrypted file plus 
`StandardKeyMetadata` bytes) and decrypt it in a test. That's the actual 
interop check.
   



##########
encryption/go_envelope_manager.go:
##########
@@ -0,0 +1,889 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package encryption
+
+import (
+       "context"
+       "crypto/aes"
+       "crypto/cipher"
+       "crypto/rand"
+       "encoding/binary"
+       "encoding/json"
+       "errors"
+       "fmt"
+       "io"
+       "io/fs"
+       "math"
+       "sync"
+
+       icebergio "github.com/apache/iceberg-go/io"
+)
+
+// Defaults for [GoEnvelopeEncryptionManager].
+const (
+       // GoEnvelopeDefaultDEKLength is the default length, in bytes, of the
+       // per-file data encryption key (DEK) generated for AES-256-GCM.
+       GoEnvelopeDefaultDEKLength = 32
+
+       // GoEnvelopeDefaultBlockSize is the default plaintext block size, in
+       // bytes, used to split a file into independently authenticated AES-GCM
+       // blocks. Blocks allow random access (Seek/ReadAt) without buffering or
+       // decrypting the whole file.
+       //
+       // This is not Java's Ciphers.PLAIN_BLOCK_SIZE (1 MiB, a compile-time
+       // constant there); see the [GoEnvelopeEncryptionManager] doc comment.
+       GoEnvelopeDefaultBlockSize = 64 * 1024
+
+       // GoEnvelopeMaxBlockSize is the largest plaintext block size accepted
+       // from the AES GCM Stream header on read, and the largest that
+       // [WithBlockSize] will configure for writing. The header's block-length
+       // field is unauthenticated, untrusted data (the Iceberg AES GCM Stream
+       // spec's "File length" note applies equally here); without a ceiling, a
+       // crafted header claiming e.g. 1<<40 bytes would force a huge
+       // allocation (see readBlock) before any authentication check can run.
+       GoEnvelopeMaxBlockSize = 128 * 1024 * 1024
+)
+
+// Sentinel errors returned by [GoEnvelopeEncryptionManager].
+var (
+       // ErrKeyIDRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile] when keyID is
+       // empty. GoEnvelopeEncryptionManager always encrypts, so it requires a
+       // KEK to wrap the generated DEK; use [PlaintextEncryptionManager] for
+       // unencrypted tables instead of passing an empty keyID here.
+       ErrKeyIDRequired = errors.New("encryption: GoEnvelopeEncryptionManager 
requires a non-empty keyID")
+
+       // ErrKeyMetadataRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when keyMetadata
+       // is empty. GoEnvelopeEncryptionManager always decrypts, so it requires
+       // the per-file key metadata produced by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile].
+       ErrKeyMetadataRequired = errors.New("encryption: 
GoEnvelopeEncryptionManager requires non-empty key metadata")
+
+       // ErrUnsupportedKeyMetadataVersion is returned when key metadata was
+       // produced by a newer, incompatible encoding version.
+       ErrUnsupportedKeyMetadataVersion = errors.New("encryption: unsupported 
key metadata version")
+
+       // ErrInvalidBlockSize is returned when a configured block size is not
+       // positive or exceeds [GoEnvelopeMaxBlockSize].
+       ErrInvalidBlockSize = errors.New("encryption: block size must be 
positive and at most GoEnvelopeMaxBlockSize")
+
+       // ErrInvalidStreamHeader is returned when the AES GCM Stream header
+       // (the "AGS1" magic and little-endian block-length fields at the start
+       // of an encrypted file, per the Iceberg AES GCM Stream spec) is
+       // missing, truncated, or specifies a block length outside the
+       // supported range.
+       ErrInvalidStreamHeader = errors.New("encryption: invalid AES GCM Stream 
header")
+
+       // ErrInvalidKeyMetadata is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when decoded key
+       // metadata fails basic sanity checks (e.g. a negative plaintext length
+       // or a missing AAD prefix). Key metadata is untrusted input on a crypto
+       // read path, so it is validated rather than trusted blindly.
+       ErrInvalidKeyMetadata = errors.New("encryption: invalid key metadata")
+
+       // ErrOutputFileClosed is returned by [goEnvelopeOutputFile.Write] when
+       // called after Close, or after a previous flush has poisoned the 
writer.
+       // It wraps [fs.ErrClosed] so callers can test with errors.Is(err, 
fs.ErrClosed).
+       ErrOutputFileClosed = fmt.Errorf("encryption: write to closed 
GoEnvelopeEncryptionManager output file: %w", fs.ErrClosed)
+
+       // ErrBlockTruncated is returned by [goEnvelopeInputFile.ReadAt] when 
the
+       // underlying storage returns fewer ciphertext bytes for a block than 
its
+       // recorded length requires, indicating the file was truncated at rest.
+       // This is distinct from [ErrCiphertextTooShort], which 
[KeyManagementClient]
+       // implementations use for a too-short wrapped key or KMS-encrypted
+       // payload; keeping them separate lets a caller tell a malformed KMS 
blob
+       // apart from a short block read.
+       ErrBlockTruncated = errors.New("encryption: block truncated: read fewer 
ciphertext bytes than expected")
+)
+
+// Constants describing the Iceberg AES GCM Stream ("AGS1") wire format used
+// for the ciphertext produced by [GoEnvelopeEncryptionManager]. See
+// https://iceberg.apache.org/gcm-stream-spec/ for the full specification.
+const (
+       // gcmStreamMagic identifies an AES GCM Stream version 1 file.
+       gcmStreamMagic = "AGS1"
+
+       // gcmStreamHeaderLength is the length, in bytes, of the magic plus the
+       // little-endian block-length field written at the start of every file.
+       gcmStreamHeaderLength = len(gcmStreamMagic) + 4
+
+       // gcmStreamNonceLength is the length, in bytes, of the random AES-GCM
+       // nonce stored at the start of every cipher block.
+       gcmStreamNonceLength = 12
+
+       // gcmStreamTagLength is the length, in bytes, of the AES-GCM
+       // authentication tag appended to every cipher block's ciphertext.
+       gcmStreamTagLength = 16
+
+       // gcmStreamBlockOverhead is the number of ciphertext bytes added to
+       // each block beyond its plaintext length (nonce + tag).
+       gcmStreamBlockOverhead = gcmStreamNonceLength + gcmStreamTagLength
+
+       // gcmStreamAADPrefixLength is the length, in bytes, of the random
+       // per-file AAD prefix generated for new output files.
+       gcmStreamAADPrefixLength = 16
+)
+
+// goEnvelopeKeyMetadataVersion is the current encoding version written by
+// [GoEnvelopeEncryptionManager]. It is bumped whenever the on-disk layout of
+// goEnvelopeKeyMetadata changes incompatibly.
+const goEnvelopeKeyMetadataVersion = 1
+
+// goEnvelopeKeyMetadata is the JSON-encoded structure stored as the opaque
+// [EncryptionKeyMetadata] for files produced by [GoEnvelopeEncryptionManager].
+//
+// This encoding is Go-specific, not Java's Avro-encoded StandardKeyMetadata;
+// see the [GoEnvelopeEncryptionManager] doc comment. Per the table spec,
+// DataFile/ManifestFile key_metadata is explicitly "implementation-specific";
+// what must be Iceberg AES GCM Stream compliant - and is - is the wire
+// format of the encrypted byte stream itself (magic, block framing, nonce
+// placement, and AAD; see the constants above).
+type goEnvelopeKeyMetadata struct {
+       Version    int    `json:"v"`
+       KeyID      string `json:"key-id"`
+       WrappedKey []byte `json:"wrapped-key"`
+
+       // BlockSize is the trusted plaintext block size the file was written
+       // with. It is validated bounded before any file I/O, and is what's
+       // actually used to size reads; the stream header's block-length field
+       // is untrusted (unauthenticated, attacker-influenced storage bytes) and
+       // is only compared against this value, never used directly to size an
+       // allocation.
+       BlockSize int64 `json:"block-size"`
+
+       // AADPrefix is combined with each block's little-endian index to form
+       // the AES GCM Stream additional authenticated data, binding every
+       // ciphertext block to this file and to its position so that blocks
+       // cannot be silently reordered, replayed from another file, or spliced
+       // in from elsewhere in the same file. It is not secret.
+       AADPrefix []byte `json:"aad-prefix"`
+
+       // PlaintextLength is the trusted total plaintext size used to compute
+       // the block count on read. Per the AES GCM Stream spec's "File length"
+       // note, a reader must use a length from a trusted source rather than
+       // the underlying storage's reported size, since storage size alone
+       // cannot distinguish a genuinely short file from one truncated by an
+       // attacker who does not also control this metadata.
+       PlaintextLength int64 `json:"plaintext-length"`
+}
+
+// GoEnvelopeEncryptionManager is a generic, format-agnostic 
[EncryptionManager]
+// that provides envelope encryption for arbitrary files (e.g. manifests,
+// manifest lists, Puffin statistics) using a [KeyManagementClient] to wrap
+// and unwrap a fresh AES-256-GCM data encryption key (DEK) per file.
+//
+// Each file is split into fixed-size plaintext blocks and written using the
+// Iceberg AES GCM Stream ("AGS1") format: a magic/block-length header
+// followed by independently authenticated blocks, each carrying its own
+// random nonce and authenticated with an AAD that binds it to the file and
+// to its position. This bounds memory usage and supports random access
+// (Seek/ReadAt) on the decrypted file without buffering or decrypting more
+// than the requested blocks.
+//
+// GoEnvelopeEncryptionManager always encrypts and always decrypts: it fails
+// closed, returning [ErrKeyIDRequired] or [ErrKeyMetadataRequired] rather
+// than silently falling back to plaintext. Use [PlaintextEncryptionManager]
+// for tables or files that are not encrypted.
+//
+// Not interoperable with Java's StandardEncryptionManager: the AGS1 byte
+// stream framing (magic, block layout, nonce placement, AAD) matches the
+// spec, but the two implementations are not cross-readable. Java's reader
+// hard-requires a 1 MiB header block size (a compile-time constant there),
+// while [GoEnvelopeDefaultBlockSize] is 64 KiB, and Java decodes key
+// metadata as a 1-byte version plus Avro-encoded StandardKeyMetadata, while
+// this type encodes key metadata as JSON (goEnvelopeKeyMetadata). Files
+// written by one are not readable by the other.
+type GoEnvelopeEncryptionManager struct {
+       kms       KeyManagementClient
+       dekLength int
+       blockSize int
+}
+
+var _ EncryptionManager = (*GoEnvelopeEncryptionManager)(nil)
+
+// GoEnvelopeManagerOption configures a [GoEnvelopeEncryptionManager] created
+// by [NewGoEnvelopeEncryptionManager].
+type GoEnvelopeManagerOption func(*GoEnvelopeEncryptionManager)
+
+// WithDEKLength overrides the default data encryption key length (in bytes).
+// Valid AES key lengths are 16, 24, or 32 bytes.
+func WithDEKLength(length int) GoEnvelopeManagerOption {
+       return func(m *GoEnvelopeEncryptionManager) { m.dekLength = length }
+}
+
+// WithBlockSize overrides the default plaintext block size (in bytes) used
+// to split files for independent block-level authentication. size must be
+// positive and at most [GoEnvelopeMaxBlockSize].
+func WithBlockSize(size int) GoEnvelopeManagerOption {
+       return func(m *GoEnvelopeEncryptionManager) { m.blockSize = size }
+}
+
+// NewGoEnvelopeEncryptionManager creates a [GoEnvelopeEncryptionManager]
+// backed by kms. kms must not be nil; NewGoEnvelopeEncryptionManager panics
+// if it is.
+func NewGoEnvelopeEncryptionManager(kms KeyManagementClient, opts 
...GoEnvelopeManagerOption) *GoEnvelopeEncryptionManager {
+       if kms == nil {
+               panic("encryption: NewGoEnvelopeEncryptionManager: kms must not 
be nil")
+       }
+
+       m := &GoEnvelopeEncryptionManager{
+               kms:       kms,
+               dekLength: GoEnvelopeDefaultDEKLength,
+               blockSize: GoEnvelopeDefaultBlockSize,
+       }
+       for _, opt := range opts {
+               opt(m)
+       }
+
+       return m
+}
+
+// NewEncryptedOutputFile creates a new AES-GCM block-encrypted output file.
+// keyID identifies the KEK used to wrap the freshly generated per-file DEK,
+// and must be non-empty; otherwise [ErrKeyIDRequired] is returned.
+func (m *GoEnvelopeEncryptionManager) NewEncryptedOutputFile(ctx 
context.Context, writer icebergio.FileWriter, keyID string) 
(EncryptedOutputFile, error) {
+       if keyID == "" {
+               return nil, ErrKeyIDRequired
+       }
+       if m.blockSize <= 0 || m.blockSize > GoEnvelopeMaxBlockSize {
+               return nil, fmt.Errorf("%w: got %d", ErrInvalidBlockSize, 
m.blockSize)
+       }
+       switch m.dekLength {
+       case 16, 24, 32:
+       default:
+               return nil, fmt.Errorf("%w: DEK length must be 16, 24, or 32 
bytes; got %d", ErrInvalidKeyLength, m.dekLength)
+       }
+
+       // The (key, nonce) uniqueness this block format relies on requires a
+       // freshly generated DEK for every file: never cache or reuse
+       // plainDEK/wrappedDEK across calls to NewEncryptedOutputFile.
+       var (
+               plainDEK, wrappedDEK []byte
+               err                  error
+       )
+       if m.kms.SupportsKeyGeneration() {
+               plainDEK, wrappedDEK, err = m.kms.GenerateKey(ctx, keyID, 
m.dekLength)
+               if err != nil {
+                       return nil, fmt.Errorf("encryption: failed to generate 
DEK: %w", err)
+               }
+       } else {
+               plainDEK = make([]byte, m.dekLength)
+               if _, err = io.ReadFull(rand.Reader, plainDEK); err != nil {
+                       return nil, fmt.Errorf("encryption: failed to generate 
DEK: %w", err)
+               }
+               if wrappedDEK, err = m.kms.WrapKey(ctx, keyID, plainDEK); err 
!= nil {
+                       return nil, fmt.Errorf("encryption: failed to wrap DEK: 
%w", err)
+               }
+       }
+
+       aead, err := newGoEnvelopeAEAD(plainDEK)
+       if err != nil {
+               return nil, err
+       }
+
+       aadPrefix := make([]byte, gcmStreamAADPrefixLength)
+       if _, err := io.ReadFull(rand.Reader, aadPrefix); err != nil {
+               return nil, fmt.Errorf("encryption: failed to generate AAD 
prefix: %w", err)
+       }
+
+       // Write the AES GCM Stream header (magic + little-endian block length)
+       // up front, before any ciphertext blocks, per the format spec.
+       header := make([]byte, gcmStreamHeaderLength)
+       copy(header, gcmStreamMagic)
+       binary.LittleEndian.PutUint32(header[len(gcmStreamMagic):], 
uint32(m.blockSize)) //nolint:gosec // bounded by ErrInvalidBlockSize above
+       if _, err := writer.Write(header); err != nil {
+               return nil, fmt.Errorf("encryption: failed to write stream 
header: %w", err)
+       }
+
+       return &goEnvelopeOutputFile{
+               FileWriter: writer,
+               aead:       aead,
+               aadPrefix:  aadPrefix,
+               blockSize:  m.blockSize,
+               keyID:      keyID,
+               wrappedKey: wrappedDEK,
+       }, nil
+}
+
+// NewDecryptedInputFile wraps file for transparent block-level AES-GCM
+// decryption. keyMetadata must be the non-empty blob produced by
+// [GoEnvelopeEncryptionManager.NewEncryptedOutputFile]; otherwise
+// [ErrKeyMetadataRequired] is returned.
+func (m *GoEnvelopeEncryptionManager) NewDecryptedInputFile(ctx 
context.Context, file icebergio.File, keyMetadata EncryptionKeyMetadata) 
(EncryptedInputFile, error) {
+       if len(keyMetadata) == 0 {
+               return nil, ErrKeyMetadataRequired
+       }
+
+       var meta goEnvelopeKeyMetadata
+       if err := json.Unmarshal(keyMetadata, &meta); err != nil {
+               return nil, fmt.Errorf("encryption: failed to decode key 
metadata: %w", err)
+       }
+       if meta.Version != goEnvelopeKeyMetadataVersion {
+               return nil, fmt.Errorf("%w: %d", 
ErrUnsupportedKeyMetadataVersion, meta.Version)
+       }
+       if meta.BlockSize <= 0 || meta.BlockSize > GoEnvelopeMaxBlockSize {
+               return nil, fmt.Errorf("%w: block-size must be positive and at 
most %d, got %d", ErrInvalidKeyMetadata, GoEnvelopeMaxBlockSize, meta.BlockSize)
+       }
+       if meta.PlaintextLength < 0 {
+               return nil, fmt.Errorf("%w: plaintext-length must be 
non-negative, got %d", ErrInvalidKeyMetadata, meta.PlaintextLength)
+       }
+       // AADPrefix is bounded to exactly the length this implementation's own
+       // writer emits (gcmStreamAADPrefixLength): key metadata is Go-specific
+       // (not cross-implementation compatible; see the type doc comment), so
+       // there is no other length a valid file could carry. Without this
+       // ceiling, an oversized prefix would be re-allocated on every
+       // gcmStreamBlockAAD call in readBlock.
+       if len(meta.AADPrefix) != gcmStreamAADPrefixLength {
+               return nil, fmt.Errorf("%w: aad-prefix must be exactly %d 
bytes, got %d", ErrInvalidKeyMetadata, gcmStreamAADPrefixLength, 
len(meta.AADPrefix))
+       }
+
+       plainDEK, err := m.kms.UnwrapKey(ctx, meta.KeyID, meta.WrappedKey)
+       if err != nil {
+               return nil, fmt.Errorf("encryption: failed to unwrap DEK: %w", 
err)
+       }
+
+       aead, err := newGoEnvelopeAEAD(plainDEK)
+       if err != nil {
+               return nil, err
+       }
+
+       // headerBlockSize is untrusted (unauthenticated bytes read from
+       // storage): it is only ever compared against the trusted, already
+       // bounded meta.BlockSize below, never used to size a read or 
allocation.
+       headerBlockSize, err := readGCMStreamHeader(file)
+       if err != nil {
+               return nil, err
+       }
+       if headerBlockSize != meta.BlockSize {
+               return nil, fmt.Errorf("%w: stream header block length %d does 
not match key metadata block-size %d", ErrInvalidStreamHeader, headerBlockSize, 
meta.BlockSize)
+       }
+
+       return &goEnvelopeInputFile{
+               underlying:      file,
+               aead:            aead,
+               aadPrefix:       meta.AADPrefix,
+               blockSize:       meta.BlockSize,
+               plaintextLength: meta.PlaintextLength,
+               keyMetadata:     keyMetadata,
+       }, nil
+}
+
+func newGoEnvelopeAEAD(key []byte) (cipher.AEAD, error) {
+       block, err := aes.NewCipher(key)
+       if err != nil {
+               return nil, fmt.Errorf("%w: %w", ErrInvalidKeyLength, err)
+       }
+       gcm, err := cipher.NewGCM(block)
+       if err != nil {
+               return nil, fmt.Errorf("encryption: failed to create GCM: %w", 
err)
+       }
+
+       return gcm, nil
+}
+
+// readGCMStreamHeader reads and validates the AES GCM Stream magic and
+// block-length header at the start of file, returning the plaintext block
+// length. The header is untrusted, unauthenticated data, so the returned
+// length is bounded to [GoEnvelopeMaxBlockSize] before any allocation sized
+// by it takes place.

Review Comment:
   These comments still describe the old design, where the header's block 
length sized allocations. Now:
   - the header value is only compared against trusted metadata (:373);
   - it's a `uint32`, so the `1<<40` example at :57 can't happen;
   - `numBlocks` (:692-696) takes `blockSize` from key metadata, not "the 
untrusted stream header".
   
   With a fixed 1 MiB block most of this goes away. Whatever remains should 
describe the header as validated only.
   



##########
encryption/go_envelope_manager.go:
##########
@@ -0,0 +1,889 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package encryption
+
+import (
+       "context"
+       "crypto/aes"
+       "crypto/cipher"
+       "crypto/rand"
+       "encoding/binary"
+       "encoding/json"
+       "errors"
+       "fmt"
+       "io"
+       "io/fs"
+       "math"
+       "sync"
+
+       icebergio "github.com/apache/iceberg-go/io"
+)
+
+// Defaults for [GoEnvelopeEncryptionManager].
+const (
+       // GoEnvelopeDefaultDEKLength is the default length, in bytes, of the
+       // per-file data encryption key (DEK) generated for AES-256-GCM.
+       GoEnvelopeDefaultDEKLength = 32
+
+       // GoEnvelopeDefaultBlockSize is the default plaintext block size, in
+       // bytes, used to split a file into independently authenticated AES-GCM
+       // blocks. Blocks allow random access (Seek/ReadAt) without buffering or
+       // decrypting the whole file.
+       //
+       // This is not Java's Ciphers.PLAIN_BLOCK_SIZE (1 MiB, a compile-time
+       // constant there); see the [GoEnvelopeEncryptionManager] doc comment.
+       GoEnvelopeDefaultBlockSize = 64 * 1024
+
+       // GoEnvelopeMaxBlockSize is the largest plaintext block size accepted
+       // from the AES GCM Stream header on read, and the largest that
+       // [WithBlockSize] will configure for writing. The header's block-length
+       // field is unauthenticated, untrusted data (the Iceberg AES GCM Stream
+       // spec's "File length" note applies equally here); without a ceiling, a
+       // crafted header claiming e.g. 1<<40 bytes would force a huge
+       // allocation (see readBlock) before any authentication check can run.
+       GoEnvelopeMaxBlockSize = 128 * 1024 * 1024
+)
+
+// Sentinel errors returned by [GoEnvelopeEncryptionManager].
+var (
+       // ErrKeyIDRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile] when keyID is
+       // empty. GoEnvelopeEncryptionManager always encrypts, so it requires a
+       // KEK to wrap the generated DEK; use [PlaintextEncryptionManager] for
+       // unencrypted tables instead of passing an empty keyID here.
+       ErrKeyIDRequired = errors.New("encryption: GoEnvelopeEncryptionManager 
requires a non-empty keyID")
+
+       // ErrKeyMetadataRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when keyMetadata
+       // is empty. GoEnvelopeEncryptionManager always decrypts, so it requires
+       // the per-file key metadata produced by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile].
+       ErrKeyMetadataRequired = errors.New("encryption: 
GoEnvelopeEncryptionManager requires non-empty key metadata")
+
+       // ErrUnsupportedKeyMetadataVersion is returned when key metadata was
+       // produced by a newer, incompatible encoding version.
+       ErrUnsupportedKeyMetadataVersion = errors.New("encryption: unsupported 
key metadata version")
+
+       // ErrInvalidBlockSize is returned when a configured block size is not
+       // positive or exceeds [GoEnvelopeMaxBlockSize].
+       ErrInvalidBlockSize = errors.New("encryption: block size must be 
positive and at most GoEnvelopeMaxBlockSize")
+
+       // ErrInvalidStreamHeader is returned when the AES GCM Stream header
+       // (the "AGS1" magic and little-endian block-length fields at the start
+       // of an encrypted file, per the Iceberg AES GCM Stream spec) is
+       // missing, truncated, or specifies a block length outside the
+       // supported range.
+       ErrInvalidStreamHeader = errors.New("encryption: invalid AES GCM Stream 
header")
+
+       // ErrInvalidKeyMetadata is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when decoded key
+       // metadata fails basic sanity checks (e.g. a negative plaintext length
+       // or a missing AAD prefix). Key metadata is untrusted input on a crypto
+       // read path, so it is validated rather than trusted blindly.
+       ErrInvalidKeyMetadata = errors.New("encryption: invalid key metadata")
+
+       // ErrOutputFileClosed is returned by [goEnvelopeOutputFile.Write] when
+       // called after Close, or after a previous flush has poisoned the 
writer.
+       // It wraps [fs.ErrClosed] so callers can test with errors.Is(err, 
fs.ErrClosed).
+       ErrOutputFileClosed = fmt.Errorf("encryption: write to closed 
GoEnvelopeEncryptionManager output file: %w", fs.ErrClosed)
+
+       // ErrBlockTruncated is returned by [goEnvelopeInputFile.ReadAt] when 
the
+       // underlying storage returns fewer ciphertext bytes for a block than 
its
+       // recorded length requires, indicating the file was truncated at rest.
+       // This is distinct from [ErrCiphertextTooShort], which 
[KeyManagementClient]
+       // implementations use for a too-short wrapped key or KMS-encrypted
+       // payload; keeping them separate lets a caller tell a malformed KMS 
blob
+       // apart from a short block read.
+       ErrBlockTruncated = errors.New("encryption: block truncated: read fewer 
ciphertext bytes than expected")
+)
+
+// Constants describing the Iceberg AES GCM Stream ("AGS1") wire format used
+// for the ciphertext produced by [GoEnvelopeEncryptionManager]. See
+// https://iceberg.apache.org/gcm-stream-spec/ for the full specification.
+const (
+       // gcmStreamMagic identifies an AES GCM Stream version 1 file.
+       gcmStreamMagic = "AGS1"
+
+       // gcmStreamHeaderLength is the length, in bytes, of the magic plus the
+       // little-endian block-length field written at the start of every file.
+       gcmStreamHeaderLength = len(gcmStreamMagic) + 4
+
+       // gcmStreamNonceLength is the length, in bytes, of the random AES-GCM
+       // nonce stored at the start of every cipher block.
+       gcmStreamNonceLength = 12
+
+       // gcmStreamTagLength is the length, in bytes, of the AES-GCM
+       // authentication tag appended to every cipher block's ciphertext.
+       gcmStreamTagLength = 16
+
+       // gcmStreamBlockOverhead is the number of ciphertext bytes added to
+       // each block beyond its plaintext length (nonce + tag).
+       gcmStreamBlockOverhead = gcmStreamNonceLength + gcmStreamTagLength
+
+       // gcmStreamAADPrefixLength is the length, in bytes, of the random
+       // per-file AAD prefix generated for new output files.
+       gcmStreamAADPrefixLength = 16
+)
+
+// goEnvelopeKeyMetadataVersion is the current encoding version written by
+// [GoEnvelopeEncryptionManager]. It is bumped whenever the on-disk layout of
+// goEnvelopeKeyMetadata changes incompatibly.
+const goEnvelopeKeyMetadataVersion = 1
+
+// goEnvelopeKeyMetadata is the JSON-encoded structure stored as the opaque
+// [EncryptionKeyMetadata] for files produced by [GoEnvelopeEncryptionManager].
+//
+// This encoding is Go-specific, not Java's Avro-encoded StandardKeyMetadata;
+// see the [GoEnvelopeEncryptionManager] doc comment. Per the table spec,
+// DataFile/ManifestFile key_metadata is explicitly "implementation-specific";
+// what must be Iceberg AES GCM Stream compliant - and is - is the wire
+// format of the encrypted byte stream itself (magic, block framing, nonce
+// placement, and AAD; see the constants above).
+type goEnvelopeKeyMetadata struct {
+       Version    int    `json:"v"`
+       KeyID      string `json:"key-id"`
+       WrappedKey []byte `json:"wrapped-key"`
+
+       // BlockSize is the trusted plaintext block size the file was written
+       // with. It is validated bounded before any file I/O, and is what's
+       // actually used to size reads; the stream header's block-length field
+       // is untrusted (unauthenticated, attacker-influenced storage bytes) and
+       // is only compared against this value, never used directly to size an
+       // allocation.
+       BlockSize int64 `json:"block-size"`
+
+       // AADPrefix is combined with each block's little-endian index to form
+       // the AES GCM Stream additional authenticated data, binding every
+       // ciphertext block to this file and to its position so that blocks
+       // cannot be silently reordered, replayed from another file, or spliced
+       // in from elsewhere in the same file. It is not secret.
+       AADPrefix []byte `json:"aad-prefix"`
+
+       // PlaintextLength is the trusted total plaintext size used to compute
+       // the block count on read. Per the AES GCM Stream spec's "File length"
+       // note, a reader must use a length from a trusted source rather than
+       // the underlying storage's reported size, since storage size alone
+       // cannot distinguish a genuinely short file from one truncated by an
+       // attacker who does not also control this metadata.
+       PlaintextLength int64 `json:"plaintext-length"`
+}
+
+// GoEnvelopeEncryptionManager is a generic, format-agnostic 
[EncryptionManager]
+// that provides envelope encryption for arbitrary files (e.g. manifests,
+// manifest lists, Puffin statistics) using a [KeyManagementClient] to wrap
+// and unwrap a fresh AES-256-GCM data encryption key (DEK) per file.
+//
+// Each file is split into fixed-size plaintext blocks and written using the
+// Iceberg AES GCM Stream ("AGS1") format: a magic/block-length header
+// followed by independently authenticated blocks, each carrying its own
+// random nonce and authenticated with an AAD that binds it to the file and
+// to its position. This bounds memory usage and supports random access
+// (Seek/ReadAt) on the decrypted file without buffering or decrypting more
+// than the requested blocks.
+//
+// GoEnvelopeEncryptionManager always encrypts and always decrypts: it fails
+// closed, returning [ErrKeyIDRequired] or [ErrKeyMetadataRequired] rather
+// than silently falling back to plaintext. Use [PlaintextEncryptionManager]
+// for tables or files that are not encrypted.
+//
+// Not interoperable with Java's StandardEncryptionManager: the AGS1 byte
+// stream framing (magic, block layout, nonce placement, AAD) matches the
+// spec, but the two implementations are not cross-readable. Java's reader
+// hard-requires a 1 MiB header block size (a compile-time constant there),
+// while [GoEnvelopeDefaultBlockSize] is 64 KiB, and Java decodes key
+// metadata as a 1-byte version plus Avro-encoded StandardKeyMetadata, while
+// this type encodes key metadata as JSON (goEnvelopeKeyMetadata). Files
+// written by one are not readable by the other.
+type GoEnvelopeEncryptionManager struct {
+       kms       KeyManagementClient
+       dekLength int
+       blockSize int
+}
+
+var _ EncryptionManager = (*GoEnvelopeEncryptionManager)(nil)
+
+// GoEnvelopeManagerOption configures a [GoEnvelopeEncryptionManager] created
+// by [NewGoEnvelopeEncryptionManager].
+type GoEnvelopeManagerOption func(*GoEnvelopeEncryptionManager)
+
+// WithDEKLength overrides the default data encryption key length (in bytes).
+// Valid AES key lengths are 16, 24, or 32 bytes.
+func WithDEKLength(length int) GoEnvelopeManagerOption {
+       return func(m *GoEnvelopeEncryptionManager) { m.dekLength = length }
+}
+
+// WithBlockSize overrides the default plaintext block size (in bytes) used
+// to split files for independent block-level authentication. size must be
+// positive and at most [GoEnvelopeMaxBlockSize].
+func WithBlockSize(size int) GoEnvelopeManagerOption {
+       return func(m *GoEnvelopeEncryptionManager) { m.blockSize = size }
+}
+
+// NewGoEnvelopeEncryptionManager creates a [GoEnvelopeEncryptionManager]
+// backed by kms. kms must not be nil; NewGoEnvelopeEncryptionManager panics
+// if it is.
+func NewGoEnvelopeEncryptionManager(kms KeyManagementClient, opts 
...GoEnvelopeManagerOption) *GoEnvelopeEncryptionManager {
+       if kms == nil {
+               panic("encryption: NewGoEnvelopeEncryptionManager: kms must not 
be nil")
+       }
+
+       m := &GoEnvelopeEncryptionManager{
+               kms:       kms,
+               dekLength: GoEnvelopeDefaultDEKLength,
+               blockSize: GoEnvelopeDefaultBlockSize,
+       }
+       for _, opt := range opts {
+               opt(m)
+       }
+
+       return m
+}
+
+// NewEncryptedOutputFile creates a new AES-GCM block-encrypted output file.
+// keyID identifies the KEK used to wrap the freshly generated per-file DEK,
+// and must be non-empty; otherwise [ErrKeyIDRequired] is returned.
+func (m *GoEnvelopeEncryptionManager) NewEncryptedOutputFile(ctx 
context.Context, writer icebergio.FileWriter, keyID string) 
(EncryptedOutputFile, error) {
+       if keyID == "" {
+               return nil, ErrKeyIDRequired
+       }
+       if m.blockSize <= 0 || m.blockSize > GoEnvelopeMaxBlockSize {
+               return nil, fmt.Errorf("%w: got %d", ErrInvalidBlockSize, 
m.blockSize)
+       }
+       switch m.dekLength {
+       case 16, 24, 32:
+       default:
+               return nil, fmt.Errorf("%w: DEK length must be 16, 24, or 32 
bytes; got %d", ErrInvalidKeyLength, m.dekLength)
+       }
+
+       // The (key, nonce) uniqueness this block format relies on requires a
+       // freshly generated DEK for every file: never cache or reuse
+       // plainDEK/wrappedDEK across calls to NewEncryptedOutputFile.
+       var (
+               plainDEK, wrappedDEK []byte
+               err                  error
+       )
+       if m.kms.SupportsKeyGeneration() {
+               plainDEK, wrappedDEK, err = m.kms.GenerateKey(ctx, keyID, 
m.dekLength)
+               if err != nil {
+                       return nil, fmt.Errorf("encryption: failed to generate 
DEK: %w", err)
+               }
+       } else {
+               plainDEK = make([]byte, m.dekLength)
+               if _, err = io.ReadFull(rand.Reader, plainDEK); err != nil {
+                       return nil, fmt.Errorf("encryption: failed to generate 
DEK: %w", err)
+               }
+               if wrappedDEK, err = m.kms.WrapKey(ctx, keyID, plainDEK); err 
!= nil {
+                       return nil, fmt.Errorf("encryption: failed to wrap DEK: 
%w", err)
+               }
+       }
+
+       aead, err := newGoEnvelopeAEAD(plainDEK)
+       if err != nil {
+               return nil, err
+       }
+
+       aadPrefix := make([]byte, gcmStreamAADPrefixLength)
+       if _, err := io.ReadFull(rand.Reader, aadPrefix); err != nil {
+               return nil, fmt.Errorf("encryption: failed to generate AAD 
prefix: %w", err)
+       }
+
+       // Write the AES GCM Stream header (magic + little-endian block length)
+       // up front, before any ciphertext blocks, per the format spec.
+       header := make([]byte, gcmStreamHeaderLength)
+       copy(header, gcmStreamMagic)
+       binary.LittleEndian.PutUint32(header[len(gcmStreamMagic):], 
uint32(m.blockSize)) //nolint:gosec // bounded by ErrInvalidBlockSize above
+       if _, err := writer.Write(header); err != nil {
+               return nil, fmt.Errorf("encryption: failed to write stream 
header: %w", err)
+       }
+
+       return &goEnvelopeOutputFile{
+               FileWriter: writer,
+               aead:       aead,
+               aadPrefix:  aadPrefix,
+               blockSize:  m.blockSize,
+               keyID:      keyID,
+               wrappedKey: wrappedDEK,
+       }, nil
+}
+
+// NewDecryptedInputFile wraps file for transparent block-level AES-GCM
+// decryption. keyMetadata must be the non-empty blob produced by
+// [GoEnvelopeEncryptionManager.NewEncryptedOutputFile]; otherwise
+// [ErrKeyMetadataRequired] is returned.
+func (m *GoEnvelopeEncryptionManager) NewDecryptedInputFile(ctx 
context.Context, file icebergio.File, keyMetadata EncryptionKeyMetadata) 
(EncryptedInputFile, error) {
+       if len(keyMetadata) == 0 {
+               return nil, ErrKeyMetadataRequired
+       }
+
+       var meta goEnvelopeKeyMetadata
+       if err := json.Unmarshal(keyMetadata, &meta); err != nil {
+               return nil, fmt.Errorf("encryption: failed to decode key 
metadata: %w", err)
+       }
+       if meta.Version != goEnvelopeKeyMetadataVersion {
+               return nil, fmt.Errorf("%w: %d", 
ErrUnsupportedKeyMetadataVersion, meta.Version)
+       }
+       if meta.BlockSize <= 0 || meta.BlockSize > GoEnvelopeMaxBlockSize {
+               return nil, fmt.Errorf("%w: block-size must be positive and at 
most %d, got %d", ErrInvalidKeyMetadata, GoEnvelopeMaxBlockSize, meta.BlockSize)
+       }
+       if meta.PlaintextLength < 0 {
+               return nil, fmt.Errorf("%w: plaintext-length must be 
non-negative, got %d", ErrInvalidKeyMetadata, meta.PlaintextLength)
+       }
+       // AADPrefix is bounded to exactly the length this implementation's own
+       // writer emits (gcmStreamAADPrefixLength): key metadata is Go-specific
+       // (not cross-implementation compatible; see the type doc comment), so
+       // there is no other length a valid file could carry. Without this
+       // ceiling, an oversized prefix would be re-allocated on every
+       // gcmStreamBlockAAD call in readBlock.
+       if len(meta.AADPrefix) != gcmStreamAADPrefixLength {
+               return nil, fmt.Errorf("%w: aad-prefix must be exactly %d 
bytes, got %d", ErrInvalidKeyMetadata, gcmStreamAADPrefixLength, 
len(meta.AADPrefix))
+       }
+
+       plainDEK, err := m.kms.UnwrapKey(ctx, meta.KeyID, meta.WrappedKey)
+       if err != nil {
+               return nil, fmt.Errorf("encryption: failed to unwrap DEK: %w", 
err)
+       }
+
+       aead, err := newGoEnvelopeAEAD(plainDEK)
+       if err != nil {
+               return nil, err
+       }
+
+       // headerBlockSize is untrusted (unauthenticated bytes read from
+       // storage): it is only ever compared against the trusted, already
+       // bounded meta.BlockSize below, never used to size a read or 
allocation.
+       headerBlockSize, err := readGCMStreamHeader(file)
+       if err != nil {
+               return nil, err
+       }
+       if headerBlockSize != meta.BlockSize {
+               return nil, fmt.Errorf("%w: stream header block length %d does 
not match key metadata block-size %d", ErrInvalidStreamHeader, headerBlockSize, 
meta.BlockSize)
+       }
+
+       return &goEnvelopeInputFile{
+               underlying:      file,
+               aead:            aead,
+               aadPrefix:       meta.AADPrefix,
+               blockSize:       meta.BlockSize,
+               plaintextLength: meta.PlaintextLength,
+               keyMetadata:     keyMetadata,
+       }, nil
+}
+
+func newGoEnvelopeAEAD(key []byte) (cipher.AEAD, error) {
+       block, err := aes.NewCipher(key)
+       if err != nil {
+               return nil, fmt.Errorf("%w: %w", ErrInvalidKeyLength, err)
+       }
+       gcm, err := cipher.NewGCM(block)
+       if err != nil {
+               return nil, fmt.Errorf("encryption: failed to create GCM: %w", 
err)
+       }
+
+       return gcm, nil
+}
+
+// readGCMStreamHeader reads and validates the AES GCM Stream magic and
+// block-length header at the start of file, returning the plaintext block
+// length. The header is untrusted, unauthenticated data, so the returned
+// length is bounded to [GoEnvelopeMaxBlockSize] before any allocation sized
+// by it takes place.
+func readGCMStreamHeader(file icebergio.File) (int64, error) {
+       header := make([]byte, gcmStreamHeaderLength)
+       n, err := file.ReadAt(header, 0)
+       if err != nil && !errors.Is(err, io.EOF) {
+               return 0, fmt.Errorf("encryption: failed to read stream header: 
%w", err)
+       }
+       if n != gcmStreamHeaderLength {
+               return 0, fmt.Errorf("%w: expected %d header bytes, got %d", 
ErrInvalidStreamHeader, gcmStreamHeaderLength, n)
+       }
+       if string(header[:len(gcmStreamMagic)]) != gcmStreamMagic {
+               return 0, fmt.Errorf("%w: missing %q magic", 
ErrInvalidStreamHeader, gcmStreamMagic)
+       }
+
+       blockSize := 
int64(binary.LittleEndian.Uint32(header[len(gcmStreamMagic):]))
+       if blockSize <= 0 || blockSize > GoEnvelopeMaxBlockSize {
+               return 0, fmt.Errorf("%w: block length %d out of supported 
range (0, %d]", ErrInvalidStreamHeader, blockSize, GoEnvelopeMaxBlockSize)
+       }
+
+       return blockSize, nil
+}
+
+// gcmStreamBlockAAD derives the AES-GCM additional authenticated data for
+// blockIndex: the per-file AAD prefix followed by the 4-byte little-endian
+// block index, per the Iceberg AES GCM Stream spec. This binds every
+// ciphertext block to this file and to its position, which matters because
+// blocks carry independent random nonces (rather than a nonce derived from
+// the block index): without the index in the AAD, an attacker able to
+// tamper with ciphertext at rest could silently reorder or splice blocks.
+func gcmStreamBlockAAD(prefix []byte, blockIndex uint32) []byte {
+       aad := make([]byte, len(prefix)+4)
+       copy(aad, prefix)
+       binary.LittleEndian.PutUint32(aad[len(prefix):], blockIndex)
+
+       return aad
+}
+
+// checkedMulInt64 returns a*b and true, or (0, false) if the multiplication
+// overflows int64.
+func checkedMulInt64(a, b int64) (int64, bool) {
+       if a == 0 || b == 0 {
+               return 0, true
+       }
+       result := a * b
+       if result/b != a {
+               return 0, false
+       }
+
+       return result, true
+}
+
+// checkedAddInt64 returns a+b and true, or (0, false) if the addition
+// overflows int64.
+func checkedAddInt64(a, b int64) (int64, bool) {
+       result := a + b
+       if (b > 0 && result < a) || (b < 0 && result > a) {
+               return 0, false
+       }
+
+       return result, true
+}
+
+// goEnvelopeOutputFile is an [EncryptedOutputFile] that seals fixed-size
+// plaintext blocks with AES-GCM as they are written, using the Iceberg AES
+// GCM Stream ("AGS1") wire format.
+type goEnvelopeOutputFile struct {
+       icebergio.FileWriter
+
+       aead       cipher.AEAD
+       aadPrefix  []byte
+       blockSize  int
+       keyID      string
+       wrappedKey []byte
+
+       buf        []byte
+       blockIndex uint32
+       written    int64
+       closed     bool
+       err        error
+
+       // underlyingClosed tracks whether FileWriter.Close has already been
+       // attempted, so a poisoned writer is closed exactly once regardless of
+       // whether the failure is first observed in Write or in Close.
+       underlyingClosed bool
+
+       keyMetadata EncryptionKeyMetadata
+}
+
+var _ EncryptedOutputFile = (*goEnvelopeOutputFile)(nil)
+
+// closeUnderlyingIgnoringError closes the underlying writer at most once.
+// The error is ignored: the caller is already reporting a more specific
+// failure (a flush or encode error), and this is best-effort cleanup so a
+// poisoned writer never leaks its underlying file descriptor or connection.
+//
+// defensive: today every call site sets f.err before calling this, and
+// Write/Close both bail out early once f.err is set, so the
+// underlyingClosed guard never actually stops a second Close in practice.
+// It stays as a hard invariant in case a future call site is added that
+// doesn't follow that pattern.
+func (f *goEnvelopeOutputFile) closeUnderlyingIgnoringError() {
+       if f.underlyingClosed {
+               return
+       }
+       f.underlyingClosed = true
+       _ = f.FileWriter.Close()
+}
+
+func (f *goEnvelopeOutputFile) Write(p []byte) (int, error) {
+       if f.err != nil {
+               return 0, f.err
+       }
+       if f.closed {
+               return 0, ErrOutputFileClosed
+       }
+
+       total := len(p)
+       consumed := 0 // bytes of p appended into f.buf so far in this call
+       accepted := 0 // bytes of p known to be durably flushed; reported on 
failure
+       for len(p) > 0 {
+               space := f.blockSize - len(f.buf)
+               n := min(space, len(p))
+               f.buf = append(f.buf, p[:n]...)
+               p = p[n:]
+               consumed += n
+               if len(f.buf) == f.blockSize {
+                       if err := f.flushBlock(); err != nil {
+                               f.err = err
+                               f.closeUnderlyingIgnoringError()
+
+                               return accepted, err
+                       }
+                       accepted = consumed
+               }
+       }
+
+       return total, nil
+}
+
+// flushBlock seals and writes the currently buffered plaintext block using a
+// fresh random nonce, per the Iceberg AES GCM Stream format. f.written is
+// only advanced once the ciphertext has actually reached the underlying
+// writer, so a failed flush never overcounts PlaintextLength.
+func (f *goEnvelopeOutputFile) flushBlock() error {
+       if f.blockIndex == math.MaxUint32 {
+               return errors.New("encryption: cannot write block: exceeded 
maximum block count")
+       }
+
+       nonce := make([]byte, gcmStreamNonceLength)
+       if _, err := io.ReadFull(rand.Reader, nonce); err != nil {
+               return fmt.Errorf("encryption: failed to generate block nonce: 
%w", err)
+       }
+
+       aad := gcmStreamBlockAAD(f.aadPrefix, f.blockIndex)
+       sealed := f.aead.Seal(nonce, nonce, f.buf, aad)
+       if _, err := f.FileWriter.Write(sealed); err != nil {
+               return fmt.Errorf("encryption: failed to write encrypted block: 
%w", err)
+       }
+       f.written += int64(len(f.buf))
+       f.blockIndex++
+       f.buf = f.buf[:0]
+
+       return nil
+}
+
+// ReadFrom copies from r, encrypting as data is written, satisfying
+// io.ReaderFrom (required by [icebergio.FileWriter]).
+func (f *goEnvelopeOutputFile) ReadFrom(r io.Reader) (int64, error) {
+       buf := make([]byte, max(32*1024, f.blockSize))
+       var total int64
+       for {
+               n, err := r.Read(buf)
+               if n > 0 {
+                       wn, werr := f.Write(buf[:n])
+                       total += int64(wn)
+                       if werr != nil {
+                               return total, werr
+                       }
+               }
+               if err == io.EOF {
+                       break
+               }
+               if err != nil {
+                       return total, err
+               }
+       }
+
+       return total, nil
+}
+
+// Close flushes any buffered partial block and finalizes the key metadata.
+// closed is only set, and keyMetadata only published, once everything -
+// including the underlying Close - has succeeded; a failed Close poisons the
+// writer (via f.err) so a retry reliably reports the same error instead of
+// masking the failure as success or exposing metadata for an output that
+// never finished.
+func (f *goEnvelopeOutputFile) Close() error {
+       if f.err != nil {
+               return f.err
+       }
+       if f.closed {
+               return nil
+       }
+
+       if len(f.buf) > 0 {

Review Comment:
   Java writes block 0 even when the plaintext is empty. `close()` always calls 
`encryptAndWriteBlock()`, which only skips an empty block when 
`currentBlockIndex != 0`. An empty Java file is therefore 8 + 28 = 36 bytes, 
and iceberg-rust rejects anything shorter (`MIN_STREAM_LENGTH`). This code 
writes an 8-byte header-only file.
   
   The spec text says the last block is non-empty, but every shipping 
implementation writes the empty block 0, so match them:
   - flush when `len(f.buf) > 0 || f.blockIndex == 0`
   - on read, accept the 36-byte empty form (encrypted length 36 → plaintext 
length 0)
   



##########
encryption/go_envelope_manager_test.go:
##########
@@ -0,0 +1,654 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package encryption_test
+
+import (
+       "bytes"
+       "encoding/base64"
+       "encoding/binary"
+       "errors"
+       "fmt"
+       "io"
+       "io/fs"
+       "math"
+       "testing"
+       "time"
+
+       "github.com/apache/iceberg-go/encryption"
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+       "golang.org/x/sync/errgroup"
+)
+
+func newTestGoEnvelopeManager(t *testing.T, opts 
...encryption.GoEnvelopeManagerOption) 
(*encryption.GoEnvelopeEncryptionManager, 
*encryption.InMemoryKeyManagementClient) {
+       t.Helper()
+       kms := encryption.NewInMemoryKeyManagementClient()
+       require.NoError(t, kms.AddKey("kek-1", bytes.Repeat([]byte{0x42}, 32)))
+
+       return encryption.NewGoEnvelopeEncryptionManager(kms, opts...), kms
+}
+
+func encryptAll(t *testing.T, mgr *encryption.GoEnvelopeEncryptionManager, 
keyID string, plaintext []byte) (ciphertext []byte, keyMetadata 
encryption.EncryptionKeyMetadata) {
+       t.Helper()
+       fw := &memFileWriter{}
+       out, err := mgr.NewEncryptedOutputFile(t.Context(), fw, keyID)
+       require.NoError(t, err)
+       _, err = out.Write(plaintext)
+       require.NoError(t, err)
+       require.NoError(t, out.Close())
+
+       return fw.Bytes(), out.KeyMetadata()
+}
+
+func decryptAll(t *testing.T, mgr *encryption.GoEnvelopeEncryptionManager, 
ciphertext []byte, keyMetadata encryption.EncryptionKeyMetadata) []byte {
+       t.Helper()
+       in, err := mgr.NewDecryptedInputFile(t.Context(), 
newMemFile(ciphertext), keyMetadata)
+       require.NoError(t, err)
+       data, err := io.ReadAll(in)
+       require.NoError(t, err)
+
+       return data
+}
+
+func TestGoEnvelopeEncryptionManager_RoundTrip_SmallerThanBlock(t *testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       plaintext := []byte("hello iceberg")
+
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", plaintext)
+       assert.NotEmpty(t, keyMetadata)
+       assert.NotEqual(t, plaintext, ciphertext, "ciphertext must differ from 
plaintext")
+
+       got := decryptAll(t, mgr, ciphertext, keyMetadata)
+       assert.Equal(t, plaintext, got)
+}
+
+func TestGoEnvelopeEncryptionManager_RoundTrip_MultiBlock(t *testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       plaintext := bytes.Repeat([]byte("0123456789abcdef"), 10)
+       plaintext = append(plaintext, []byte("partial-tail")...)
+
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", plaintext)
+       got := decryptAll(t, mgr, ciphertext, keyMetadata)
+       assert.Equal(t, plaintext, got)
+}
+
+func TestGoEnvelopeEncryptionManager_RoundTrip_ExactBlockMultiple(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       plaintext := bytes.Repeat([]byte("0123456789abcdef"), 4)
+
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", plaintext)
+       got := decryptAll(t, mgr, ciphertext, keyMetadata)
+       assert.Equal(t, plaintext, got)
+}
+
+func TestGoEnvelopeEncryptionManager_RoundTrip_Empty(t *testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", nil)
+       got := decryptAll(t, mgr, ciphertext, keyMetadata)
+       assert.Empty(t, got)
+}
+
+func TestGoEnvelopeEncryptionManager_RandomAccess(t *testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       plaintext := bytes.Repeat([]byte("0123456789abcdef"), 10)
+       plaintext = append(plaintext, []byte("tail-bytes")...)
+
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", plaintext)
+       in, err := mgr.NewDecryptedInputFile(t.Context(), 
newMemFile(ciphertext), keyMetadata)
+       require.NoError(t, err)
+
+       // Read a slice spanning a block boundary.
+       buf := make([]byte, 20)
+       n, err := in.ReadAt(buf, 10)
+       require.NoError(t, err)
+       assert.Equal(t, plaintext[10:30], buf[:n])
+
+       // Seek + Read.
+       pos, err := in.Seek(5, io.SeekStart)
+       require.NoError(t, err)
+       assert.Equal(t, int64(5), pos)
+
+       rest, err := io.ReadAll(in)
+       require.NoError(t, err)
+       assert.Equal(t, plaintext[5:], rest)
+}
+
+func TestGoEnvelopeEncryptionManager_TamperedBlockFailsAuthentication(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       plaintext := bytes.Repeat([]byte("0123456789abcdef"), 4)
+
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", plaintext)
+       // Flip a bit inside the first block's ciphertext, past the 8-byte AES
+       // GCM Stream header and 12-byte nonce; flipping the header instead
+       // would fail earlier, at stream-header validation.
+       ciphertext[8+12] ^= 0xFF
+
+       in, err := mgr.NewDecryptedInputFile(t.Context(), 
newMemFile(ciphertext), keyMetadata)
+       require.NoError(t, err)
+
+       _, err = io.ReadAll(in)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrAuthenticationFailed))
+}
+
+func TestGoEnvelopeEncryptionManager_ReorderedBlocksFailAuthentication(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       plaintext := bytes.Repeat([]byte("0123456789abcdef"), 2) // exactly two 
full 16-byte blocks
+
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", plaintext)
+
+       // The AES GCM Stream header (magic + block length) is 8 bytes; each
+       // cipher block here is 16 (block size) + 12 (nonce) + 16 (tag) = 44
+       // bytes. Swap the two blocks: with a random nonce per block, only the
+       // AAD (which binds a block to its position) can catch this.
+       const headerLen, blockPhysicalLen = 8, 44
+       block1 := append([]byte(nil), 
ciphertext[headerLen:headerLen+blockPhysicalLen]...)
+       block2 := append([]byte(nil), 
ciphertext[headerLen+blockPhysicalLen:headerLen+2*blockPhysicalLen]...)
+       copy(ciphertext[headerLen:headerLen+blockPhysicalLen], block2)
+       
copy(ciphertext[headerLen+blockPhysicalLen:headerLen+2*blockPhysicalLen], 
block1)
+
+       in, err := mgr.NewDecryptedInputFile(t.Context(), 
newMemFile(ciphertext), keyMetadata)
+       require.NoError(t, err)
+
+       _, err = io.ReadAll(in)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrAuthenticationFailed), 
"swapping blocks must be detected: the AAD binds each block to its position")
+}
+
+func TestGoEnvelopeEncryptionManager_StreamHeaderMagicMismatchRejected(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", []byte("hello 
iceberg"))
+       copy(ciphertext[:4], "XXXX")
+
+       _, err := mgr.NewDecryptedInputFile(t.Context(), 
newMemFile(ciphertext), keyMetadata)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidStreamHeader))
+}
+
+func TestGoEnvelopeEncryptionManager_StreamHeaderZeroBlockLengthRejected(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", []byte("hello 
iceberg"))
+       binary.LittleEndian.PutUint32(ciphertext[4:8], 0)
+
+       _, err := mgr.NewDecryptedInputFile(t.Context(), 
newMemFile(ciphertext), keyMetadata)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidStreamHeader))
+}
+
+func 
TestGoEnvelopeEncryptionManager_StreamHeaderBlockLengthExceedsMaxRejected(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", []byte("hello 
iceberg"))
+       binary.LittleEndian.PutUint32(ciphertext[4:8], math.MaxUint32)
+
+       _, err := mgr.NewDecryptedInputFile(t.Context(), 
newMemFile(ciphertext), keyMetadata)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidStreamHeader))
+}
+
+func TestGoEnvelopeEncryptionManager_OutputFile_EmptyKeyIDRejected(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t)
+       _, err := mgr.NewEncryptedOutputFile(t.Context(), &memFileWriter{}, "")
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrKeyIDRequired))
+}
+
+func TestGoEnvelopeEncryptionManager_InputFile_EmptyKeyMetadataRejected(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t)
+       _, err := mgr.NewDecryptedInputFile(t.Context(), newMemFile(nil), nil)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrKeyMetadataRequired))
+}
+
+func TestGoEnvelopeEncryptionManager_UnknownKeyIDPropagatesFromKMS(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t)
+       _, err := mgr.NewEncryptedOutputFile(t.Context(), &memFileWriter{}, 
"does-not-exist")
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrUnknownKeyID))
+}
+
+func TestGoEnvelopeEncryptionManager_MalformedKeyMetadataRejected(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t)
+       _, err := mgr.NewDecryptedInputFile(t.Context(), newMemFile(nil), 
encryption.EncryptionKeyMetadata("not json"))
+       require.Error(t, err)
+}
+
+func TestGoEnvelopeEncryptionManager_ZeroBlockSizeRejectedOnWrite(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(0))
+       _, err := mgr.NewEncryptedOutputFile(t.Context(), &memFileWriter{}, 
"kek-1")
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidBlockSize))
+}
+
+func TestGoEnvelopeEncryptionManager_NegativeBlockSizeRejectedOnWrite(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(-1))
+       _, err := mgr.NewEncryptedOutputFile(t.Context(), &memFileWriter{}, 
"kek-1")
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidBlockSize))
+}
+
+func TestGoEnvelopeEncryptionManager_BlockSizeExceedingMaxRejectedOnWrite(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, 
encryption.WithBlockSize(encryption.GoEnvelopeMaxBlockSize+1))
+       _, err := mgr.NewEncryptedOutputFile(t.Context(), &memFileWriter{}, 
"kek-1")
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidBlockSize))
+}
+
+func 
TestGoEnvelopeEncryptionManager_NegativePlaintextLengthMetadataRejectedOnRead(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t)
+       meta := 
[]byte(`{"v":1,"key-id":"kek-1","wrapped-key":"AA==","block-size":16,"aad-prefix":"MTIzNDU2Nzg5MDEyMzQ1Ng==","plaintext-length":-1}`)
+       _, err := mgr.NewDecryptedInputFile(t.Context(), newMemFile(nil), meta)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidKeyMetadata))
+       assert.ErrorContains(t, err, "plaintext-length must be non-negative")
+}
+
+func TestGoEnvelopeEncryptionManager_EmptyAADPrefixMetadataRejectedOnRead(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t)
+       meta := 
[]byte(`{"v":1,"key-id":"kek-1","wrapped-key":"AA==","block-size":16,"aad-prefix":"","plaintext-length":10}`)
+       _, err := mgr.NewDecryptedInputFile(t.Context(), newMemFile(nil), meta)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidKeyMetadata))
+       assert.ErrorContains(t, err, "aad-prefix must be exactly 16 bytes")
+}
+
+func 
TestGoEnvelopeEncryptionManager_OversizedAADPrefixMetadataRejectedOnRead(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t)
+       // 17 zero bytes, base64-encoded: one byte longer than 
gcmStreamAADPrefixLength.
+       oversized := base64.StdEncoding.EncodeToString(make([]byte, 17))
+       meta := fmt.Appendf(nil, 
`{"v":1,"key-id":"kek-1","wrapped-key":"AA==","block-size":16,"aad-prefix":"%s","plaintext-length":10}`,
 oversized)
+       _, err := mgr.NewDecryptedInputFile(t.Context(), newMemFile(nil), meta)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidKeyMetadata))
+       assert.ErrorContains(t, err, "aad-prefix must be exactly 16 bytes")
+}
+
+func TestGoEnvelopeEncryptionManager_ZeroBlockSizeMetadataRejectedOnRead(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t)
+       meta := 
[]byte(`{"v":1,"key-id":"kek-1","wrapped-key":"AA==","block-size":0,"aad-prefix":"MTIzNDU2Nzg5MDEyMzQ1Ng==","plaintext-length":10}`)
+       _, err := mgr.NewDecryptedInputFile(t.Context(), newMemFile(nil), meta)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidKeyMetadata))
+       assert.ErrorContains(t, err, "block-size must be positive")
+}
+
+func 
TestGoEnvelopeEncryptionManager_BlockSizeExceedingMaxMetadataRejectedOnRead(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t)
+       meta := fmt.Appendf(nil, 
`{"v":1,"key-id":"kek-1","wrapped-key":"AA==","block-size":%d,"aad-prefix":"MTIzNDU2Nzg5MDEyMzQ1Ng==","plaintext-length":10}`,
 encryption.GoEnvelopeMaxBlockSize+1)
+       _, err := mgr.NewDecryptedInputFile(t.Context(), newMemFile(nil), meta)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidKeyMetadata))
+       assert.ErrorContains(t, err, "block-size must be positive")
+}
+
+func TestGoEnvelopeEncryptionManager_InvalidDEKLengthRejected(t *testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithDEKLength(15))
+       _, err := mgr.NewEncryptedOutputFile(t.Context(), &memFileWriter{}, 
"kek-1")
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrInvalidKeyLength))
+}
+
+// shortReadFile is an icebergio.File stub whose ReadAt returns fewer bytes
+// than requested along with io.EOF once truncateAt is reached, mimicking a
+// real backend (S3, local fs) reading a genuinely truncated file. A plain
+// bytes.Reader always fills the requested slice, which is why the round-trip
+// tests never exercise this path.
+type shortReadFile struct {
+       *memFile
+       truncateAt int64
+}
+
+func (f *shortReadFile) ReadAt(p []byte, off int64) (int, error) {
+       if off >= f.truncateAt {
+               return 0, io.EOF
+       }
+       if off+int64(len(p)) > f.truncateAt {
+               p = p[:f.truncateAt-off]
+       }
+       n, err := f.memFile.ReadAt(p, off)
+       if err == nil && int64(n) < int64(len(p)) {
+               err = io.EOF
+       }
+
+       return n, err
+}
+
+func 
TestGoEnvelopeEncryptionManager_TruncatedBackendReadReportsBlockTruncated(t 
*testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       plaintext := bytes.Repeat([]byte("0123456789abcdef"), 4)
+
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", plaintext)
+       truncated := &shortReadFile{memFile: newMemFile(ciphertext), 
truncateAt: int64(len(ciphertext) - 5)}
+
+       in, err := mgr.NewDecryptedInputFile(t.Context(), truncated, 
keyMetadata)
+       require.NoError(t, err)
+
+       _, err = io.ReadAll(in)
+       require.Error(t, err)
+       assert.True(t, errors.Is(err, encryption.ErrBlockTruncated))
+       assert.False(t, errors.Is(err, encryption.ErrCiphertextTooShort), "a 
truncated block read is distinct from a too-short KMS-wrapped key/payload")
+       assert.False(t, errors.Is(err, encryption.ErrAuthenticationFailed), "a 
truncated read must not be misreported as tampering")
+}
+
+func TestGoEnvelopeEncryptionManager_EmptyFile_ZeroLengthReadAt(t *testing.T) {
+       mgr, _ := newTestGoEnvelopeManager(t, encryption.WithBlockSize(16))
+       ciphertext, keyMetadata := encryptAll(t, mgr, "kek-1", nil)
+
+       in, err := mgr.NewDecryptedInputFile(t.Context(), 
newMemFile(ciphertext), keyMetadata)
+       require.NoError(t, err)
+
+       n, err := in.ReadAt(nil, 0)
+       assert.NoError(t, err, "a zero-length ReadAt on an empty file should 
behave like bytes.Reader, not report io.EOF")

Review Comment:
   Returning `(0, nil)` here is fine, but the reason in the message is wrong. 
`bytes.Reader.ReadAt` returns `(0, io.EOF)` whenever `off >= len`, whatever 
`len(p)` is:
   
   ```go
   if off >= int64(len(r.s)) {
        return 0, io.EOF
   }
   ```
   
   The behavior here matches `os.File.ReadAt` instead. Reword the message, e.g. 
"zero-length ReadAt must return (0, nil), like os.File". (@laskoviymishka: same 
correction to the premise of the 08-05 zero-length-read thread.)
   



##########
encryption/go_envelope_manager.go:
##########
@@ -0,0 +1,889 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package encryption
+
+import (
+       "context"
+       "crypto/aes"
+       "crypto/cipher"
+       "crypto/rand"
+       "encoding/binary"
+       "encoding/json"
+       "errors"
+       "fmt"
+       "io"
+       "io/fs"
+       "math"
+       "sync"
+
+       icebergio "github.com/apache/iceberg-go/io"
+)
+
+// Defaults for [GoEnvelopeEncryptionManager].
+const (
+       // GoEnvelopeDefaultDEKLength is the default length, in bytes, of the
+       // per-file data encryption key (DEK) generated for AES-256-GCM.
+       GoEnvelopeDefaultDEKLength = 32
+
+       // GoEnvelopeDefaultBlockSize is the default plaintext block size, in
+       // bytes, used to split a file into independently authenticated AES-GCM
+       // blocks. Blocks allow random access (Seek/ReadAt) without buffering or
+       // decrypting the whole file.
+       //
+       // This is not Java's Ciphers.PLAIN_BLOCK_SIZE (1 MiB, a compile-time
+       // constant there); see the [GoEnvelopeEncryptionManager] doc comment.
+       GoEnvelopeDefaultBlockSize = 64 * 1024
+
+       // GoEnvelopeMaxBlockSize is the largest plaintext block size accepted
+       // from the AES GCM Stream header on read, and the largest that
+       // [WithBlockSize] will configure for writing. The header's block-length
+       // field is unauthenticated, untrusted data (the Iceberg AES GCM Stream
+       // spec's "File length" note applies equally here); without a ceiling, a
+       // crafted header claiming e.g. 1<<40 bytes would force a huge
+       // allocation (see readBlock) before any authentication check can run.
+       GoEnvelopeMaxBlockSize = 128 * 1024 * 1024
+)
+
+// Sentinel errors returned by [GoEnvelopeEncryptionManager].
+var (
+       // ErrKeyIDRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile] when keyID is
+       // empty. GoEnvelopeEncryptionManager always encrypts, so it requires a
+       // KEK to wrap the generated DEK; use [PlaintextEncryptionManager] for
+       // unencrypted tables instead of passing an empty keyID here.
+       ErrKeyIDRequired = errors.New("encryption: GoEnvelopeEncryptionManager 
requires a non-empty keyID")
+
+       // ErrKeyMetadataRequired is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when keyMetadata
+       // is empty. GoEnvelopeEncryptionManager always decrypts, so it requires
+       // the per-file key metadata produced by
+       // [GoEnvelopeEncryptionManager.NewEncryptedOutputFile].
+       ErrKeyMetadataRequired = errors.New("encryption: 
GoEnvelopeEncryptionManager requires non-empty key metadata")
+
+       // ErrUnsupportedKeyMetadataVersion is returned when key metadata was
+       // produced by a newer, incompatible encoding version.
+       ErrUnsupportedKeyMetadataVersion = errors.New("encryption: unsupported 
key metadata version")
+
+       // ErrInvalidBlockSize is returned when a configured block size is not
+       // positive or exceeds [GoEnvelopeMaxBlockSize].
+       ErrInvalidBlockSize = errors.New("encryption: block size must be 
positive and at most GoEnvelopeMaxBlockSize")
+
+       // ErrInvalidStreamHeader is returned when the AES GCM Stream header
+       // (the "AGS1" magic and little-endian block-length fields at the start
+       // of an encrypted file, per the Iceberg AES GCM Stream spec) is
+       // missing, truncated, or specifies a block length outside the
+       // supported range.
+       ErrInvalidStreamHeader = errors.New("encryption: invalid AES GCM Stream 
header")
+
+       // ErrInvalidKeyMetadata is returned by
+       // [GoEnvelopeEncryptionManager.NewDecryptedInputFile] when decoded key
+       // metadata fails basic sanity checks (e.g. a negative plaintext length
+       // or a missing AAD prefix). Key metadata is untrusted input on a crypto
+       // read path, so it is validated rather than trusted blindly.
+       ErrInvalidKeyMetadata = errors.New("encryption: invalid key metadata")
+
+       // ErrOutputFileClosed is returned by [goEnvelopeOutputFile.Write] when
+       // called after Close, or after a previous flush has poisoned the 
writer.
+       // It wraps [fs.ErrClosed] so callers can test with errors.Is(err, 
fs.ErrClosed).
+       ErrOutputFileClosed = fmt.Errorf("encryption: write to closed 
GoEnvelopeEncryptionManager output file: %w", fs.ErrClosed)
+
+       // ErrBlockTruncated is returned by [goEnvelopeInputFile.ReadAt] when 
the
+       // underlying storage returns fewer ciphertext bytes for a block than 
its
+       // recorded length requires, indicating the file was truncated at rest.
+       // This is distinct from [ErrCiphertextTooShort], which 
[KeyManagementClient]
+       // implementations use for a too-short wrapped key or KMS-encrypted
+       // payload; keeping them separate lets a caller tell a malformed KMS 
blob
+       // apart from a short block read.
+       ErrBlockTruncated = errors.New("encryption: block truncated: read fewer 
ciphertext bytes than expected")
+)
+
+// Constants describing the Iceberg AES GCM Stream ("AGS1") wire format used
+// for the ciphertext produced by [GoEnvelopeEncryptionManager]. See
+// https://iceberg.apache.org/gcm-stream-spec/ for the full specification.
+const (
+       // gcmStreamMagic identifies an AES GCM Stream version 1 file.
+       gcmStreamMagic = "AGS1"
+
+       // gcmStreamHeaderLength is the length, in bytes, of the magic plus the
+       // little-endian block-length field written at the start of every file.
+       gcmStreamHeaderLength = len(gcmStreamMagic) + 4
+
+       // gcmStreamNonceLength is the length, in bytes, of the random AES-GCM
+       // nonce stored at the start of every cipher block.
+       gcmStreamNonceLength = 12
+
+       // gcmStreamTagLength is the length, in bytes, of the AES-GCM
+       // authentication tag appended to every cipher block's ciphertext.
+       gcmStreamTagLength = 16
+
+       // gcmStreamBlockOverhead is the number of ciphertext bytes added to
+       // each block beyond its plaintext length (nonce + tag).
+       gcmStreamBlockOverhead = gcmStreamNonceLength + gcmStreamTagLength
+
+       // gcmStreamAADPrefixLength is the length, in bytes, of the random
+       // per-file AAD prefix generated for new output files.
+       gcmStreamAADPrefixLength = 16
+)
+
+// goEnvelopeKeyMetadataVersion is the current encoding version written by
+// [GoEnvelopeEncryptionManager]. It is bumped whenever the on-disk layout of
+// goEnvelopeKeyMetadata changes incompatibly.
+const goEnvelopeKeyMetadataVersion = 1
+
+// goEnvelopeKeyMetadata is the JSON-encoded structure stored as the opaque
+// [EncryptionKeyMetadata] for files produced by [GoEnvelopeEncryptionManager].
+//
+// This encoding is Go-specific, not Java's Avro-encoded StandardKeyMetadata;
+// see the [GoEnvelopeEncryptionManager] doc comment. Per the table spec,
+// DataFile/ManifestFile key_metadata is explicitly "implementation-specific";
+// what must be Iceberg AES GCM Stream compliant - and is - is the wire
+// format of the encrypted byte stream itself (magic, block framing, nonce
+// placement, and AAD; see the constants above).
+type goEnvelopeKeyMetadata struct {
+       Version    int    `json:"v"`
+       KeyID      string `json:"key-id"`
+       WrappedKey []byte `json:"wrapped-key"`
+
+       // BlockSize is the trusted plaintext block size the file was written
+       // with. It is validated bounded before any file I/O, and is what's
+       // actually used to size reads; the stream header's block-length field
+       // is untrusted (unauthenticated, attacker-influenced storage bytes) and
+       // is only compared against this value, never used directly to size an
+       // allocation.
+       BlockSize int64 `json:"block-size"`
+
+       // AADPrefix is combined with each block's little-endian index to form
+       // the AES GCM Stream additional authenticated data, binding every
+       // ciphertext block to this file and to its position so that blocks
+       // cannot be silently reordered, replayed from another file, or spliced
+       // in from elsewhere in the same file. It is not secret.
+       AADPrefix []byte `json:"aad-prefix"`
+
+       // PlaintextLength is the trusted total plaintext size used to compute
+       // the block count on read. Per the AES GCM Stream spec's "File length"
+       // note, a reader must use a length from a trusted source rather than
+       // the underlying storage's reported size, since storage size alone
+       // cannot distinguish a genuinely short file from one truncated by an
+       // attacker who does not also control this metadata.
+       PlaintextLength int64 `json:"plaintext-length"`
+}
+
+// GoEnvelopeEncryptionManager is a generic, format-agnostic 
[EncryptionManager]
+// that provides envelope encryption for arbitrary files (e.g. manifests,
+// manifest lists, Puffin statistics) using a [KeyManagementClient] to wrap
+// and unwrap a fresh AES-256-GCM data encryption key (DEK) per file.
+//
+// Each file is split into fixed-size plaintext blocks and written using the
+// Iceberg AES GCM Stream ("AGS1") format: a magic/block-length header
+// followed by independently authenticated blocks, each carrying its own
+// random nonce and authenticated with an AAD that binds it to the file and
+// to its position. This bounds memory usage and supports random access
+// (Seek/ReadAt) on the decrypted file without buffering or decrypting more
+// than the requested blocks.
+//
+// GoEnvelopeEncryptionManager always encrypts and always decrypts: it fails
+// closed, returning [ErrKeyIDRequired] or [ErrKeyMetadataRequired] rather
+// than silently falling back to plaintext. Use [PlaintextEncryptionManager]
+// for tables or files that are not encrypted.
+//
+// Not interoperable with Java's StandardEncryptionManager: the AGS1 byte
+// stream framing (magic, block layout, nonce placement, AAD) matches the
+// spec, but the two implementations are not cross-readable. Java's reader
+// hard-requires a 1 MiB header block size (a compile-time constant there),
+// while [GoEnvelopeDefaultBlockSize] is 64 KiB, and Java decodes key
+// metadata as a 1-byte version plus Avro-encoded StandardKeyMetadata, while
+// this type encodes key metadata as JSON (goEnvelopeKeyMetadata). Files
+// written by one are not readable by the other.
+type GoEnvelopeEncryptionManager struct {
+       kms       KeyManagementClient
+       dekLength int
+       blockSize int
+}
+
+var _ EncryptionManager = (*GoEnvelopeEncryptionManager)(nil)
+
+// GoEnvelopeManagerOption configures a [GoEnvelopeEncryptionManager] created
+// by [NewGoEnvelopeEncryptionManager].
+type GoEnvelopeManagerOption func(*GoEnvelopeEncryptionManager)
+
+// WithDEKLength overrides the default data encryption key length (in bytes).
+// Valid AES key lengths are 16, 24, or 32 bytes.
+func WithDEKLength(length int) GoEnvelopeManagerOption {
+       return func(m *GoEnvelopeEncryptionManager) { m.dekLength = length }
+}
+
+// WithBlockSize overrides the default plaintext block size (in bytes) used
+// to split files for independent block-level authentication. size must be
+// positive and at most [GoEnvelopeMaxBlockSize].
+func WithBlockSize(size int) GoEnvelopeManagerOption {
+       return func(m *GoEnvelopeEncryptionManager) { m.blockSize = size }
+}
+
+// NewGoEnvelopeEncryptionManager creates a [GoEnvelopeEncryptionManager]
+// backed by kms. kms must not be nil; NewGoEnvelopeEncryptionManager panics
+// if it is.
+func NewGoEnvelopeEncryptionManager(kms KeyManagementClient, opts 
...GoEnvelopeManagerOption) *GoEnvelopeEncryptionManager {
+       if kms == nil {
+               panic("encryption: NewGoEnvelopeEncryptionManager: kms must not 
be nil")
+       }
+
+       m := &GoEnvelopeEncryptionManager{
+               kms:       kms,
+               dekLength: GoEnvelopeDefaultDEKLength,
+               blockSize: GoEnvelopeDefaultBlockSize,
+       }
+       for _, opt := range opts {
+               opt(m)
+       }
+
+       return m
+}
+
+// NewEncryptedOutputFile creates a new AES-GCM block-encrypted output file.
+// keyID identifies the KEK used to wrap the freshly generated per-file DEK,
+// and must be non-empty; otherwise [ErrKeyIDRequired] is returned.
+func (m *GoEnvelopeEncryptionManager) NewEncryptedOutputFile(ctx 
context.Context, writer icebergio.FileWriter, keyID string) 
(EncryptedOutputFile, error) {
+       if keyID == "" {
+               return nil, ErrKeyIDRequired
+       }
+       if m.blockSize <= 0 || m.blockSize > GoEnvelopeMaxBlockSize {
+               return nil, fmt.Errorf("%w: got %d", ErrInvalidBlockSize, 
m.blockSize)
+       }
+       switch m.dekLength {
+       case 16, 24, 32:
+       default:
+               return nil, fmt.Errorf("%w: DEK length must be 16, 24, or 32 
bytes; got %d", ErrInvalidKeyLength, m.dekLength)
+       }
+
+       // The (key, nonce) uniqueness this block format relies on requires a
+       // freshly generated DEK for every file: never cache or reuse
+       // plainDEK/wrappedDEK across calls to NewEncryptedOutputFile.
+       var (
+               plainDEK, wrappedDEK []byte
+               err                  error
+       )
+       if m.kms.SupportsKeyGeneration() {
+               plainDEK, wrappedDEK, err = m.kms.GenerateKey(ctx, keyID, 
m.dekLength)
+               if err != nil {
+                       return nil, fmt.Errorf("encryption: failed to generate 
DEK: %w", err)
+               }
+       } else {
+               plainDEK = make([]byte, m.dekLength)
+               if _, err = io.ReadFull(rand.Reader, plainDEK); err != nil {
+                       return nil, fmt.Errorf("encryption: failed to generate 
DEK: %w", err)
+               }
+               if wrappedDEK, err = m.kms.WrapKey(ctx, keyID, plainDEK); err 
!= nil {
+                       return nil, fmt.Errorf("encryption: failed to wrap DEK: 
%w", err)
+               }
+       }

Review Comment:
   With the raw DEK in `StandardKeyMetadata`, this block and the `UnwrapKey` at 
:356 go away. Generate the DEK locally with `crypto/rand`, the way Java's 
`StandardEncryptedOutputFile.keyMetadata()` does.
   
   Wrapping each file's DEK costs a KMS round trip for every file opened, and 
no cache can help across files. Java and Rust call the KMS once per KEK and 
cache it for an hour. In that model the DEK is protected by encrypting the 
thing that holds the key metadata (manifest, manifest list), with the KEK layer 
stored in the table's `encryption-keys`. That layer and the unwrap cache are 
the follow-up PR, and it has to land before EncryptingFileIO writes anything. 
Nothing calls this manager yet, so a raw DEK in this slice is acceptable.
   
   @laskoviymishka this replaces the per-file wrapping you reviewed on 08-05 
and approved on 08-27 and 09-08 with the Java/Rust key hierarchy.
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to