Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 22 additions & 8 deletions puffin/puffin_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,15 @@ import (
"github.com/pierrec/lz4/v4"
)

// ErrNotPuffinFile is returned by NewReader when the input does not carry a
// Puffin container: the file is too short to hold the header magic, or its
// leading bytes are not the Puffin magic. A file that starts with the magic
// but is truncated or has a broken footer is a damaged Puffin file and does
// NOT wrap this error. Callers that can still consume the payload by other
// means (e.g. a deletion vector addressed by content_offset) may test for it
// with errors.Is.
var ErrNotPuffinFile = errors.New("puffin: not a puffin file")

// ReaderAtSeeker combines io.ReaderAt and io.Seeker for reading Puffin files.
// This interface is implemented by *os.File, *bytes.Reader, and similar types.
type ReaderAtSeeker interface {
Expand Down Expand Up @@ -105,20 +114,25 @@ func NewReader(r ReaderAtSeeker, opts ...ReaderOption) (*Reader, error) {
return nil, fmt.Errorf("puffin: detect file size: %w", err)
}

// Minimum size: header magic + footer magic + footer trailer
// [Magic] + zero for blob + [Magic] + [FooterPayloadSize (assuming ~0)] + [Flags] + [Magic]
minSize := int64(MagicSize + MagicSize + footerTrailerSize)
if size < minSize {
return nil, fmt.Errorf("puffin: file too small (%d bytes, minimum %d)", size, minSize)
// Validate header magic first: only a file that cannot show the Puffin
// magic is "not a Puffin file". A file that starts with the magic but is
// truncated is a damaged Puffin file and must not be mistaken for one.
if size < int64(MagicSize) {
return nil, fmt.Errorf("%w: file too small to hold header magic (%d bytes)", ErrNotPuffinFile, size)
}

// Validate header magic
var headerMagic [MagicSize]byte
if _, err := r.ReadAt(headerMagic[:], 0); err != nil {
return nil, fmt.Errorf("puffin: read header magic: %w", err)
}
if !bytes.Equal(headerMagic[:], magic[:]) {
return nil, errors.New("puffin: invalid header magic")
return nil, fmt.Errorf("%w: invalid header magic", ErrNotPuffinFile)
}

// Minimum size: header magic + footer magic + footer trailer
// [Magic] + zero for blob + [Magic] + [FooterPayloadSize (assuming ~0)] + [Flags] + [Magic]
minSize := int64(MagicSize + MagicSize + footerTrailerSize)
if size < minSize {
return nil, fmt.Errorf("puffin: file too small (%d bytes, minimum %d)", size, minSize)
}

pr := &Reader{
Expand Down
7 changes: 5 additions & 2 deletions puffin/puffin_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -653,8 +653,11 @@ func TestReaderInvalidFile(t *testing.T) {
data func() []byte
wantErr string
}{
// file too small: Minimum valid puffin file has header magic + footer, rejects truncated files.
{"file too small", func() []byte { return []byte("tiny") }, "too small"},
// file too small: Minimum valid puffin file has header magic + footer, rejects truncated files
// (magic is valid here, so this is a damaged Puffin file, not ErrNotPuffinFile).
{"file too small", func() []byte { return []byte("PFA1tiny") }, "too small"},
// too small to hold the magic: cannot be identified as Puffin at all.
{"too small for magic", func() []byte { return []byte("ti") }, "too small to hold header magic"},
// invalid header magic: First 4 bytes must be 'PFA1' to identify puffin format.
{"invalid header magic", func() []byte {
d := validFile()
Expand Down
81 changes: 81 additions & 0 deletions table/dv/deletion_vector.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import (
"math"
"slices"
"strconv"
"sync"

"github.com/apache/iceberg-go"
iceberginternal "github.com/apache/iceberg-go/internal"
Expand Down Expand Up @@ -209,6 +210,14 @@ func ReadDV(fs iceio.IO, dvFile iceberg.DataFile) (*RoaringPositionBitmap, error
return nil, err
}
defer f.Close()
if reader == nil { // not a Puffin container: read the blob directly at content_offset
bitmaps, err := readBareDVs(f, []iceberg.DataFile{dvFile})
if err != nil {
return nil, err
}

return bitmaps[0], nil
}

_, _, manifestReferencedDataFile, contentOffset, contentSize := iceberginternal.BorrowedDataFilePointers(dvFile)
offset, size := *contentOffset, *contentSize
Expand Down Expand Up @@ -252,6 +261,9 @@ func ReadDVs(fs iceio.IO, dvFiles []iceberg.DataFile) ([]*RoaringPositionBitmap,
return nil, err
}
defer f.Close()
if reader == nil { // not a Puffin container: read the blobs directly at content_offset
return readBareDVs(f, dvFiles)
}

blobsByOffset := indexBlobMetadataByOffset(reader.Blobs())
type dvBlobRead struct {
Expand Down Expand Up @@ -327,6 +339,68 @@ func ReadDVs(fs iceio.IO, dvFiles []iceberg.DataFile) ([]*RoaringPositionBitmap,
return bitmaps, nil
}

// bareDVWarned records the bare (non-Puffin) DV files already reported, so a
// snapshot with thousands of such files logs each path once.
var bareDVWarned sync.Map

// readBareDVs reads deletion vectors from an open file that is not a Puffin
// container: the deletion-vector-v1 blobs are addressed directly by the
// manifest's content_offset / content_size_in_bytes, exactly as the Java
// reference reader (BaseDeleteLoader.readDV) does, which never consults the
// Puffin footer. Databricks writes DVs for IcebergCompatV3 (UniForm) tables
// this way: a Delta deletion_vector_*.bin file with a one-byte version
// prefix and no Puffin header or footer.
//
// Without footer metadata the blob's type, referenced data file and
// cardinality property cannot be cross-checked; the blob's own length,
// magic and CRC-32 are still verified by DeserializeDV, and the decoded
// cardinality is validated against the manifest record_count.
//
// dvFiles must already have passed validateDVFile (both callers validate
// before dispatching here) and must all point at f. Blobs are read in
// content_offset order to avoid backward seeks; results keep dvFiles order.
func readBareDVs(f iceio.File, dvFiles []iceberg.DataFile) ([]*RoaringPositionBitmap, error) {
filePath := dvFiles[0].FilePath()

order := make([]int, len(dvFiles))
for i := range order {
order[i] = i
}
offsetOf := func(i int) int64 {
_, _, _, contentOffset, _ := iceberginternal.BorrowedDataFilePointers(dvFiles[i])

return *contentOffset
}
slices.SortFunc(order, func(a, b int) int { return cmp.Compare(offsetOf(a), offsetOf(b)) })

bitmaps := make([]*RoaringPositionBitmap, len(dvFiles))
for _, i := range order {
dvFile := dvFiles[i]
_, _, _, contentOffset, contentSize := iceberginternal.BorrowedDataFilePointers(dvFile)
offset, size := *contentOffset, *contentSize

data := make([]byte, size)
if _, err := f.ReadAt(data, offset); err != nil {
return nil, fmt.Errorf("%w: DV file %s is not a Puffin container; direct read of %d bytes at offset %d: %w",
ErrInvalidDeletionVector, filePath, size, offset, err)
}

bitmap, err := DeserializeDV(data, dvFile.Count())
if err != nil {
return nil, fmt.Errorf("%w: DV file %s is not a Puffin container; blob at offset %d: %w",
ErrInvalidDeletionVector, filePath, offset, err)
}
bitmaps[i] = bitmap

if _, seen := bareDVWarned.LoadOrStore(filePath, struct{}{}); !seen {
slog.Warn("DV file is not a Puffin container; reading deletion-vector-v1 blobs directly at content_offset, footer metadata validation skipped",
"dv_file", filePath)
}
}

return bitmaps, nil
}

func validateDVFile(dvFile iceberg.DataFile) error {
if dvFile.FileFormat() != iceberg.PuffinFile {
return fmt.Errorf("expected PUFFIN format for deletion vector, got %s", dvFile.FileFormat())
Expand All @@ -348,13 +422,20 @@ func validateDVFile(dvFile iceberg.DataFile) error {
return nil
}

// openDVReader opens the DV file and its Puffin footer. When the file is not
// a Puffin container (puffin.ErrNotPuffinFile) it returns a nil reader and
// the still-open file, so the caller can read bare blobs without a second
// Open round-trip; the caller owns closing f whenever err is nil.
func openDVReader(fs iceio.IO, filePath string) (*puffin.Reader, iceio.File, error) {
f, err := fs.Open(filePath)
if err != nil {
return nil, nil, fmt.Errorf("open DV file %s: %w", filePath, err)
}

reader, err := puffin.NewReader(f)
if errors.Is(err, puffin.ErrNotPuffinFile) {
return nil, f, nil
}
if err != nil {
_ = f.Close()

Expand Down
120 changes: 119 additions & 1 deletion table/dv/deletion_vector_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import (
"os"
"path/filepath"
"strconv"
"strings"
"testing"

"github.com/apache/iceberg-go"
Expand Down Expand Up @@ -878,7 +879,124 @@ func TestReadDVInvalidPuffin(t *testing.T) {

offset, size := int64(4), int64(16)
_, err := ReadDV(iceio.LocalFS{}, newDVTestFile(path, 0, &offset, &size))
assert.ErrorContains(t, err, "create puffin reader")
require.ErrorIs(t, err, ErrInvalidDeletionVector)
assert.ErrorContains(t, err, "not a Puffin container")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

TestReadDVInvalidPuffin now feeds a 17-byte file, which is the too-small to bare-path case, so nothing here pins the tighter invariant: a file that starts with real Puffin magic but has a broken/truncated footer should fail in readFooter with an error that does NOT wrap ErrNotPuffinFile, i.e. the bare fallback must not activate. I'd add that case so a future change that accidentally wraps footer errors in ErrNotPuffinFile can't silently route a corrupt Puffin file to the bare reader.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in 12b915c: TestReadDVTruncatedPuffinDoesNotFallBack covers magic-only, truncated-footer and corrupt-footer files — each fails in the Puffin reader, errors.Is(err, puffin.ErrNotPuffinFile) is false, and the error does not carry "not a Puffin container".

}

// Why: the bare-blob fallback must stay pinned to inputs that are genuinely
// not Puffin. A file that starts with the Puffin magic but is truncated or has
// a broken footer is a damaged Puffin file; routing it to the bare reader would
// hide a partial upload behind a "format mismatch" diagnostic.
// Condition: valid header magic, footer missing or corrupt.
// Assertion: ReadDV fails in the Puffin reader, and the error does not wrap
// puffin.ErrNotPuffinFile.
func TestReadDVTruncatedPuffinDoesNotFallBack(t *testing.T) {
dir := t.TempDir()
dvBlobBytes := readDVTestData(t, "small-alternating-values-position-index.bin")
goodPath, meta := writePuffinWithDVBlob(t, dir, dvBlobBytes)
good, err := os.ReadFile(goodPath)
require.NoError(t, err)

cases := map[string][]byte{
"magic only": good[:4],
"truncated footer": good[:len(good)-6],
"corrupt footer": append(append([]byte{}, good[:len(good)-12]...), []byte("xxxxxxxxxxxx")...),
}
for name, raw := range cases {
t.Run(name, func(t *testing.T) {
path := filepath.Join(dir, strings.ReplaceAll(name, " ", "_")+".puffin")
require.NoError(t, os.WriteFile(path, raw, 0o644))
offset, size := meta.Offset, meta.Length
_, err := ReadDV(iceio.LocalFS{}, newDVTestFile(path, 5, &offset, &size))
require.Error(t, err)
assert.False(t, errors.Is(err, puffin.ErrNotPuffinFile), "damaged Puffin file must not be treated as bare: %v", err)
assert.ErrorContains(t, err, "create puffin reader")
assert.NotContains(t, err.Error(), "not a Puffin container")
})
}
}

// Why: Databricks writes deletion vectors for IcebergCompatV3 (UniForm) tables
// as a Delta deletion_vector_*.bin — a one-byte version prefix followed by
// deletion-vector-v1 blobs, with no Puffin header or footer — and the manifest
// addresses the blob with content_offset / content_size_in_bytes. The Java
// reference reader reads these directly; so must we.
// Condition: the DV file has no Puffin container but the manifest range holds a valid blob.
// Assertion: ReadDV/ReadDVs decode the blob and validate cardinality against record_count.
func TestReadDVBareBlobWithoutPuffinContainer(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "deletion_vector_0001.bin")

first := NewRoaringPositionBitmap()
first.Set(1)
first.Set(9)
firstData, err := SerializeDV(first)
require.NoError(t, err)
second := NewRoaringPositionBitmap()
second.Set(7)
secondData, err := SerializeDV(second)
require.NoError(t, err)

raw := append([]byte{0x01}, firstData...) // Delta DV file version byte, then blobs back to back
secondOffset := int64(len(raw))
raw = append(raw, secondData...)
require.NoError(t, os.WriteFile(path, raw, 0o644))

firstOffset, firstSize := int64(1), int64(len(firstData))
bm, err := ReadDV(iceio.LocalFS{}, newDVTestFile(path, 2, &firstOffset, &firstSize))
require.NoError(t, err)
assert.Equal(t, int64(2), bm.Cardinality())
assert.True(t, bm.Contains(1))
assert.True(t, bm.Contains(9))

secondSize := int64(len(secondData))
files := []iceberg.DataFile{
newDVTestFile(path, 2, &firstOffset, &firstSize),
newDVTestFile(path, 1, &secondOffset, &secondSize),
}
bitmaps, err := ReadDVs(iceio.LocalFS{}, files)
require.NoError(t, err)
require.Len(t, bitmaps, 2)
assert.Equal(t, int64(2), bitmaps[0].Cardinality())
assert.True(t, bitmaps[1].Contains(7))

t.Run("cardinality still validated against record_count", func(t *testing.T) {
_, err := ReadDV(iceio.LocalFS{}, newDVTestFile(path, 3, &firstOffset, &firstSize))
require.ErrorIs(t, err, ErrInvalidDeletionVector)
assert.ErrorContains(t, err, "cardinality mismatch")
})

t.Run("range beyond file", func(t *testing.T) {
badOffset, badSize := int64(1), int64(len(raw)+8)
_, err := ReadDV(iceio.LocalFS{}, newDVTestFile(path, 2, &badOffset, &badSize))
require.ErrorIs(t, err, ErrInvalidDeletionVector)
assert.ErrorContains(t, err, "direct read")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The subtests cover ReadAt failing (range beyond file) and post-deserialize cardinality mismatch, but not the in-between: ReadAt succeeds and DeserializeDV rejects the blob for a DV-level reason (bad magic / CRC), which is the blob at offset wrap in readBareDVs. I'd add a subtest that flips the last 4 CRC bytes and asserts ErrorIs(err, ErrInvalidDeletionVector) plus the "blob at offset" message.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in 12b915c: added the "corrupt blob CRC" subtest in TestReadDVBareBlobWithoutPuffinContainer — flips the blob's CRC bytes, asserts ErrorIs(err, ErrInvalidDeletionVector), "blob at offset" and "CRC mismatch".

})

t.Run("corrupt blob CRC", func(t *testing.T) {
bad := append([]byte{}, raw...)
for i := int(firstOffset+firstSize) - 4; i < int(firstOffset+firstSize); i++ {
bad[i] ^= 0xFF
}
badPath := filepath.Join(dir, "deletion_vector_crc.bin")
require.NoError(t, os.WriteFile(badPath, bad, 0o644))
_, err := ReadDV(iceio.LocalFS{}, newDVTestFile(badPath, 2, &firstOffset, &firstSize))
require.ErrorIs(t, err, ErrInvalidDeletionVector)
assert.ErrorContains(t, err, "blob at offset")
assert.ErrorContains(t, err, "CRC mismatch")
})

t.Run("blobs read in offset order, results in input order", func(t *testing.T) {
reversed := []iceberg.DataFile{
newDVTestFile(path, 1, &secondOffset, &secondSize),
newDVTestFile(path, 2, &firstOffset, &firstSize),
}
bitmaps, err := ReadDVs(iceio.LocalFS{}, reversed)
require.NoError(t, err)
require.Len(t, bitmaps, 2)
assert.True(t, bitmaps[0].Contains(7))
assert.Equal(t, int64(2), bitmaps[1].Cardinality())
})
}

// Why: offset, size, and cardinality cannot prove that the selected Puffin blob
Expand Down
Loading