Skip to content
Draft
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
4 changes: 2 additions & 2 deletions catalog/glue/glue_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -911,7 +911,7 @@ func TestGluePurgeTableSwallowsPurgeFilesError(t *testing.T) {
assert.True(dropCalled.Load())
assert.Positive(removeCalls.Load())
assert.False(removeBeforeDrop.Load(), "PurgeTable should drop the catalog entry before removing files")
file, err := failingFS.Open(dataFile)
file, err := failingFS.Open(context.Background(), dataFile)
assert.NoError(err, "data file should remain when FileIO remove fails")
assert.NotNil(file)
assert.NoError(file.Close())
Expand Down Expand Up @@ -950,7 +950,7 @@ type failRemoveIO struct {
onRemove func()
}

func (f failRemoveIO) Remove(string) error {
func (f failRemoveIO) Remove(context.Context, string) error {
if f.onRemove != nil {
f.onRemove()
}
Expand Down
33 changes: 17 additions & 16 deletions catalog/hadoop/hadoop.go
Original file line number Diff line number Diff line change
Expand Up @@ -427,7 +427,8 @@ func (c *Catalog) writeVersionHint(ident table.Identifier, version int) {

if err := c.filesystem.Rename(tempPath, hintPath); err != nil {
log.Printf("hadoop catalog: failed to rename version hint: %v", err)
_ = c.filesystem.Remove(tempPath)
// Best-effort cleanup of the version-hint temp file; no caller context.
_ = c.filesystem.Remove(context.Background(), tempPath)
}
}

Expand Down Expand Up @@ -575,12 +576,12 @@ func (c *Catalog) CreateTable(ctx context.Context, ident table.Identifier, sc *i
tempPath := joinPath(c.isLocal, metaDir, uuid.New().String()+".metadata.json")

if err := internal.WriteTableMetadata(metadata, c.filesystem, tempPath, compression); err != nil {
_ = c.filesystem.Remove(tempPath)
_ = c.filesystem.Remove(ctx, tempPath)

return nil, fmt.Errorf("hadoop catalog: failed to write table metadata: %w", err)
}

if err := c.commitMetadataFile(ident, version, tempPath, metaPath, catalog.ErrTableAlreadyExists); err != nil {
if err := c.commitMetadataFile(ctx, ident, version, tempPath, metaPath, catalog.ErrTableAlreadyExists); err != nil {
return nil, err
}

Expand Down Expand Up @@ -711,12 +712,12 @@ func (c *Catalog) CommitTable(ctx context.Context, ident table.Identifier, reqs
tempPath := joinPath(c.isLocal, metaDir, uuid.New().String()+".metadata.json")

if err := internal.WriteTableMetadata(updated, c.filesystem, tempPath, compression); err != nil {
_ = c.filesystem.Remove(tempPath)
_ = c.filesystem.Remove(ctx, tempPath)

return nil, "", fmt.Errorf("hadoop catalog: failed to write table metadata: %w", err)
}

if err := c.commitMetadataFile(ident, newVersion, tempPath, newMetaPath, table.ErrCommitFailed); err != nil {
if err := c.commitMetadataFile(ctx, ident, newVersion, tempPath, newMetaPath, table.ErrCommitFailed); err != nil {
return nil, "", err
}

Expand All @@ -726,25 +727,25 @@ func (c *Catalog) CommitTable(ctx context.Context, ident table.Identifier, reqs
return updated, newMetaPath, nil
}

func (c *Catalog) commitMetadataFile(ident table.Identifier, version int, tempPath, metaPath string, conflictErr error) error {
func (c *Catalog) commitMetadataFile(ctx context.Context, ident table.Identifier, version int, tempPath, metaPath string, conflictErr error) error {
claimPath := c.metadataVersionClaimPath(ident, version)
for {
if err := c.filesystem.RenameNoReplace(tempPath, claimPath); err != nil {
if !errors.Is(err, fs.ErrExist) {
_ = c.filesystem.Remove(tempPath)
_ = c.filesystem.Remove(ctx, tempPath)

return fmt.Errorf("hadoop catalog: failed to claim metadata version: %w", err)
}

existingPath, exists, err := c.metadataVersionLocation(ident, version)
if err != nil {
_ = c.filesystem.Remove(tempPath)
_ = c.filesystem.Remove(ctx, tempPath)

return fmt.Errorf("hadoop catalog: failed to inspect metadata directory for version %d: %w",
version, err)
}
if exists {
_ = c.filesystem.Remove(tempPath)
_ = c.filesystem.Remove(ctx, tempPath)

return fmt.Errorf("%w: metadata file already exists for table %s: %s",
conflictErr, strings.Join(ident, "."), existingPath)
Expand All @@ -756,19 +757,19 @@ func (c *Catalog) commitMetadataFile(ident table.Identifier, version int, tempPa
continue
}

_ = c.filesystem.Remove(tempPath)
_ = c.filesystem.Remove(ctx, tempPath)

return fmt.Errorf("hadoop catalog: failed to inspect stale metadata claim %s: %w", claimPath, err)
}
if time.Since(claimInfo.ModTime()) < metadataClaimStaleAfter {
_ = c.filesystem.Remove(tempPath)
_ = c.filesystem.Remove(ctx, tempPath)

return fmt.Errorf("%w: metadata version already claimed for table %s: %s",
conflictErr, strings.Join(ident, "."), claimPath)
}

if err := c.filesystem.Remove(claimPath); err != nil && !errors.Is(err, fs.ErrNotExist) {
_ = c.filesystem.Remove(tempPath)
if err := c.filesystem.Remove(ctx, claimPath); err != nil && !errors.Is(err, fs.ErrNotExist) {
_ = c.filesystem.Remove(ctx, tempPath)

return fmt.Errorf("hadoop catalog: failed to clear stale metadata claim %s: %w", claimPath, err)
}
Expand All @@ -782,7 +783,7 @@ func (c *Catalog) commitMetadataFile(ident table.Identifier, version int, tempPa
removeClaim := true
defer func() {
if removeClaim {
_ = c.filesystem.Remove(claimPath)
_ = c.filesystem.Remove(ctx, claimPath)
}
}()

Expand Down Expand Up @@ -956,7 +957,7 @@ func (c *Catalog) CreateNamespace(_ context.Context, ns table.Identifier, props
return nil
}

func (c *Catalog) DropNamespace(_ context.Context, ns table.Identifier) error {
func (c *Catalog) DropNamespace(ctx context.Context, ns table.Identifier) error {
if err := validateIdentifier(ns); err != nil {
return err
}
Expand Down Expand Up @@ -998,7 +999,7 @@ func (c *Catalog) DropNamespace(_ context.Context, ns table.Identifier) error {
return fmt.Errorf("%w: %s", catalog.ErrNamespaceNotEmpty, strings.Join(ns, "."))
}

return c.filesystem.Remove(path)
return c.filesystem.Remove(ctx, path)
}

func (c *Catalog) CheckNamespaceExists(_ context.Context, ns table.Identifier) (bool, error) {
Expand Down
6 changes: 3 additions & 3 deletions catalog/hadoop/hadoop_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -123,11 +123,11 @@ var _ HadoopCatalogFS = (*stubHadoopCatalogFS)(nil)

type stubIO struct{}

func (stubIO) Open(string) (icebergio.File, error) {
func (stubIO) Open(context.Context, string) (icebergio.File, error) {
return nil, fs.ErrNotExist
}

func (stubIO) Remove(string) error {
func (stubIO) Remove(context.Context, string) error {
return nil
}

Expand Down Expand Up @@ -1857,7 +1857,7 @@ func (s *HadoopCatalogTestSuite) TestCommitMetadataFileFailsClosedOnMetadataScan
s.cat.filesystem = originalFS
}()

err = s.cat.commitMetadataFile(ident, 1, tempPath, metaPath, table.ErrCommitFailed)
err = s.cat.commitMetadataFile(context.Background(), ident, 1, tempPath, metaPath, table.ErrCommitFailed)
s.Require().Error(err)
s.Contains(err.Error(), "failed to inspect metadata directory for version 1")
s.ErrorIs(err, fs.ErrPermission)
Expand Down
2 changes: 1 addition & 1 deletion catalog/hive/hive.go
Original file line number Diff line number Diff line change
Expand Up @@ -764,7 +764,7 @@ func (c *Catalog) DropView(ctx context.Context, identifier table.Identifier) err
return fmt.Errorf("failed to load filesystem for view metadata: %w", err)
}

if err := fs.Remove(metadataLocation); err != nil {
if err := fs.Remove(ctx, metadataLocation); err != nil {
return fmt.Errorf("failed to remove view metadata file at %s: %w", metadataLocation, err)
}

Expand Down
4 changes: 2 additions & 2 deletions catalog/hive/hive_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -960,7 +960,7 @@ func TestHivePurgeTableSwallowsPurgeFilesError(t *testing.T) {
assert.True(dropCalled.Load())
assert.Positive(removeCalls.Load())
assert.False(removeBeforeDrop.Load(), "PurgeTable should drop the catalog entry before removing files")
file, err := failingFS.Open(dataFile)
file, err := failingFS.Open(context.Background(), dataFile)
assert.NoError(err, "data file should remain when FileIO remove fails")
assert.NotNil(file)
assert.NoError(file.Close())
Expand Down Expand Up @@ -997,7 +997,7 @@ type failRemoveIO struct {
onRemove func()
}

func (f failRemoveIO) Remove(string) error {
func (f failRemoveIO) Remove(context.Context, string) error {
if f.onRemove != nil {
f.onRemove()
}
Expand Down
6 changes: 3 additions & 3 deletions catalog/rest/rest_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -382,7 +382,7 @@ func (s *RestIntegrationSuite) TestWriteCommitTable() {
s.Require().NoError(err)
s.Require().NoError(pqarrow.WriteTable(table, fw, table.NumRows(),
nil, pqarrow.DefaultWriterProps()))
defer mustFS(s.T(), tbl).Remove(pqfile)
defer mustFS(s.T(), tbl).Remove(s.ctx, pqfile)

txn := tbl.NewTransaction()
s.Require().NoError(txn.AddFiles(s.ctx, []string{pqfile}, nil, false))
Expand Down Expand Up @@ -458,10 +458,10 @@ func (s *RestIntegrationSuite) TestMultiTableCommit() {

// Write a parquet data file for each table.
pq1 := s.writeParquetFile(tbl1, tableSchemaSimple, "multi-txn-1")
defer mustFS(s.T(), tbl1).Remove(pq1)
defer mustFS(s.T(), tbl1).Remove(s.ctx, pq1)

pq2 := s.writeParquetFile(tbl2, tableSchemaSimple, "multi-txn-2")
defer mustFS(s.T(), tbl2).Remove(pq2)
defer mustFS(s.T(), tbl2).Remove(s.ctx, pq2)

// Build transactions that add data files to each table.
tx1 := tbl1.NewTransaction()
Expand Down
2 changes: 1 addition & 1 deletion catalog/rest/scan_planning_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -262,7 +262,7 @@ func (s *RestIntegrationSuite) createScanPlanningTable(cat *rest.Catalog, ident
s.Require().NoError(err)
s.Require().NoError(pqarrow.WriteTable(arrowTable, file, arrowTable.NumRows(),
nil, pqarrow.DefaultWriterProps()))
s.T().Cleanup(func() { s.Require().NoError(mustFS(s.T(), tbl).Remove(dataPath)) })
s.T().Cleanup(func() { s.Require().NoError(mustFS(s.T(), tbl).Remove(s.ctx, dataPath)) })

txn := tbl.NewTransaction()
s.Require().NoError(txn.AddFiles(s.ctx, []string{dataPath}, nil, false))
Expand Down
2 changes: 1 addition & 1 deletion catalog/rest/scan_planning_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1833,7 +1833,7 @@ func TestPlanScopedIOAlreadyExpiredCredentials(t *testing.T) {

fs, err := p.Load(context.Background())
require.NoError(t, err)
_, err = fs.Open("file:///bucket/data.parquet")
_, err = fs.Open(context.Background(), "file:///bucket/data.parquet")
require.ErrorIs(t, err, ErrVendedCredentialsExpired)
}

Expand Down
8 changes: 4 additions & 4 deletions catalog/rest/vended_creds.go
Original file line number Diff line number Diff line change
Expand Up @@ -296,22 +296,22 @@ func newPrefixScopedIO(ctx context.Context, baseProps iceberg.Properties, creden
}
}

func (p *prefixScopedIO) Open(name string) (iceio.File, error) {
func (p *prefixScopedIO) Open(ctx context.Context, name string) (iceio.File, error) {
fs, err := p.filesystemFor(name)
if err != nil {
return nil, err
}

return fs.Open(name)
return fs.Open(ctx, name)
}

func (p *prefixScopedIO) Remove(name string) error {
func (p *prefixScopedIO) Remove(ctx context.Context, name string) error {
fs, err := p.filesystemFor(name)
if err != nil {
return err
}

return fs.Remove(name)
return fs.Remove(ctx, name)
}

func (p *prefixScopedIO) filesystemFor(name string) (iceio.IO, error) {
Expand Down
4 changes: 2 additions & 2 deletions catalog/sql/sql.go
Original file line number Diff line number Diff line change
Expand Up @@ -751,7 +751,7 @@ func removeUncommittedMetadata(ctx context.Context, metadataLocation string, loa
return fmt.Errorf("failed to load filesystem while removing uncommitted metadata %s: %w", metadataLocation, err)
}

if err := fs.Remove(metadataLocation); err != nil {
if err := fs.Remove(ctx, metadataLocation); err != nil {
return fmt.Errorf("failed to remove uncommitted metadata %s: %w", metadataLocation, err)
}

Expand Down Expand Up @@ -1825,7 +1825,7 @@ func (c *Catalog) DropView(ctx context.Context, identifier table.Identifier) err
return err
}

_ = fs.Remove(metadataLocation)
_ = fs.Remove(ctx, metadataLocation)
}

return nil
Expand Down
2 changes: 1 addition & 1 deletion catalog/sql/sql_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -292,7 +292,7 @@ func (s *SQLIntegrationSuite) TestWriteCommitTable() {
s.Require().NoError(pqarrow.WriteTable(table, fw, table.NumRows(),
nil, pqarrow.DefaultWriterProps()))
defer func(fs io.IO, name string) {
err = fs.Remove(name)
err = fs.Remove(s.ctx, name)
s.Require().NoError(err)
}(fs, pqfile)

Expand Down
2 changes: 1 addition & 1 deletion cmd/iceberg/clean_orphan_files_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,7 +95,7 @@ func TestRunCleanOrphanFilesPreviewsPlanBeforeDeletion(t *testing.T) {
assert.False(t, deleted.DryRun)
assert.Equal(t, 1, deleted.OrphanFileCount)
assert.Equal(t, orphanPath, deleted.OrphanFiles[0].Path)
_, err = memFS.Open(orphanPath)
_, err = memFS.Open(context.Background(), orphanPath)
assert.ErrorIs(t, err, fs.ErrNotExist)
}

Expand Down
5 changes: 3 additions & 2 deletions internal/mock_fs.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ package internal

import (
"bytes"
"context"
"errors"
sio "io"
"io/fs"
Expand All @@ -31,7 +32,7 @@ type MockFS struct {
mock.Mock
}

func (m *MockFS) Open(name string) (io.File, error) {
func (m *MockFS) Open(_ context.Context, name string) (io.File, error) {
args := m.Called(name)

return args.Get(0).(io.File), args.Error(1)
Expand All @@ -47,7 +48,7 @@ func (m *MockFS) WriteFile(name string, content []byte) error {
return m.Called(name, content).Error(0)
}

func (m *MockFS) Remove(name string) error {
func (m *MockFS) Remove(_ context.Context, name string) error {
return m.Called(name).Error(0)
}

Expand Down
10 changes: 6 additions & 4 deletions io/gocloud/blobfs/blob.go
Original file line number Diff line number Diff line change
Expand Up @@ -285,7 +285,9 @@ func directoryName(key string) string {
return pathpkg.Base(key)
}

func (bfs *FileIO) Open(path string) (icebergio.File, error) {
func (bfs *FileIO) Open(_ context.Context, path string) (icebergio.File, error) {
// The returned File is read after Open returns, so its reads stay bound to
// the FileIO's lifetime context (bfs.ctx) rather than this per-call ctx.
originalPath := path
var err error
path, err = bfs.preprocess(path)
Expand All @@ -306,17 +308,17 @@ func (bfs *FileIO) Open(path string) (icebergio.File, error) {
return &blobOpenFile{Reader: r, name: name, key: key, b: bfs, ctx: bfs.ctx}, nil
}

func (bfs *FileIO) Remove(name string) error {
func (bfs *FileIO) Remove(ctx context.Context, name string) error {
key, err := bfs.preprocess(name)
if err != nil {
return &fs.PathError{Op: "remove", Path: name, Err: err}
}

if err := bfs.Delete(bfs.ctx, key); err != nil {
if err := bfs.Delete(ctx, key); err != nil {
if gcerrors.Code(err) == gcerrors.NotFound {
marker := directoryMarker(key)
if marker != "" {
if markerErr := bfs.Delete(bfs.ctx, marker); markerErr == nil {
if markerErr := bfs.Delete(ctx, marker); markerErr == nil {
return nil
} else if gcerrors.Code(markerErr) != gcerrors.NotFound {
return &fs.PathError{Op: "remove", Path: name, Err: markerErr}
Expand Down
Loading
Loading