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 f3045ff16 fix(table): preserve branch/tag rentention settings on
snapshot ref updates (#1601)
f3045ff16 is described below
commit f3045ff16a9b8be5c05660d4da951b4f282d14a4
Author: Badal Prasad Singh <[email protected]>
AuthorDate: Fri Jul 31 00:27:35 2026 +0530
fix(table): preserve branch/tag rentention settings on snapshot ref updates
(#1601)
## Description
Fixes #1600, silent loss of _snapshot-ref_ retention settings
(`min-snapshots-to-keep`, `max-snapshot-age-ms`, `max-ref-age-ms`) on
every commit and rollback.
Have added the necessary tests for it.
---------
Signed-off-by: badalprasadsingh <[email protected]>
---
table/metadata.go | 18 ++++
table/metadata_builder_internal_test.go | 166 ++++++++++++++++++++++++++++++++
table/snapshot_producers.go | 12 ++-
table/transaction.go | 2 +-
table/transaction_internal_test.go | 51 ++++++++++
5 files changed, 243 insertions(+), 6 deletions(-)
diff --git a/table/metadata.go b/table/metadata.go
index 4aa87c94b..7b7a18749 100644
--- a/table/metadata.go
+++ b/table/metadata.go
@@ -963,6 +963,24 @@ func (b *MetadataBuilder) SetSnapshotRef(
return nil
}
+func (b *MetadataBuilder) NewRetainingSnapshotRefUpdate(name string,
snapshotID int64, refType RefType) *setSnapshotRefUpdate {
+ var maxRefAgeMs, maxSnapshotAgeMs int64
+ var minSnapshotsToKeep int
+ if existing, ok := b.refs[name]; ok && existing.SnapshotRefType ==
refType {
+ if existing.MaxRefAgeMs != nil {
+ maxRefAgeMs = *existing.MaxRefAgeMs
+ }
+ if existing.MaxSnapshotAgeMs != nil {
+ maxSnapshotAgeMs = *existing.MaxSnapshotAgeMs
+ }
+ if existing.MinSnapshotsToKeep != nil {
+ minSnapshotsToKeep = *existing.MinSnapshotsToKeep
+ }
+ }
+
+ return NewSetSnapshotRefUpdate(name, snapshotID, refType, maxRefAgeMs,
maxSnapshotAgeMs, minSnapshotsToKeep)
+}
+
func (b *MetadataBuilder) RemoveSnapshotRef(name string) error {
if _, found := b.refs[name]; !found {
return fmt.Errorf("%w: snapshot ref not found: %s",
iceberg.ErrInvalidArgument, name)
diff --git a/table/metadata_builder_internal_test.go
b/table/metadata_builder_internal_test.go
index 3b6baee5b..b3b37be9e 100644
--- a/table/metadata_builder_internal_test.go
+++ b/table/metadata_builder_internal_test.go
@@ -482,6 +482,172 @@ func TestSetRef(t *testing.T) {
require.Len(t, builder.snapshotLog, 1)
}
+func TestSetSnapshotRefUpdateApplyPreservesRetention(t *testing.T) {
+ builder := builderWithoutChanges(2)
+ schemaID := 0
+ snapshot1 := Snapshot{
+ SnapshotID: 1,
+ ParentSnapshotID: nil,
+ SequenceNumber: 0,
+ TimestampMs: builder.base.LastUpdatedMillis() + 1,
+ ManifestList: "/snap-1.avro",
+ Summary: &Summary{Operation: OpAppend},
+ SchemaID: &schemaID,
+ }
+ parentID := int64(1)
+ snapshot2 := Snapshot{
+ SnapshotID: 2,
+ ParentSnapshotID: &parentID,
+ SequenceNumber: 1,
+ TimestampMs: builder.base.LastUpdatedMillis() + 2,
+ ManifestList: "/snap-2.avro",
+ Summary: &Summary{Operation: OpAppend},
+ SchemaID: &schemaID,
+ }
+
+ const (
+ minKeep = 5
+ maxSnapAgeMs = int64(172800000) // 2 days
+ maxRefAgeMsIn = int64(604800000) // 7 days
+ )
+
+ require.NoError(t, builder.AddSnapshot(&snapshot1))
+ require.NoError(t, builder.SetSnapshotRef(
+ MainBranch, 1, BranchRef,
+ WithMinSnapshotsToKeep(minKeep),
+ WithMaxSnapshotAgeMs(maxSnapAgeMs),
+ WithMaxRefAgeMs(maxRefAgeMsIn),
+ ))
+ require.NoError(t, builder.AddSnapshot(&snapshot2))
+
+ // build the update the way the commit/rollback producers do.
+ upd := builder.NewRetainingSnapshotRefUpdate(MainBranch, 2, BranchRef)
+
+ require.Equal(t, minKeep, upd.MinSnapshotsToKeep)
+ require.Equal(t, maxSnapAgeMs, upd.MaxSnapshotAgeMs)
+ require.Equal(t, maxRefAgeMsIn, upd.MaxRefAgeMs)
+
+ require.NoError(t, upd.Apply(&builder))
+ ref := builder.refs[MainBranch]
+ require.Equal(t, int64(2), ref.SnapshotID, "snapshot pointer should
advance")
+ require.NotNil(t, ref.MinSnapshotsToKeep, "min-snapshots-to-keep must
not be wiped out")
+ require.Equal(t, minKeep, *ref.MinSnapshotsToKeep)
+ require.NotNil(t, ref.MaxSnapshotAgeMs, "max-snapshot-age-ms must not
be wiped out")
+ require.Equal(t, maxSnapAgeMs, *ref.MaxSnapshotAgeMs)
+ require.NotNil(t, ref.MaxRefAgeMs, "max-ref-age-ms must not be wiped
out")
+ require.Equal(t, maxRefAgeMsIn, *ref.MaxRefAgeMs)
+}
+
+func TestSetSnapshotRefUpdateApplyClearsRetentionWhenAbsent(t *testing.T) {
+ builder := builderWithoutChanges(2)
+ schemaID := 0
+ snapshot1 := Snapshot{
+ SnapshotID: 1,
+ ParentSnapshotID: nil,
+ SequenceNumber: 0,
+ TimestampMs: builder.base.LastUpdatedMillis() + 1,
+ ManifestList: "/snap-1.avro",
+ Summary: &Summary{Operation: OpAppend},
+ SchemaID: &schemaID,
+ }
+ parentID := int64(1)
+ snapshot2 := Snapshot{
+ SnapshotID: 2,
+ ParentSnapshotID: &parentID,
+ SequenceNumber: 1,
+ TimestampMs: builder.base.LastUpdatedMillis() + 2,
+ ManifestList: "/snap-2.avro",
+ Summary: &Summary{Operation: OpAppend},
+ SchemaID: &schemaID,
+ }
+ require.NoError(t, builder.AddSnapshot(&snapshot1))
+ require.NoError(t, builder.SetSnapshotRef(MainBranch, 1, BranchRef,
WithMinSnapshotsToKeep(5)))
+ require.NoError(t, builder.AddSnapshot(&snapshot2))
+
+ require.NoError(t, NewSetSnapshotRefUpdate(MainBranch, 2, BranchRef, 0,
0, 0).Apply(&builder))
+ ref := builder.refs[MainBranch]
+ require.Equal(t, int64(2), ref.SnapshotID)
+ require.Nil(t, ref.MinSnapshotsToKeep, "bare set-ref must clear
retention (pure replace)")
+ require.Nil(t, ref.MaxSnapshotAgeMs)
+ require.Nil(t, ref.MaxRefAgeMs)
+}
+
+func TestSetSnapshotRefBranchToTagDropsAllRetention(t *testing.T) {
+ builder := builderWithoutChanges(2)
+ schemaID := 0
+ snapshot := Snapshot{
+ SnapshotID: 1,
+ ParentSnapshotID: nil,
+ SequenceNumber: 0,
+ TimestampMs: builder.base.LastUpdatedMillis() + 1,
+ ManifestList: "/snap-1.avro",
+ Summary: &Summary{Operation: OpAppend},
+ SchemaID: &schemaID,
+ }
+ require.NoError(t, builder.AddSnapshot(&snapshot))
+ require.NoError(t, builder.SetSnapshotRef(
+ MainBranch, 1, BranchRef,
+ WithMinSnapshotsToKeep(5),
+ WithMaxSnapshotAgeMs(int64(172800000)),
+ WithMaxRefAgeMs(int64(604800000)),
+ ))
+
+ // The retaining helper must not carry retention across a type change.
+ upd := builder.NewRetainingSnapshotRefUpdate(MainBranch, 1, TagRef)
+ require.Zero(t, upd.MinSnapshotsToKeep)
+ require.Zero(t, upd.MaxSnapshotAgeMs)
+ require.Zero(t, upd.MaxRefAgeMs)
+
+ require.NoError(t, builder.SetSnapshotRef(MainBranch, 1, TagRef))
+ ref := builder.refs[MainBranch]
+ require.Equal(t, TagRef, ref.SnapshotRefType)
+ require.Nil(t, ref.MinSnapshotsToKeep)
+ require.Nil(t, ref.MaxSnapshotAgeMs)
+ require.Nil(t, ref.MaxRefAgeMs)
+}
+
+func TestNewRetainingSnapshotRefUpdateTagPreservesMaxRefAge(t *testing.T) {
+ builder := builderWithoutChanges(2)
+ schemaID := 0
+ snapshot1 := Snapshot{
+ SnapshotID: 1,
+ ParentSnapshotID: nil,
+ SequenceNumber: 0,
+ TimestampMs: builder.base.LastUpdatedMillis() + 1,
+ ManifestList: "/snap-1.avro",
+ Summary: &Summary{Operation: OpAppend},
+ SchemaID: &schemaID,
+ }
+ parentID := int64(1)
+ snapshot2 := Snapshot{
+ SnapshotID: 2,
+ ParentSnapshotID: &parentID,
+ SequenceNumber: 1,
+ TimestampMs: builder.base.LastUpdatedMillis() + 2,
+ ManifestList: "/snap-2.avro",
+ Summary: &Summary{Operation: OpAppend},
+ SchemaID: &schemaID,
+ }
+
+ const maxRefAgeMsIn = int64(604800000)
+ require.NoError(t, builder.AddSnapshot(&snapshot1))
+ require.NoError(t, builder.SetSnapshotRef("release", 1, TagRef,
WithMaxRefAgeMs(maxRefAgeMsIn)))
+ require.NoError(t, builder.AddSnapshot(&snapshot2))
+
+ upd := builder.NewRetainingSnapshotRefUpdate("release", 2, TagRef)
+ require.Equal(t, maxRefAgeMsIn, upd.MaxRefAgeMs)
+ require.Zero(t, upd.MinSnapshotsToKeep)
+ require.Zero(t, upd.MaxSnapshotAgeMs)
+
+ require.NoError(t, upd.Apply(&builder))
+ ref := builder.refs["release"]
+ require.Equal(t, int64(2), ref.SnapshotID)
+ require.NotNil(t, ref.MaxRefAgeMs)
+ require.Equal(t, maxRefAgeMsIn, *ref.MaxRefAgeMs)
+ require.Nil(t, ref.MinSnapshotsToKeep)
+ require.Nil(t, ref.MaxSnapshotAgeMs)
+}
+
func TestSetRefRejectsInvalidTypeAndTagRetention(t *testing.T) {
builder := builderWithoutChanges(2)
diff --git a/table/snapshot_producers.go b/table/snapshot_producers.go
index fc6d1d80a..b305d3f49 100644
--- a/table/snapshot_producers.go
+++ b/table/snapshot_producers.go
@@ -1365,11 +1365,13 @@ func (sp *snapshotProducer)
commitManifests(newManifests, addedContent []iceberg
return []Update{
addSnap,
- // Use 0 (not -1) for the optional fields so they are
omitted by
- // `omitempty` in JSON marshalling. -1 is a sentinel
meaning
- // "no limit" internally, but strict catalogs such as
AWS S3 Tables
- // reject a payload that explicitly contains negative
values.
- NewSetSnapshotRefUpdate(branch, sp.snapshotID,
BranchRef, 0, 0, 0),
+ // Carry over the branch's existing retention settings
so advancing
+ // the ref on commit does not silently discard them.
The update
+ // encodes exactly the current ref's retention
(settings the branch
+ // lacks stay 0 and are dropped by the `omitempty`
tags); the catalog
+ // applies a set-snapshot-ref as a pure replace, so
this fully
+ // determines the resulting ref rather than merging
with the old one.
+ sp.txn.meta.NewRetainingSnapshotRefUpdate(branch,
sp.snapshotID, BranchRef),
}, []Requirement{
AssertRefSnapshotID(branch,
sp.txn.meta.currentSnapshotID),
}, nil
diff --git a/table/transaction.go b/table/transaction.go
index 527d42443..6e91a35ca 100644
--- a/table/transaction.go
+++ b/table/transaction.go
@@ -313,7 +313,7 @@ func (t *Transaction) RollbackToSnapshot(snapshotID int64)
error {
snapshotID, cs.SnapshotID)
}
- update := NewSetSnapshotRefUpdate(MainBranch, snapshotID, BranchRef, 0,
0, 0)
+ update := meta.NewRetainingSnapshotRefUpdate(MainBranch, snapshotID,
BranchRef)
req := AssertRefSnapshotID(MainBranch, &cs.SnapshotID)
return t.apply([]Update{update}, []Requirement{req})
diff --git a/table/transaction_internal_test.go
b/table/transaction_internal_test.go
index 3d65f8410..3f10e509d 100644
--- a/table/transaction_internal_test.go
+++ b/table/transaction_internal_test.go
@@ -479,6 +479,57 @@ func
TestTransactionApplyDedupesIdenticalNonRefRequirements(t *testing.T) {
require.Equal(t, 1, specAsserts)
}
+func TestRollbackToSnapshotPreservesRetention(t *testing.T) {
+ txn, _ := createTestTransactionWithMemIO(t, *iceberg.UnpartitionedSpec)
+ now := time.Now().UnixMilli()
+
+ const (
+ minKeep = 5
+ maxSnapAgeMs = int64(172800000) // 2 days
+ maxRefAgeMs = int64(604800000) // 7 days
+ )
+
+ require.NoError(t, txn.meta.AddSnapshot(&Snapshot{
+ SnapshotID: 10,
+ SequenceNumber: 1,
+ ManifestList:
"mem://default/table-location/metadata/manifest-10.avro",
+ Summary: &Summary{Operation: OpAppend},
+ TimestampMs: now,
+ }))
+ require.NoError(t, txn.meta.SetSnapshotRef(
+ MainBranch, 10, BranchRef,
+ WithMinSnapshotsToKeep(minKeep),
+ WithMaxSnapshotAgeMs(maxSnapAgeMs),
+ WithMaxRefAgeMs(maxRefAgeMs),
+ ))
+
+ require.NoError(t, txn.meta.AddSnapshot(&Snapshot{
+ SnapshotID: 20,
+ ParentSnapshotID: transactionTestPtr(int64(10)),
+ SequenceNumber: 2,
+ ManifestList:
"mem://default/table-location/metadata/manifest-20.avro",
+ Summary: &Summary{Operation: OpAppend},
+ TimestampMs: now + 1,
+ }))
+ require.NoError(t, txn.meta.SetSnapshotRef(
+ MainBranch, 20, BranchRef,
+ WithMinSnapshotsToKeep(minKeep),
+ WithMaxSnapshotAgeMs(maxSnapAgeMs),
+ WithMaxRefAgeMs(maxRefAgeMs),
+ ))
+
+ require.NoError(t, txn.RollbackToSnapshot(10))
+
+ ref := txn.meta.refs[MainBranch]
+ require.Equal(t, int64(10), ref.SnapshotID, "rollback should move main
to the ancestor")
+ require.NotNil(t, ref.MinSnapshotsToKeep, "min-snapshots-to-keep must
survive rollback")
+ require.Equal(t, minKeep, *ref.MinSnapshotsToKeep)
+ require.NotNil(t, ref.MaxSnapshotAgeMs, "max-snapshot-age-ms must
survive rollback")
+ require.Equal(t, maxSnapAgeMs, *ref.MaxSnapshotAgeMs)
+ require.NotNil(t, ref.MaxRefAgeMs, "max-ref-age-ms must survive
rollback")
+ require.Equal(t, maxRefAgeMs, *ref.MaxRefAgeMs)
+}
+
func newTransactionWithSnapshotRefs(t *testing.T) *Transaction {
t.Helper()