diff --git a/puffin/puffin_reader.go b/puffin/puffin_reader.go index ab2c6c0e6..bbcdff586 100644 --- a/puffin/puffin_reader.go +++ b/puffin/puffin_reader.go @@ -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 { @@ -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{ diff --git a/puffin/puffin_test.go b/puffin/puffin_test.go index 0132e1f6a..907c8212d 100644 --- a/puffin/puffin_test.go +++ b/puffin/puffin_test.go @@ -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() diff --git a/table/dv/deletion_vector.go b/table/dv/deletion_vector.go index abaa707a6..206b91338 100644 --- a/table/dv/deletion_vector.go +++ b/table/dv/deletion_vector.go @@ -28,6 +28,7 @@ import ( "math" "slices" "strconv" + "sync" "github.com/apache/iceberg-go" iceberginternal "github.com/apache/iceberg-go/internal" @@ -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 @@ -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 { @@ -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()) @@ -348,6 +422,10 @@ 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 { @@ -355,6 +433,9 @@ func openDVReader(fs iceio.IO, filePath string) (*puffin.Reader, iceio.File, err } reader, err := puffin.NewReader(f) + if errors.Is(err, puffin.ErrNotPuffinFile) { + return nil, f, nil + } if err != nil { _ = f.Close() diff --git a/table/dv/deletion_vector_test.go b/table/dv/deletion_vector_test.go index bd9814e65..f75594d08 100644 --- a/table/dv/deletion_vector_test.go +++ b/table/dv/deletion_vector_test.go @@ -27,6 +27,7 @@ import ( "os" "path/filepath" "strconv" + "strings" "testing" "github.com/apache/iceberg-go" @@ -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") +} + +// 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") + }) + + 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