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 211799b28 fix(hive): release acquired locks with a cleanup context 
(#1605)
211799b28 is described below

commit 211799b2864689719f8423e511785ea0c8db5b12
Author: Minh Vu <[email protected]>
AuthorDate: Fri Jul 31 16:23:55 2026 +0200

    fix(hive): release acquired locks with a cleanup context (#1605)
    
    ## Summary
    
    - release acquired Hive locks with a live cleanup context
    - bound cleanup with the existing five-second timeout
    - return unlock failures alongside the operation result
    - apply the cleanup path to drop, rename, and commit
    
    ## Why
    
    These operations deferred `Release` with the caller's context. If that
    context was canceled after lock acquisition, Hive could reject the
    unlock immediately, and the error was discarded. This could leave other
    writers blocked until the lock was cleared externally.
    
    ## Testing
    
    - `go test ./catalog/hive`
---
 catalog/hive/hive.go      | 30 ++++++++++++++++++++++++------
 catalog/hive/hive_test.go |  7 +++++--
 catalog/hive/lock.go      | 11 +++++++++++
 catalog/hive/lock_test.go | 29 +++++++++++++++++++++++++++++
 4 files changed, 69 insertions(+), 8 deletions(-)

diff --git a/catalog/hive/hive.go b/catalog/hive/hive.go
index 3f19d6c27..3e108c1e1 100644
--- a/catalog/hive/hive.go
+++ b/catalog/hive/hive.go
@@ -417,7 +417,7 @@ func sqlFromVersion(v *view.Version) (string, error) {
        return "", errors.New("view version has no SQL representation")
 }
 
-func (c *Catalog) DropTable(ctx context.Context, identifier table.Identifier) 
error {
+func (c *Catalog) DropTable(ctx context.Context, identifier table.Identifier) 
(err error) {
        database, tableName, err := identifierToTableName(identifier)
        if err != nil {
                return err
@@ -428,7 +428,13 @@ func (c *Catalog) DropTable(ctx context.Context, 
identifier table.Identifier) er
                return fmt.Errorf("%w: failed to acquire lock for %s.%s: %w", 
table.ErrCommitFailed, database, tableName, err)
        }
        defer func() {
-               _ = lock.Release(ctx)
+               if releaseErr := lock.releaseForCleanup(ctx); releaseErr != nil 
{
+                       if err != nil {
+                               err = errors.Join(err, releaseErr)
+                       } else {
+                               log.Printf("WARNING: failed to release Hive 
lock after dropping %s.%s: %v", database, tableName, releaseErr)
+                       }
+               }
        }()
 
        // Re-read after acquiring the lock so drop cannot act on table state 
that a
@@ -464,7 +470,7 @@ func (c *Catalog) PurgeTable(ctx context.Context, 
identifier table.Identifier) e
        return nil
 }
 
-func (c *Catalog) RenameTable(ctx context.Context, from, to table.Identifier) 
(*table.Table, error) {
+func (c *Catalog) RenameTable(ctx context.Context, from, to table.Identifier) 
(_ *table.Table, err error) {
        fromDB, fromTable, err := identifierToTableName(from)
        if err != nil {
                return nil, err
@@ -502,7 +508,13 @@ func (c *Catalog) RenameTable(ctx context.Context, from, 
to table.Identifier) (*
                        table.ErrCommitFailed, fromDB, fromTable, toDB, 
toTable, err)
        }
        defer func() {
-               _ = lock.Release(ctx)
+               if releaseErr := lock.releaseForCleanup(ctx); releaseErr != nil 
{
+                       if err != nil {
+                               err = errors.Join(err, releaseErr)
+                       } else {
+                               log.Printf("WARNING: failed to release Hive 
lock after renaming %s.%s to %s.%s: %v", fromDB, fromTable, toDB, toTable, 
releaseErr)
+                       }
+               }
        }()
 
        hiveTbl, err := c.getIcebergTable(ctx, fromDB, fromTable)
@@ -533,7 +545,7 @@ func (c *Catalog) RenameTable(ctx context.Context, from, to 
table.Identifier) (*
        return c.LoadTable(ctx, to)
 }
 
-func (c *Catalog) CommitTable(ctx context.Context, identifier 
table.Identifier, requirements []table.Requirement, updates []table.Update) 
(table.Metadata, string, error) {
+func (c *Catalog) CommitTable(ctx context.Context, identifier 
table.Identifier, requirements []table.Requirement, updates []table.Update) (_ 
table.Metadata, _ string, err error) {
        database, tableName, err := identifierToTableName(identifier)
        if err != nil {
                return nil, "", err
@@ -548,7 +560,13 @@ func (c *Catalog) CommitTable(ctx context.Context, 
identifier table.Identifier,
                        table.ErrCommitFailed, database, tableName, err)
        }
        defer func() {
-               _ = lock.Release(ctx)
+               if releaseErr := lock.releaseForCleanup(ctx); releaseErr != nil 
{
+                       if err != nil {
+                               err = errors.Join(err, releaseErr)
+                       } else {
+                               log.Printf("WARNING: failed to release Hive 
lock after committing %s.%s: %v", database, tableName, releaseErr)
+                       }
+               }
        }()
 
        currentHiveTbl, err := c.client.GetTable(ctx, database, tableName)
diff --git a/catalog/hive/hive_test.go b/catalog/hive/hive_test.go
index 040a78b62..d741a9566 100644
--- a/catalog/hive/hive_test.go
+++ b/catalog/hive/hive_test.go
@@ -743,7 +743,9 @@ func TestHiveCommitTableSynchronizesHMSMetadata(t 
*testing.T) {
                iceberg.NestedField{ID: 1, Name: "renamed", Type: 
iceberg.PrimitiveTypes.String},
        )
        mockClient := &mockHiveClient{}
-       expectImmediateTableLock(mockClient, 1)
+       mockClient.On("Lock", mock.Anything, 
mock.AnythingOfType("*hive_metastore.LockRequest")).
+               Return(&hive_metastore.LockResponse{Lockid: 1, State: 
hive_metastore.LockState_ACQUIRED}, nil).Once()
+       mockClient.On("Unlock", mock.Anything, 
int64(1)).Return(errors.New("unlock failed")).Once()
        mockClient.On("GetTable", mock.Anything, "test_database", "test_table").
                Return(existing, nil).Once()
        var altered *hive_metastore.Table
@@ -753,7 +755,7 @@ func TestHiveCommitTableSynchronizesHMSMetadata(t 
*testing.T) {
                }).Return(nil).Once()
        cat := NewCatalogWithClient(mockClient, iceberg.Properties{})
 
-       _, metadataLocation, err := cat.CommitTable(
+       metadata, metadataLocation, err := cat.CommitTable(
                ctx,
                TableIdentifier("test_database", "test_table"),
                nil,
@@ -771,6 +773,7 @@ func TestHiveCommitTableSynchronizesHMSMetadata(t 
*testing.T) {
        )
 
        require.NoError(t, err)
+       require.NotNil(t, metadata)
        require.NotNil(t, altered)
        require.Equal(t, metadataLocation, 
altered.Parameters[MetadataLocationKey])
        require.Equal(t, oldMetadataLocation, 
altered.Parameters[PreviousMetadataLocationKey])
diff --git a/catalog/hive/lock.go b/catalog/hive/lock.go
index a39d5e0f3..540ff70dc 100644
--- a/catalog/hive/lock.go
+++ b/catalog/hive/lock.go
@@ -164,6 +164,17 @@ func (l *HiveLock) Release(ctx context.Context) error {
        return l.client.Unlock(ctx, l.lockId)
 }
 
+func (l *HiveLock) releaseForCleanup(ctx context.Context) error {
+       cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 
pendingLockCleanupTimeout)
+       defer cancel()
+
+       if err := l.Release(cleanupCtx); err != nil {
+               return fmt.Errorf("failed to release acquired lock %d: %w", 
l.lockId, err)
+       }
+
+       return nil
+}
+
 func (l *HiveLock) LockID() int64 {
        if l == nil {
                return 0
diff --git a/catalog/hive/lock_test.go b/catalog/hive/lock_test.go
index 64d1a1423..6ec0c2eb4 100644
--- a/catalog/hive/lock_test.go
+++ b/catalog/hive/lock_test.go
@@ -323,6 +323,35 @@ func TestReleaseLock(t *testing.T) {
        mockClient.AssertExpectations(t)
 }
 
+func TestReleaseLockForCleanupUsesLiveBoundedContext(t *testing.T) {
+       mockClient := new(mockHiveClient)
+       ctx, cancel := context.WithCancel(context.Background())
+       cancel()
+       lock := &HiveLock{client: mockClient, lockId: 1000}
+
+       mockClient.On("Unlock", mock.MatchedBy(func(cleanupCtx context.Context) 
bool {
+               _, hasDeadline := cleanupCtx.Deadline()
+
+               return cleanupCtx.Err() == nil && hasDeadline
+       }), int64(1000)).Return(nil).Once()
+
+       require.NoError(t, lock.releaseForCleanup(ctx))
+       mockClient.AssertExpectations(t)
+}
+
+func TestReleaseLockForCleanupWrapsFailure(t *testing.T) {
+       mockClient := new(mockHiveClient)
+       unlockErr := errors.New("unlock failed")
+       lock := &HiveLock{client: mockClient, lockId: 1001}
+       mockClient.On("Unlock", mock.Anything, 
int64(1001)).Return(unlockErr).Once()
+
+       err := lock.releaseForCleanup(context.Background())
+
+       require.ErrorIs(t, err, unlockErr)
+       require.ErrorContains(t, err, "failed to release acquired lock 1001")
+       mockClient.AssertExpectations(t)
+}
+
 func TestCalculateBackoff(t *testing.T) {
        minWait := 100 * time.Millisecond
        maxWait := 1 * time.Second

Reply via email to