This is an automated email from the ASF dual-hosted git repository.
laskoviymishka pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-go.git
The following commit(s) were added to refs/heads/main by this push:
new 3ff8d6f35 fix(puffin): poison writers after write failures (#1594)
3ff8d6f35 is described below
commit 3ff8d6f350fc3cfd3dce725a407cff23bf7021fa
Author: Minh Vu <[email protected]>
AuthorDate: Fri Jul 31 16:23:29 2026 +0200
fix(puffin): poison writers after write failures (#1594)
## What changed
Record the first physical write failure on a Puffin writer and reject
every subsequent mutation, blob write, or finish attempt while
preserving the original cause. Route blob and footer writes through the
shared failure-tracking path.
Add fault-injection coverage for partial failures during blob, footer
magic, footer payload, and trailer writes, plus short writes that return
no underlying error. The tests verify that no additional writes occur
afterward.
## Why
An `io.Writer` may consume bytes and return an error or report a short
write. Reusing the Puffin writer afterward left its logical offset
behind the physical stream and could emit corrupt metadata.
## Testing
- `go test ./puffin -count=1`
---------
Signed-off-by: Minh Vu <[email protected]>
---
puffin/puffin_test.go | 81 +++++++++++++++++++++++++++++++++++++++++++++++++
puffin/puffin_writer.go | 37 +++++++++++++++++++---
2 files changed, 114 insertions(+), 4 deletions(-)
diff --git a/puffin/puffin_test.go b/puffin/puffin_test.go
index 689fb1948..4c186ebb4 100644
--- a/puffin/puffin_test.go
+++ b/puffin/puffin_test.go
@@ -19,6 +19,7 @@ package puffin_test
import (
"bytes"
+ "errors"
"math"
"os"
"path"
@@ -39,6 +40,25 @@ func newWriter() (*puffin.Writer, *bytes.Buffer) {
return w, buf
}
+type partialFailWriter struct {
+ bytes.Buffer
+ failCall int
+ calls int
+ err error
+}
+
+func (w *partialFailWriter) Write(data []byte) (int, error) {
+ w.calls++
+ if w.calls == w.failCall {
+ partial := len(data) / 2
+ _, _ = w.Buffer.Write(data[:partial])
+
+ return partial, w.err
+ }
+
+ return w.Buffer.Write(data)
+}
+
func newReader(t *testing.T, buf *bytes.Buffer) *puffin.Reader {
r, err := puffin.NewReader(bytes.NewReader(buf.Bytes()))
require.NoError(t, err)
@@ -53,6 +73,67 @@ func defaultBlobInput() puffin.BlobMetadataInput {
}
}
+func TestWriterRejectsOperationsAfterPartialWriteFailure(t *testing.T) {
+ tests := []struct {
+ name string
+ failCall int
+ failsDuringAdd bool
+ }{
+ // NewWriter writes first, AddBlob writes second, and Finish
writes calls 3-5.
+ {name: "blob", failCall: 2, failsDuringAdd: true},
+ {name: "footer magic", failCall: 3},
+ {name: "footer payload", failCall: 4},
+ {name: "footer trailer", failCall: 5},
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ writeErr := errors.New("injected partial write")
+ output := &partialFailWriter{failCall: tt.failCall,
err: writeErr}
+ writer, err := puffin.NewWriter(output)
+ require.NoError(t, err)
+
+ _, err = writer.AddBlob(defaultBlobInput(),
[]byte("first blob payload"))
+ if tt.failsDuringAdd {
+ require.ErrorIs(t, err, writeErr)
+ } else {
+ require.NoError(t, err)
+ require.ErrorIs(t, writer.Finish(), writeErr)
+ }
+
+ callsAfterFailure := output.calls
+ bytesAfterFailure := output.Len()
+
+ _, err = writer.AddBlob(defaultBlobInput(),
[]byte("second blob payload"))
+ assert.ErrorIs(t, err, writeErr)
+ assert.ErrorIs(t, writer.Finish(), writeErr)
+ assert.ErrorIs(t,
writer.AddProperties(map[string]string{"key": "value"}), writeErr)
+ assert.ErrorIs(t, writer.ClearProperties(), writeErr)
+ assert.ErrorIs(t, writer.SetCreatedBy("test"), writeErr)
+ assert.Equal(t, callsAfterFailure, output.calls)
+ assert.Equal(t, bytesAfterFailure, output.Len())
+ })
+ }
+}
+
+func TestWriterRejectsOperationsAfterShortWrite(t *testing.T) {
+ output := &partialFailWriter{failCall: 2}
+ writer, err := puffin.NewWriter(output)
+ require.NoError(t, err)
+
+ _, err = writer.AddBlob(defaultBlobInput(), []byte("first blob
payload"))
+ require.ErrorContains(t, err, "short write")
+
+ callsAfterFailure := output.calls
+ _, err = writer.AddBlob(defaultBlobInput(), []byte("second blob
payload"))
+ assert.ErrorContains(t, err, "short write")
+ assert.ErrorContains(t, writer.Finish(), "short write")
+ assert.ErrorContains(t, writer.AddProperties(map[string]string{"key":
"value"}), "short write")
+ assert.ErrorContains(t, writer.ClearProperties(), "short write")
+ assert.ErrorContains(t, writer.SetCreatedBy("test"), "short write")
+ assert.Equal(t, callsAfterFailure, output.calls)
+}
+
func validFile() []byte {
w, buf := newWriter()
w.Finish()
diff --git a/puffin/puffin_writer.go b/puffin/puffin_writer.go
index 45da4195a..16230d71d 100644
--- a/puffin/puffin_writer.go
+++ b/puffin/puffin_writer.go
@@ -47,12 +47,16 @@ import (
// return err
// }
// return w.Finish()
+//
+// Writer cannot be reused after a physical write fails. Once poisoned, all
+// subsequent operations return the original write error without changing
state.
type Writer struct {
w io.Writer
offset int64
blobs []BlobMetadata
props map[string]string
done bool
+ failed error
createdBy string
}
@@ -92,6 +96,9 @@ func (w *Writer) AddProperties(props map[string]string) error
{
if w.done {
return errors.New("puffin: cannot set properties: writer
already finalized")
}
+ if w.failed != nil {
+ return fmt.Errorf("puffin: cannot set properties: writer
previously failed: %w", w.failed)
+ }
for k, v := range props {
w.props[k] = v
}
@@ -104,6 +111,9 @@ func (w *Writer) ClearProperties() error {
if w.done {
return errors.New("puffin: cannot clear properties: writer
already finalized")
}
+ if w.failed != nil {
+ return fmt.Errorf("puffin: cannot clear properties: writer
previously failed: %w", w.failed)
+ }
w.props = make(map[string]string)
return nil
@@ -115,6 +125,9 @@ func (w *Writer) SetCreatedBy(createdBy string) error {
if w.done {
return errors.New("puffin: cannot set created-by: writer
already finalized")
}
+ if w.failed != nil {
+ return fmt.Errorf("puffin: cannot set created-by: writer
previously failed: %w", w.failed)
+ }
if createdBy == "" {
return errors.New("puffin: cannot set created-by: value cannot
be empty")
}
@@ -130,6 +143,9 @@ func (w *Writer) AddBlob(input BlobMetadataInput, data
[]byte) (BlobMetadata, er
if w.done {
return BlobMetadata{}, errors.New("puffin: cannot add blob:
writer already finalized")
}
+ if w.failed != nil {
+ return BlobMetadata{}, fmt.Errorf("puffin: cannot add blob:
writer previously failed: %w", w.failed)
+ }
if input.Type == "" {
return BlobMetadata{}, errors.New("puffin: cannot add blob:
type is required")
}
@@ -184,7 +200,7 @@ func (w *Writer) AddBlob(input BlobMetadataInput, data
[]byte) (BlobMetadata, er
Properties: properties,
}
- if err := writeAll(w.w, data); err != nil {
+ if err := w.write(data); err != nil {
return BlobMetadata{}, fmt.Errorf("puffin: write blob: %w", err)
}
@@ -204,6 +220,9 @@ func (w *Writer) Finish() error {
if w.done {
return errors.New("puffin: cannot finish: writer already
finalized")
}
+ if w.failed != nil {
+ return fmt.Errorf("puffin: cannot finish: writer previously
failed: %w", w.failed)
+ }
// Build footer
blobs := w.blobs
@@ -237,12 +256,12 @@ func (w *Writer) Finish() error {
}
// Write footer start magic
- if err := writeAll(w.w, magic[:]); err != nil {
+ if err := w.write(magic[:]); err != nil {
return fmt.Errorf("puffin: write footer magic: %w", err)
}
// Write footer payload
- if err := writeAll(w.w, payload); err != nil {
+ if err := w.write(payload); err != nil {
return fmt.Errorf("puffin: write footer payload: %w", err)
}
@@ -252,7 +271,7 @@ func (w *Writer) Finish() error {
binary.LittleEndian.PutUint32(trailer[4:8], 0) // flags = 0
(uncompressed)
copy(trailer[8:12], magic[:])
- if err := writeAll(w.w, trailer[:]); err != nil {
+ if err := w.write(trailer[:]); err != nil {
return fmt.Errorf("puffin: write footer trailer: %w", err)
}
@@ -261,6 +280,16 @@ func (w *Writer) Finish() error {
return nil
}
+func (w *Writer) write(data []byte) error {
+ if err := writeAll(w.w, data); err != nil {
+ w.failed = err
+
+ return err
+ }
+
+ return nil
+}
+
// writeAll writes all bytes to w or returns an error.
// Handles the io.Writer contract where Write can return n < len(data) without
error.
func writeAll(w io.Writer, data []byte) error {