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 027b116da fix(partitions): validate decoded partition specs (#1593)
027b116da is described below

commit 027b116dae5863a509bfdc37c533c43e4a006583
Author: Minh Vu <[email protected]>
AuthorDate: Fri Jul 31 16:16:09 2026 +0200

    fix(partitions): validate decoded partition specs (#1593)
    
    ## What changed
    
    Apply structural partition-field and partition-spec validation during
    JSON decoding. Reject invalid spec IDs, missing or non-positive source
    IDs, empty and duplicate names, duplicate field IDs, redundant
    definitions, and invalid transforms.
    
    The validation remains independent of the current schema so historical
    specs may still reference dropped source columns.
    
    ## Why
    
    The JSON path directly populated internal fields and bypassed invariants
    enforced by the constructor. Malformed standalone specs and historical
    table specs could therefore enter the model and be serialized again.
    
    ## Testing
    
    - `go test . -run
    
'Test(PartitionFieldUnmarshalJSON|PartitionSpecUnmarshalRejectsInvalidStructure|DeserializePartitionSpec)'
    -count=1`\n- `go test ./table -run
    'TestRejectStructurallyInvalidHistoricalPartitionSpec$' -count=1`
    
    ---------
    
    Signed-off-by: Minh Vu <[email protected]>
---
 partitions.go                   | 64 +++++++++++++++++++++++++++---
 partitions_test.go              | 87 +++++++++++++++++++++++++++++++++++++----
 table/metadata_internal_test.go | 19 +++++++++
 3 files changed, 156 insertions(+), 14 deletions(-)

diff --git a/partitions.go b/partitions.go
index 6f04299ae..7f2e79ea2 100644
--- a/partitions.go
+++ b/partitions.go
@@ -119,9 +119,11 @@ func (p *PartitionField) UnmarshalJSON(b []byte) error {
 
        if _, ok := raw["source-id"]; ok {
                if _, ok := raw["source-ids"]; ok {
-                       return errors.New("partition field cannot contain both 
source-id and source-ids")
+                       return fmt.Errorf("%w: partition field cannot contain 
both source-id and source-ids", ErrInvalidPartitionSpec)
                }
        }
+       _, hasSourceID := raw["source-id"]
+       _, hasSourceIDs := raw["source-ids"]
 
        aux := struct {
                SourceID        int    `json:"source-id"`
@@ -138,15 +140,36 @@ func (p *PartitionField) UnmarshalJSON(b []byte) error {
        p.FieldID = aux.FieldID
        p.Name = aux.Name
 
-       if len(aux.SourceIDs) > 0 {
+       var err error
+       if p.Transform, err = ParseTransform(aux.TransformString); err != nil {
+               return fmt.Errorf("%w: %w", ErrInvalidPartitionSpec, err)
+       }
+       if err := validateTransform(p.Transform); err != nil {
+               return fmt.Errorf("%w: %w", ErrInvalidPartitionSpec, err)
+       }
+
+       if hasSourceIDs && len(aux.SourceIDs) == 0 {
+               return fmt.Errorf("%w: partition source-ids cannot be empty", 
ErrInvalidPartitionSpec)
+       }
+       if !hasSourceID && !hasSourceIDs {
+               if _, isVoid := p.Transform.(VoidTransform); !isVoid {
+                       return fmt.Errorf("%w: partition field requires 
source-id or source-ids", ErrInvalidPartitionSpec)
+               }
+               // Preserve compatibility with historical source-less void 
tombstones.
+               p.SourceIDs = []int{0}
+       } else if len(aux.SourceIDs) > 0 {
                p.SourceIDs = aux.SourceIDs
        } else {
                p.SourceIDs = []int{aux.SourceID}
        }
-
-       var err error
-       if p.Transform, err = ParseTransform(aux.TransformString); err != nil {
-               return err
+       for _, sourceID := range p.SourceIDs {
+               _, isVoid := p.Transform.(VoidTransform)
+               if sourceID <= 0 && (!isVoid || hasSourceID || hasSourceIDs) {
+                       return fmt.Errorf("%w: partition source ID must be 
positive: %d", ErrInvalidPartitionSpec, sourceID)
+               }
+       }
+       if p.Name == "" {
+               return fmt.Errorf("%w: partition name cannot be empty", 
ErrInvalidPartitionSpec)
        }
 
        return nil
@@ -439,6 +462,9 @@ func (ps *PartitionSpec) UnmarshalJSON(b []byte) error {
        if err := json.Unmarshal(b, &aux); err != nil {
                return err
        }
+       if aux.ID < 0 {
+               return fmt.Errorf("%w: spec ID must be non-negative: %d", 
ErrInvalidPartitionSpec, aux.ID)
+       }
 
        fields := make([]PartitionField, len(aux.Fields))
        for i, rawField := range aux.Fields {
@@ -459,6 +485,9 @@ func (ps *PartitionSpec) UnmarshalJSON(b []byte) error {
                        return err
                }
        }
+       if err := validatePartitionFields(fields); err != nil {
+               return err
+       }
 
        ps.id, ps.fields = aux.ID, fields
        if err := ps.assignPartitionFieldIds(nil); err != nil {
@@ -469,6 +498,29 @@ func (ps *PartitionSpec) UnmarshalJSON(b []byte) error {
        return nil
 }
 
+func validatePartitionFields(fields []PartitionField) error {
+       names := make(map[string]struct{}, len(fields))
+       definitions := make(map[string]struct{}, len(fields))
+       for _, field := range fields {
+               if _, ok := names[field.Name]; ok {
+                       return fmt.Errorf("%w: duplicate partition name: %s", 
ErrInvalidPartitionSpec, field.Name)
+               }
+               names[field.Name] = struct{}{}
+               if _, isVoid := field.Transform.(VoidTransform); isVoid {
+                       continue
+               }
+
+               definition := fmt.Sprintf("%v:%s", field.SourceIDs, 
field.Transform)
+               if _, ok := definitions[definition]; ok {
+                       return fmt.Errorf("%w: redundant partition field for 
source IDs %v and transform %s",
+                               ErrInvalidPartitionSpec, field.SourceIDs, 
field.Transform)
+               }
+               definitions[definition] = struct{}{}
+       }
+
+       return nil
+}
+
 func (ps *PartitionSpec) initialize() {
        ps.sourceIdToFields = make(map[int][]PartitionField)
 
diff --git a/partitions_test.go b/partitions_test.go
index 6b4ef641a..b5c0ece37 100644
--- a/partitions_test.go
+++ b/partitions_test.go
@@ -635,10 +635,11 @@ func TestPartitionFieldUnmarshalJSON(t *testing.T) {
                var field iceberg.PartitionField
                err := json.Unmarshal([]byte(jsonData), &field)
                require.Error(t, err)
-               assert.EqualError(t, err, "partition field cannot contain both 
source-id and source-ids")
+               require.ErrorIs(t, err, iceberg.ErrInvalidPartitionSpec)
+               assert.ErrorContains(t, err, "partition field cannot contain 
both source-id and source-ids")
        })
 
-       t.Run("unmarshal with no source id", func(t *testing.T) {
+       t.Run("unmarshal source-less void tombstone", func(t *testing.T) {
                jsonData := `
                {
                        "field-id": 1003,
@@ -646,11 +647,81 @@ func TestPartitionFieldUnmarshalJSON(t *testing.T) {
                        "name": "void"
                }`
                var field iceberg.PartitionField
-               err := json.Unmarshal([]byte(jsonData), &field)
-               require.NoError(t, err)
-               assert.Zero(t, field.SourceID())
-               assert.Equal(t, 1003, field.FieldID)
-               assert.Equal(t, "void", field.Name)
-               assert.Equal(t, iceberg.VoidTransform{}, field.Transform)
+               require.NoError(t, json.Unmarshal([]byte(jsonData), &field))
+               assert.Equal(t, 0, field.SourceID())
        })
+
+       t.Run("unmarshal void with source id", func(t *testing.T) {
+               var field iceberg.PartitionField
+               require.NoError(t, json.Unmarshal([]byte(`{
+                       "source-id": 1, "field-id": 1003, "transform": "void", 
"name": "void"
+               }`), &field))
+               assert.Equal(t, 1, field.SourceID())
+       })
+}
+
+func TestPartitionSpecUnmarshalRejectsInvalidStructure(t *testing.T) {
+       tests := []struct {
+               name    string
+               data    string
+               message string
+       }{
+               {
+                       name:    "negative spec ID",
+                       data:    `{"spec-id":-1,"fields":[]}`,
+                       message: "spec ID must be non-negative",
+               },
+               {
+                       name:    "zero source ID",
+                       data:    
`{"spec-id":0,"fields":[{"source-id":0,"field-id":1000,"name":"part","transform":"identity"}]}`,
+                       message: "source ID must be positive",
+               },
+               {
+                       name:    "empty source IDs",
+                       data:    
`{"spec-id":0,"fields":[{"source-ids":[],"field-id":1000,"name":"part","transform":"identity"}]}`,
+                       message: "source-ids cannot be empty",
+               },
+               {
+                       name:    "empty name",
+                       data:    
`{"spec-id":0,"fields":[{"source-id":1,"field-id":1000,"name":"","transform":"identity"}]}`,
+                       message: "partition name cannot be empty",
+               },
+               {
+                       name:    "duplicate names",
+                       data:    
`{"spec-id":0,"fields":[{"source-id":1,"field-id":1000,"name":"part","transform":"identity"},{"source-id":2,"field-id":1001,"name":"part","transform":"identity"}]}`,
+                       message: "duplicate partition name",
+               },
+               {
+                       name:    "duplicate field IDs",
+                       data:    
`{"spec-id":0,"fields":[{"source-id":1,"field-id":1000,"name":"first","transform":"identity"},{"source-id":2,"field-id":1000,"name":"second","transform":"identity"}]}`,
+                       message: "duplicate field ID provided",
+               },
+               {
+                       name:    "redundant fields",
+                       data:    
`{"spec-id":0,"fields":[{"source-id":1,"field-id":1000,"name":"first","transform":"identity"},{"source-id":1,"field-id":1001,"name":"second","transform":"identity"}]}`,
+                       message: "redundant partition field",
+               },
+               {
+                       name:    "invalid bucket count",
+                       data:    
`{"spec-id":0,"fields":[{"source-id":1,"field-id":1000,"name":"part","transform":"bucket[0]"}]}`,
+                       message: "invalid transform syntax",
+               },
+       }
+
+       for _, tt := range tests {
+               t.Run(tt.name, func(t *testing.T) {
+                       var spec iceberg.PartitionSpec
+                       err := json.Unmarshal([]byte(tt.data), &spec)
+                       require.ErrorIs(t, err, iceberg.ErrInvalidPartitionSpec)
+                       assert.ErrorContains(t, err, tt.message)
+               })
+       }
+}
+
+func TestPartitionSpecUnmarshalAllowsRepeatedVoidTransforms(t *testing.T) {
+       data := 
`{"spec-id":0,"fields":[{"source-id":1,"field-id":1000,"name":"first","transform":"void"},{"source-id":1,"field-id":1001,"name":"second","transform":"void"}]}`
+
+       var spec iceberg.PartitionSpec
+       require.NoError(t, json.Unmarshal([]byte(data), &spec))
+       assert.Equal(t, 2, spec.NumFields())
 }
diff --git a/table/metadata_internal_test.go b/table/metadata_internal_test.go
index bfe821968..087bf6d73 100644
--- a/table/metadata_internal_test.go
+++ b/table/metadata_internal_test.go
@@ -602,6 +602,25 @@ func TestRejectDuplicateSnapshotIDs(t *testing.T) {
        }
 }
 
+func TestRejectStructurallyInvalidHistoricalPartitionSpec(t *testing.T) {
+       var metadata map[string]any
+       decoder := json.NewDecoder(strings.NewReader(ExampleTableMetadataV2))
+       decoder.UseNumber()
+       require.NoError(t, decoder.Decode(&metadata))
+
+       specs := metadata["partition-specs"].([]any)
+       metadata["partition-specs"] = append(slices.Clone(specs), 
map[string]any{
+               "spec-id": json.Number("-1"),
+               "fields":  []any{},
+       })
+       data, err := json.Marshal(metadata)
+       require.NoError(t, err)
+
+       _, err = ParseMetadataBytes(data)
+       require.ErrorIs(t, err, iceberg.ErrInvalidPartitionSpec)
+       assert.ErrorContains(t, err, "spec ID must be non-negative")
+}
+
 func TestSortOrderNotFound(t *testing.T) {
        metadataSortOrderNotFound := `{
         "format-version": 2,

Reply via email to